diff --git a/.changeset/external-trace-id-per-run.md b/.changeset/external-trace-id-per-run.md new file mode 100644 index 00000000000..65d0ca98b8e --- /dev/null +++ b/.changeset/external-trace-id-per-run.md @@ -0,0 +1,5 @@ +--- +"@trigger.dev/core": patch +--- + +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 0b3a66a87b4..2212a274305 100644 --- a/packages/core/src/v3/otel/tracingSDK.ts +++ b/packages/core/src/v3/otel/tracingSDK.ts @@ -162,12 +162,13 @@ export class TracingSDK { ) ); - const externalTraceId = idGenerator.generateTraceId(); + // Shared by every wrapper below so a run's spans and logs agree on the id. + 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" ), @@ -179,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)) ); } @@ -231,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" @@ -246,7 +247,7 @@ export class TracingSDK { } ) : new SimpleLogRecordProcessor( - new ExternalLogRecordExporterWrapper(externalLogExporter, externalTraceId) + new ExternalLogRecordExporterWrapper(externalLogExporter, fallbackTraceId) ) ); } @@ -393,10 +394,52 @@ 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. + * + * 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 seenEpoch: number; + + constructor( + private seed: string, + private traceIdGenerator: Pick = idGenerator + ) { + this.traceId = seed; + this.seenEpoch = traceContext.getTraceContextEpoch(); + } + + 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 epoch = traceContext.getTraceContextEpoch(); + + if (epoch !== this.seenEpoch) { + this.seenEpoch = epoch; + this.traceId = this.traceIdGenerator.generateTraceId(); + } + + return this.traceId; + } +} + export class ExternalSpanExporterWrapper { constructor( private underlyingExporter: SpanExporter, - private externalTraceId: string + private fallback: FallbackExternalTraceId ) {} private transformSpan(span: ReadableSpan): ReadableSpan | undefined { @@ -404,10 +447,11 @@ 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.forCurrentRun(); const isExternallySampled = externalTraceContext ? isTraceFlagSampled(externalTraceContext.traceFlags) - : !!this.externalTraceId; + : !!fallbackTraceId; if (!isExternallySampled) { return; @@ -417,9 +461,7 @@ export class ExternalSpanExporterWrapper { return; } - const externalTraceId = externalTraceContext - ? externalTraceContext.traceId - : this.externalTraceId; + const externalTraceId = externalTraceContext ? externalTraceContext.traceId : fallbackTraceId; const isAttemptSpan = span.attributes[SemanticInternalAttributes.SPAN_ATTEMPT]; @@ -477,18 +519,19 @@ export class ExternalSpanExporterWrapper { } } -class ExternalLogRecordExporterWrapper { +export class ExternalLogRecordExporterWrapper { constructor( private underlyingExporter: LogRecordExporter, - private externalTraceId: string + private fallback: FallbackExternalTraceId ) {} export(logs: any[], resultCallback: (result: any) => void): void { const externalTraceContext = traceContext.getExternalTraceContext(); + const fallbackTraceId = this.fallback.forCurrentRun(); const isExternallySampled = externalTraceContext ? isTraceFlagSampled(externalTraceContext.traceFlags) - : !!this.externalTraceId; + : !!fallbackTraceId; if (!isExternallySampled) { this.underlyingExporter.export([], resultCallback); @@ -496,7 +539,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 +562,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/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 9b51653a1ec..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 } 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,10 +67,40 @@ 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; 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); }); @@ -62,7 +110,7 @@ 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(SEED)); manager.traceContext = { external: { traceparent: TRACEPARENT_RUN_B } }; @@ -77,4 +125,139 @@ 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(); + const idGenerator = makeIdGenerator(); + + manager.traceContext = { traceparent: TRACEPARENT_RUN_A }; + + const wrapper = new ExternalSpanExporterWrapper( + exporter, + new FallbackExternalTraceId(SEED, idGenerator) + ); + + wrapper.export([createAttemptSpan()], () => {}); + + // A second run on the same warm process. + 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(SEED); + expect(runBTraceId).not.toBe(runATraceId); + expect(runBTraceId).toBe("00000000000000000000000000000001"); + }); + + // 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) + ); + + 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(SEED, idGenerator) + ); + + wrapper.export([createAttemptSpan()], () => {}); + wrapper.export([createAttemptSpan()], () => {}); + + expect(captured[1]![0]!.spanContext().traceId).toBe(captured[0]![0]!.spanContext().traceId); + expect(idGenerator.count).toBe(0); + }); + + it("leaves external export off when no external trace id was configured", () => { + const { exporter, captured } = makeCapturingExporter(); + const idGenerator = makeIdGenerator(); + + manager.traceContext = { traceparent: TRACEPARENT_RUN_A }; + + const wrapper = new ExternalSpanExporterWrapper( + exporter, + new FallbackExternalTraceId("", 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); + }); + + // 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(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()], () => {}); + logWrapper.export([createLogRecord()], () => {}); + + expect(logs.captured[0]![0]!.spanContext!.traceId).toBe( + spans.captured[0]![0]!.spanContext().traceId + ); + expect(idGenerator.count).toBe(1); + }); + + // 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(); + const idGenerator = makeIdGenerator(); + + traceContext.disable(); + + const wrapper = new ExternalSpanExporterWrapper( + exporter, + new FallbackExternalTraceId(SEED, idGenerator) + ); + + wrapper.export([createAttemptSpan()], () => {}); + wrapper.export([createAttemptSpan()], () => {}); + + expect(captured[1]![0]!.spanContext().traceId).toBe(captured[0]![0]!.spanContext().traceId); + expect(idGenerator.count).toBe(0); + }); });