From 2bc20c875d4a50df80980140c2765ed7543b41f2 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Marcus=20Nerl=C3=B8e?= Date: Fri, 7 Aug 2026 10:53:59 +0200 Subject: [PATCH 1/3] fix(core): mint the fallback external trace id per run MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Runs that carry no external trace context (schedules, task-to-task triggers) fall back to a trace id generated once in the TracingSDK constructor. With `experimental_processKeepAlive` the TracingSDK outlives the run, so every run on a warm process was exported to the external OTLP endpoint under that one id — merging unrelated runs into a single trace. This is the same warm-start hazard c043c4a6a fixed for the external context path, which read the context live but deliberately left the fallback captured at construction. Remint the fallback when the trace context manager's context object is reassigned, which is the run boundary. An empty configured id still means external export is off and is left alone rather than switched on. The test harness needed a fix too: `setGlobalManager` delegates to `registerGlobal`, which ignores a second registration, so every test after the first was mutating the first test's manager. Co-Authored-By: Claude Opus 5 (1M context) --- .changeset/external-trace-id-per-run.md | 5 ++ packages/core/src/v3/otel/tracingSDK.ts | 80 +++++++++++++++--- .../test/externalSpanExporterWrapper.test.ts | 81 +++++++++++++++++++ 3 files changed, 153 insertions(+), 13 deletions(-) create mode 100644 .changeset/external-trace-id-per-run.md diff --git a/.changeset/external-trace-id-per-run.md b/.changeset/external-trace-id-per-run.md new file mode 100644 index 00000000000..f6e7c335c44 --- /dev/null +++ b/.changeset/external-trace-id-per-run.md @@ -0,0 +1,5 @@ +--- +"@trigger.dev/core": patch +--- + +Mint the fallback external trace id per run rather than once per `TracingSDK`. Runs that carry no external trace context fall back to a generated trace id, and with `experimental_processKeepAlive` the `TracingSDK` outlives the run — so every run on a warm process was exported to the external OTLP endpoint under one shared trace id, merging unrelated runs into a single trace. diff --git a/packages/core/src/v3/otel/tracingSDK.ts b/packages/core/src/v3/otel/tracingSDK.ts index 0b3a66a87b4..e3dd023b3ca 100644 --- a/packages/core/src/v3/otel/tracingSDK.ts +++ b/packages/core/src/v3/otel/tracingSDK.ts @@ -393,21 +393,67 @@ function setLogLevel(level: TracingDiagnosticLogLevel) { diag.setLogger(new DiagConsoleLogger(), diagLogLevel); } +/** + * The external trace id used by runs that carry no external trace context, + * minted once per run. + * + * It has to change per run for the same reason the wrappers read the external + * context live: with `processKeepAlive` the `TracingSDK` — and so the wrappers + * — outlive the run, so an id captured at construction merges every run on the + * process into one trace. The manager's trace context object is reassigned per + * run, which makes its identity the run boundary. + */ +class FallbackExternalTraceId { + private traceId: string; + private seenTraceContext: unknown; + + constructor( + private seed: string, + private traceIdGenerator: Pick = idGenerator + ) { + this.traceId = seed; + this.seenTraceContext = traceContext.getTraceContext(); + } + + get(): string { + // An empty seed means external export is disabled — leave it that way + // rather than minting an id and switching the feature on. + if (!this.seed) { + return this.seed; + } + + const currentTraceContext = traceContext.getTraceContext(); + + if (currentTraceContext !== this.seenTraceContext) { + this.seenTraceContext = currentTraceContext; + this.traceId = this.traceIdGenerator.generateTraceId(); + } + + return this.traceId; + } +} + export class ExternalSpanExporterWrapper { + private fallback: FallbackExternalTraceId; + constructor( private underlyingExporter: SpanExporter, - private externalTraceId: string - ) {} + externalTraceId: string, + traceIdGenerator?: Pick + ) { + this.fallback = new FallbackExternalTraceId(externalTraceId, traceIdGenerator); + } private transformSpan(span: ReadableSpan): ReadableSpan | undefined { // Read external context live, so per-run reassignment of // standardTraceContextManager.traceContext is honoured on warm-started // workers that reuse a single TracingSDK across runs. const externalTraceContext = traceContext.getExternalTraceContext(); + const fallbackTraceId = this.fallback.get(); const isExternallySampled = externalTraceContext ? isTraceFlagSampled(externalTraceContext.traceFlags) - : !!this.externalTraceId; + : !!fallbackTraceId; if (!isExternallySampled) { return; @@ -419,7 +465,7 @@ export class ExternalSpanExporterWrapper { const externalTraceId = externalTraceContext ? externalTraceContext.traceId - : this.externalTraceId; + : fallbackTraceId; const isAttemptSpan = span.attributes[SemanticInternalAttributes.SPAN_ATTEMPT]; @@ -478,17 +524,23 @@ export class ExternalSpanExporterWrapper { } class ExternalLogRecordExporterWrapper { + private fallback: FallbackExternalTraceId; + constructor( private underlyingExporter: LogRecordExporter, - private externalTraceId: string - ) {} + externalTraceId: string, + traceIdGenerator?: Pick + ) { + this.fallback = new FallbackExternalTraceId(externalTraceId, traceIdGenerator); + } export(logs: any[], resultCallback: (result: any) => void): void { const externalTraceContext = traceContext.getExternalTraceContext(); + const fallbackTraceId = this.fallback.get(); const isExternallySampled = externalTraceContext ? isTraceFlagSampled(externalTraceContext.traceFlags) - : !!this.externalTraceId; + : !!fallbackTraceId; if (!isExternallySampled) { this.underlyingExporter.export([], resultCallback); @@ -496,7 +548,9 @@ class ExternalLogRecordExporterWrapper { return; } - const modifiedLogs = logs.map((log) => this.transformLogRecord(log, externalTraceContext)); + const modifiedLogs = logs.map((log) => + this.transformLogRecord(log, externalTraceContext, fallbackTraceId) + ); this.underlyingExporter.export(modifiedLogs, resultCallback); } @@ -517,13 +571,13 @@ class ExternalLogRecordExporterWrapper { logRecord: ReadableLogRecord, externalTraceContext: | { traceId: string; spanId: string; tracestate?: string; traceFlags: number } - | undefined + | undefined, + fallbackTraceId: string ): ReadableLogRecord { // Capture externalTraceId for use within the proxy's scope. - // Use externalTraceContext.traceId if available, otherwise fall back to generated externalTraceId - const externalTraceId = externalTraceContext - ? externalTraceContext.traceId - : this.externalTraceId; + // Use externalTraceContext.traceId if available, otherwise fall back to the + // per-run generated id. + const externalTraceId = externalTraceContext ? externalTraceContext.traceId : fallbackTraceId; // If there's no spanContext, or if the externalTraceId is not set, return the original logRecord. if (!logRecord.spanContext || !externalTraceId) { diff --git a/packages/core/test/externalSpanExporterWrapper.test.ts b/packages/core/test/externalSpanExporterWrapper.test.ts index 9b51653a1ec..8880daad92e 100644 --- a/packages/core/test/externalSpanExporterWrapper.test.ts +++ b/packages/core/test/externalSpanExporterWrapper.test.ts @@ -53,6 +53,10 @@ describe("ExternalSpanExporterWrapper warm-start regression", () => { let manager: StandardTraceContextManager; beforeEach(() => { + // `setGlobalManager` delegates to `registerGlobal`, which ignores a second + // registration — without disabling first, every test after the first would + // keep mutating the first test's manager. + traceContext.disable(); manager = new StandardTraceContextManager(); traceContext.setGlobalManager(manager); }); @@ -77,4 +81,81 @@ describe("ExternalSpanExporterWrapper warm-start regression", () => { expect(span.parentSpanContext?.traceId).toBe("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"); expect(span.spanContext().traceId).toBe("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"); }); + + // Runs triggered internally — a schedule, or one task triggering another — + // carry no external trace context and so take the generated fallback. That id + // was captured at construction, which on a warm-started worker meant every run + // on the process shared a single trace id. + it("mints a new fallback trace id per run when there is no external context", () => { + const { exporter, captured } = makeCapturingExporter(); + + let generated = 0; + const idGenerator = { + generateTraceId: () => `${++generated}`.padStart(32, "0"), + }; + + manager.traceContext = { traceparent: TRACEPARENT_RUN_A }; + + const wrapper = new ExternalSpanExporterWrapper( + exporter, + "ffffffffffffffffffffffffffffffff", + idGenerator + ); + + wrapper.export([createAttemptSpan()], () => {}); + + // A second run on the same warm process: the manager is reassigned, so the + // fallback has to be reminted. + manager.traceContext = { traceparent: TRACEPARENT_RUN_B }; + + wrapper.export([createAttemptSpan()], () => {}); + + const runATraceId = captured[0]![0]!.spanContext().traceId; + const runBTraceId = captured[1]![0]!.spanContext().traceId; + + expect(runATraceId).toBe("ffffffffffffffffffffffffffffffff"); + expect(runBTraceId).not.toBe(runATraceId); + expect(runBTraceId).toBe("00000000000000000000000000000001"); + }); + + it("keeps one fallback trace id across every export within a run", () => { + const { exporter, captured } = makeCapturingExporter(); + + let generated = 0; + const idGenerator = { + generateTraceId: () => `${++generated}`.padStart(32, "0"), + }; + + manager.traceContext = { traceparent: TRACEPARENT_RUN_A }; + + const wrapper = new ExternalSpanExporterWrapper( + exporter, + "ffffffffffffffffffffffffffffffff", + idGenerator + ); + + wrapper.export([createAttemptSpan()], () => {}); + wrapper.export([createAttemptSpan()], () => {}); + + expect(captured[1]![0]!.spanContext().traceId).toBe(captured[0]![0]!.spanContext().traceId); + expect(generated).toBe(0); + }); + + it("leaves external export off when no external trace id was configured", () => { + const { exporter, captured } = makeCapturingExporter(); + + const idGenerator = { + generateTraceId: () => "00000000000000000000000000000001", + }; + + manager.traceContext = { traceparent: TRACEPARENT_RUN_A }; + + const wrapper = new ExternalSpanExporterWrapper(exporter, "", idGenerator); + + wrapper.export([createAttemptSpan()], () => {}); + + // Minting an id here would switch external export on for a deployment that + // never asked for it. + expect(captured[0]).toHaveLength(0); + }); }); From 81cf2a585d22952863261969e2a630a73bed406a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Marcus=20Nerl=C3=B8e?= Date: Fri, 7 Aug 2026 11:07:58 +0200 Subject: [PATCH 2/3] fix(core): share one fallback trace id across span and log exporters Address review feedback on the per-run fallback. Giving each exporter wrapper its own FallbackExternalTraceId reintroduced the problem it was meant to fix, one signal down: before, every wrapper got the same generated string, so a run's spans and logs agreed. With per-wrapper state each one reminted independently, so from the second run on a warm process the logs carried a different trace id than the spans and stopped correlating. Construct one instance in the TracingSDK and pass it to every span and log wrapper. Also stop treating an empty trace context as a run boundary. The noop manager returns a fresh object on every call, so its identity always differs and would remint on every export batch, shattering one run's trace into many. Not reachable today (the wrappers are only built where a StandardTraceContextManager is registered) but the invariant was implicit. Rewrite the changeset for users per AGENTS.md, and format with oxfmt. Co-Authored-By: Claude Opus 5 (1M context) --- .changeset/external-trace-id-per-run.md | 2 +- packages/core/src/v3/otel/tracingSDK.ts | 53 ++++++++------ .../test/externalSpanExporterWrapper.test.ts | 73 +++++++++++++++++-- 3 files changed, 96 insertions(+), 32 deletions(-) diff --git a/.changeset/external-trace-id-per-run.md b/.changeset/external-trace-id-per-run.md index f6e7c335c44..65d0ca98b8e 100644 --- a/.changeset/external-trace-id-per-run.md +++ b/.changeset/external-trace-id-per-run.md @@ -2,4 +2,4 @@ "@trigger.dev/core": patch --- -Mint the fallback external trace id per run rather than once per `TracingSDK`. Runs that carry no external trace context fall back to a generated trace id, and with `experimental_processKeepAlive` the `TracingSDK` outlives the run — so every run on a warm process was exported to the external OTLP endpoint under one shared trace id, merging unrelated runs into a single trace. +Runs that don't continue an incoming trace are no longer merged into one trace when they execute on the same warm worker process. Each run now appears as its own trace in your external observability tool, so per-run cost and latency attribution works again. diff --git a/packages/core/src/v3/otel/tracingSDK.ts b/packages/core/src/v3/otel/tracingSDK.ts index e3dd023b3ca..49706c4eccf 100644 --- a/packages/core/src/v3/otel/tracingSDK.ts +++ b/packages/core/src/v3/otel/tracingSDK.ts @@ -162,7 +162,8 @@ export class TracingSDK { ) ); - const externalTraceId = idGenerator.generateTraceId(); + // Shared by every wrapper below so a run's spans and logs agree on the id. + const externalTraceId = new FallbackExternalTraceId(idGenerator.generateTraceId()); for (const exporter of config.exporters ?? []) { spanProcessors.push( @@ -393,6 +394,19 @@ function setLogLevel(level: TracingDiagnosticLogLevel) { diag.setLogger(new DiagConsoleLogger(), diagLogLevel); } +/** + * Identity of the current run's trace context, or undefined when no run is + * active. + * + * An empty context is not a run: the noop manager returns a fresh `{}` on every + * call, so treating it as a run would mint a new id on every export. + */ +function currentRunTraceContext(): object | undefined { + const current = traceContext.getTraceContext(); + + return current && Object.keys(current).length > 0 ? current : undefined; +} + /** * The external trace id used by runs that carry no external trace context, * minted once per run. @@ -402,17 +416,20 @@ function setLogLevel(level: TracingDiagnosticLogLevel) { * — outlive the run, so an id captured at construction merges every run on the * process into one trace. The manager's trace context object is reassigned per * run, which makes its identity the run boundary. + * + * One instance is shared by every wrapper, so a run's spans and logs still + * agree on the id after a remint. */ -class FallbackExternalTraceId { +export class FallbackExternalTraceId { private traceId: string; - private seenTraceContext: unknown; + private seenTraceContext: object | undefined; constructor( private seed: string, private traceIdGenerator: Pick = idGenerator ) { this.traceId = seed; - this.seenTraceContext = traceContext.getTraceContext(); + this.seenTraceContext = currentRunTraceContext(); } get(): string { @@ -422,10 +439,10 @@ class FallbackExternalTraceId { return this.seed; } - const currentTraceContext = traceContext.getTraceContext(); + const current = currentRunTraceContext(); - if (currentTraceContext !== this.seenTraceContext) { - this.seenTraceContext = currentTraceContext; + if (current && current !== this.seenTraceContext) { + this.seenTraceContext = current; this.traceId = this.traceIdGenerator.generateTraceId(); } @@ -434,15 +451,10 @@ class FallbackExternalTraceId { } export class ExternalSpanExporterWrapper { - private fallback: FallbackExternalTraceId; - constructor( private underlyingExporter: SpanExporter, - externalTraceId: string, - traceIdGenerator?: Pick - ) { - this.fallback = new FallbackExternalTraceId(externalTraceId, traceIdGenerator); - } + private fallback: FallbackExternalTraceId + ) {} private transformSpan(span: ReadableSpan): ReadableSpan | undefined { // Read external context live, so per-run reassignment of @@ -463,9 +475,7 @@ export class ExternalSpanExporterWrapper { return; } - const externalTraceId = externalTraceContext - ? externalTraceContext.traceId - : fallbackTraceId; + const externalTraceId = externalTraceContext ? externalTraceContext.traceId : fallbackTraceId; const isAttemptSpan = span.attributes[SemanticInternalAttributes.SPAN_ATTEMPT]; @@ -524,15 +534,10 @@ export class ExternalSpanExporterWrapper { } class ExternalLogRecordExporterWrapper { - private fallback: FallbackExternalTraceId; - constructor( private underlyingExporter: LogRecordExporter, - externalTraceId: string, - traceIdGenerator?: Pick - ) { - this.fallback = new FallbackExternalTraceId(externalTraceId, traceIdGenerator); - } + private fallback: FallbackExternalTraceId + ) {} export(logs: any[], resultCallback: (result: any) => void): void { const externalTraceContext = traceContext.getExternalTraceContext(); diff --git a/packages/core/test/externalSpanExporterWrapper.test.ts b/packages/core/test/externalSpanExporterWrapper.test.ts index 8880daad92e..f1e0655d518 100644 --- a/packages/core/test/externalSpanExporterWrapper.test.ts +++ b/packages/core/test/externalSpanExporterWrapper.test.ts @@ -1,7 +1,7 @@ import { SpanKind, SpanStatusCode, TraceFlags } from "@opentelemetry/api"; import type { ReadableSpan, SpanExporter } from "@opentelemetry/sdk-trace-node"; import { beforeEach, describe, expect, it } from "vitest"; -import { ExternalSpanExporterWrapper } from "../src/v3/otel/tracingSDK.js"; +import { ExternalSpanExporterWrapper, FallbackExternalTraceId } from "../src/v3/otel/tracingSDK.js"; import { SemanticInternalAttributes } from "../src/v3/semanticInternalAttributes.js"; import { traceContext } from "../src/v3/trace-context-api.js"; import { StandardTraceContextManager } from "../src/v3/traceContext/manager.js"; @@ -66,7 +66,10 @@ describe("ExternalSpanExporterWrapper warm-start regression", () => { manager.traceContext = { external: { traceparent: TRACEPARENT_RUN_A } }; - const wrapper = new ExternalSpanExporterWrapper(exporter, "ffffffffffffffffffffffffffffffff"); + const wrapper = new ExternalSpanExporterWrapper( + exporter, + new FallbackExternalTraceId("ffffffffffffffffffffffffffffffff") + ); manager.traceContext = { external: { traceparent: TRACEPARENT_RUN_B } }; @@ -98,8 +101,7 @@ describe("ExternalSpanExporterWrapper warm-start regression", () => { const wrapper = new ExternalSpanExporterWrapper( exporter, - "ffffffffffffffffffffffffffffffff", - idGenerator + new FallbackExternalTraceId("ffffffffffffffffffffffffffffffff", idGenerator) ); wrapper.export([createAttemptSpan()], () => {}); @@ -130,8 +132,7 @@ describe("ExternalSpanExporterWrapper warm-start regression", () => { const wrapper = new ExternalSpanExporterWrapper( exporter, - "ffffffffffffffffffffffffffffffff", - idGenerator + new FallbackExternalTraceId("ffffffffffffffffffffffffffffffff", idGenerator) ); wrapper.export([createAttemptSpan()], () => {}); @@ -150,7 +151,10 @@ describe("ExternalSpanExporterWrapper warm-start regression", () => { manager.traceContext = { traceparent: TRACEPARENT_RUN_A }; - const wrapper = new ExternalSpanExporterWrapper(exporter, "", idGenerator); + const wrapper = new ExternalSpanExporterWrapper( + exporter, + new FallbackExternalTraceId("", idGenerator) + ); wrapper.export([createAttemptSpan()], () => {}); @@ -158,4 +162,59 @@ describe("ExternalSpanExporterWrapper warm-start regression", () => { // never asked for it. expect(captured[0]).toHaveLength(0); }); + + // The TracingSDK shares one FallbackExternalTraceId across every span and log + // wrapper. Giving each wrapper its own would remint them independently, so + // from the second run on, a run's logs would carry a different trace id than + // its spans and stop correlating in the external backend. + it("gives every wrapper sharing one fallback the same id after a remint", () => { + const first = makeCapturingExporter(); + const second = makeCapturingExporter(); + + let generated = 0; + const idGenerator = { + generateTraceId: () => `${++generated}`.padStart(32, "0"), + }; + + manager.traceContext = { traceparent: TRACEPARENT_RUN_A }; + + const fallback = new FallbackExternalTraceId("ffffffffffffffffffffffffffffffff", idGenerator); + const spanWrapper = new ExternalSpanExporterWrapper(first.exporter, fallback); + const otherWrapper = new ExternalSpanExporterWrapper(second.exporter, fallback); + + manager.traceContext = { traceparent: TRACEPARENT_RUN_B }; + + spanWrapper.export([createAttemptSpan()], () => {}); + otherWrapper.export([createAttemptSpan()], () => {}); + + expect(second.captured[0]![0]!.spanContext().traceId).toBe( + first.captured[0]![0]!.spanContext().traceId + ); + expect(generated).toBe(1); + }); + + // The noop manager returns a fresh `{}` on every call, so using its identity + // as the run boundary would remint on every export and shatter one run's + // trace into many. + it("holds the id when no trace context manager is registered", () => { + const { exporter, captured } = makeCapturingExporter(); + + let generated = 0; + const idGenerator = { + generateTraceId: () => `${++generated}`.padStart(32, "0"), + }; + + traceContext.disable(); + + const wrapper = new ExternalSpanExporterWrapper( + exporter, + new FallbackExternalTraceId("ffffffffffffffffffffffffffffffff", idGenerator) + ); + + wrapper.export([createAttemptSpan()], () => {}); + wrapper.export([createAttemptSpan()], () => {}); + + expect(captured[1]![0]!.spanContext().traceId).toBe(captured[0]![0]!.spanContext().traceId); + expect(generated).toBe(0); + }); }); From f837ccc5c58bfd9535f0c577eadd8633cf64bc94 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Marcus=20Nerl=C3=B8e?= Date: Fri, 7 Aug 2026 12:15:43 +0200 Subject: [PATCH 3/3] refactor(core): let the trace context manager mark the run boundary Detecting the run boundary by reference-identity of the trace context was a heuristic, and guarding it against the noop manager's fresh `{}` meant testing the context for emptiness. That traded an unreachable bug for a reachable one: `traceContext` is `z.record(z.unknown())`, so a run whose context is empty would stop reminting and silently merge back into the previous run's trace. Replace the inference with a fact. `StandardTraceContextManager.traceContext` becomes an accessor pair that advances an epoch whenever the context is replaced, which is exactly what starting a run does, so no call site changes. The noop manager reports a constant epoch, so with no manager registered there are no boundaries to react to and nothing churns. Drops the emptiness heuristic entirely, and covers the case it would have broken. Also renames `get()` to `forCurrentRun()` and the shared instance to `fallbackTraceId`, since the old name read as a string, and exports the log wrapper so a test can prove a run's spans and logs stay on one id. Co-Authored-By: Claude Opus 5 (1M context) --- packages/core/src/v3/otel/tracingSDK.ts | 44 ++--- packages/core/src/v3/traceContext/api.ts | 9 ++ packages/core/src/v3/traceContext/manager.ts | 19 ++- packages/core/src/v3/traceContext/types.ts | 6 + .../test/externalSpanExporterWrapper.test.ts | 151 +++++++++++------- 5 files changed, 145 insertions(+), 84 deletions(-) diff --git a/packages/core/src/v3/otel/tracingSDK.ts b/packages/core/src/v3/otel/tracingSDK.ts index 49706c4eccf..2212a274305 100644 --- a/packages/core/src/v3/otel/tracingSDK.ts +++ b/packages/core/src/v3/otel/tracingSDK.ts @@ -163,12 +163,12 @@ export class TracingSDK { ); // Shared by every wrapper below so a run's spans and logs agree on the id. - const externalTraceId = new FallbackExternalTraceId(idGenerator.generateTraceId()); + const fallbackTraceId = new FallbackExternalTraceId(idGenerator.generateTraceId()); for (const exporter of config.exporters ?? []) { spanProcessors.push( getEnvVar("TRIGGER_OTEL_BATCH_PROCESSING_ENABLED") === "1" - ? new BatchSpanProcessor(new ExternalSpanExporterWrapper(exporter, externalTraceId), { + ? new BatchSpanProcessor(new ExternalSpanExporterWrapper(exporter, fallbackTraceId), { maxExportBatchSize: parseInt( getEnvVar("TRIGGER_OTEL_SPAN_MAX_EXPORT_BATCH_SIZE") ?? "64" ), @@ -180,7 +180,7 @@ export class TracingSDK { ), maxQueueSize: parseInt(getEnvVar("TRIGGER_OTEL_SPAN_MAX_QUEUE_SIZE") ?? "512"), }) - : new SimpleSpanProcessor(new ExternalSpanExporterWrapper(exporter, externalTraceId)) + : new SimpleSpanProcessor(new ExternalSpanExporterWrapper(exporter, fallbackTraceId)) ); } @@ -232,7 +232,7 @@ export class TracingSDK { logProcessors.push( getEnvVar("TRIGGER_OTEL_BATCH_PROCESSING_ENABLED") === "1" ? new BatchLogRecordProcessor( - new ExternalLogRecordExporterWrapper(externalLogExporter, externalTraceId), + new ExternalLogRecordExporterWrapper(externalLogExporter, fallbackTraceId), { maxExportBatchSize: parseInt( getEnvVar("TRIGGER_OTEL_LOG_MAX_EXPORT_BATCH_SIZE") ?? "64" @@ -247,7 +247,7 @@ export class TracingSDK { } ) : new SimpleLogRecordProcessor( - new ExternalLogRecordExporterWrapper(externalLogExporter, externalTraceId) + new ExternalLogRecordExporterWrapper(externalLogExporter, fallbackTraceId) ) ); } @@ -394,19 +394,6 @@ function setLogLevel(level: TracingDiagnosticLogLevel) { diag.setLogger(new DiagConsoleLogger(), diagLogLevel); } -/** - * Identity of the current run's trace context, or undefined when no run is - * active. - * - * An empty context is not a run: the noop manager returns a fresh `{}` on every - * call, so treating it as a run would mint a new id on every export. - */ -function currentRunTraceContext(): object | undefined { - const current = traceContext.getTraceContext(); - - return current && Object.keys(current).length > 0 ? current : undefined; -} - /** * The external trace id used by runs that carry no external trace context, * minted once per run. @@ -414,35 +401,34 @@ function currentRunTraceContext(): object | undefined { * It has to change per run for the same reason the wrappers read the external * context live: with `processKeepAlive` the `TracingSDK` — and so the wrappers * — outlive the run, so an id captured at construction merges every run on the - * process into one trace. The manager's trace context object is reassigned per - * run, which makes its identity the run boundary. + * process into one trace. * * One instance is shared by every wrapper, so a run's spans and logs still * agree on the id after a remint. */ export class FallbackExternalTraceId { private traceId: string; - private seenTraceContext: object | undefined; + private seenEpoch: number; constructor( private seed: string, private traceIdGenerator: Pick = idGenerator ) { this.traceId = seed; - this.seenTraceContext = currentRunTraceContext(); + this.seenEpoch = traceContext.getTraceContextEpoch(); } - get(): string { + forCurrentRun(): string { // An empty seed means external export is disabled — leave it that way // rather than minting an id and switching the feature on. if (!this.seed) { return this.seed; } - const current = currentRunTraceContext(); + const epoch = traceContext.getTraceContextEpoch(); - if (current && current !== this.seenTraceContext) { - this.seenTraceContext = current; + if (epoch !== this.seenEpoch) { + this.seenEpoch = epoch; this.traceId = this.traceIdGenerator.generateTraceId(); } @@ -461,7 +447,7 @@ export class ExternalSpanExporterWrapper { // standardTraceContextManager.traceContext is honoured on warm-started // workers that reuse a single TracingSDK across runs. const externalTraceContext = traceContext.getExternalTraceContext(); - const fallbackTraceId = this.fallback.get(); + const fallbackTraceId = this.fallback.forCurrentRun(); const isExternallySampled = externalTraceContext ? isTraceFlagSampled(externalTraceContext.traceFlags) @@ -533,7 +519,7 @@ export class ExternalSpanExporterWrapper { } } -class ExternalLogRecordExporterWrapper { +export class ExternalLogRecordExporterWrapper { constructor( private underlyingExporter: LogRecordExporter, private fallback: FallbackExternalTraceId @@ -541,7 +527,7 @@ class ExternalLogRecordExporterWrapper { export(logs: any[], resultCallback: (result: any) => void): void { const externalTraceContext = traceContext.getExternalTraceContext(); - const fallbackTraceId = this.fallback.get(); + const fallbackTraceId = this.fallback.forCurrentRun(); const isExternallySampled = externalTraceContext ? isTraceFlagSampled(externalTraceContext.traceFlags) diff --git a/packages/core/src/v3/traceContext/api.ts b/packages/core/src/v3/traceContext/api.ts index b4d0074314b..c8105ea56e5 100644 --- a/packages/core/src/v3/traceContext/api.ts +++ b/packages/core/src/v3/traceContext/api.ts @@ -10,6 +10,11 @@ class NoopTraceContextManager implements TraceContextManager { return {}; } + // Never advances: with no manager registered there are no runs to separate. + getTraceContextEpoch() { + return 0; + } + reset() {} getExternalTraceContext() { @@ -57,6 +62,10 @@ export class TraceContextAPI implements TraceContextManager { return this.#getManager().getTraceContext(); } + public getTraceContextEpoch() { + return this.#getManager().getTraceContextEpoch(); + } + public getExternalTraceContext() { return this.#getManager().getExternalTraceContext(); } diff --git a/packages/core/src/v3/traceContext/manager.ts b/packages/core/src/v3/traceContext/manager.ts index ebf8f9b53ec..d0d5cc44707 100644 --- a/packages/core/src/v3/traceContext/manager.ts +++ b/packages/core/src/v3/traceContext/manager.ts @@ -4,12 +4,29 @@ import { parseTraceParent } from "@opentelemetry/core"; import type { TraceContextManager } from "./types.js"; export class StandardTraceContextManager implements TraceContextManager { - public traceContext: Record = {}; + #traceContext: Record = {}; + #epoch = 0; + + // An accessor rather than a plain field so that replacing the context, which + // is what starting a run does, is what advances the epoch. Call sites are + // unchanged. + get traceContext(): Record { + return this.#traceContext; + } + + set traceContext(value: Record) { + this.#traceContext = value; + this.#epoch++; + } getTraceContext() { return this.traceContext; } + getTraceContextEpoch() { + return this.#epoch; + } + reset() { this.traceContext = {}; } diff --git a/packages/core/src/v3/traceContext/types.ts b/packages/core/src/v3/traceContext/types.ts index 065cf73a2af..43834b55c4f 100644 --- a/packages/core/src/v3/traceContext/types.ts +++ b/packages/core/src/v3/traceContext/types.ts @@ -2,6 +2,12 @@ import type { Context } from "@opentelemetry/api"; export interface TraceContextManager { getTraceContext(): Record; + /** + * Increments every time the trace context is replaced, which on a worker that + * reuses one process across runs is the run boundary. Long-lived consumers + * compare it to tell "still the same run" from "a new run started". + */ + getTraceContextEpoch(): number; extractContext(): Context; reset(): void; getExternalTraceContext(): diff --git a/packages/core/test/externalSpanExporterWrapper.test.ts b/packages/core/test/externalSpanExporterWrapper.test.ts index f1e0655d518..87e9f596dba 100644 --- a/packages/core/test/externalSpanExporterWrapper.test.ts +++ b/packages/core/test/externalSpanExporterWrapper.test.ts @@ -1,13 +1,19 @@ import { SpanKind, SpanStatusCode, TraceFlags } from "@opentelemetry/api"; +import type { LogRecordExporter, ReadableLogRecord } from "@opentelemetry/sdk-logs"; import type { ReadableSpan, SpanExporter } from "@opentelemetry/sdk-trace-node"; import { beforeEach, describe, expect, it } from "vitest"; -import { ExternalSpanExporterWrapper, FallbackExternalTraceId } from "../src/v3/otel/tracingSDK.js"; +import { + ExternalLogRecordExporterWrapper, + ExternalSpanExporterWrapper, + FallbackExternalTraceId, +} from "../src/v3/otel/tracingSDK.js"; import { SemanticInternalAttributes } from "../src/v3/semanticInternalAttributes.js"; import { traceContext } from "../src/v3/trace-context-api.js"; import { StandardTraceContextManager } from "../src/v3/traceContext/manager.js"; const TRACEPARENT_RUN_A = "00-aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa-1111111111111111-01"; const TRACEPARENT_RUN_B = "00-bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb-2222222222222222-01"; +const SEED = "ffffffffffffffffffffffffffffffff"; function createAttemptSpan(): ReadableSpan { const spanCtx = { @@ -36,6 +42,18 @@ function createAttemptSpan(): ReadableSpan { } as unknown as ReadableSpan; } +function createLogRecord(): ReadableLogRecord { + return { + body: "hello", + attributes: {}, + spanContext: { + traceId: "cccccccccccccccccccccccccccccccc", + spanId: "3333333333333333", + traceFlags: TraceFlags.SAMPLED, + }, + } as unknown as ReadableLogRecord; +} + function makeCapturingExporter(): { exporter: SpanExporter; captured: ReadableSpan[][] } { const captured: ReadableSpan[][] = []; const exporter: SpanExporter = { @@ -49,6 +67,32 @@ function makeCapturingExporter(): { exporter: SpanExporter; captured: ReadableSp return { exporter, captured }; } +function makeCapturingLogExporter(): { + exporter: LogRecordExporter; + captured: ReadableLogRecord[][]; +} { + const captured: ReadableLogRecord[][] = []; + const exporter: LogRecordExporter = { + export: (records, cb) => { + captured.push(records); + cb({ code: 0 } as any); + }, + shutdown: () => Promise.resolve(), + }; + return { exporter, captured }; +} + +/** Yields 000…001, 000…002, … so a reminted id is identifiable by its ordinal. */ +function makeIdGenerator() { + let generated = 0; + return { + generateTraceId: () => `${++generated}`.padStart(32, "0"), + get count() { + return generated; + }, + }; +} + describe("ExternalSpanExporterWrapper warm-start regression", () => { let manager: StandardTraceContextManager; @@ -66,10 +110,7 @@ describe("ExternalSpanExporterWrapper warm-start regression", () => { manager.traceContext = { external: { traceparent: TRACEPARENT_RUN_A } }; - const wrapper = new ExternalSpanExporterWrapper( - exporter, - new FallbackExternalTraceId("ffffffffffffffffffffffffffffffff") - ); + const wrapper = new ExternalSpanExporterWrapper(exporter, new FallbackExternalTraceId(SEED)); manager.traceContext = { external: { traceparent: TRACEPARENT_RUN_B } }; @@ -91,23 +132,18 @@ describe("ExternalSpanExporterWrapper warm-start regression", () => { // on the process shared a single trace id. it("mints a new fallback trace id per run when there is no external context", () => { const { exporter, captured } = makeCapturingExporter(); - - let generated = 0; - const idGenerator = { - generateTraceId: () => `${++generated}`.padStart(32, "0"), - }; + const idGenerator = makeIdGenerator(); manager.traceContext = { traceparent: TRACEPARENT_RUN_A }; const wrapper = new ExternalSpanExporterWrapper( exporter, - new FallbackExternalTraceId("ffffffffffffffffffffffffffffffff", idGenerator) + new FallbackExternalTraceId(SEED, idGenerator) ); wrapper.export([createAttemptSpan()], () => {}); - // A second run on the same warm process: the manager is reassigned, so the - // fallback has to be reminted. + // A second run on the same warm process. manager.traceContext = { traceparent: TRACEPARENT_RUN_B }; wrapper.export([createAttemptSpan()], () => {}); @@ -115,39 +151,55 @@ describe("ExternalSpanExporterWrapper warm-start regression", () => { const runATraceId = captured[0]![0]!.spanContext().traceId; const runBTraceId = captured[1]![0]!.spanContext().traceId; - expect(runATraceId).toBe("ffffffffffffffffffffffffffffffff"); + expect(runATraceId).toBe(SEED); expect(runBTraceId).not.toBe(runATraceId); expect(runBTraceId).toBe("00000000000000000000000000000001"); }); - it("keeps one fallback trace id across every export within a run", () => { + // The run boundary is the manager being handed a new context, not that + // context having any particular content. A run whose trace context is empty + // is still a different run. + it("mints a new fallback trace id for a run whose trace context is empty", () => { const { exporter, captured } = makeCapturingExporter(); + const idGenerator = makeIdGenerator(); + + manager.traceContext = { traceparent: TRACEPARENT_RUN_A }; + + const wrapper = new ExternalSpanExporterWrapper( + exporter, + new FallbackExternalTraceId(SEED, idGenerator) + ); - let generated = 0; - const idGenerator = { - generateTraceId: () => `${++generated}`.padStart(32, "0"), - }; + wrapper.export([createAttemptSpan()], () => {}); + + manager.traceContext = {}; + + wrapper.export([createAttemptSpan()], () => {}); + + expect(captured[1]![0]!.spanContext().traceId).not.toBe(captured[0]![0]!.spanContext().traceId); + }); + + it("keeps one fallback trace id across every export within a run", () => { + const { exporter, captured } = makeCapturingExporter(); + const idGenerator = makeIdGenerator(); manager.traceContext = { traceparent: TRACEPARENT_RUN_A }; const wrapper = new ExternalSpanExporterWrapper( exporter, - new FallbackExternalTraceId("ffffffffffffffffffffffffffffffff", idGenerator) + new FallbackExternalTraceId(SEED, idGenerator) ); wrapper.export([createAttemptSpan()], () => {}); wrapper.export([createAttemptSpan()], () => {}); expect(captured[1]![0]!.spanContext().traceId).toBe(captured[0]![0]!.spanContext().traceId); - expect(generated).toBe(0); + expect(idGenerator.count).toBe(0); }); it("leaves external export off when no external trace id was configured", () => { const { exporter, captured } = makeCapturingExporter(); - - const idGenerator = { - generateTraceId: () => "00000000000000000000000000000001", - }; + const idGenerator = makeIdGenerator(); manager.traceContext = { traceparent: TRACEPARENT_RUN_A }; @@ -163,58 +215,49 @@ describe("ExternalSpanExporterWrapper warm-start regression", () => { expect(captured[0]).toHaveLength(0); }); - // The TracingSDK shares one FallbackExternalTraceId across every span and log - // wrapper. Giving each wrapper its own would remint them independently, so - // from the second run on, a run's logs would carry a different trace id than - // its spans and stop correlating in the external backend. - it("gives every wrapper sharing one fallback the same id after a remint", () => { - const first = makeCapturingExporter(); - const second = makeCapturingExporter(); - - let generated = 0; - const idGenerator = { - generateTraceId: () => `${++generated}`.padStart(32, "0"), - }; + // The TracingSDK shares one FallbackExternalTraceId across its span and log + // wrappers. Giving each its own would remint them independently, so from the + // second run on, a run's logs would carry a different trace id than its spans + // and stop correlating in the external backend. + it("keeps a run's spans and logs on the same id after a remint", () => { + const spans = makeCapturingExporter(); + const logs = makeCapturingLogExporter(); + const idGenerator = makeIdGenerator(); manager.traceContext = { traceparent: TRACEPARENT_RUN_A }; - const fallback = new FallbackExternalTraceId("ffffffffffffffffffffffffffffffff", idGenerator); - const spanWrapper = new ExternalSpanExporterWrapper(first.exporter, fallback); - const otherWrapper = new ExternalSpanExporterWrapper(second.exporter, fallback); + const fallback = new FallbackExternalTraceId(SEED, idGenerator); + const spanWrapper = new ExternalSpanExporterWrapper(spans.exporter, fallback); + const logWrapper = new ExternalLogRecordExporterWrapper(logs.exporter, fallback); manager.traceContext = { traceparent: TRACEPARENT_RUN_B }; spanWrapper.export([createAttemptSpan()], () => {}); - otherWrapper.export([createAttemptSpan()], () => {}); + logWrapper.export([createLogRecord()], () => {}); - expect(second.captured[0]![0]!.spanContext().traceId).toBe( - first.captured[0]![0]!.spanContext().traceId + expect(logs.captured[0]![0]!.spanContext!.traceId).toBe( + spans.captured[0]![0]!.spanContext().traceId ); - expect(generated).toBe(1); + expect(idGenerator.count).toBe(1); }); - // The noop manager returns a fresh `{}` on every call, so using its identity - // as the run boundary would remint on every export and shatter one run's - // trace into many. + // With no manager registered the epoch is a constant, so there are no run + // boundaries to react to and the id must hold rather than churn per export. it("holds the id when no trace context manager is registered", () => { const { exporter, captured } = makeCapturingExporter(); - - let generated = 0; - const idGenerator = { - generateTraceId: () => `${++generated}`.padStart(32, "0"), - }; + const idGenerator = makeIdGenerator(); traceContext.disable(); const wrapper = new ExternalSpanExporterWrapper( exporter, - new FallbackExternalTraceId("ffffffffffffffffffffffffffffffff", idGenerator) + new FallbackExternalTraceId(SEED, idGenerator) ); wrapper.export([createAttemptSpan()], () => {}); wrapper.export([createAttemptSpan()], () => {}); expect(captured[1]![0]!.spanContext().traceId).toBe(captured[0]![0]!.spanContext().traceId); - expect(generated).toBe(0); + expect(idGenerator.count).toBe(0); }); });