diff --git a/.changeset/external-trace-id-per-run.md b/.changeset/external-trace-id-per-run.md new file mode 100644 index 0000000000..e074b53cd8 --- /dev/null +++ b/.changeset/external-trace-id-per-run.md @@ -0,0 +1,5 @@ +--- +"@trigger.dev/core": patch +--- + +Unrelated runs are no longer merged into a single trace in your external observability tool when they happen to execute on the same warm worker process. A run and the runs it triggers still share one trace, so a run tree stays together. diff --git a/packages/core/src/v3/otel/tracingSDK.ts b/packages/core/src/v3/otel/tracingSDK.ts index 0f4ea82a22..fcf49ac4da 100644 --- a/packages/core/src/v3/otel/tracingSDK.ts +++ b/packages/core/src/v3/otel/tracingSDK.ts @@ -163,12 +163,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 fallbackTraceIds = new FallbackExternalTraceIds(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, fallbackTraceIds), { maxExportBatchSize: parseInt( getEnvVar("TRIGGER_OTEL_SPAN_MAX_EXPORT_BATCH_SIZE") ?? "64" ), @@ -180,7 +181,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, fallbackTraceIds)) ); } @@ -232,7 +233,7 @@ export class TracingSDK { logProcessors.push( getEnvVar("TRIGGER_OTEL_BATCH_PROCESSING_ENABLED") === "1" ? new BatchLogRecordProcessor( - new ExternalLogRecordExporterWrapper(externalLogExporter, externalTraceId), + new ExternalLogRecordExporterWrapper(externalLogExporter, fallbackTraceIds), { maxExportBatchSize: parseInt( getEnvVar("TRIGGER_OTEL_LOG_MAX_EXPORT_BATCH_SIZE") ?? "64" @@ -247,7 +248,7 @@ export class TracingSDK { } ) : new SimpleLogRecordProcessor( - new ExternalLogRecordExporterWrapper(externalLogExporter, externalTraceId) + new ExternalLogRecordExporterWrapper(externalLogExporter, fallbackTraceIds) ) ); } @@ -424,10 +425,81 @@ function setLogLevel(level: TracingDiagnosticLogLevel) { diag.setLogger(new DiagConsoleLogger(), diagLogLevel); } +/** Only the current run and the tail of recently ended ones can still export. */ +export const MAX_TRACKED_INTERNAL_TRACES = 64; + +/** + * External trace ids for runs that carry no external trace context — with + * `processKeepAlive` the `TracingSDK` outlives the run, so an id captured at + * construction merges every run on the process into one trace. + * + * A record's id comes from its own internal trace id rather than from whatever + * run is current when the exporter is called. Batch processors drain + * asynchronously, so a run's records are routinely exported after the next run + * has started, and reading ambient state then would stamp them with the wrong + * run's id. It also makes a run's spans and logs agree without coordinating. + * + * Granularity therefore follows the internal trace, not the run: a run tree + * shares one internal trace, so a parent and the runs it triggers land on one + * external trace together, which is the grouping you want. + */ +export class FallbackExternalTraceIds { + private readonly byInternalTrace = new Map(); + + constructor( + private seed: string, + private traceIdGenerator: Pick = idGenerator + ) {} + + /** False when no external trace id was configured, i.e. external export is off. */ + get enabled(): boolean { + return !!this.seed; + } + + forInternalTrace(internalTraceId: string): 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 known = this.byInternalTrace.get(internalTraceId); + + if (known) { + // Re-insert so the map is ordered by last use rather than first. A run + // that is still exporting keeps its id even if enough unrelated traces + // appear alongside it to fill the map, which would otherwise split it + // across two external traces. + this.byInternalTrace.delete(internalTraceId); + this.byInternalTrace.set(internalTraceId, known); + + return known; + } + + // The first run reuses the id generated at construction, so the configured + // seed is not thrown away. + const traceId = + this.byInternalTrace.size === 0 ? this.seed : this.traceIdGenerator.generateTraceId(); + + this.byInternalTrace.set(internalTraceId, traceId); + + if (this.byInternalTrace.size > MAX_TRACKED_INTERNAL_TRACES) { + // Map iterates in insertion order, so this drops the least recently used. + const stalest = this.byInternalTrace.keys().next().value; + + if (stalest !== undefined) { + this.byInternalTrace.delete(stalest); + } + } + + return traceId; + } +} + export class ExternalSpanExporterWrapper { constructor( private underlyingExporter: SpanExporter, - private externalTraceId: string + private fallback: FallbackExternalTraceIds ) {} private transformSpan(span: ReadableSpan): ReadableSpan | undefined { @@ -438,7 +510,7 @@ export class ExternalSpanExporterWrapper { const isExternallySampled = externalTraceContext ? isTraceFlagSampled(externalTraceContext.traceFlags) - : !!this.externalTraceId; + : this.fallback.enabled; if (!isExternallySampled) { return; @@ -450,7 +522,7 @@ export class ExternalSpanExporterWrapper { const externalTraceId = externalTraceContext ? externalTraceContext.traceId - : this.externalTraceId; + : this.fallback.forInternalTrace(span.spanContext().traceId); const isAttemptSpan = span.attributes[SemanticInternalAttributes.SPAN_ATTEMPT]; @@ -508,10 +580,10 @@ export class ExternalSpanExporterWrapper { } } -class ExternalLogRecordExporterWrapper { +export class ExternalLogRecordExporterWrapper { constructor( private underlyingExporter: LogRecordExporter, - private externalTraceId: string + private fallback: FallbackExternalTraceIds ) {} export(logs: any[], resultCallback: (result: any) => void): void { @@ -519,7 +591,7 @@ class ExternalLogRecordExporterWrapper { const isExternallySampled = externalTraceContext ? isTraceFlagSampled(externalTraceContext.traceFlags) - : !!this.externalTraceId; + : this.fallback.enabled; if (!isExternallySampled) { this.underlyingExporter.export([], resultCallback); @@ -550,14 +622,20 @@ class ExternalLogRecordExporterWrapper { | { traceId: string; spanId: string; tracestate?: string; traceFlags: number } | undefined ): ReadableLogRecord { - // Capture externalTraceId for use within the proxy's scope. - // Use externalTraceContext.traceId if available, otherwise fall back to generated externalTraceId + // Without a spanContext there is no internal trace id to key the fallback + // on, and nothing to rewrite. + if (!logRecord.spanContext) { + return logRecord; + } + + // Capture externalTraceId for use within the proxy's scope. Use + // externalTraceContext.traceId if available, otherwise the id belonging to + // the run this record came from. const externalTraceId = externalTraceContext ? externalTraceContext.traceId - : this.externalTraceId; + : this.fallback.forInternalTrace(logRecord.spanContext.traceId); - // If there's no spanContext, or if the externalTraceId is not set, return the original logRecord. - if (!logRecord.spanContext || !externalTraceId) { + if (!externalTraceId) { return logRecord; } diff --git a/packages/core/test/externalSpanExporterWrapper.test.ts b/packages/core/test/externalSpanExporterWrapper.test.ts index 9b51653a1e..cadc816583 100644 --- a/packages/core/test/externalSpanExporterWrapper.test.ts +++ b/packages/core/test/externalSpanExporterWrapper.test.ts @@ -1,17 +1,27 @@ 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, + FallbackExternalTraceIds, + MAX_TRACKED_INTERNAL_TRACES, +} 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"; +// Every span and log record of one run shares the run's internal trace id. +const INTERNAL_TRACE_RUN_A = "cccccccccccccccccccccccccccccccc"; +const INTERNAL_TRACE_RUN_B = "dddddddddddddddddddddddddddddddd"; -function createAttemptSpan(): ReadableSpan { +function createAttemptSpan(internalTraceId = INTERNAL_TRACE_RUN_A): ReadableSpan { const spanCtx = { - traceId: "cccccccccccccccccccccccccccccccc", + traceId: internalTraceId, spanId: "3333333333333333", traceFlags: TraceFlags.SAMPLED, }; @@ -36,6 +46,18 @@ function createAttemptSpan(): ReadableSpan { } as unknown as ReadableSpan; } +function createLogRecord(internalTraceId = INTERNAL_TRACE_RUN_A): ReadableLogRecord { + return { + body: "hello", + attributes: {}, + spanContext: { + traceId: internalTraceId, + spanId: "3333333333333333", + traceFlags: TraceFlags.SAMPLED, + }, + } as unknown as ReadableLogRecord; +} + function makeCapturingExporter(): { exporter: SpanExporter; captured: ReadableSpan[][] } { const captured: ReadableSpan[][] = []; const exporter: SpanExporter = { @@ -49,10 +71,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 +114,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 FallbackExternalTraceIds(SEED)); manager.traceContext = { external: { traceparent: TRACEPARENT_RUN_B } }; @@ -77,4 +129,146 @@ 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("gives each run its own fallback trace id when there is no external context", () => { + const { exporter, captured } = makeCapturingExporter(); + const idGenerator = makeIdGenerator(); + + const wrapper = new ExternalSpanExporterWrapper( + exporter, + new FallbackExternalTraceIds(SEED, idGenerator) + ); + + wrapper.export([createAttemptSpan(INTERNAL_TRACE_RUN_A)], () => {}); + // A second run on the same warm process. + wrapper.export([createAttemptSpan(INTERNAL_TRACE_RUN_B)], () => {}); + + 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"); + }); + + it("keeps one fallback trace id across every export within a run", () => { + const { exporter, captured } = makeCapturingExporter(); + const idGenerator = makeIdGenerator(); + + const wrapper = new ExternalSpanExporterWrapper( + exporter, + new FallbackExternalTraceIds(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); + }); + + // Batch processors drain asynchronously, so a run's records are routinely + // exported after the next run has already started. Deciding the id from + // ambient state at that moment would stamp the earlier run's records with the + // later run's id, merging exactly the traces this is meant to separate. + // + // Drives the span and log wrappers together: the TracingSDK shares one + // instance between them, and a run's spans and logs have to land on one trace. + it("stamps records with their own run's id even when exported after the next run started", () => { + const spans = makeCapturingExporter(); + const logs = makeCapturingLogExporter(); + const idGenerator = makeIdGenerator(); + + const fallback = new FallbackExternalTraceIds(SEED, idGenerator); + const spanWrapper = new ExternalSpanExporterWrapper(spans.exporter, fallback); + const logWrapper = new ExternalLogRecordExporterWrapper(logs.exporter, fallback); + + // Run B is underway and has already exported. Its ambient context has no + // `external` key, which is what a run on the fallback path looks like, so + // `getExternalTraceContext()` stays undefined throughout: the point is that + // the run currently in scope must not influence the records below at all. + spanWrapper.export([createAttemptSpan(INTERNAL_TRACE_RUN_B)], () => {}); + manager.traceContext = { traceparent: TRACEPARENT_RUN_B }; + + // Run A's queued records only drain now. + spanWrapper.export([createAttemptSpan(INTERNAL_TRACE_RUN_A)], () => {}); + logWrapper.export([createLogRecord(INTERNAL_TRACE_RUN_A)], () => {}); + + const runBTraceId = spans.captured[0]![0]!.spanContext().traceId; + const lateRunASpanId = spans.captured[1]![0]!.spanContext().traceId; + const lateRunALogId = logs.captured[0]![0]!.spanContext!.traceId; + + expect(lateRunASpanId).not.toBe(runBTraceId); + expect(lateRunALogId).toBe(lateRunASpanId); + }); + + it("leaves external export off when no external trace id was configured", () => { + const { exporter, captured } = makeCapturingExporter(); + const idGenerator = makeIdGenerator(); + + const wrapper = new ExternalSpanExporterWrapper( + exporter, + new FallbackExternalTraceIds("", 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); + }); + + // Instrumentation can start root spans outside a run's async context, each + // its own internal trace, so a run can be alive while the map churns. Evicting + // by insertion order would drop the run still using its id and split it across + // two external traces. + it("keeps the id of a run that is still exporting while other traces fill the map", () => { + const { exporter, captured } = makeCapturingExporter(); + + const wrapper = new ExternalSpanExporterWrapper( + exporter, + new FallbackExternalTraceIds(SEED, makeIdGenerator()) + ); + + const liveRun = "aa000000000000000000000000000000"; + wrapper.export([createAttemptSpan(liveRun)], () => {}); + + for (let i = 0; i < MAX_TRACKED_INTERNAL_TRACES * 2; i++) { + wrapper.export([createAttemptSpan(`bb${`${i}`.padStart(30, "0")}`)], () => {}); + // The run is still going, so it keeps exporting alongside the noise. + wrapper.export([createAttemptSpan(liveRun)], () => {}); + } + + expect(captured.at(-1)![0]!.spanContext().traceId).toBe(captured[0]![0]!.spanContext().traceId); + }); + + // A warm process is long-lived, so the map that remembers each run's id has + // to be bounded rather than growing for the life of the worker. + it("bounds how many runs it remembers", () => { + const { exporter, captured } = makeCapturingExporter(); + const idGenerator = makeIdGenerator(); + + const wrapper = new ExternalSpanExporterWrapper( + exporter, + new FallbackExternalTraceIds(SEED, idGenerator) + ); + + const firstRun = "aa000000000000000000000000000000"; + wrapper.export([createAttemptSpan(firstRun)], () => {}); + + for (let i = 0; i < MAX_TRACKED_INTERNAL_TRACES; i++) { + wrapper.export([createAttemptSpan(`bb${`${i}`.padStart(30, "0")}`)], () => {}); + } + + // Evicted, so it is treated as a run never seen before. + wrapper.export([createAttemptSpan(firstRun)], () => {}); + + expect(captured.at(-1)![0]!.spanContext().traceId).not.toBe( + captured[0]![0]!.spanContext().traceId + ); + }); });