GraphRunReport — il risultato di una esecuzione
Il servizio è stateless e agnostico: esegue una pipeline e restituisce un GraphRunReport con tutti i dettagli del caso. Non possiede DB, non scrive da nessuna parte, non fa push. Cosa farne (archiviare, alert, ignorare) è deciso dal chiamante.
Contratto di I/O (CLI / container)
stdout= un solo oggetto JSON: il report. Parsabile conjq.stderr= log umani / progress (rigapipeline.completed, ecc.).- exit code:
0sestatusèsuccess/partial;≠0sefailed. - Anche in caso di errore viene emesso un report (
status: "failed").
etl run langfuse-sessions-daily > report.json 2> run.log # separa report e log
etl run langfuse-sessions-daily | jq '.status, .written' # consuma al voloCome libreria, graph(...).build().run(ctx) restituisce direttamente il report e lancia sull'errore fatale — contratto invariato. Da ADR-0038 l'errore lanciato porta con sé un report completo (topologia reale e conteggi parziali), recuperabile con graphRunReportOf(err):
try {
report = await pipeline.run(ctx);
} catch (err) {
// L'errore mantiene il suo tipo concreto (ValidationError, ExtractError, …).
report =
graphRunReportOf(err) ??
buildFailedGraphReport(ctx, name, err, startedAtMs);
}buildFailedGraphReport resta il fallback per i fallimenti sollevati prima che il grafo esista (es. wiring dei connettori in run.ts), dove una topologia da riportare non c'è davvero.
Schema (reportSchemaVersion: 4)
| Campo | Tipo | Note |
|---|---|---|
reportSchemaVersion | number | versione dello schema del report |
runId | string | id univoco della run |
pipeline | string | nome pipeline |
status | success | partial | failed | partial = girata con scarti |
ok | boolean | true se non abortita |
env | string | ambiente logico (ETL_ENV) |
version | string | null | release che ha eseguito (ETL_RELEASE, es. tag immagine) |
startedAt / finishedAt | string (ISO 8601) | |
durationMs | number | |
nodes | NodeReport[] | topologia reale del grafo eseguito — vedi § NodeReport (mai credenziali nei nomi) |
edges | { from: string; to: string }[] | archi del grafo eseguito |
params | object | parametri della run (es. watermark). No segreti |
extracted | number | record letti dal mondo esterno: somma dei soli nodi source (niente doppio conteggio su un fan-out) |
written | number | record scritti verso l'esterno: somma dei nodi sink. Su un fan-out lo stesso record scritto su due sink conta due volte |
rejected | number | record isolati a dead-letter, su tutti i nodi |
rejectedReasons | object | categorie di scarto → conteggio, aggregate su tutti i nodi |
metrics | object | contatori (standard + eventuali custom degli stage) |
nodeReports | Record<string, unknown> | mappa nodeId → esito di writeReport (o il default { name, success } se il nodo non lo implementa) — vedi ADR-0018 |
error | object | null | { type, message, stage, nodeId } se failed (no PII). nodeId è il nodo che ha fallito, non l'ultimo attraversato dall'errore mentre risaliva lo stream — vedi ADR-0038 |
NodeReport — cosa c'è per ogni nodo
| Campo | Tipo | Note |
|---|---|---|
id / kind / name | string | identità del nodo nel grafo eseguito |
extracted / written | number | record entrati / usciti dal nodo |
rejected / rejectedReasons | number / object | record isolati a dead-letter e loro categorie |
caught / caughtErrors | number / CaughtErrorSummary[] | solo sui nodi try-catch (ADR-0029) |
startedAt | string | null | primo record attraversato (ISO 8601). null = il nodo non ha mai lavorato |
finishedAt | string | null | completamento del nodo (ISO 8601). null = il nodo non si è mai concluso (run abortita) |
durationMs | number | tempo di attraversamento — risponde a «qual è il collo di bottiglia» |
I tre campi di tempo sono l'aggiunta dello schema v4 (additivi: nessun campo rimosso o rinominato, i consumer v3 continuano a funzionare).
startedAt: null non è un dato mancante: significa che il nodo non ha mai attraversato un record — input vuoto, filtro che scarta tutto a monte, o run abortita prima che il nodo partisse. È il corollario documentato su ProgressReporter.onNodeStart. In quel caso durationMs è 0: una durata non misurata si dichiara zero, non si deriva da un solo estremo.
Run fallita: topologia e conteggi parziali (ADR-0038)
Un report failed prodotto da una run che è partita contiene la topologia reale e i conteggi accumulati fino all'abort — è proprio la run che più serve ispezionare, e lasciarla come guscio vuoto era il difetto corretto in ADR-0038.
Un report failed prodotto da buildFailedGraphReport (fallimento prima del grafo) ha invece nodes: [] / edges: [], perché non c'è alcuna topologia eseguita da riportare.
nodeReports e writeReport (ADR-0018)
Ogni nodo del grafo (source/transform/lookup/join/sink/try-catch) può ricevere un'opzione writeReport: (ctx, stats) => Record<string, unknown> in fase di costruzione: GraphBuilder.source/transform/lookup/join/sink/tryCatch(..., { writeReport }). stats espone i conteggi del nodo (extracted/written/ rejected/rejectedReasons) così un nodo può riassumere i volumi senza scrivere il payload grezzo (es. { recordsFetched: 1_000_000 } invece del milione di record). Se writeReport è assente o lancia, il runner scrive il default { name: node.name, success }, dove success è true se il nodo non ha abortito la run (uno scarto isolato a dead-letter non conta come fallimento).
processes.output_data riceve gli stessi report incrementalmente, durante l'esecuzione (non solo a fine run): vedi src/worker/progress-reporter.ts e ADR-0017/0018. Accanto a nodeReports viaggia runningNodes: string[] — i nodi che stanno elaborando adesso (partiti e non ancora completati), da ProgressReporter.onNodeStart. È ciò che permette alla UI di distinguere un nodo che lavora da uno in attesa, e quindi di mostrare due rami concorrenti. Assente sulle run concluse. Vedi ADR-0037.
nodes ed edges sono la topologia reale del grafo eseguito: un nodo per ogni source/transform/lookup/join/sink/try-catch dichiarato nel builder. Non esiste più un nodo sintetico che collassa la catena — vedi ADR-0036.
Scarti: report vs payload
Il report contiene solo conteggi e categorie di scarto (niente payload → niente PII). I record scartati completi vanno nel dead-letter sink, che è scelto dal chiamante (file, tabella, coda). Così il servizio resta agnostico e il report è sicuro da archiviare.
Il DeadLetter porta quindi campi che il report non ha, perché contengono payload reale:
| Campo | Quando | Nel report? |
|---|---|---|
raw | catture per-record (transform isolate, sink con loadOne) | Solo con includeRawInReport sul guard (ADR-0034) |
lastConsumed | nodi non isolabili, durante il consumo (ADR-0033) | Mai |
guardSnapshots | record d'ingresso da cui deriva il fallimento, risalendo derivedFrom (ADR-0035) | Mai |
Sintetici e quindi sicuri da esporre: error.context.recordId (id del record entrato nel nodo che ha fallito) e error.context.derivedFrom (id dei contribuenti, quando il fallimento avviene in emissione di un nodo stateful). Sono id, non dati.
Nota: guardSnapshots è additivo e opzionale, quindi non ha inciso su reportSchemaVersion.
Esempio (successo)
{
"reportSchemaVersion": 4,
"runId": "ebe7d7cc-…",
"pipeline": "langfuse-sessions-daily",
"status": "success",
"ok": true,
"env": "dev",
"version": "20260603-200239",
"startedAt": "2026-06-03T19:02:10.878Z",
"finishedAt": "2026-06-03T19:02:10.927Z",
"durationMs": 49,
"nodes": [
{
"id": "sessions",
"kind": "source",
"name": "trace_sessions",
"extracted": 74,
"written": 0,
"rejected": 0,
"rejectedReasons": {},
"caught": 0,
"caughtErrors": [],
"startedAt": "2026-06-03T19:02:10.880Z",
"finishedAt": "2026-06-03T19:02:10.902Z",
"durationMs": 22
},
{
"id": "raw-contract",
"kind": "transform",
"name": "raw-contract",
"extracted": 74,
"written": 74,
"rejected": 0,
"rejectedReasons": {},
"caught": 0,
"caughtErrors": [],
"startedAt": "2026-06-03T19:02:10.881Z",
"finishedAt": "2026-06-03T19:02:10.903Z",
"durationMs": 22
},
{
"id": "aggregate-daily",
"kind": "transform",
"name": "aggregate-daily",
"extracted": 74,
"written": 14,
"rejected": 0,
"rejectedReasons": {},
"caught": 0,
"caughtErrors": [],
"startedAt": "2026-06-03T19:02:10.882Z",
"finishedAt": "2026-06-03T19:02:10.918Z",
"durationMs": 36
},
{
"id": "sqlserver:dbo.fact_session_daily",
"kind": "sink",
"name": "sqlserver:dbo.fact_session_daily",
"extracted": 14,
"written": 14,
"rejected": 0,
"rejectedReasons": {},
"caught": 0,
"caughtErrors": [],
"startedAt": "2026-06-03T19:02:10.883Z",
"finishedAt": "2026-06-03T19:02:10.926Z",
"durationMs": 43
}
],
"edges": [
{ "from": "sessions", "to": "raw-contract" },
{ "from": "raw-contract", "to": "aggregate-daily" },
{ "from": "aggregate-daily", "to": "sqlserver:dbo.fact_session_daily" }
],
"params": {},
"extracted": 74,
"written": 14,
"rejected": 0,
"rejectedReasons": {},
"metrics": { "extracted": 74, "written": 14 },
"nodeReports": {
"sessions": { "name": "trace_sessions", "success": true },
"raw-contract": { "name": "raw-contract", "success": true },
"aggregate-daily": { "name": "aggregate-daily", "success": true },
"sqlserver:dbo.fact_session_daily": {
"recordsWritten": 14,
"success": true
}
},
"error": null
}Esempio (fallimento durante la run)
La topologia c'è, e i conteggi dicono fin dove la run era arrivata: la source aveva già letto 40 record, il sink non ne aveva ancora scritto nessuno, e il nodo a valle non si è mai concluso (finishedAt: null).
{
"reportSchemaVersion": 4,
"status": "failed",
"ok": false,
"extracted": 40,
"written": 0,
"nodes": [
{
"id": "sessions",
"kind": "source",
"name": "trace_sessions",
"extracted": 40,
"written": 40,
"rejected": 0,
"rejectedReasons": {},
"caught": 0,
"caughtErrors": [],
"startedAt": "2026-06-03T19:02:10.880Z",
"finishedAt": "2026-06-03T19:02:12.104Z",
"durationMs": 1224
},
{
"id": "sqlserver:dbo.fact_session_daily",
"kind": "sink",
"name": "sqlserver:dbo.fact_session_daily",
"extracted": 40,
"written": 0,
"rejected": 0,
"rejectedReasons": {},
"caught": 0,
"caughtErrors": [],
"startedAt": "2026-06-03T19:02:10.884Z",
"finishedAt": null,
"durationMs": 0
}
],
"edges": [{ "from": "sessions", "to": "sqlserver:dbo.fact_session_daily" }],
"error": {
"type": "LoadError",
"message": "Failed to connect …",
"stage": "fact_session_daily",
"nodeId": "sqlserver:dbo.fact_session_daily"
}
}Esempio (fallimento prima della run)
Wiring dei connettori fallito: nessun grafo è mai partito, quindi non c'è topologia da riportare. È il caso coperto da buildFailedGraphReport.
{
"reportSchemaVersion": 4,
"status": "failed",
"ok": false,
"extracted": 0,
"written": 0,
"nodes": [],
"edges": [],
"error": {
"type": "ConfigError",
"message": "PGHOST non configurato",
"stage": null,
"nodeId": null
}
}
