diff --git a/.changeset/otel-linked-trace-mode.md b/.changeset/otel-linked-trace-mode.md new file mode 100644 index 0000000000..ff792d53f0 --- /dev/null +++ b/.changeset/otel-linked-trace-mode.md @@ -0,0 +1,21 @@ +--- +'@workflow/core': minor +'workflow': minor +'@workflow/world-vercel': minor +'@workflow/utils': minor +--- + +Add `WORKFLOW_TRACE_MODE` with a new `linked` default: each workflow/step invocation span is now its own trace root with span links to the delivery and run-origin contexts, instead of one trace spanning the entire run. world-vercel now explicitly injects W3C `traceparent`/`tracestate`/`baggage` headers on outgoing workflow-server requests. + +Span names are also friendlier: workflow and step spans now use the short function name (e.g. `workflow.execute processOrder`, `step.execute chargeCard`, `workflow.start processOrder`) instead of the uppercase prefixes and full machine names (`WORKFLOW_V2 workflow//./src/jobs/order//processOrder`). The full name remains available in the `workflow.name` / `step.name` span attributes, and new `workflowDisplayName` / `stepDisplayName` helpers are exported from `@workflow/utils`. + +Behavioral changes to telemetry under the new default (set `WORKFLOW_TRACE_MODE=continuous` to restore the previous trace shape exactly; the span-name change applies in both modes): + +- A run no longer shares one trace ID: the trace of the request that called `start()` no longer contains the workflow's execution spans — navigate via span links or the `workflow.run.id` attribute instead. +- Sampling decisions are made independently per invocation root (previously one parent-based decision covered the whole run), and the number of root spans/traces increases to one per invocation. +- `workflow.execute`/`step.execute` invocation spans (formerly `WORKFLOW_V2`/`STEP`) become parentless roots, which changes parent/child-based queries and service-map edges. +- Re-enqueued queue messages forward the original run-origin trace carrier unchanged, rather than each invocation's current context. +- Queries or dashboards matching the old `WORKFLOW_V2 ...`/`STEP ...` span names must switch to the new names. +- The queue-delivered `workflow.execute` span kind changed from `internal` to `consumer`, matching the queue-delivered `step.execute` span (this applies in both modes). + +Existing attributes and baggage keys are unchanged, and everything remains a no-op when no OpenTelemetry SDK is registered. diff --git a/docs/content/docs/v5/observability/index.mdx b/docs/content/docs/v5/observability/index.mdx index 58af35f4ae..2bb5142c73 100644 --- a/docs/content/docs/v5/observability/index.mdx +++ b/docs/content/docs/v5/observability/index.mdx @@ -67,6 +67,9 @@ When deployed to Vercel, workflow data is [encrypted end-to-end](/docs/how-it-wo ## More Observability Features + + Distributed tracing with OpenTelemetry for workflow runs, steps, and queue deliveries. + Attach experimental metadata to workflow runs for observability. diff --git a/docs/content/docs/v5/observability/meta.json b/docs/content/docs/v5/observability/meta.json index 41d8063ef6..dd2d7afe24 100644 --- a/docs/content/docs/v5/observability/meta.json +++ b/docs/content/docs/v5/observability/meta.json @@ -1,4 +1,7 @@ { "title": "Observability", - "pages": ["attributes"] + "pages": [ + "tracing", + "attributes" + ] } diff --git a/docs/content/docs/v5/observability/tracing.mdx b/docs/content/docs/v5/observability/tracing.mdx new file mode 100644 index 0000000000..4472cc077d --- /dev/null +++ b/docs/content/docs/v5/observability/tracing.mdx @@ -0,0 +1,106 @@ +--- +title: Tracing +description: Distributed tracing with OpenTelemetry for workflow runs, steps, and queue deliveries. +type: guide +summary: Trace workflow execution end to end with OpenTelemetry. +prerequisites: + - /docs/foundations/workflows-and-steps +related: + - /docs/observability + - /docs/observability/attributes + - /docs/how-it-works/event-sourcing +--- + +The Workflow SDK is instrumented with [OpenTelemetry](https://opentelemetry.io) out of the box. It emits spans for workflow starts, every workflow and step invocation, and the HTTP calls it makes to the workflow backend — and it propagates trace context across queue deliveries so a run remains traceable end to end. + +The SDK only depends on the OpenTelemetry **API**, never on an SDK or exporter. If your application does not register an OpenTelemetry SDK, all tracing code is a silent no-op with no overhead and no behavior change. + +## Enabling tracing + +Register any OpenTelemetry Node SDK in your application. On Vercel with Next.js, the simplest setup is [`@vercel/otel`](https://vercel.com/docs/observability/otel-overview) in `instrumentation.ts`: + +```typescript title="instrumentation.ts" lineNumbers +import { registerOTel } from "@vercel/otel" + +export function register() { + registerOTel({ serviceName: "my-app" }) +} +``` + +No workflow-specific configuration is required. As soon as a tracer provider and propagator are registered, the SDK's spans, context propagation, and span links activate automatically. + +## Spans + +| Span name | Kind | Emitted when | +| --- | --- | --- | +| `workflow.start ` | internal | `start()` is called in your application code | +| `workflow.execute ` | consumer (root) | a queue delivery invokes the workflow — replay, orchestration, and inline steps run under it | +| `step.execute ` | internal (inline) / consumer + root (queue-delivered) | a step function executes | +| `http ` | client | the SDK calls the workflow backend (event reads/writes) | + +`` is the short function name (for example `processOrder`); the full machine name, including the source module, is available in the `workflow.name` / `step.name` attributes. + +## Key attributes + +| Attribute | Description | +| --- | --- | +| `workflow.run.id` | The run ID (`wrun_...`). Present on every workflow and step span — the primary key for finding all spans of a run. | +| `workflow.name` | The workflow function name. | +| `workflow.trace.mode` | The active trace mode (`linked` or `continuous`). | +| `workflow.trace.propagated` | Whether the invocation received trace context from the queue message. | +| `workflow.queue.overhead_ms` | Time between the message being enqueued and the handler starting — queue dwell plus any cold start. | + +## Trace shape: one trace per invocation + +A single workflow run can span hours or days across many separate function invocations: every step completion, `sleep()` wake-up, and retry is a new queue delivery. Stitching all of that into one trace produces giant, slow-loading traces that most tracing backends truncate. + +Instead, the SDK creates **one bounded trace per invocation**. Each `workflow.execute` (or background `step.execute`) span starts a new trace root and attaches two **span links**: + +- a link to the **enqueue site** — the span that queued the message which triggered this invocation, and +- a link to the **run origin** — the trace in which `start()` was originally called. + +A span link is OpenTelemetry's relationship for "causally related, but in a different trace." It is the standard pattern for asynchronous messaging, where producing and consuming a message can be separated by arbitrary time. + +```mermaid +flowchart LR + O["start() request trace"] + A["invocation 1"] + B["invocation 2"] + C["invocation 3 ..."] + A -. "link" .-> O + B -. "link" .-> O + C -. "link" .-> O + B -. "link" .-> A + C -. "link" .-> B + + style O fill:#a78bfa,stroke:#8b5cf6,color:#000 +``` + +Each invocation links back to the trace that enqueued it and to the run origin. + +To see a whole run, query by attribute rather than by trace ID — for example `workflow.run.id = wrun_...` in your tracing backend — or follow the span links between invocation traces. + +## Trace modes + +The `WORKFLOW_TRACE_MODE` environment variable controls the shape: + +| Mode | Behavior | +| --- | --- | +| `linked` (default) | Each invocation is its own trace root with span links to the enqueue site and the run origin. Traces stay small; sampling is decided per invocation. | +| `continuous` | The run-origin context becomes the **parent** of every invocation, so the entire run shares one trace ID. | + + +This is a behavior change from v4, which always used `continuous`-style tracing. If you have dashboards or queries that assume one trace ID per run, either update them to use `workflow.run.id` and span links, or set `WORKFLOW_TRACE_MODE=continuous` to restore the previous shape. Note that in `linked` mode each invocation root makes its own sampling decision, and the number of root spans increases to one per invocation. + + +## Context propagation + +When tracing is enabled, the SDK propagates [W3C Trace Context](https://www.w3.org/TR/trace-context/) on its outbound calls: + +- **Backend requests** carry `traceparent`, `tracestate`, and `baggage` headers, so backend spans can join your trace. +- **Queue messages** carry the run-origin trace context in the message payload, and the queue re-delivers the producer's context to the workflow handler, where it becomes the enqueue-site span link. +- **Baggage** carries `workflow.run_id` and `workflow.name` entries during workflow execution, allowing downstream services you call from steps to tag their own telemetry with the run ID. + + +Baggage entries set by your application are propagated as a `baggage` HTTP header on the SDK's backend requests, like any other OpenTelemetry-instrumented HTTP call. Avoid placing sensitive values in baggage. + diff --git a/packages/core/package.json b/packages/core/package.json index 333b87ac27..e95dbb2d33 100644 --- a/packages/core/package.json +++ b/packages/core/package.json @@ -108,6 +108,9 @@ }, "devDependencies": { "@opentelemetry/api": "1.9.0", + "@opentelemetry/context-async-hooks": "1.30.1", + "@opentelemetry/core": "1.30.1", + "@opentelemetry/sdk-trace-base": "1.30.1", "@types/debug": "4.1.12", "@types/node": "catalog:", "@types/seedrandom": "3.0.8", diff --git a/packages/core/src/runtime-trace-mode.test.ts b/packages/core/src/runtime-trace-mode.test.ts new file mode 100644 index 0000000000..f0a6cfb067 --- /dev/null +++ b/packages/core/src/runtime-trace-mode.test.ts @@ -0,0 +1,359 @@ +import { + context, + trace as otelTrace, + propagation, + type Span, + SpanKind, +} from '@opentelemetry/api'; +import { AsyncLocalStorageContextManager } from '@opentelemetry/context-async-hooks'; +import { W3CTraceContextPropagator } from '@opentelemetry/core'; +import { + BasicTracerProvider, + InMemorySpanExporter, + type ReadableSpan, + SimpleSpanProcessor, +} from '@opentelemetry/sdk-trace-base'; +import { + type Event, + SPEC_VERSION_CURRENT, + type WorkflowRun, +} from '@workflow/world'; +import { + afterAll, + afterEach, + beforeAll, + describe, + expect, + it, + vi, +} from 'vitest'; +import { runtimeLogger } from './logger.js'; +import { setWorld } from './runtime/world.js'; +import { workflowEntrypoint } from './runtime.js'; +import { dehydrateWorkflowArguments } from './serialization.js'; +import { getNextTraceCarrier, getWorkflowTraceMode } from './telemetry.js'; + +vi.mock('@vercel/functions', () => ({ + waitUntil: vi.fn((p: Promise) => { + p.catch(() => {}); + }), +})); + +// Run-origin trace context, as carried in queue messages. Uses a fixed, +// valid W3C traceparent so assertions are deterministic. +const ORIGIN_TRACE_ID = '0af7651916cd43dd8448eb211c80319c'; +const ORIGIN_SPAN_ID = 'b7ad6b7169203331'; +const ORIGIN_CARRIER = { + traceparent: `00-${ORIGIN_TRACE_ID}-${ORIGIN_SPAN_ID}-01`, +}; + +const exporter = new InMemorySpanExporter(); +const provider = new BasicTracerProvider(); +const contextManager = new AsyncLocalStorageContextManager(); + +beforeAll(() => { + provider.addSpanProcessor(new SimpleSpanProcessor(exporter)); + contextManager.enable(); + context.setGlobalContextManager(contextManager); + propagation.setGlobalPropagator(new W3CTraceContextPropagator()); + otelTrace.setGlobalTracerProvider(provider); +}); + +afterAll(async () => { + await provider.shutdown(); + context.disable(); + propagation.disable(); + otelTrace.disable(); +}); + +afterEach(() => { + exporter.reset(); + setWorld(undefined); + vi.unstubAllEnvs(); + vi.clearAllMocks(); +}); + +const getWorkflowTransformCode = (workflowName: string) => + `;globalThis.__private_workflows = new Map(); + globalThis.__private_workflows.set(${JSON.stringify(workflowName)}, ${workflowName});`; + +const simpleWorkflow = `async function workflow() { + return 'done'; + }${getWorkflowTransformCode('workflow')}`; + +async function makeRunningRun(runId: string): Promise { + return { + runId, + workflowName: 'workflow', + status: 'running', + input: await dehydrateWorkflowArguments([], runId, undefined, []), + createdAt: new Date('2024-01-01T00:00:00.000Z'), + updatedAt: new Date('2024-01-01T00:00:00.000Z'), + startedAt: new Date('2024-01-01T00:00:00.000Z'), + deploymentId: 'test-deployment', + }; +} + +/** + * Drives the workflow queue handler once with the given trace carrier, + * inside an active "delivery" span (simulating the span the platform/queue + * consumer creates around the delivery request). Returns the finished + * WORKFLOW_V2 span, the delivery span, and all queued messages. + */ +async function driveHandler(opts: { + runId: string; + workflowCode: string; + traceCarrier?: Record; +}) { + const workflowRun = await makeRunningRun(opts.runId); + const queuedMessages: any[] = []; + + const eventsCreate = vi.fn(async (_runId: string, data: any) => { + if (data.eventType === 'run_started') { + return { run: workflowRun, events: [] as Event[] }; + } + return { + event: { + eventId: `event-${Math.random()}`, + runId: workflowRun.runId, + createdAt: new Date(), + ...data, + }, + }; + }); + + setWorld({ + specVersion: SPEC_VERSION_CURRENT, + createQueueHandler: vi.fn( + ( + _prefix: string, + handler: (message: unknown, metadata: unknown) => Promise + ) => { + return async () => { + await handler( + { + runId: workflowRun.runId, + requestedAt: new Date('2024-01-01T00:00:00.000Z'), + traceCarrier: opts.traceCarrier, + }, + { + requestId: 'req_test', + attempt: 1, + queueName: '__wkf_workflow_workflow', + messageId: 'msg_test', + } + ); + return new Response(null, { status: 204 }); + }; + } + ), + events: { + create: eventsCreate, + list: vi.fn(async () => ({ + data: [] as Event[], + hasMore: false, + cursor: 'cursor_test', + })), + }, + runs: { + get: vi.fn(async () => workflowRun), + }, + queue: vi.fn(async (_queueName: string, message: unknown) => { + queuedMessages.push(message); + return { messageId: null }; + }), + getEncryptionKeyForRun: vi.fn(async () => undefined), + } as any); + + const handler = workflowEntrypoint(opts.workflowCode); + + // Invoke inside an active "delivery" span so linkToCurrentContext() + // observes a live delivery context, as it would in production. + const tracer = otelTrace.getTracer('test'); + let deliverySpan!: Span; + await tracer.startActiveSpan('queue delivery', async (span) => { + deliverySpan = span; + try { + await handler(new Request('https://example.test')); + } finally { + span.end(); + } + }); + + const workflowSpan = exporter + .getFinishedSpans() + .find((s) => s.name === 'workflow.execute workflow'); + + return { workflowSpan, deliverySpan, queuedMessages }; +} + +function linkTraceIds(span: ReadableSpan | undefined): string[] { + return (span?.links ?? []).map((l) => l.context.traceId); +} + +describe('getWorkflowTraceMode', () => { + it('defaults to linked when WORKFLOW_TRACE_MODE is unset', () => { + vi.stubEnv('WORKFLOW_TRACE_MODE', ''); + expect(getWorkflowTraceMode()).toBe('linked'); + }); + + it('returns continuous when WORKFLOW_TRACE_MODE=continuous', () => { + vi.stubEnv('WORKFLOW_TRACE_MODE', 'continuous'); + expect(getWorkflowTraceMode()).toBe('continuous'); + }); + + it('warns once for unrecognized values and falls back to linked', () => { + const warnSpy = vi + .spyOn(runtimeLogger, 'warn') + .mockImplementation(() => {}); + vi.stubEnv('WORKFLOW_TRACE_MODE', 'continous'); + + expect(getWorkflowTraceMode()).toBe('linked'); + expect(getWorkflowTraceMode()).toBe('linked'); + + expect(warnSpy).toHaveBeenCalledTimes(1); + const message = warnSpy.mock.calls[0][0] as string; + expect(message).toContain('"continous"'); + expect(message).toContain('"linked"'); + expect(message).toContain('"continuous"'); + warnSpy.mockRestore(); + }); +}); + +describe('workflowEntrypoint trace modes', () => { + it('linked (default): creates the WORKFLOW_V2 span as a new root with links to delivery and run-origin contexts', async () => { + const { workflowSpan, deliverySpan } = await driveHandler({ + runId: 'wrun_trace_linked', + workflowCode: simpleWorkflow, + traceCarrier: ORIGIN_CARRIER, + }); + + expect(workflowSpan).toBeDefined(); + // New root: no parent, and a fresh trace distinct from both the + // delivery trace and the run-origin trace. + expect(workflowSpan?.parentSpanId).toBeUndefined(); + expect(workflowSpan?.spanContext().traceId).not.toBe(ORIGIN_TRACE_ID); + expect(workflowSpan?.spanContext().traceId).not.toBe( + deliverySpan.spanContext().traceId + ); + + // Links to BOTH the delivery context and the run-origin context. + expect(workflowSpan?.links).toHaveLength(2); + expect(linkTraceIds(workflowSpan)).toContain( + deliverySpan.spanContext().traceId + ); + expect(linkTraceIds(workflowSpan)).toContain(ORIGIN_TRACE_ID); + + expect(workflowSpan?.attributes['workflow.trace.mode']).toBe('linked'); + expect(workflowSpan?.attributes['workflow.trace.propagated']).toBe(true); + + // Queue-delivered invocation spans use the CONSUMER kind, matching + // queue-delivered step.execute spans. + expect(workflowSpan?.kind).toBe(SpanKind.CONSUMER); + }); + + it('linked: treats an empty trace carrier ({}) like an absent one', async () => { + const { workflowSpan, deliverySpan } = await driveHandler({ + runId: 'wrun_trace_linked_empty_carrier', + workflowCode: simpleWorkflow, + traceCarrier: {}, + }); + + expect(workflowSpan).toBeDefined(); + expect(workflowSpan?.parentSpanId).toBeUndefined(); + // Only the delivery link — no origin link is derived from `{}`. + expect(workflowSpan?.links).toHaveLength(1); + expect(linkTraceIds(workflowSpan)).toContain( + deliverySpan.spanContext().traceId + ); + // An empty carrier does not count as propagated trace context. + expect(workflowSpan?.attributes['workflow.trace.propagated']).toBe(false); + }); + + it('linked: without an incoming carrier, still creates a root span with only the delivery link', async () => { + const { workflowSpan, deliverySpan } = await driveHandler({ + runId: 'wrun_trace_linked_no_carrier', + workflowCode: simpleWorkflow, + traceCarrier: undefined, + }); + + expect(workflowSpan).toBeDefined(); + expect(workflowSpan?.parentSpanId).toBeUndefined(); + expect(workflowSpan?.links).toHaveLength(1); + expect(linkTraceIds(workflowSpan)).toContain( + deliverySpan.spanContext().traceId + ); + expect(workflowSpan?.attributes['workflow.trace.propagated']).toBe(false); + }); + + it('continuous: preserves the legacy shape — parented to the run-origin context with a delivery link', async () => { + vi.stubEnv('WORKFLOW_TRACE_MODE', 'continuous'); + + const { workflowSpan, deliverySpan } = await driveHandler({ + runId: 'wrun_trace_continuous', + workflowCode: simpleWorkflow, + traceCarrier: ORIGIN_CARRIER, + }); + + expect(workflowSpan).toBeDefined(); + // Same trace as the run origin, parented to the carrier's span. + expect(workflowSpan?.spanContext().traceId).toBe(ORIGIN_TRACE_ID); + expect(workflowSpan?.parentSpanId).toBe(ORIGIN_SPAN_ID); + + // Only the delivery link (no self-link to the origin). + expect(workflowSpan?.links).toHaveLength(1); + expect(linkTraceIds(workflowSpan)).toContain( + deliverySpan.spanContext().traceId + ); + + expect(workflowSpan?.attributes['workflow.trace.mode']).toBe('continuous'); + }); +}); + +// The carrier put on re-enqueued messages is produced by getNextTraceCarrier. +// (The combined V2 handler now executes a single owned step inline rather +// than queueing it, so this invariant is exercised directly on the helper +// instead of by inspecting a queued step message.) +describe('getNextTraceCarrier (re-enqueue carrier semantics)', () => { + it('linked: forwards a usable run-origin carrier unchanged', async () => { + // The original carrier flows forward unchanged so the run-origin + // identity is preserved for future links. + const next = await getNextTraceCarrier('linked', ORIGIN_CARRIER); + expect(next).toEqual(ORIGIN_CARRIER); + }); + + it('linked: replaces an empty carrier with the current invocation context', async () => { + const tracer = otelTrace.getTracer('test'); + let activeTraceId = ''; + const next = await tracer.startActiveSpan('inv', async (span) => { + activeTraceId = span.spanContext().traceId; + try { + return await getNextTraceCarrier('linked', {}); + } finally { + span.end(); + } + }); + // The useless `{}` is not forwarded — this invocation serializes its + // own context, becoming the de-facto run origin for future links. + expect(next.traceparent).toContain(activeTraceId); + expect(next.traceparent).not.toBe(ORIGIN_CARRIER.traceparent); + }); + + it('continuous: serializes the current invocation context, not the origin carrier', async () => { + const tracer = otelTrace.getTracer('test'); + let activeTraceId = ''; + const next = await tracer.startActiveSpan('inv', async (span) => { + activeTraceId = span.spanContext().traceId; + try { + return await getNextTraceCarrier('continuous', ORIGIN_CARRIER); + } finally { + span.end(); + } + }); + // Continuous mode always serializes the current context rather than + // forwarding the incoming carrier. + expect(next.traceparent).toContain(activeTraceId); + expect(next.traceparent).not.toBe(ORIGIN_CARRIER.traceparent); + expect(activeTraceId).not.toBe(ORIGIN_TRACE_ID); + }); +}); diff --git a/packages/core/src/runtime.ts b/packages/core/src/runtime.ts index 9d9b6351a3..b8ed8c7e55 100644 --- a/packages/core/src/runtime.ts +++ b/packages/core/src/runtime.ts @@ -9,7 +9,10 @@ import { RunExpiredError, WorkflowRuntimeError, } from '@workflow/errors'; -import { parseWorkflowName } from '@workflow/utils/parse-name'; +import { + parseWorkflowName, + workflowDisplayName, +} from '@workflow/utils/parse-name'; import { type Event, getQueueTopicPrefix, @@ -53,8 +56,11 @@ import { dehydrateRunError } from './serialization.js'; import { remapErrorStack } from './source-map.js'; import * as Attribute from './telemetry/semantic-conventions.js'; import { - linkToCurrentContext, - serializeTraceCarrier, + buildInvocationSpanLinks, + getNextTraceCarrier, + getSpanKind, + getWorkflowTraceMode, + isUsableTraceCarrier, trace, withTraceContext, withWorkflowBaggage, @@ -246,13 +252,21 @@ export function workflowEntrypoint( const { runId, - traceCarrier: traceContext, + traceCarrier: incomingTraceCarrier, requestedAt, stepId: incomingStepId, stepName: incomingStepName, replayDivergence, runInput, } = WorkflowInvokePayloadSchema.parse(message_); + // `start()` always attaches a trace carrier, but + // serializeTraceCarrier() returns `{}` when no OTEL SDK is registered + // or no span is active — treat an empty carrier the same as an + // absent one so linked mode falls back to a fresh origin instead of + // forwarding a useless `{}` forever. + const traceContext = isUsableTraceCarrier(incomingTraceCarrier) + ? incomingTraceCarrier + : undefined; const { requestId } = metadata; const workflowName = metadata.queueName.slice(workflowPrefix.length); @@ -320,7 +334,27 @@ export function workflowEntrypoint( return; } - const spanLinks = await linkToCurrentContext(); + // --- Trace correlation mode --- + // 'linked' (default): the WORKFLOW_V2 span below starts a NEW root + // trace, with span links to (a) the incoming delivery context and + // (b) the run-origin context from the message's trace carrier. This + // bounds each trace to a single invocation instead of stitching an + // entire (potentially hours-long) run into one giant trace. + // 'continuous': legacy behavior — the restored run-origin context + // becomes the parent of this invocation's spans. + const traceMode = getWorkflowTraceMode(); + + // Trace carrier to attach to messages this invocation enqueues — + // see getNextTraceCarrier for the linked/continuous semantics. + const nextTraceCarrier = (): Promise> => + getNextTraceCarrier(traceMode, traceContext); + + // Span links to the incoming delivery context and (in linked mode) + // the run-origin context from the trace carrier. + const spanLinks = await buildInvocationSpanLinks( + traceMode, + traceContext + ); // --- Replay budget bookkeeping --- // The replay budget bounds the *non-step* portion of a single @@ -348,14 +382,25 @@ export function workflowEntrypoint( // single step longer than the budget. const replayBudget = new ReplayBudget(); - return await withTraceContext(traceContext, async () => { + // In linked mode the run-origin context is NOT restored as the + // active (parent) context — passing `undefined` makes + // withTraceContext a passthrough, and the WORKFLOW_V2 span below + // becomes a new trace root carrying span links instead. + const parentTraceCarrier = + traceMode === 'continuous' ? traceContext : undefined; + // Queue-delivered invocation: CONSUMER kind, matching the + // queue-delivered step.execute span. + const spanKind = await getSpanKind('CONSUMER'); + return await withTraceContext(parentTraceCarrier, async () => { return await withWorkflowBaggage( { workflowRunId: runId, workflowName }, async () => { const world = await getWorld(); return trace( - `WORKFLOW_V2 ${workflowName}`, - { links: spanLinks }, + `workflow.execute ${workflowDisplayName(workflowName)}`, + traceMode === 'linked' + ? { kind: spanKind, links: spanLinks, root: true } + : { kind: spanKind, links: spanLinks }, async (span) => { span?.setAttributes({ ...Attribute.WorkflowName(workflowName), @@ -367,6 +412,7 @@ export function workflowEntrypoint( ...getQueueOverhead({ requestedAt }), ...Attribute.WorkflowRunId(runId), ...Attribute.WorkflowTracePropagated(!!traceContext), + ...Attribute.WorkflowTraceMode(traceMode), }); const invocationStartTime = Date.now(); @@ -440,7 +486,7 @@ export function workflowEntrypoint( getWorkflowQueueName(workflowName, namespace), { runId, - traceCarrier: await serializeTraceCarrier(), + traceCarrier: await nextTraceCarrier(), requestedAt: new Date(), } ); @@ -702,7 +748,7 @@ export function workflowEntrypoint( getWorkflowQueueName(workflowName, namespace), { runId, - traceCarrier: await serializeTraceCarrier(), + traceCarrier: await nextTraceCarrier(), requestedAt: new Date(), } ); @@ -1149,7 +1195,7 @@ export function workflowEntrypoint( // suspension passes — see // runtime/wait-continuation.ts for the full // delay/key selection rationale. - const traceCarrier = await serializeTraceCarrier(); + const traceCarrier = await nextTraceCarrier(); const dispatches: Promise[] = []; for (const step of pendingSteps) { if ( @@ -1238,8 +1284,7 @@ export function workflowEntrypoint( // Any pending wait timer was already enqueued as part // of the unified dispatch above, so we can return // unconditionally here. - const retryTraceCarrier = - await serializeTraceCarrier(); + const retryTraceCarrier = await nextTraceCarrier(); await queueMessage( world, getWorkflowQueueName(workflowName, namespace), @@ -1290,7 +1335,7 @@ export function workflowEntrypoint( getWorkflowQueueName(workflowName, namespace), { runId, - traceCarrier: await serializeTraceCarrier(), + traceCarrier: await nextTraceCarrier(), requestedAt: new Date(), } ); @@ -1324,7 +1369,7 @@ export function workflowEntrypoint( getWorkflowQueueName(workflowName, namespace), { runId, - traceCarrier: await serializeTraceCarrier(), + traceCarrier: await nextTraceCarrier(), requestedAt: new Date(), replayDivergence: { eventId: err.eventId, diff --git a/packages/core/src/runtime/resume-hook.ts b/packages/core/src/runtime/resume-hook.ts index 3c5a653d76..f0e8f466fd 100644 --- a/packages/core/src/runtime/resume-hook.ts +++ b/packages/core/src/runtime/resume-hook.ts @@ -21,7 +21,7 @@ import { } from '../serialization.js'; import { WEBHOOK_RESPONSE_WRITABLE } from '../symbols.js'; import * as Attribute from '../telemetry/semantic-conventions.js'; -import { getSpanContextForTraceCarrier, trace } from '../telemetry.js'; +import { linkToTraceCarrier, trace } from '../telemetry.js'; import { getWorldLazy } from './get-world-lazy.js'; import { getWorkflowQueueName } from './helpers.js'; import { safeWaitUntil, waitedUntil } from './wait-until.js'; @@ -185,13 +185,13 @@ export async function resumeHook( ...Attribute.WorkflowName(workflowRun.workflowName), }); - const traceCarrier = workflowRun.executionContext?.traceCarrier; - - if (traceCarrier) { - const context = await getSpanContextForTraceCarrier(traceCarrier); - if (context) { - span?.addLink?.({ context }); - } + // Link to the run-origin context from the workflow run's stored + // trace carrier (skipped when absent or invalid). + const originLink = await linkToTraceCarrier( + workflowRun.executionContext?.traceCarrier + ); + if (originLink) { + span?.addLink?.(originLink); } // Re-trigger the workflow against the deployment ID associated diff --git a/packages/core/src/runtime/start.ts b/packages/core/src/runtime/start.ts index c6f6c60ef2..76effc2c57 100644 --- a/packages/core/src/runtime/start.ts +++ b/packages/core/src/runtime/start.ts @@ -4,6 +4,7 @@ import { WorkflowRuntimeError, WorkflowWorldError, } from '@workflow/errors'; +import { workflowDisplayName } from '@workflow/utils/parse-name'; import type { WorkflowInvokePayload, World } from '@workflow/world'; import { isLegacySpecVersion, @@ -165,7 +166,8 @@ export async function start( ); } - return trace(`workflow.start ${workflowName}`, async (span) => { + const spanName = `workflow.start ${workflowDisplayName(workflowName)}`; + return trace(spanName, async (span) => { span?.setAttributes({ ...Attribute.WorkflowName(workflowName), ...Attribute.WorkflowOperation('start'), diff --git a/packages/core/src/runtime/step-executor.ts b/packages/core/src/runtime/step-executor.ts index 1cfe99a78b..e901d93c23 100644 --- a/packages/core/src/runtime/step-executor.ts +++ b/packages/core/src/runtime/step-executor.ts @@ -8,7 +8,7 @@ import { TooEarlyError, WorkflowRuntimeError, } from '@workflow/errors'; -import { pluralize } from '@workflow/utils'; +import { pluralize, stepDisplayName } from '@workflow/utils'; import type { World } from '@workflow/world'; import { SPEC_VERSION_CURRENT } from '@workflow/world'; import type { CryptoKey } from '../encryption.js'; @@ -80,7 +80,8 @@ export async function executeStep( } = params; const isVercel = process.env.VERCEL_URL !== undefined; - return trace(`STEP ${stepName}`, {}, async (span) => { + const spanName = `step.execute ${stepDisplayName(stepName)}`; + return trace(spanName, {}, async (span) => { span?.setAttributes({ ...Attribute.StepName(stepName), ...Attribute.WorkflowName(workflowName), diff --git a/packages/core/src/runtime/step-handler.test.ts b/packages/core/src/runtime/step-handler.test.ts index 278f5a6b0e..3afd079902 100644 --- a/packages/core/src/runtime/step-handler.test.ts +++ b/packages/core/src/runtime/step-handler.test.ts @@ -95,7 +95,15 @@ vi.mock('../telemetry.js', () => ({ }), withTraceContext: vi.fn((_ctx: unknown, fn: () => unknown) => fn()), getSpanKind: vi.fn().mockResolvedValue(undefined), + getWorkflowTraceMode: vi.fn(() => 'linked'), + isUsableTraceCarrier: vi.fn( + (carrier?: Record) => + carrier !== undefined && Object.keys(carrier).length > 0 + ), + getNextTraceCarrier: vi.fn().mockResolvedValue({}), + buildInvocationSpanLinks: vi.fn().mockResolvedValue([]), linkToCurrentContext: vi.fn().mockResolvedValue([]), + linkToTraceCarrier: vi.fn().mockResolvedValue(undefined), withWorkflowBaggage: vi.fn((_attrs: unknown, fn: () => unknown) => fn()), })); diff --git a/packages/core/src/runtime/step-handler.ts b/packages/core/src/runtime/step-handler.ts index 201d93500b..659bbd6768 100644 --- a/packages/core/src/runtime/step-handler.ts +++ b/packages/core/src/runtime/step-handler.ts @@ -9,7 +9,7 @@ import { WorkflowRuntimeError, WorkflowWorldError, } from '@workflow/errors'; -import { formatStepName, pluralize } from '@workflow/utils'; +import { formatStepName, pluralize, stepDisplayName } from '@workflow/utils'; import { getPort } from '@workflow/utils/get-port'; import { getQueueTopicPrefix, @@ -31,9 +31,11 @@ import { import { contextStorage } from '../step/context-storage.js'; import * as Attribute from '../telemetry/semantic-conventions.js'; import { + buildInvocationSpanLinks, + getNextTraceCarrier, getSpanKind, - linkToCurrentContext, - serializeTraceCarrier, + getWorkflowTraceMode, + isUsableTraceCarrier, trace, withTraceContext, } from '../telemetry.js'; @@ -84,10 +86,29 @@ function createStepHandler(namespace?: string) { workflowRunId, workflowStartedAt, stepId, - traceCarrier: traceContext, + traceCarrier: incomingTraceCarrier, requestedAt, } = StepInvokePayloadSchema.parse(message_); const { requestId } = metadata; + // serializeTraceCarrier() returns `{}` when no OTEL SDK is registered + // or no span is active — treat an empty carrier the same as an absent + // one so linked mode falls back to a fresh origin instead of + // forwarding a useless `{}` forever. + const traceContext = isUsableTraceCarrier(incomingTraceCarrier) + ? incomingTraceCarrier + : undefined; + + // --- Trace correlation mode --- + // 'linked' (default): the STEP span below starts a NEW root trace with + // span links to the delivery context and the run-origin context from + // the message's trace carrier. 'continuous': legacy behavior — the + // restored run-origin context parents this invocation's spans. + const traceMode = getWorkflowTraceMode(); + + // Trace carrier to attach to messages this invocation enqueues — see + // getNextTraceCarrier for the linked/continuous semantics. + const nextTraceCarrier = (): Promise> => + getNextTraceCarrier(traceMode, traceContext); // --- Max delivery check --- // Enforce max delivery limit before any infrastructure calls. @@ -141,7 +162,7 @@ function createStepHandler(namespace?: string) { getWorkflowQueueName(workflowName, resolvedNamespace), { runId: workflowRunId, - traceCarrier: await serializeTraceCarrier(), + traceCarrier: await nextTraceCarrier(), requestedAt: new Date(), } ); @@ -168,9 +189,16 @@ function createStepHandler(namespace?: string) { return; } - const spanLinks = await linkToCurrentContext(); - // Execute step within the propagated trace context - return await withTraceContext(traceContext, async () => { + // Span links to the incoming delivery context and (in linked mode) + // the run-origin context from the trace carrier. + const spanLinks = await buildInvocationSpanLinks(traceMode, traceContext); + + // Execute step within the propagated trace context (continuous mode + // only — in linked mode the STEP span below becomes a new trace root + // carrying span links instead, so withTraceContext is a passthrough). + const parentTraceCarrier = + traceMode === 'continuous' ? traceContext : undefined; + return await withTraceContext(parentTraceCarrier, async () => { // Extract the step name from the topic name const stepName = metadata.queueName.slice(stepPrefix.length); const world = await getWorld(); @@ -192,8 +220,10 @@ function createStepHandler(namespace?: string) { ]); return trace( - `STEP ${stepName}`, - { kind: spanKind, links: spanLinks }, + `step.execute ${stepDisplayName(stepName)}`, + traceMode === 'linked' + ? { kind: spanKind, links: spanLinks, root: true } + : { kind: spanKind, links: spanLinks }, async (span) => { span?.setAttributes({ ...Attribute.StepName(stepName), @@ -216,6 +246,7 @@ function createStepHandler(namespace?: string) { ...Attribute.WorkflowRunId(workflowRunId), ...Attribute.StepId(stepId), ...Attribute.StepTracePropagated(!!traceContext), + ...Attribute.WorkflowTraceMode(traceMode), }); // step_started validates state and returns the step entity, so no separate @@ -289,7 +320,7 @@ function createStepHandler(namespace?: string) { getWorkflowQueueName(workflowName, resolvedNamespace), { runId: workflowRunId, - traceCarrier: await serializeTraceCarrier(), + traceCarrier: await nextTraceCarrier(), requestedAt: new Date(), } ); @@ -393,7 +424,7 @@ function createStepHandler(namespace?: string) { getWorkflowQueueName(workflowName, resolvedNamespace), { runId: workflowRunId, - traceCarrier: await serializeTraceCarrier(), + traceCarrier: await nextTraceCarrier(), requestedAt: new Date(), } ); @@ -493,7 +524,7 @@ function createStepHandler(namespace?: string) { getWorkflowQueueName(workflowName, resolvedNamespace), { runId: workflowRunId, - traceCarrier: await serializeTraceCarrier(), + traceCarrier: await nextTraceCarrier(), requestedAt: new Date(), } ); @@ -545,7 +576,7 @@ function createStepHandler(namespace?: string) { getWorkflowQueueName(workflowName, resolvedNamespace), { runId: workflowRunId, - traceCarrier: await serializeTraceCarrier(), + traceCarrier: await nextTraceCarrier(), requestedAt: new Date(), } ); @@ -988,7 +1019,7 @@ function createStepHandler(namespace?: string) { getWorkflowQueueName(workflowName, resolvedNamespace), { runId: workflowRunId, - traceCarrier: await serializeTraceCarrier(), + traceCarrier: await nextTraceCarrier(), requestedAt: new Date(), } ); @@ -1049,7 +1080,7 @@ function createStepHandler(namespace?: string) { } throw err; }), - serializeTraceCarrier(), + nextTraceCarrier(), ]); if (stepCompleted409) { diff --git a/packages/core/src/telemetry.ts b/packages/core/src/telemetry.ts index 1b26294b5e..120a8f7a70 100644 --- a/packages/core/src/telemetry.ts +++ b/packages/core/src/telemetry.ts @@ -9,6 +9,105 @@ import * as Attr from './telemetry/semantic-conventions.js'; // Trace Context Propagation Utilities // ============================================================ +/** + * Controls how workflow/step queue-handler spans relate to the run-origin + * trace context carried in queue messages: + * + * - `'linked'` (default): each invocation starts a NEW root trace; the + * run-origin context (and the incoming delivery context) are attached as + * span links. This keeps traces bounded per invocation instead of + * stitching a long-lived run into one giant trace. + * - `'continuous'`: the restored run-origin context becomes the parent of + * the invocation span (legacy behavior), producing a single trace that + * spans the entire run. + */ +export type WorkflowTraceMode = 'linked' | 'continuous'; + +/** Unrecognized `WORKFLOW_TRACE_MODE` values we already warned about. */ +const warnedUnrecognizedTraceModes = new Set(); + +/** + * Resolves the active trace mode from the `WORKFLOW_TRACE_MODE` env var. + * Defaults to `'linked'`; any value other than `'continuous'` selects it. + * Unrecognized non-empty values (e.g. typos like `continous`) emit a + * one-time warning so silent fallback to `linked` is at least visible. + */ +export function getWorkflowTraceMode(): WorkflowTraceMode { + const value = process.env.WORKFLOW_TRACE_MODE; + if (value === 'continuous') return 'continuous'; + if (value && value !== 'linked' && !warnedUnrecognizedTraceModes.has(value)) { + warnedUnrecognizedTraceModes.add(value); + runtimeLogger.warn( + `Unrecognized WORKFLOW_TRACE_MODE value "${value}"; expected "linked" or "continuous". Falling back to "linked".` + ); + } + return 'linked'; +} + +/** + * Returns whether a serialized trace carrier is usable, i.e. present and + * non-empty. `serializeTraceCarrier()` returns `{}` when no OTEL SDK is + * registered or no span is active, and `start()` always attaches the + * carrier to the first queue message — so an empty carrier must be treated + * the same as an absent one wherever the trace-mode logic branches. + */ +export function isUsableTraceCarrier( + carrier: Record | undefined +): carrier is Record { + return carrier !== undefined && Object.keys(carrier).length > 0; +} + +/** + * Returns the trace carrier to attach to messages the current invocation + * enqueues. In `linked` mode the ORIGINAL run-origin carrier is forwarded + * unchanged (when usable) so every future invocation links back to the same + * origin; otherwise — `continuous` mode, or no usable incoming carrier — + * the current (active) context is serialized, so the trace keeps chaining + * (continuous) or the first instrumented invocation becomes the de-facto + * origin (linked). + */ +export function getNextTraceCarrier( + traceMode: WorkflowTraceMode, + incomingCarrier: Record | undefined +): Promise> { + return traceMode === 'linked' && isUsableTraceCarrier(incomingCarrier) + ? Promise.resolve(incomingCarrier) + : serializeTraceCarrier(); +} + +/** + * Builds the span links for a queue-delivered invocation span + * (`workflow.execute` / `step.execute`): + * + * - a link to the incoming delivery context (the active span extracted from + * the queue delivery request, i.e. the enqueue site once the queue + * propagates producer context), and + * - in `linked` mode, a link to the run-origin context from the message's + * trace carrier (skipped when absent, empty, invalid, or identical to the + * delivery context). + */ +export async function buildInvocationSpanLinks( + traceMode: WorkflowTraceMode, + incomingCarrier: Record | undefined +): Promise { + const deliveryLinks = await linkToCurrentContext(); + if (traceMode !== 'linked') return deliveryLinks; + const originLink = await linkToTraceCarrier( + isUsableTraceCarrier(incomingCarrier) ? incomingCarrier : undefined + ); + if ( + !originLink || + deliveryLinks?.some( + (link) => + link.context.traceId === originLink.context.traceId && + link.context.spanId === originLink.context.spanId + ) + ) { + return deliveryLinks; + } + return [...(deliveryLinks ?? []), originLink]; +} + /** * Serializes the current trace context into a format that can be passed through queues * @returns A record of strings representing the trace context @@ -189,6 +288,25 @@ export function linkToCurrentContext(): Promise<[api.Link] | undefined> { }); } +/** + * Builds a span link pointing at the span context embedded in a serialized + * trace carrier (e.g. the run-origin context flowing through queue + * messages). Returns `undefined` when OTEL is unavailable, the carrier is + * absent, or it does not contain a valid span context. + */ +export async function linkToTraceCarrier( + carrier: Record | undefined +): Promise { + if (!carrier) return; + const [context, otel] = await Promise.all([ + getSpanContextForTraceCarrier(carrier), + OtelApi.value, + ]); + if (!context || !otel) return; + if (!otel.trace.isSpanContextValid(context)) return; + return { context }; +} + // ============================================================ // Baggage Propagation Utilities // ============================================================ diff --git a/packages/core/src/telemetry/semantic-conventions.ts b/packages/core/src/telemetry/semantic-conventions.ts index 58cfd9cda0..ca3fea17da 100644 --- a/packages/core/src/telemetry/semantic-conventions.ts +++ b/packages/core/src/telemetry/semantic-conventions.ts @@ -92,6 +92,11 @@ export const WorkflowTracePropagated = SemanticConvention( 'workflow.trace.propagated' ); +/** Active trace-correlation mode for this invocation (linked or continuous) */ +export const WorkflowTraceMode = SemanticConvention<'linked' | 'continuous'>( + 'workflow.trace.mode' +); + /** Name of the error that caused workflow failure */ export const WorkflowErrorName = SemanticConvention( 'workflow.error.name' diff --git a/packages/utils/src/index.ts b/packages/utils/src/index.ts index 7128e83132..5d485cf047 100644 --- a/packages/utils/src/index.ts +++ b/packages/utils/src/index.ts @@ -5,6 +5,8 @@ export { parseClassName, parseStepName, parseWorkflowName, + stepDisplayName, + workflowDisplayName, } from './parse-name.js'; export { once, type PromiseWithResolvers, withResolvers } from './promise.js'; export { parseDurationToDate } from './time.js'; diff --git a/packages/utils/src/parse-name.test.ts b/packages/utils/src/parse-name.test.ts index e87a3aa4f1..0ff8f6f8f7 100644 --- a/packages/utils/src/parse-name.test.ts +++ b/packages/utils/src/parse-name.test.ts @@ -5,6 +5,8 @@ import { parseClassName, parseStepName, parseWorkflowName, + stepDisplayName, + workflowDisplayName, } from './parse-name'; describe('parseWorkflowName', () => { @@ -260,3 +262,68 @@ describe('formatStepName / formatWorkflowName', () => { ); }); }); + +describe('workflowDisplayName / stepDisplayName', () => { + test('returns the short name for raw machine names', () => { + expect( + workflowDisplayName('workflow//./src/jobs/order//processOrder') + ).toBe('processOrder'); + expect(stepDisplayName('step//./src/jobs/order//chargeCard')).toBe( + 'chargeCard' + ); + }); + + test('returns the short name for module-specifier names', () => { + expect(workflowDisplayName('workflow//@myorg/shared@1.2.3//sync')).toBe( + 'sync' + ); + }); + + test('recovers the function name from queue-sanitized names', () => { + // `workflow//./src/jobs/order//processOrder` after the queue's + // `replace(/[^A-Za-z0-9-_]/g, '-')` sanitization: + expect( + workflowDisplayName('workflow----src-jobs-order--processOrder') + ).toBe('processOrder'); + expect(stepDisplayName('step----src-jobs-order--chargeCard')).toBe( + 'chargeCard' + ); + }); + + test('sanitized nested functions use the leaf name', () => { + expect( + stepDisplayName('step----src-jobs-order--processOrder-innerStep') + ).toBe('innerStep'); + }); + + test('sanitized default exports map to the module short name', () => { + // `workflow//./src/jobs/order//default` after sanitization — mirror + // parseName's default-export handling so the same workflow doesn't + // display as `order` in workflow.start but `default` in + // workflow.execute. + expect(workflowDisplayName('workflow----src-jobs-order--default')).toBe( + 'order' + ); + expect(stepDisplayName('step----src-jobs-order--default')).toBe('order'); + // `__default` survives sanitization (underscores are preserved). + expect(workflowDisplayName('workflow----src-jobs-order--__default')).toBe( + 'order' + ); + }); + + test('sanitized names with `$` degrade to the trailing segment (best effort)', () => { + // `$` is a valid JS identifier character but sanitizes to `-`, so + // `process$Order` is indistinguishable from a nested function — the + // best-effort recovery returns just `Order`. Accepted limitation. + expect(stepDisplayName('step----src-jobs-order--process-Order')).toBe( + 'Order' + ); + }); + + test('falls back to the input when unrecognized', () => { + expect(workflowDisplayName('my-plain-name')).toBe('my-plain-name'); + expect(stepDisplayName('workflow--wrong-tag--fn')).toBe( + 'workflow--wrong-tag--fn' + ); + }); +}); diff --git a/packages/utils/src/parse-name.ts b/packages/utils/src/parse-name.ts index 9be071d93b..43f84429d1 100644 --- a/packages/utils/src/parse-name.ts +++ b/packages/utils/src/parse-name.ts @@ -106,6 +106,55 @@ export function formatWorkflowName(name: string): string { return formatParsedName(parseWorkflowName(name), name); } +/** + * Best-effort short display name for spans and UI labels. Accepts either the + * raw machine name (`workflow//./src/jobs/order//processOrder`) or the + * queue-sanitized form (`workflow----src-jobs-order--processOrder`, where + * every non-alphanumeric character was replaced with `-`) and returns just + * the function name (`processOrder`). Falls back to the input unchanged when + * neither form is recognized. + */ +export function workflowDisplayName(name: string): string { + return ( + parseWorkflowName(name)?.shortName ?? + shortNameFromSanitized('workflow', name) ?? + name + ); +} + +/** See {@link workflowDisplayName} — the step-name equivalent. */ +export function stepDisplayName(name: string): string { + return ( + parseStepName(name)?.shortName ?? + shortNameFromSanitized('step', name) ?? + name + ); +} + +function shortNameFromSanitized(tag: string, name: string): string | null { + if (!name.startsWith(`${tag}--`)) return null; + // The `//` separators became `--`, and within the function-name part any + // nested-function `/` became `-`. Function names are mostly dash-free, so + // the innermost name is the last dash-free segment. This is best-effort: + // `$` is a valid JS identifier character but is also sanitized to `-`, so + // a name like `process$Order` displays as `Order` — accepted limitation. + const segments = name.split('--').filter(Boolean); + const functionPart = segments.at(-1); + let shortName = functionPart?.split('-').filter(Boolean).at(-1) ?? ''; + // Mirror parseName's default-export handling: display default exports as + // the module's short name (the module part is the second-to-last `--` + // segment; its last `-` segment is the module short name), so the same + // workflow doesn't render as e.g. `order` from the raw name but + // `default` from the sanitized name. + if (['default', '__default'].includes(shortName)) { + const moduleShortName = segments.at(-2)?.split('-').filter(Boolean).at(-1); + if (moduleShortName && moduleShortName !== tag) { + shortName = moduleShortName; + } + } + return shortName || null; +} + function formatParsedName( parsed: { shortName: string; diff --git a/packages/world-vercel/package.json b/packages/world-vercel/package.json index be93526134..9aefe2854a 100644 --- a/packages/world-vercel/package.json +++ b/packages/world-vercel/package.json @@ -52,6 +52,9 @@ }, "devDependencies": { "@opentelemetry/api": "1.9.0", + "@opentelemetry/context-async-hooks": "1.30.1", + "@opentelemetry/core": "1.30.1", + "@opentelemetry/sdk-trace-base": "1.30.1", "@types/node": "catalog:", "@workflow/tsconfig": "workspace:*", "genversion": "3.2.0", diff --git a/packages/world-vercel/src/telemetry.ts b/packages/world-vercel/src/telemetry.ts index b80b4b34ea..6c254a7981 100644 --- a/packages/world-vercel/src/telemetry.ts +++ b/packages/world-vercel/src/telemetry.ts @@ -87,6 +87,27 @@ export async function getSpanKind( return otel.SpanKind[field]; } +/** + * Injects the active trace context into the given request headers using the + * registered propagator (typically W3C `traceparent`/`tracestate` plus + * `baggage`). Call inside an active client span so the receiving server can + * parent its spans to it. + * + * No-ops when `@opentelemetry/api` is unavailable or no SDK/propagator is + * registered (the default no-op propagator injects nothing). + */ +export async function injectTraceContextIntoHeaders( + headers: Headers +): Promise { + const otel = await getOtelApi(); + if (!otel) return; + const carrier: Record = {}; + otel.propagation.inject(otel.context.active(), carrier); + for (const [key, value] of Object.entries(carrier)) { + headers.set(key, value); + } +} + // Semantic conventions for World/Storage tracing // Standard OTEL conventions: https://opentelemetry.io/docs/specs/semconv/http/http-spans/ function SemanticConvention(...names: string[]) { diff --git a/packages/world-vercel/src/trace-propagation.test.ts b/packages/world-vercel/src/trace-propagation.test.ts new file mode 100644 index 0000000000..3e0bfea3dd --- /dev/null +++ b/packages/world-vercel/src/trace-propagation.test.ts @@ -0,0 +1,115 @@ +import { context, trace as otelTrace, propagation } from '@opentelemetry/api'; +import { AsyncLocalStorageContextManager } from '@opentelemetry/context-async-hooks'; +import { W3CTraceContextPropagator } from '@opentelemetry/core'; +import { + BasicTracerProvider, + InMemorySpanExporter, + SimpleSpanProcessor, +} from '@opentelemetry/sdk-trace-base'; +import { encode } from 'cbor-x'; +import { + afterAll, + afterEach, + beforeAll, + describe, + expect, + it, + vi, +} from 'vitest'; +import { z } from 'zod'; +import { injectTraceContextIntoHeaders } from './telemetry.js'; +import { makeRequest } from './utils.js'; + +vi.mock('@vercel/oidc', () => ({ + getVercelOidcToken: vi.fn().mockRejectedValue(new Error('no OIDC')), +})); + +const exporter = new InMemorySpanExporter(); +const provider = new BasicTracerProvider(); +const contextManager = new AsyncLocalStorageContextManager(); + +beforeAll(() => { + provider.addSpanProcessor(new SimpleSpanProcessor(exporter)); + contextManager.enable(); + context.setGlobalContextManager(contextManager); + propagation.setGlobalPropagator(new W3CTraceContextPropagator()); + otelTrace.setGlobalTracerProvider(provider); +}); + +afterAll(async () => { + await provider.shutdown(); + context.disable(); + propagation.disable(); + otelTrace.disable(); +}); + +afterEach(() => { + exporter.reset(); + vi.unstubAllGlobals(); +}); + +/** Minimal 2xx CBOR response, mirroring utils.test.ts. */ +function cborResponse(data: unknown) { + const bytes = encode(data); + return { + ok: true, + status: 200, + statusText: 'OK', + headers: { + get: (k: string) => + k.toLowerCase() === 'content-type' ? 'application/cbor' : null, + }, + arrayBuffer: async () => + bytes.buffer.slice(bytes.byteOffset, bytes.byteOffset + bytes.byteLength), + }; +} + +describe('injectTraceContextIntoHeaders', () => { + it('injects traceparent for the active span', async () => { + const tracer = otelTrace.getTracer('test'); + await tracer.startActiveSpan('client', async (span) => { + const headers = new Headers(); + await injectTraceContextIntoHeaders(headers); + expect(headers.get('traceparent')).toBe( + `00-${span.spanContext().traceId}-${span.spanContext().spanId}-01` + ); + span.end(); + }); + }); + + it('is a no-op when there is no active span context', async () => { + const headers = new Headers(); + await injectTraceContextIntoHeaders(headers); + expect(headers.get('traceparent')).toBeNull(); + }); +}); + +describe('makeRequest trace propagation', () => { + const schema = z.object({ value: z.string() }); + + it('sends traceparent on the outgoing workflow-server request, parented to the client span', async () => { + const fetchMock = vi.fn().mockResolvedValue(cborResponse({ value: 'ok' })); + vi.stubGlobal('fetch', fetchMock); + + const result = await makeRequest({ + endpoint: '/v3/runs/wrun_test/events', + options: { method: 'GET' }, + schema, + }); + expect(result).toEqual({ value: 'ok' }); + + const request = fetchMock.mock.calls[0][0] as Request; + const traceparent = request.headers.get('traceparent'); + expect(traceparent).toMatch(/^00-[0-9a-f]{32}-[0-9a-f]{16}-0[01]$/); + + // The injected context must be the `http GET` CLIENT span created by + // makeRequest, so the server's spans become its children. + const clientSpan = exporter + .getFinishedSpans() + .find((s) => s.name === 'http GET'); + expect(clientSpan).toBeDefined(); + expect(traceparent).toBe( + `00-${clientSpan?.spanContext().traceId}-${clientSpan?.spanContext().spanId}-01` + ); + }); +}); diff --git a/packages/world-vercel/src/utils.ts b/packages/world-vercel/src/utils.ts index afa5fe8086..0f1f6d789d 100644 --- a/packages/world-vercel/src/utils.ts +++ b/packages/world-vercel/src/utils.ts @@ -18,6 +18,7 @@ import { getSpanKind, HttpRequestMethod, HttpResponseStatusCode, + injectTraceContextIntoHeaders, PeerService, RpcService, RpcSystem, @@ -332,6 +333,13 @@ export async function makeRequest({ headers.set('Accept', 'application/cbor'); + // Explicitly propagate the active trace context (traceparent / + // tracestate / baggage) onto the outgoing request so workflow-server + // can parent its spans to this client span — without relying on the + // customer app having undici auto-instrumentation. No-ops when no + // OTEL SDK is registered. + await injectTraceContextIntoHeaders(headers); + // Encode body as CBOR if data is provided let body: Buffer | undefined; if (data !== undefined) { diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index ef054c8265..a24df0bc4a 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -610,6 +610,15 @@ importers: '@opentelemetry/api': specifier: 1.9.0 version: 1.9.0 + '@opentelemetry/context-async-hooks': + specifier: 1.30.1 + version: 1.30.1(@opentelemetry/api@1.9.0) + '@opentelemetry/core': + specifier: 1.30.1 + version: 1.30.1(@opentelemetry/api@1.9.0) + '@opentelemetry/sdk-trace-base': + specifier: 1.30.1 + version: 1.30.1(@opentelemetry/api@1.9.0) '@types/debug': specifier: 4.1.12 version: 4.1.12 @@ -1507,6 +1516,15 @@ importers: '@opentelemetry/api': specifier: 1.9.0 version: 1.9.0 + '@opentelemetry/context-async-hooks': + specifier: 1.30.1 + version: 1.30.1(@opentelemetry/api@1.9.0) + '@opentelemetry/core': + specifier: 1.30.1 + version: 1.30.1(@opentelemetry/api@1.9.0) + '@opentelemetry/sdk-trace-base': + specifier: 1.30.1 + version: 1.30.1(@opentelemetry/api@1.9.0) '@types/node': specifier: 'catalog:' version: 22.19.0 @@ -5025,6 +5043,12 @@ packages: resolution: {integrity: sha512-gLyJlPHPZYdAk1JENA9LeHejZe1Ti77/pTeFm/nMXmQH/HFZlcS/O2XJB+L8fkbrNSqhdtlvjBVjxwUYanNH5Q==} engines: {node: '>=8.0.0'} + '@opentelemetry/context-async-hooks@1.30.1': + resolution: {integrity: sha512-s5vvxXPVdjqS3kTLKMeBMvop9hbWkwzBpu+mUO2M7sZtlkyDJGwFe33wRKnbaYDo8ExRVBIIdwIGrqpxHuKttA==} + engines: {node: '>=14'} + peerDependencies: + '@opentelemetry/api': '>=1.0.0 <1.10.0' + '@opentelemetry/core@1.30.1': resolution: {integrity: sha512-OOCM2C/QIURhJMuKaekP3TRBxBKxG/TWWA0TL2J6nXUtDnuCtccy49LUJF8xPFXMX+0LMcxFpCo8M9cGY1W6rQ==} engines: {node: '>=14'} @@ -21179,6 +21203,10 @@ snapshots: '@opentelemetry/api@1.9.1': {} + '@opentelemetry/context-async-hooks@1.30.1(@opentelemetry/api@1.9.0)': + dependencies: + '@opentelemetry/api': 1.9.0 + '@opentelemetry/core@1.30.1(@opentelemetry/api@1.9.0)': dependencies: '@opentelemetry/api': 1.9.0