🧱 2d.3 / Tier B — el brazo nativo PRESERVANTE del pipeline.output (re-deploy, merge, delta, model-node)
Entregable de approach (2026-07-09). Cierra lo que 2d.2 dejó fail-loud: los writes de
pipeline.outputque preservan o mezclan con el target existente (re-deploys), el delta nativo, y el write-arm del model-node. Companion de paso6-write-approach.md · registro vivo port-migration.md · spectro 2d 2d-pipeline-output-flip.md.Estado de entrada (2d.2 CERRADO): el brazo iceberg del deploy escribe nativo SOLO en el caso replace-puro (primer deploy, target vacío, estrategia con
computeReplaceSql). Verificado por el canary VERDE E2E (scripts/dataspaces/pipeline-output-canary.ts). Todo lo demás lanza "Tier B — deferred to 2d.3" — nunca un write silencioso incorrecto.defaultMode=pg→ todo inerte hasta el flip por-dataset.🟢 PROGRESO 2d.3 (2026-07-09) — SLICE 1 (probe#0 + B1 + B2) + SLICE 2 (B3) CERRADOS:
- probe #0: first-vs-re-deploy por el oráculo de frescura (
icebergFreshness.hasSucceededSnapshot, ya NO leedataset_rows→ materialiser a 0 accesos crudos).- B1/B2:
WriteStrategy.nativeReplayMode(snapshot_replace/default→upsert,append_always→append); ml-runner mergea server-side, sin target-read.- B3
append_new:WriteStrategy.nativeAppendNew+buildAppendNewFilter(native-output-write.ts) — anti-join client-side leyendo los PKs del target POR LA FACADE (handle.stream({columns:[pk]}), set acotado al cap) +rowFilterenwriteOutputToTrident→appendsolo los nuevos. Opción (A) del §3 elegida.- B4
delta nativo(opción A):computeReplaceSqldelta-aware en las 3 estrategias (añadeAND ${deltaPredicate('$1','$2','$3')}cuandoctx.sourceDelta, mirror deapply) + quitar el throw del brazo iceberg. El dispatch de re-deploy (upsert/append/append_new) aplica la ventana;windowAllowsDeltagarantiza que es aditiva. Asunción documentada (TODO Fase 5): el read de la ventana sigue endataset_rows(ledger PG). Ver 2d3-b4-delta-native.md.- B5
model-node:writeRowsToTrident(nuevo, hermano en-memoria dewriteOutputToTrident) + el loop bufferea las ventanas → UN soloreplace(acumulación) + la correlación source↔predicción va como columna explícita__source_row_index__(ml-runner poseerow_index) + truncate-on-reuse native → no-op. Ver 2d3-b5-model-node.md.- Verificado E2E (canary §5a-e): upsert + append + append_new + delta + model (2 ventanas → 1 snapshot,
__source_row_index__preservado).🎉 TIER B CERRADO (2026-07-09).
pipeline.outputnativo cubre el espectro completo del deploy + model-node: primer-deploy (replace) + re-deploy (upsert/append/anti-join) + delta (incremental) + model (predicciones).defaultMode=pg→ todo opt-in por-dataset (canary). Diferido a Karma/Fase 5: read nativo de la ventana delta (B4 §5), unificación de la correlación model (B5 §5).
1. Qué queda fail-loud hoy (el inventario exacto)
Del gate en dataset-output-materialiser.ts (brazo iceberg: de sink.run) y del model-node:
| # | Superficie | Condición que lanza hoy | Semántica que falta |
|---|---|---|---|
snapshot_replace re-deploy | hasPriorRowsupsert nativo | ✅ merge-por-PK preservante (ml-runner upsert) | |
append_always re-deploy | hasPriorRowsappend nativo | ✅ append acumulativo (ml-runner append) | |
append_new | computeReplaceSqlappend | ✅ anti-join client-side (PKs del target por facade) + append | |
delta nativo (sourceDelta) | ✅ opción A (read ventana=PG, write=re-deploy dispatch); asunción dataset_rows-no-drenado (TODO Fase 5) | ||
model-node (2 sink.run) | __source_row_index__ col | ✅ correlación por columna + acumulación por ventanas + truncate→no-op |
El probe del gate es el bug latente #0. El gate decide first-deploy vs Tier-B con:
SELECT count(*) FROM dataset_rows WHERE dataset_id = $1 -- ⚠️ lee PG, no el target nativo
Para un target nativo con dataset_rows drenado (Fase 5) esto da 0 → el gate creería "primer deploy" y haría un REPLACE que borra el histórico nativo acumulado. Hoy sólo es seguro porque la flip-doc manda NO drenar dataset_rows durante la ventana canary (el count espeja el snapshot). La pieza #0 de Tier B = mover el probe del target al lado nativo (oráculo de frescura / count por la facade), no a dataset_rows.
2. El insight — no todo Tier B necesita "leer el target"
ml-runner ya expone write_mode ∈ {replace, upsert, append} (ingest-client.ts; upsert = table.upsert(join_cols=pk), merge-on-read, update matched + insert new, nunca delete). Mapeando las estrategias a esos primitivos, la mayoría de Tier B NO requiere un target-read en Node — ml-runner hace el merge server-side:
| Estrategia (re-deploy) | Write-mode nativo | ¿Target-read en Node? |
|---|---|---|
B1 snapshot_replace preservante | upsert (join = PK) | ❌ No — el merge lo hace ml-runner |
B2 append_always | append | ❌ No — append puro |
B3 append_new (anti-join) | append de solo-nuevos | ✅ Sí (o primitiva nueva, ver §3) |
Es decir: B1 y B2 se cierran cambiando el writeMode hardcodeado 'replace' de writeOutputToTrident por el modo de la estrategia — el compute es la misma proyección (data jsonb por fila); solo cambia el modo de commit. B3 es el único que necesita conocer los PK existentes.
⚠️ Matiz de
upsert(B1):table.upsertNO borra filas ausentes del nuevo compute. Eso ES la semántica "preservante" desnapshot_replace(rowsPreserved= filas previas no-conflictivas). Confirmado alineado conWriteResult. Si alguna vez se quisiera un replace que TAMBIÉN borra ausentes por PK, eso sería una primitivamerge-delete(no existe en el writer nativo hoy → Karma).
3. B3 append_new — las dos rutas (decisión)
append_new = insertar filas cuyo PK no existe en el target (skip colisiones, sin update). upsert no sirve (actualizaría las matched). Opciones:
- (A) Anti-join client-side por la facade (recomendada para el primer corte): leer los PK del target nativo por
handle.read({ columns: [pk], includeIdentity:false })(o unaggregate/distinctsobre el PK), materializar el set en memoria acotada, yAND NOT (pk = ANY($existing))en elcomputeSql; luegowriteMode: 'append'. Reusa el read nativo YA verificado; cero cambio en ml-runner. Coste: el set de PK cabe en memoria (cap como el de escritura,PIPELINE_OUTPUT_NATIVE_MAX_ROWS); un target enorme lanza claro (→ ruta B). - (B) Primitiva
append_newen ml-runner (anti-join server-side, PyIceberg): unwrite_mode: 'append_new'que hace el anti-join por PK dentro del runner. Más limpio a escala, pero es superficie nueva en el runner → diferir a menos que (A) tope el cap. Karma lo subsume.
Recomendación: (A) ahora (contenido, sin tocar el runner); reevaluar (B) si el cap de PK molesta.
4. B4 delta nativo
📄 Entregable de espectro completo: 2d3-b4-delta-native.md (2026-07-09) — los dos ejes ortogonales (READ de la ventana = ledger PG, difícil; WRITE = ya hecho por B1/B2/B3), el espectro del read (A interim PG / B native-scan / C Karma), los cambios mínimos (delta-aware
computeReplaceSql+ quitar 1 throw), la asuncióndataset_rows-no-drenado, y el plan de canary. Recomendación: opción A ahora.
El delta hoy fuerza PG (sourceDelta → throw). El compute delta ya narra su source por sourceTable (linchpin) y la ventana (low,high]. Nativizarlo = correr ese mismo compute-de-ventana y aplicarlo con el write-mode de la estrategia:
- delta sobre
snapshot_replace→upsert(las filas de la ventana mergean por PK). - delta sobre
append_always→append(las filas de la ventana se añaden). - delta sobre
append_new→ anti-join (§3) restringido a la ventana.
Prerequisito: el probe #0 (target por facade) y B1/B2/B3. El delta es una restricción de source (ya resuelta por el linchpin), no un write-mode nuevo → cae casi solo una vez B1-B3 están. La ventana sigue leyéndose de PG (el ledger de watermark es control-plane PG); solo cambia dónde aterrizan las filas.
5. B5 model-node write-arm
El model-node (model-node-materialiser.ts) tiene 2 sink.run fail-loud. Su compute correlaciona predicciones con la fila fuente por __source_row_key (§9 del spectro doc). Nativizar:
- El source read ya va por el linchpin (P0.2). El write produce
{data jsonb}por fila (predicción + join key) →writeOutputToTridentcon el write-mode del nodo (típicamente replace/upsert). - Abierto: la correlación
__source_row_keyasume identidad estable del source. Con un source nativo, la identidad la posee ml-runner (__row_id/__row_index) → el join key debe derivarse de esa identidad, no de unrow_indexPG. Verificar que__source_row_keysobrevive el round-trip nativo (candidato a un test dedicado antes de flipear model-node).
6. Cambios de código (mínimos, ordenados)
- Probe #0 → facade. En el gate, reemplazar
count(*) FROM dataset_rowspor un check de target nativo:icebergFreshness(producedDatasetId).hasSucceededSnapshot(existe snapshot ⇒ re-deploy) o unhandle.rowCount. Elimina el bug latente del drenado. WriteStrategyinterface. GeneralizarcomputeReplaceSql?→computeNativeSql?(ctx): { sql; params; writeMode: 'replace'|'upsert'|'append' }(o añadir hermano) para que cada estrategia declare su write-mode nativo.snapshot_replace=upsert(PK),append_always=append,append_new=append+anti-join (§3),default=replace.writeOutputToTrident. Ya aceptawriteMode— dejar de hardcodear'replace'en el brazo del deploy; pasarcompute.writeMode. (El helper no cambia.)- B3 anti-join. Helper
readTargetPrimaryKeys(datasetId, pk)por la facade + inyecciónNOT (pk = ANY($n))en el compute. - model-node. Cablear sus 2
sink.runiceberg-arm análogo al deploy, tras verificar__source_row_key.
Invariante que se mantiene: todo bajo defaultMode=pg → inerte; el flip es por-dataset (canary). Un caso aún no cableado sigue fail-loud, nunca write silencioso.
7. Decisiones abiertas (con recomendación)
| Decisión | Recomendación |
|---|---|
| Probe de re-deploy | Oráculo de frescura (hasSucceededSnapshot), no count(*) dataset_rows. Cierra el bug del drenado. |
append_new | Anti-join client-side por la facade (A) primero; primitiva runner (B) solo si el cap de PK topa. |
snapshot_replace preservante | upsert nativo (ml-runner merge por PK) — sin target-read. |
| Atomicidad re-deploy | withNativeControlPlane sigue all-or-nothing por snapshot; re-deploy idempotente. Sin ledger nuevo (Karma-era). |
| ¿Karma? | No es prerequisito. upsert/append/anti-join(client) cubren Tier B con el runner de hoy. Karma subsume el compute+merge después. |
8. Verificación — extender el canary
scripts/dataspaces/pipeline-output-canary.ts prueba hoy el replace-puro. Tier B añade casos (mismo patrón: setup en main.test → write → control-plane → paridad PG↔Iceberg, limpieza):
- B1 re-deploy
upsert: escribir snapshot inicial → re-escribir un delta con PKs solapados + nuevos → leer back → paridad contra el merge esperado (update matched + insert new + keep resto). - B2
appendre-deploy: dos writes → count acumulado + paridad. - B3
append_new: target con PKs 3 + compute 5 → solo 5 añadidas; paridad +rowsSkipped. - B4 delta: ventana acota el compute; paridad contra el compute-de-ventana en PG.
- B5 model-node:
__source_row_keysobrevive el round-trip nativo.
DoD de 2d.3: los 5 casos verdes E2E en main.test + el gate ya no lee dataset_rows + port-migration.md/paso6-write-approach.md marcan Tier B cerrado.
Documento de approach. Grounded en dataset-output-materialiser.ts (gate + brazo iceberg), dataset-write-strategies/* (WriteStrategy + computeReplaceSql), iceberg-native-write.ts (write-modes), ingest-client.ts (upsert/append). Mantener al día conforme cada caso Tier B aterrice.