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/project-step-provenance.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
'@workflow/web-shared': patch
---

Allow trace callers to add product-specific attributes to event-derived step spans.
7 changes: 6 additions & 1 deletion packages/web-shared/src/components/trace-viewer.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import {
} from './sidebar/sidebar-data-context';
import { TraceViewerSkeleton } from './trace-viewer/components/trace-viewer-skeleton';
import { TraceViewer as TraceViewerComponent } from './trace-viewer/trace-viewer';
import type { GetStepAttributes } from './workflow-traces/trace-span-construction';

const TraceViewer = ({
run,
Expand All @@ -17,6 +18,7 @@ const TraceViewer = ({
hasMore,
isLoadingMore,
loading = false,
getStepAttributes,
}: {
run: WorkflowRun;
events: Event[];
Expand All @@ -25,6 +27,8 @@ const TraceViewer = ({
hasMore?: boolean;
isLoadingMore?: boolean;
loading?: boolean;
/** Adds product-specific attributes to event-derived step span data. */
getStepAttributes?: GetStepAttributes;
}) => {
const trace: TraceWithMeta | undefined = useMemo(() => {
if (!run?.runId) {
Expand All @@ -35,9 +39,10 @@ const TraceViewer = ({
// repeats with the whole log in hand.
return buildTrace(run, events, new Date(), {
isCompleteHistory: !hasMore,
getStepAttributes,
});
// eslint-disable-next-line react-hooks/exhaustive-deps -- `new Date()` is intentionally not a dep
}, [run, events, hasMore]);
}, [run, events, hasMore, getStepAttributes]);

// The sidebar shows one entity's slice of the log, so it takes the trace's
// answer rather than recomputing one from the slice.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -150,6 +150,10 @@ export function waitToSpan(
};
}

export type GetStepAttributes = (
events: Event[]

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

AI Review: This callback is meant to read product-specific enrichment fields, but its input is fixed to the base Event[]. The new test already has to cast to access externalAttemptId, and Front will need the same workaround for vercelId / computeInstanceId. Could we make GetStepAttributes, buildTrace, and TraceViewer generic over TEvent extends Event so consumers retain their enriched event type?

) => Record<string, unknown> | undefined;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

AI Review: The return type suggests any product-specific key can be surfaced, but the sidebar drops keys that are not registered in attributeToDisplayFn. The current Vercel fields happen to be registered, so this works for this use case; please either document that display registration is also required or expose a way for consumers to provide display metadata.


export const stepEventsToStepEntity = (
events: Event[]
): {
Expand Down Expand Up @@ -231,7 +235,11 @@ export const stepEventsToStepEntity = (
/**
* Converts step events to an OpenTelemetry Span
*/
export function stepToSpan(stepEvents: Event[], maxEndTime: Date): Span | null {
export function stepToSpan(
stepEvents: Event[],
maxEndTime: Date,
getStepAttributes?: GetStepAttributes
): Span | null {
const step = stepEventsToStepEntity(stepEvents);
if (!step) {
return null;
Expand All @@ -242,7 +250,11 @@ export function stepToSpan(stepEvents: Event[], maxEndTime: Date): Span | null {

const attributes = {
resource: 'step' as const,
data: step,
data: {
...getStepAttributes?.(stepEvents),
// Canonical event-derived fields cannot be overridden by extensions.
...step,
},
};

const resource = 'step';
Expand Down
43 changes: 42 additions & 1 deletion packages/web-shared/src/lib/trace-builder.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,11 @@ let nextId = 0;

function event(
eventType: EventType,
options: { correlationId?: string; at: number }
options: {
correlationId?: string;
at: number;
externalAttemptId?: string;
}
): Event {
nextId += 1;
return {
Expand All @@ -20,6 +24,7 @@ function event(
createdAt: new Date(BASE_TIME + options.at * 1000),
occurredAt: new Date(BASE_TIME + options.at * 1000),
eventData: eventType === 'step_created' ? { stepName: 'doWork' } : {},
externalAttemptId: options.externalAttemptId,
} as unknown as Event;
}

Expand All @@ -31,6 +36,42 @@ const run = {
} as unknown as WorkflowRun;

describe('buildTrace', () => {
it('adds caller-derived attributes to step span data', () => {
const events = [
event('run_created', { at: 0 }),
event('run_started', { at: 0 }),
event('step_created', { correlationId: 'step_a', at: 1 }),
event('step_started', {
correlationId: 'step_a',
at: 2,
externalAttemptId: 'attempt_first',
}),
event('step_retrying', { correlationId: 'step_a', at: 3 }),
event('step_started', {
correlationId: 'step_a',
at: 4,
externalAttemptId: 'attempt_latest',
}),
];

const trace = buildTrace(run, events, new Date(BASE_TIME + 5000), {
getStepAttributes(stepEvents) {
const latestStart = stepEvents
.slice()
.reverse()
.find((candidate) => candidate.eventType === 'step_started') as
| (Event & { externalAttemptId?: string })
| undefined;
return { externalAttemptId: latestStart?.externalAttemptId };
},
});
const stepSpan = trace.spans.find((span) => span.resource === 'step');

expect(stepSpan?.attributes.data).toMatchObject({

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

AI Review: Could we add a collision assertion here? The implementation intentionally spreads canonical step fields after extension attributes, so a test where the callback returns stepId or status would lock in the non-overridable guarantee documented in the implementation.

externalAttemptId: 'attempt_latest',
});
});

it('ends a step span on the terminal event the run acted on', () => {
const events = [
event('run_created', { at: 0 }),
Expand Down
17 changes: 13 additions & 4 deletions packages/web-shared/src/lib/trace-builder.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ import {
type WorkflowRun,
} from '@workflow/world';
import {
type GetStepAttributes,
getEventTimestamp,
hookToSpan,
runToSpan,
Expand Down Expand Up @@ -137,7 +138,8 @@ function buildSpans(
run: WorkflowRun,
groupedEvents: GroupedEvents,
now: Date,
latestKnownTime: Date
latestKnownTime: Date,
getStepAttributes?: GetStepAttributes
) {
// Active child spans cap at latestKnownTime so they don't extend into
// unknown territory. Even when the run is completed, we may not have loaded
Expand All @@ -146,7 +148,7 @@ function buildSpans(
const runMaxEnd = run.completedAt ?? now;

const stepSpans = Array.from(groupedEvents.eventsByStepId.values())
.map((events) => stepToSpan(events, childMaxEnd))
.map((events) => stepToSpan(events, childMaxEnd, getStepAttributes))
.filter((span): span is Span => span !== null);

const hookSpans = Array.from(groupedEvents.hookEvents.values())
Expand Down Expand Up @@ -208,7 +210,13 @@ export function buildTrace(
* from the only copy the caller was given, and dropping the wrong one moves
* a span. See {@link findDuplicateEventIds}.
*/
{ isCompleteHistory = false }: { isCompleteHistory?: boolean } = {}
{
isCompleteHistory = false,
getStepAttributes,
}: {
isCompleteHistory?: boolean;
getStepAttributes?: GetStepAttributes;
} = {}
): TraceWithMeta {
// Span geometry comes from what the run acted on. A repeat of a class the
// log already records is read past by every replay, and letting one through
Expand All @@ -232,7 +240,8 @@ export function buildTrace(
run,
groupedEvents,
now,
latestKnownTime
latestKnownTime,
getStepAttributes
);
const sortedCascadingSpans = cascadeSpans(runSpan, spans);

Expand Down
Loading