src/orchestrator/telemetry.ts — canonical telemetry export contract
Purpose
Section titled “Purpose”THE versioned, serializable contract every telemetry consumer speaks
(TELEMETRY_SCHEMA_VERSION = 3). Exporters (otel, a custom sink) receive
these records instead of re-deriving facts from the rendering-oriented
WireEvent stream — cacheSource is derived once, git/CI/host context
is pre-folded, bigint wallclock spans are decimal strings.
Public surface
Section titled “Public surface”TelemetryRecord— per-event union:run.start/task.start/task.log/task.end/run.end.RunSummaryRecord— one per run:RunContextRecord+ totals + per-taskTaskTelemetry[]. What every telemetry sink receives at end of run.abortedCount(v3, item 851) counts the tasks a shutdown signal or an embedder’s abort killed; they are not in the task list, which holds real runs only, so without it a stopped run read as a failure with nothing failed.assembleRunSummary(run, tasks, timing)— builds theRunSummaryRecordfrom the context and theTaskTelemetry[], the per-task tallies derived fromtasks, so a distributed run and a local one produce the same summary.deriveCacheSource(status)— theCacheSource:'local'/'remote'for the two hits,'miss'forsuccess/failed,'none'forskipped/aborted; never null.taskTelemetryOf(outcome)— the one projection of aTaskOutcomeintoTaskTelemetry, used by the streamingtask.endrecord and the summary’stasks[]alike, so the two cannot drift (item 660).createTelemetrySource({ sinks, run, warn?, owners? })→ aTelemetrySource— projects the bus once and fans out to sinks.runis theRunContextRecordstamped onrun.start;ownersnames a nameless sink by its plugin.TelemetrySink— what atelemetryplugin returns: an optionalname, awantslist of record kinds (the source checks it BEFORE projecting, so a sink pays nothing for kinds it declines), andonRecord/onRunSummary/flush. Every one is crash-isolated andflushis deadline-bounded. A warning names the sink by itsname, else by its plugin’s (org/p, ororg/p #2for the second of a list), else by its place in the source’s list (#2; item 1027).TelemetryContext— what the hook is handed:workspaceRoot,cacheDir(a STRING, not a Cache handle — a sink cannot reach the cache) andwarn.TelemetrySource— the live projection: a bussubscriber,emitSummaryandflush.CacheSource—'miss' | 'local' | 'remote' | 'none', the cache axis every record carries.TASK_STATUSES,isPassStatus(status),isCacheHit(status)— the task axis and the two predicates every surface shares.TELEMETRY_SCHEMA_VERSION— bumped when a record’s shape changes, so a receiver can refuse what it cannot read.
Invariants
Section titled “Invariants”- Observe-only by construction: sinks receive immutable records and a read-only context — no bus, no Cache, no path back into scheduling.
- Crash isolation: a throwing sink is disabled for the run, never propagates.
task.logis opt-in viaTelemetrySink.wants(default excludes it).- Version bumps are additive-or-bump: consumers reject unknown majors.