Skip to content

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 con jq.
  • stderr = log umani / progress (riga pipeline.completed, ecc.).
  • exit code: 0 se status è success/partial; ≠0 se failed.
  • Anche in caso di errore viene emesso un report (status: "failed").
bash
etl run langfuse-sessions-daily > report.json 2> run.log   # separa report e log
etl run langfuse-sessions-daily | jq '.status, .written'   # consuma al volo

Come 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):

ts
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)

CampoTipoNote
reportSchemaVersionnumberversione dello schema del report
runIdstringid univoco della run
pipelinestringnome pipeline
statussuccess | partial | failedpartial = girata con scarti
okbooleantrue se non abortita
envstringambiente logico (ETL_ENV)
versionstring | nullrelease che ha eseguito (ETL_RELEASE, es. tag immagine)
startedAt / finishedAtstring (ISO 8601)
durationMsnumber
nodesNodeReport[]topologia reale del grafo eseguito — vedi § NodeReport (mai credenziali nei nomi)
edges{ from: string; to: string }[]archi del grafo eseguito
paramsobjectparametri della run (es. watermark). No segreti
extractednumberrecord letti dal mondo esterno: somma dei soli nodi source (niente doppio conteggio su un fan-out)
writtennumberrecord scritti verso l'esterno: somma dei nodi sink. Su un fan-out lo stesso record scritto su due sink conta due volte
rejectednumberrecord isolati a dead-letter, su tutti i nodi
rejectedReasonsobjectcategorie di scarto → conteggio, aggregate su tutti i nodi
metricsobjectcontatori (standard + eventuali custom degli stage)
nodeReportsRecord<string, unknown>mappa nodeId → esito di writeReport (o il default { name, success } se il nodo non lo implementa) — vedi ADR-0018
errorobject | 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

CampoTipoNote
id / kind / namestringidentità del nodo nel grafo eseguito
extracted / writtennumberrecord entrati / usciti dal nodo
rejected / rejectedReasonsnumber / objectrecord isolati a dead-letter e loro categorie
caught / caughtErrorsnumber / CaughtErrorSummary[]solo sui nodi try-catch (ADR-0029)
startedAtstring | nullprimo record attraversato (ISO 8601). null = il nodo non ha mai lavorato
finishedAtstring | nullcompletamento del nodo (ISO 8601). null = il nodo non si è mai concluso (run abortita)
durationMsnumbertempo 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:

CampoQuandoNel report?
rawcatture per-record (transform isolate, sink con loadOne)Solo con includeRawInReport sul guard (ADR-0034)
lastConsumednodi non isolabili, durante il consumo (ADR-0033)Mai
guardSnapshotsrecord 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)

json
{
  "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).

json
{
  "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.

json
{
  "reportSchemaVersion": 4,
  "status": "failed",
  "ok": false,
  "extracted": 0,
  "written": 0,
  "nodes": [],
  "edges": [],
  "error": {
    "type": "ConfigError",
    "message": "PGHOST non configurato",
    "stage": null,
    "nodeId": null
  }
}

Noeva è un marchio registrato di 4D S.R.L.