Published

🧱 2d.3 / Tier B — el brazo nativo PRESERVANTE del pipeline.output (re-deploy, merge, delta, model-node)

Connect any source, model it as an ontology, transform it, and operationalize it, analytics, automation and machine learning, under one governed, self-hostable roof. --- Most teams stitch the...

🧱 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.output que 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 lee dataset_rows → materialiser a 0 accesos crudos).
  • B1/B2: WriteStrategy.nativeReplayMode (snapshot_replace/defaultupsert, append_alwaysappend); 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) + rowFilter en writeOutputToTridentappend solo los nuevos. Opción (A) del §3 elegida.
  • B4 delta nativo (opción A): computeReplaceSql delta-aware en las 3 estrategias (añade AND ${deltaPredicate('$1','$2','$3')} cuando ctx.sourceDelta, mirror de apply) + quitar el throw del brazo iceberg. El dispatch de re-deploy (upsert/append/append_new) aplica la ventana; windowAllowsDelta garantiza que es aditiva. Asunción documentada (TODO Fase 5): el read de la ventana sigue en dataset_rows (ledger PG). Ver 2d3-b4-delta-native.md.
  • B5 model-node: writeRowsToTrident (nuevo, hermano en-memoria de writeOutputToTrident) + el loop bufferea las ventanas → UN solo replace (acumulación) + la correlación source↔predicción va como columna explícita __source_row_index__ (ml-runner posee row_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.output nativo 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:

#SuperficieCondición que lanza hoySemántica que falta
B1snapshot_replace re-deployhasPriorRowsupsert nativo✅ merge-por-PK preservante (ml-runner upsert)
B2append_always re-deployhasPriorRowsappend nativo✅ append acumulativo (ml-runner append)
B3append_newsin computeReplaceSql → anti-join + append✅ anti-join client-side (PKs del target por facade) + append
B4delta nativo (sourceDelta)throw → compute delta-aware + dispatch✅ opción A (read ventana=PG, write=re-deploy dispatch); asunción dataset_rows-no-drenado (TODO Fase 5)
B5model-node (2 sink.run)fail-loud → buffer + single replace + __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 preservanteupsert (join = PK)❌ No — el merge lo hace ml-runner
B2 append_alwaysappend❌ No — append puro
B3 append_new (anti-join)append de solo-nuevos (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.upsert NO borra filas ausentes del nuevo compute. Eso ES la semántica "preservante" de snapshot_replace (rowsPreserved = filas previas no-conflictivas). Confirmado alineado con WriteResult. Si alguna vez se quisiera un replace que TAMBIÉN borra ausentes por PK, eso sería una primitiva merge-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 un aggregate/distinct sobre el PK), materializar el set en memoria acotada, y AND NOT (pk = ANY($existing)) en el computeSql; luego writeMode: '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_new en ml-runner (anti-join server-side, PyIceberg): un write_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ón dataset_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_replaceupsert (las filas de la ventana mergean por PK).
  • delta sobre append_alwaysappend (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) → writeOutputToTrident con el write-mode del nodo (típicamente replace/upsert).
  • Abierto: la correlación __source_row_key asume 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 un row_index PG. Verificar que __source_row_key sobrevive el round-trip nativo (candidato a un test dedicado antes de flipear model-node).

6. Cambios de código (mínimos, ordenados)

  1. Probe #0 → facade. En el gate, reemplazar count(*) FROM dataset_rows por un check de target nativo: icebergFreshness(producedDatasetId).hasSucceededSnapshot (existe snapshot ⇒ re-deploy) o un handle.rowCount. Elimina el bug latente del drenado.
  2. WriteStrategy interface. Generalizar computeReplaceSql?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.
  3. writeOutputToTrident. Ya acepta writeMode — dejar de hardcodear 'replace' en el brazo del deploy; pasar compute.writeMode. (El helper no cambia.)
  4. B3 anti-join. Helper readTargetPrimaryKeys(datasetId, pk) por la facade + inyección NOT (pk = ANY($n)) en el compute.
  5. model-node. Cablear sus 2 sink.run iceberg-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ónRecomendación
Probe de re-deployOráculo de frescura (hasSucceededSnapshot), no count(*) dataset_rows. Cierra el bug del drenado.
append_newAnti-join client-side por la facade (A) primero; primitiva runner (B) solo si el cap de PK topa.
snapshot_replace preservanteupsert nativo (ml-runner merge por PK) — sin target-read.
Atomicidad re-deploywithNativeControlPlane 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 append re-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_key sobrevive 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.