Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/nest-linked-invocations-under-delivery.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
'@workflow/core': patch
---

Refine `WORKFLOW_TRACE_MODE=linked` (the default) so each queue-delivered `workflow.execute` / `step.execute` span nests under its local delivery context instead of starting a new trace root.
40 changes: 18 additions & 22 deletions packages/core/src/runtime-trace-mode.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -221,28 +221,29 @@ describe('getWorkflowTraceMode', () => {
});

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 () => {
it('linked (default): nests under the delivery context with a link to the run-origin context', 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(
// Child of the local delivery (flow-route) span — same trace, so one
// invocation is a single bounded trace rather than a new root.
expect(workflowSpan?.parentSpanId).toBe(deliverySpan.spanContext().spanId);
expect(workflowSpan?.spanContext().traceId).toBe(
deliverySpan.spanContext().traceId
);

// Links to BOTH the delivery context and the run-origin context.
expect(workflowSpan?.links).toHaveLength(2);
expect(linkTraceIds(workflowSpan)).toContain(
// Single link to the run-origin context (NOT a parent) — connecting this
// bounded invocation trace back to where the run was started.
expect(workflowSpan?.links).toHaveLength(1);
expect(linkTraceIds(workflowSpan)).toContain(ORIGIN_TRACE_ID);
// The delivery context is the parent now, so it is not also a link.
expect(linkTraceIds(workflowSpan)).not.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);
Expand All @@ -260,29 +261,24 @@ describe('workflowEntrypoint trace modes', () => {
});

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
);
// Still nested under the delivery context...
expect(workflowSpan?.parentSpanId).toBe(deliverySpan.spanContext().spanId);
// ...but no run-origin link is derived from `{}`.
expect(workflowSpan?.links ?? []).toHaveLength(0);
// 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 () => {
it('linked: without an incoming carrier, nests under the delivery context with no links', 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?.parentSpanId).toBe(deliverySpan.spanContext().spanId);
expect(workflowSpan?.links ?? []).toHaveLength(0);
expect(workflowSpan?.attributes['workflow.trace.propagated']).toBe(false);
});

Expand Down
21 changes: 11 additions & 10 deletions packages/core/src/runtime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -393,11 +393,13 @@ export function workflowEntrypoint(
}

// --- 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.
// 'linked' (default): the workflow.execute span below stays a CHILD
// of the local delivery (flow-route) context, so one invocation —
// route handler, workflow replay, inline steps, event writes — is a
// single bounded trace. The run-origin context travels as a span
// LINK (not a parent), and re-enqueues forward the original carrier
// unchanged, so a (potentially hours-long) run is never stitched
// into one giant trace across invocations.
// 'continuous': legacy behavior — the restored run-origin context
// becomes the parent of this invocation's spans.
const traceMode = getWorkflowTraceMode();
Expand Down Expand Up @@ -442,8 +444,9 @@ export function workflowEntrypoint(

// 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.
// withTraceContext a passthrough, so the workflow.execute span below
// stays a child of the local delivery (flow-route) context and the
// run-origin travels as a span link instead.
const parentTraceCarrier =
traceMode === 'continuous' ? traceContext : undefined;
// Queue-delivered invocation: CONSUMER kind, matching the
Expand All @@ -456,9 +459,7 @@ export function workflowEntrypoint(
const world = await getWorld();
return trace(
`workflow.execute ${workflowDisplayName(workflowName)}`,
traceMode === 'linked'
? { kind: spanKind, links: spanLinks, root: true }
: { kind: spanKind, links: spanLinks },
{ kind: spanKind, links: spanLinks },
async (span) => {
span?.setAttributes({
...Attribute.WorkflowName(workflowName),
Expand Down
14 changes: 7 additions & 7 deletions packages/core/src/runtime/step-handler.ts
Original file line number Diff line number Diff line change
Expand Up @@ -190,13 +190,15 @@ function createStepHandler(namespace?: string) {
return;
}

// Span links to the incoming delivery context and (in linked mode)
// the run-origin context from the trace carrier.
// In linked mode the only span link is to the run-origin context; the
// step.execute span stays a child of the local delivery (flow-route)
// context. In continuous mode the link points at the delivery context.
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).
// only — in linked mode the run-origin is NOT restored as the parent,
// so withTraceContext is a passthrough and the step.execute span stays
// a child of the local delivery context, linked to the run origin).
const parentTraceCarrier =
traceMode === 'continuous' ? traceContext : undefined;
return await withTraceContext(parentTraceCarrier, async () => {
Expand All @@ -222,9 +224,7 @@ function createStepHandler(namespace?: string) {

return trace(
`step.execute ${stepDisplayName(stepName)}`,
traceMode === 'linked'
? { kind: spanKind, links: spanLinks, root: true }
: { kind: spanKind, links: spanLinks },
{ kind: spanKind, links: spanLinks },
async (span) => {
span?.setAttributes({
...Attribute.StepName(stepName),
Expand Down
30 changes: 11 additions & 19 deletions packages/core/src/telemetry.ts
Original file line number Diff line number Diff line change
Expand Up @@ -79,33 +79,25 @@ export function getNextTraceCarrier(
* 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).
* - In `linked` mode the invocation span is a CHILD of the local delivery
* (flow-route) context, so the only link is to the run-origin context
* from the message's trace carrier — connecting this bounded per-invocation
* trace back to where the run was started. The run-origin context is a
* link, never a parent, and re-enqueues forward the original carrier
* unchanged, so the whole run is never stitched into one giant trace.
* The link is skipped when the carrier is absent, empty, or invalid.
* - In `continuous` mode the run-origin context is restored as the parent
* instead, and the link points at the incoming delivery context.
*/
export async function buildInvocationSpanLinks(
traceMode: WorkflowTraceMode,
incomingCarrier: Record<string, string> | undefined
): Promise<api.Link[] | undefined> {
const deliveryLinks = await linkToCurrentContext();
if (traceMode !== 'linked') return deliveryLinks;
if (traceMode !== 'linked') return linkToCurrentContext();
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];
return originLink ? [originLink] : undefined;
}

/**
Expand Down
Loading