ADR-0016: pipeline a grafo (DAG) — nodi/edge validati a runtime, non a compile-time
- Stato: Accettata
- Data: 2026-07-09
Contesto
Le pipeline erano modellate come una catena lineare (pipeline().from().pipe()...to()): un solo source, una sequenza di Transform, un solo sink. Questo non copre casi reali frequenti — join a chiave fra due sorgenti (es. ordini + anagrafica clienti) e fan-out di uno stesso stream su più destinazioni indipendenti (es. fatto arricchito + riepilogo aggregato). Serviva un vero DAG: nodi source | transform | join | sink collegati da edge, con dead-letter isolabile per nodo.
Decisione
Aggiunto un DSL a grafo (graph(name).source().transform().join().sink().build(), src/core/graph-builder.ts) che costruisce una GraphSpec (nodi + edge, src/core/graph.ts) eseguita da graph-runner.ts. La struttura del grafo — id duplicati, edge verso nodi sconosciuti, arità delle porte build/probe di un join, nodi source con edge in ingresso, nodi sink con edge in uscita, nodi irraggiungibili, fan-in non passante per un join, cicli — è validata a runtime, alla chiamata di .build() (validateGraph()), non a compile-time via il sistema di tipi.
Alternative scartate:
- Validazione a compile-time via tipi TypeScript — builder con generics che tracciano gli id dei nodi dichiarati, per rifiutare a compile-time un edge verso un id inesistente o un ordine errato di
.join(). Più sicura in teoria, ma con un fluent builder a id stringa arbitrari il costo in complessità dei tipi (generics ricorsivi su union di stringhe letterali) è alto e fragile — messaggi d'errore TypeScript illeggibili, inference che si rompe facilmente — a fronte di un grafo che, in questo dominio, si costruisce una volta all'avvio del processo e non in un hot path: il costo di una validazione runtime eager (fail-fast al primo.build(), prima di qualunque esecuzione) è trascurabile e il messaggio d'errore (GraphValidationError, italiano, con l'id del nodo/edge incriminato) è immediatamente più leggibile di un errore di tipo. - Windowed join (finestre temporali/di conteggio per unire due stream continui) — scartato perché fuori scopo: le pipeline di questo framework sono batch (source/sink array o file/DB interi, non stream infiniti), quindi basta una hash join classica con il lato
buildmaterializzato per intero (vedi anti-pattern nella skilletl-pipeline). Un windowed join aggiungerebbe stato temporale e complessità di watermarking che qui non servono. - Dead-letter unico per l'intera pipeline (un solo sink di scarto condiviso da tutti i nodi) — scartato in favore del dead-letter per nodo, scelta esplicita fatta in fase di brainstorming: con un solo dead-letter condiviso si perderebbe l'origine dello scarto (quale nodo ha rifiutato il record) a meno di annotarla nel payload, e non sarebbe possibile abilitare il dead-letter solo su alcuni nodi. Il costo è che ogni nodo isolabile richiede un sink opzionale esplicito nella firma della pipeline (vedi
OrdersEnrichedDailyOptionsinsrc/pipelines/orders-enriched-daily/pipeline.ts), invece di un unico parametro.
Conseguenze
- Pro: builder semplice da leggere e scrivere, errori di grafo con messaggio umano (id del nodo, motivo) invece di errori di tipo generici; aggiungere un nuovo kind di nodo non richiede ristrutturare i generics del builder.
- Contro: un grafo mal formato (id sbagliato, join senza
probe, ciclo) fallisce solo quando il modulo che chiama.build()viene eseguito (es. al primopnpm run etl run <pipeline>, o in un test), non all'editing/compilazione. Mitigato da: (1)validateGraph()fallisce eager su.build(), prima di qualunque esecuzione di record — non serve arrivare a runtime "in produzione" per scoprire il problema, basta costruire il grafo; (2) ogni pipeline a grafo ha un test che chiama.build()(vediorders-enriched-daily/pipeline.test.ts), quindi l'errore emerge comunque in CI prima del merge, non in produzione. RunnablePipeline.run(src/pipelines/registry.ts) è stato allargato aPromise<RunReport | GraphRunReport>perchéGraphPipeline.run()restituisce un report strutturalmente diverso (nodi/edge multipli, niente aggregato top-levelextracted/written/rejected— quel concetto è solo della pipeline lineare a singolo source/sink). CLI e registry leggono solo i campi comuni (status, serializzazione JSON), quindi l'unione è un allargamento sicuro, non un compromesso sul contratto.
Riferimenti
src/core/graph.ts(GraphSpec,validateGraph),src/core/graph-builder.ts(graph()DSL),src/core/graph-runner.ts(esecuzione:tee()per il fan-out,hashJoin()per il join a chiave).- Esempio end-to-end:
src/pipelines/orders-enriched-daily/, docs/pipelines/orders-enriched-daily.md. - Piano di implementazione: docs/plans/2026-07-09-pipeline-dag-nodes-edges.md.

