diff --git a/.changeset/retain-workflow-vm.md b/.changeset/retain-workflow-vm.md new file mode 100644 index 0000000000..7d3ad50076 --- /dev/null +++ b/.changeset/retain-workflow-vm.md @@ -0,0 +1,6 @@ +--- +'@workflow/core': patch +'workflow': patch +--- + +Retain workflow execution across inline steps within one invocation. diff --git a/packages/core/src/events-consumer.test.ts b/packages/core/src/events-consumer.test.ts index ecf828f728..3e320b16c7 100644 --- a/packages/core/src/events-consumer.test.ts +++ b/packages/core/src/events-consumer.test.ts @@ -45,6 +45,15 @@ describe('EventsConsumer', () => { expect(consumer.eventIndex).toBe(0); expect(consumer.callbacks).toEqual([]); }); + + it('should own its event array', () => { + const events = [createMockEvent()]; + const consumer = new EventsConsumer(events, defaultOptions); + + events.push(createMockEvent({ id: 'event-2' })); + + expect(consumer.events).toHaveLength(1); + }); }); describe('subscribe', () => { @@ -82,6 +91,24 @@ describe('EventsConsumer', () => { expect(callback).toHaveBeenCalledWith(event); expect(callback).toHaveBeenCalledTimes(1); }); + + it('should consume appended events', async () => { + const event = createMockEvent(); + const consumer = new EventsConsumer([], defaultOptions); + const callback = vi + .fn() + .mockImplementation((value: Event | null) => + value ? EventConsumerResult.Finished : EventConsumerResult.NotConsumed + ); + consumer.subscribe(callback); + await waitForNextTick(); + + consumer.append([event]); + await waitForNextTick(); + + expect(callback).toHaveBeenLastCalledWith(event); + expect(consumer.eventIndex).toBe(1); + }); }); describe('consume (implicit)', () => { diff --git a/packages/core/src/events-consumer.ts b/packages/core/src/events-consumer.ts index 4b7cf0742e..7b4afd6a9a 100644 --- a/packages/core/src/events-consumer.ts +++ b/packages/core/src/events-consumer.ts @@ -79,13 +79,18 @@ export class EventsConsumer { private unconsumedCheckVersion = 0; constructor(events: Event[], options: EventsConsumerOptions) { - this.events = events; + this.events = [...events]; this.eventIndex = 0; this.onConsumedEvent = options.onConsumedEvent; this.onUnconsumedEvent = options.onUnconsumedEvent; this.getPromiseQueue = options.getPromiseQueue; } + append(events: Event[]): void { + this.events.push(...events); + process.nextTick(this.consume); + } + /** * Registers a callback function to be called after an event has been consumed * by a different callback. The callback can return: diff --git a/packages/core/src/runtime.test.ts b/packages/core/src/runtime.test.ts index e687ae44de..dd7df78732 100644 --- a/packages/core/src/runtime.test.ts +++ b/packages/core/src/runtime.test.ts @@ -10,6 +10,7 @@ import { } from '@workflow/world'; import { ulid } from 'ulid'; import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import { runtimeLogger } from './logger.js'; import { registerStepFunction } from './private.js'; import { REPLAY_DIVERGENCE_MAX_RETRIES } from './runtime/constants.js'; import { setWorld } from './runtime/world.js'; @@ -701,6 +702,9 @@ describe('workflowEntrypoint replay guards', () => { }); it('replays attribute events before executing a step that loses the same race', async () => { + const debug = vi + .spyOn(runtimeLogger, 'debug') + .mockImplementation(() => undefined); const ops: Promise[] = []; const workflowRun: WorkflowRun = { runId: 'wrun_attribute_step_race', @@ -759,6 +763,11 @@ describe('workflowEntrypoint replay guards', () => { expect(createdEvents).not.toContainEqual( expect.objectContaining({ eventType: 'step_started' }) ); + const executionModes = debug.mock.calls + .filter(([message]) => message === 'Starting workflow execution') + .map(([, context]) => context?.executionMode); + expect(executionModes).toEqual(['replay', 'replay']); + debug.mockRestore(); }); it('fails the run when the World rejects an attr_set event as invalid', async () => { diff --git a/packages/core/src/runtime.ts b/packages/core/src/runtime.ts index 4076949432..c703ce7a44 100644 --- a/packages/core/src/runtime.ts +++ b/packages/core/src/runtime.ts @@ -90,7 +90,11 @@ import { } from './telemetry.js'; import { getErrorName, getErrorStack, normalizeUnknownError } from './types.js'; import { buildWorkflowSuspensionMessage } from './util.js'; -import { runWorkflow } from './workflow.js'; +import { + executeWorkflow, + type WorkflowExecutionResult, + type WorkflowSession, +} from './workflow.js'; export type { Event, WorkflowRun }; export { WorkflowSuspension } from './global.js'; @@ -306,6 +310,25 @@ function hasOpenHookOrWait(events: Event[]): boolean { return false; } +/** + * Retain only pure step boundaries with no out-of-band continuation source. + * Attributes require replay; hooks and waits can wake another invocation. + */ +function canRetainWorkflowSession( + suspension: WorkflowSuspension, + events: Event[] +): boolean { + return ( + suspension.stepCount > 0 && + suspension.hookCount === 0 && + suspension.waitCount === 0 && + suspension.attributeCount === 0 && + suspension.hookDisposedCount === 0 && + suspension.abortCount === 0 && + !hasOpenHookOrWait(events) + ); +} + /** * Creates a single route which handles workflow execution requests, * executing steps inline when possible to reduce function invocations @@ -1100,6 +1123,13 @@ export function workflowEntrypoint( encryptionKey ); + let workflowExecution: + | { readonly type: 'replay' } + | { + readonly type: 'retained'; + readonly session: WorkflowSession; + } = { type: 'replay' }; + // Main replay loop // biome-ignore lint/correctness/noConstantCondition: intentional loop while (true) { @@ -1414,37 +1444,76 @@ export function workflowEntrypoint( // point and the inline executeStep mutates eventsCursor. preInlineWriteCursor = eventsCursor; - // Replay workflow - runtimeLogger.debug('Starting workflow replay', { + let executionMode = workflowExecution.type; + runtimeLogger.debug('Starting workflow execution', { workflowRunId: runId, loopIteration, eventCount: events.length, + executionMode, }); replayStart = Date.now(); - // Start every missing decrypt/decompress operation before - // VM setup. Web Crypto work can overlap bundle evaluation; - // consumers still deserialize and resolve in event order. - const payloadPrewarm = replayPayloadCache.prewarm( - workflowRun, - events - ); - const result = await runWorkflow( - workflowCode, - workflowRun, - events, - encryptionKey, - replayPayloadCache, - // Turbo: the end-of-run drain inside runWorkflow commits - // fire-and-forget `*_created` events before the terminal - // `awaitRunReady()` below, so gate those writes on the - // backgrounded run_started too. Undefined outside turbo. - runReadyBarrier - ); - await payloadPrewarm; - runtimeLogger.debug('Workflow replay completed', { + let workflowResult: WorkflowExecutionResult = { + type: 'replay', + }; + if (workflowExecution.type === 'retained') { + workflowResult = await executeWorkflow({ + type: 'resume', + session: workflowExecution.session, + events, + }); + } + + if (workflowResult.type === 'replay') { + executionMode = 'replay'; + workflowExecution = { type: 'replay' }; + // Start every missing decrypt/decompress operation + // before VM setup. Web Crypto work can overlap bundle + // evaluation; consumers still deserialize and resolve + // in event order. + const payloadPrewarm = replayPayloadCache.prewarm( + workflowRun, + events + ); + workflowResult = await executeWorkflow({ + type: 'replay', + workflowCode, + workflowRun, + events, + encryptionKey, + replayPayloadCache, + // Turbo: the end-of-run drain inside workflow + // execution commits fire-and-forget `*_created` + // events before the terminal `awaitRunReady()` below. + runReadyBarrier, + }); + await payloadPrewarm; + } + + if (workflowResult.type === 'suspended') { + workflowExecution = canRetainWorkflowSession( + workflowResult.suspension, + events + ) + ? { + type: 'retained', + session: workflowResult.session, + } + : { type: 'replay' }; + throw workflowResult.suspension; + } + + if (workflowResult.type === 'replay') { + throw new Error( + 'Invariant violation: fresh workflow execution requested another replay' + ); + } + + const result = workflowResult.output; + runtimeLogger.debug('Workflow execution completed', { workflowRunId: runId, loopIteration, replayMs: Date.now() - replayStart, + executionMode, }); // Workflow completed. Send the snapshot but do NOT diff --git a/packages/core/src/telemetry.ts b/packages/core/src/telemetry.ts index da69d2756d..0c3a32f9d5 100644 --- a/packages/core/src/telemetry.ts +++ b/packages/core/src/telemetry.ts @@ -249,7 +249,10 @@ export async function trace( code: otel.SpanStatusCode.ERROR, message: (e as Error).message, }); - applyWorkflowSuspensionToSpan(e, otel, span); + if (WorkflowSuspension.is(e)) { + span.setStatus({ code: otel.SpanStatusCode.OK }); + applyWorkflowSuspensionToSpan(e, span); + } throw e; } finally { span.end(); @@ -275,19 +278,12 @@ export async function recordElapsedSpan( } /** - * Applies workflow suspension attributes to the given span if the error is a WorkflowSuspension - * which is technically not an error, but an algebraic effect indicating suspension. + * Applies the workflow suspension algebraic effect to an active span. */ -function applyWorkflowSuspensionToSpan( - error: unknown, - otel: typeof api, +export function applyWorkflowSuspensionToSpan( + error: WorkflowSuspension, span: api.Span -) { - if (!error || !WorkflowSuspension.is(error)) { - return; - } - - span.setStatus({ code: otel.SpanStatusCode.OK }); +): void { span.setAttributes({ ...Attr.WorkflowSuspensionState('suspended'), ...Attr.WorkflowSuspensionStepCount(error.stepCount), diff --git a/packages/core/src/telemetry/semantic-conventions.ts b/packages/core/src/telemetry/semantic-conventions.ts index bcc6d1dabf..0793b26d87 100644 --- a/packages/core/src/telemetry/semantic-conventions.ts +++ b/packages/core/src/telemetry/semantic-conventions.ts @@ -77,6 +77,11 @@ export const WorkflowEventsCount = SemanticConvention( 'workflow.events.count' ); +/** Whether workflow execution starts with replay or resumes a retained VM */ +export const WorkflowExecutionMode = SemanticConvention<'replay' | 'retained'>( + 'workflow.execution.mode' +); + /** Number of arguments passed to the workflow */ export const WorkflowArgumentsCount = SemanticConvention( 'workflow.arguments.count' diff --git a/packages/core/src/workflow-session-telemetry.test.ts b/packages/core/src/workflow-session-telemetry.test.ts new file mode 100644 index 0000000000..49058a4585 --- /dev/null +++ b/packages/core/src/workflow-session-telemetry.test.ts @@ -0,0 +1,69 @@ +import { trace as otelTrace } from '@opentelemetry/api'; +import { + BasicTracerProvider, + InMemorySpanExporter, + SimpleSpanProcessor, +} from '@opentelemetry/sdk-trace-base'; +import type { WorkflowRun } from '@workflow/world'; +import { afterAll, afterEach, beforeAll, describe, expect, it } from 'vitest'; +import { ReplayPayloadCache } from './replay-payload-cache.js'; +import { dehydrateWorkflowArguments } from './serialization.js'; +import { executeWorkflow } from './workflow.js'; + +const exporter = new InMemorySpanExporter(); +const provider = new BasicTracerProvider(); + +beforeAll(() => { + provider.addSpanProcessor(new SimpleSpanProcessor(exporter)); + otelTrace.setGlobalTracerProvider(provider); +}); + +afterAll(async () => { + await provider.shutdown(); + otelTrace.disable(); +}); + +afterEach(() => { + exporter.reset(); +}); + +describe('retained workflow telemetry', () => { + it('records suspension attributes on workflow.run spans', async () => { + const runId = 'wrun_retained_telemetry'; + const workflowRun: WorkflowRun = { + 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', + }; + const workflowCode = ` + const step = globalThis[Symbol.for("WORKFLOW_USE_STEP")]("step"); + async function workflow() { await step(); } + globalThis.__private_workflows = new Map([["workflow", workflow]]);`; + + const result = await executeWorkflow({ + type: 'replay', + workflowCode, + workflowRun, + events: [], + encryptionKey: undefined, + replayPayloadCache: new ReplayPayloadCache(undefined), + }); + expect(result.type).toBe('suspended'); + + const span = exporter + .getFinishedSpans() + .find((candidate) => candidate.name === 'workflow.run workflow'); + expect(span?.attributes).toMatchObject({ + 'workflow.execution.mode': 'replay', + 'workflow.suspension.state': 'suspended', + 'workflow.suspension.step_count': 1, + 'workflow.suspension.hook_count': 0, + 'workflow.suspension.wait_count': 0, + }); + }); +}); diff --git a/packages/core/src/workflow.test.ts b/packages/core/src/workflow.test.ts index fc21423518..66cc295e27 100644 --- a/packages/core/src/workflow.test.ts +++ b/packages/core/src/workflow.test.ts @@ -6,6 +6,7 @@ import { monotonicFactory } from 'ulid'; import { afterEach, assert, describe, expect, it, vi } from 'vitest'; import { DEFERRED_CHECK_DELAY_MS } from './events-consumer.js'; import type { WorkflowSuspension } from './global.js'; +import { ReplayPayloadCache } from './replay-payload-cache.js'; import { setWorld } from './runtime/world.js'; import { dehydrateStepReturnValue, @@ -13,7 +14,7 @@ import { hydrateWorkflowReturnValue, } from './serialization.js'; import { createContext } from './vm/index.js'; -import { runWorkflow } from './workflow.js'; +import { executeWorkflow, runWorkflow } from './workflow.js'; // No encryption key = encryption disabled const noEncryptionKey = undefined; @@ -221,6 +222,345 @@ describe('runWorkflow', () => { ).toEqual(3); }); + it('should retain workflow execution across sequential step suspensions', async () => { + const ops: Promise[] = []; + const workflowRun: WorkflowRun = { + runId: 'wrun_retained', + workflowName: 'workflow', + status: 'running', + input: await dehydrateWorkflowArguments( + [], + 'wrun_retained', + noEncryptionKey, + ops + ), + 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', + }; + const workflowCode = ` + const add = globalThis[Symbol.for("WORKFLOW_USE_STEP")]("add"); + async function workflow() { + console.log("retained:entered"); + const first = await add(1, 2); + console.log("retained:continued"); + const second = await add(first, 3); + return second; + } + ${getWorkflowTransformCode('workflow')}`; + const log = vi.spyOn(console, 'log').mockImplementation(() => undefined); + const events: Event[] = []; + let eventNumber = 0; + const appendStepEvents = async ( + correlationId: string, + result: number + ): Promise => { + const createdAt = new Date(`2024-01-01T00:00:0${eventNumber + 1}.000Z`); + events.push( + { + eventId: `event-${eventNumber++}`, + runId: workflowRun.runId, + eventType: 'step_created', + correlationId, + eventData: { stepName: 'add' }, + createdAt, + }, + { + eventId: `event-${eventNumber++}`, + runId: workflowRun.runId, + eventType: 'step_started', + correlationId, + eventData: { stepName: 'add' }, + createdAt, + }, + { + eventId: `event-${eventNumber++}`, + runId: workflowRun.runId, + eventType: 'step_completed', + correlationId, + eventData: { + stepName: 'add', + result: await dehydrateStepReturnValue( + result, + workflowRun.runId, + noEncryptionKey, + ops + ), + }, + createdAt, + } + ); + }; + + const first = await executeWorkflow({ + type: 'replay', + workflowCode, + workflowRun, + events, + encryptionKey: noEncryptionKey, + replayPayloadCache: new ReplayPayloadCache(noEncryptionKey), + }); + assert(first.type === 'suspended'); + const firstStep = first.suspension.steps[0]; + assert(firstStep?.type === 'step'); + + await appendStepEvents(firstStep.correlationId, 3); + + const second = await executeWorkflow({ + type: 'resume', + session: first.session, + events, + }); + assert(second.type === 'suspended'); + expect(second.session).toBe(first.session); + const secondStep = second.suspension.steps[0]; + assert(secondStep?.type === 'step'); + + await appendStepEvents(secondStep.correlationId, 6); + + const completed = await executeWorkflow({ + type: 'resume', + session: second.session, + events, + }); + assert(completed.type === 'completed'); + + expect( + await hydrateWorkflowReturnValue( + completed.output as any, + workflowRun.runId, + noEncryptionKey, + ops + ) + ).toBe(6); + expect( + log.mock.calls + .map(([message]) => message) + .filter((message) => String(message).startsWith('retained:')) + ).toEqual(['retained:entered', 'retained:continued']); + log.mockRestore(); + }); + + it('returns a workflow that completes while its retained session is suspended', async () => { + const ops: Promise[] = []; + const workflowRun: WorkflowRun = { + runId: 'wrun_retained_async_completion', + workflowName: 'workflow', + status: 'running', + input: await dehydrateWorkflowArguments( + [], + 'wrun_retained_async_completion', + noEncryptionKey, + ops + ), + 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', + }; + const workflowCode = ` + const step = globalThis[Symbol.for("WORKFLOW_USE_STEP")]("step"); + async function digestRepeatedly() { + const bytes = new Uint8Array(8 * 1024 * 1024); + for (let index = 0; index < 8; index++) { + await crypto.subtle.digest("SHA-256", bytes); + } + } + async function workflow() { + await Promise.race([step(), digestRepeatedly()]); + console.log("retained:digest-completed"); + return "digest completed"; + } + ${getWorkflowTransformCode('workflow')}`; + const events: Event[] = [ + { + eventId: 'event-run-created', + runId: workflowRun.runId, + eventType: 'run_created', + createdAt: new Date('2024-01-01T00:00:00.000Z'), + }, + ]; + const log = vi.spyOn(console, 'log').mockImplementation(() => undefined); + + const suspended = await executeWorkflow({ + type: 'replay', + workflowCode, + workflowRun, + events, + encryptionKey: noEncryptionKey, + replayPayloadCache: new ReplayPayloadCache(noEncryptionKey), + }); + assert(suspended.type === 'suspended'); + const pendingStep = suspended.suspension.steps[0]; + assert(pendingStep?.type === 'step'); + + const rewrittenSession = await executeWorkflow({ + type: 'replay', + workflowCode, + workflowRun, + events, + encryptionKey: noEncryptionKey, + replayPayloadCache: new ReplayPayloadCache(noEncryptionKey), + }); + assert(rewrittenSession.type === 'suspended'); + + await vi.waitFor( + () => { + expect(log).toHaveBeenCalledTimes(2); + }, + { timeout: 10_000, interval: 10 } + ); + const rewrittenEvents = [ + { ...events[0], eventId: 'rewritten-event' } as Event, + { ...events[0], eventId: 'new-event' } as Event, + ]; + expect( + await executeWorkflow({ + type: 'resume', + session: rewrittenSession.session, + events: rewrittenEvents, + }) + ).toEqual({ type: 'replay' }); + + const createdAt = new Date('2024-01-01T00:00:01.000Z'); + events.push( + { + eventId: 'event-step-created', + runId: workflowRun.runId, + eventType: 'step_created', + correlationId: pendingStep.correlationId, + eventData: { stepName: pendingStep.stepName }, + createdAt, + }, + { + eventId: 'event-step-started', + runId: workflowRun.runId, + eventType: 'step_started', + correlationId: pendingStep.correlationId, + eventData: { stepName: pendingStep.stepName }, + createdAt, + }, + { + eventId: 'event-step-completed', + runId: workflowRun.runId, + eventType: 'step_completed', + correlationId: pendingStep.correlationId, + eventData: { + stepName: pendingStep.stepName, + result: await dehydrateStepReturnValue( + undefined, + workflowRun.runId, + noEncryptionKey, + ops + ), + }, + createdAt, + } + ); + + const completed = await executeWorkflow({ + type: 'resume', + session: suspended.session, + events, + }); + assert(completed.type === 'completed'); + expect( + await executeWorkflow({ + type: 'resume', + session: rewrittenSession.session, + events, + }) + ).toEqual({ type: 'replay' }); + expect( + await hydrateWorkflowReturnValue( + completed.output as any, + workflowRun.runId, + noEncryptionKey, + ops + ) + ).toBe('digest completed'); + log.mockRestore(); + }); + + it('does not drain operations from a discarded retained session', async () => { + const ops: Promise[] = []; + const workflowRun: WorkflowRun = { + runId: 'wrun_retained_discarded', + workflowName: 'workflow', + status: 'running', + input: await dehydrateWorkflowArguments( + [], + 'wrun_retained_discarded', + noEncryptionKey, + ops + ), + 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', + }; + const workflowCode = ` + const step = globalThis[Symbol.for("WORKFLOW_USE_STEP")]("step"); + const setAttributes = globalThis[Symbol.for("WORKFLOW_SET_ATTRIBUTES")]; + async function digestRepeatedly() { + const bytes = new Uint8Array(8 * 1024 * 1024); + for (let index = 0; index < 8; index++) { + await crypto.subtle.digest("SHA-256", bytes); + } + } + async function workflow() { + await Promise.race([step(), digestRepeatedly()]); + void setAttributes([{ key: "stale", value: true }]); + console.log("retained:discarded-completed"); + return "discarded"; + } + ${getWorkflowTransformCode('workflow')}`; + const create = vi.fn(); + const log = vi.spyOn(console, 'log').mockImplementation(() => undefined); + setWorld({ + specVersion: SPEC_VERSION_CURRENT, + events: { create }, + streams: { write: vi.fn(), close: vi.fn() }, + } as any); + + const suspended = await executeWorkflow({ + type: 'replay', + workflowCode, + workflowRun, + events: [], + encryptionKey: noEncryptionKey, + replayPayloadCache: new ReplayPayloadCache(noEncryptionKey), + }); + assert(suspended.type === 'suspended'); + + await vi.waitFor( + () => { + expect(log).toHaveBeenCalledWith('retained:discarded-completed'); + }, + { timeout: 10_000, interval: 10 } + ); + await new Promise((resolve) => setTimeout(resolve, 250)); + + expect( + await executeWorkflow({ + type: 'resume', + session: suspended.session, + events: [ + { + eventId: 'event-run-created', + runId: workflowRun.runId, + eventType: 'run_created', + createdAt: new Date('2024-01-01T00:00:00.000Z'), + }, + ], + }) + ).toEqual({ type: 'replay' }); + expect(create).not.toHaveBeenCalled(); + setWorld(undefined); + log.mockRestore(); + }); + it('regenerates step correlation IDs independent of startedAt (turbo replay-stability)', async () => { // Turbo's first delivery synthesizes `startedAt` from the local clock, // while later (non-turbo) deliveries load the server-canonical `startedAt`. diff --git a/packages/core/src/workflow.ts b/packages/core/src/workflow.ts index 4595f6a17c..1fe6668470 100644 --- a/packages/core/src/workflow.ts +++ b/packages/core/src/workflow.ts @@ -4,7 +4,11 @@ import { WorkflowNotRegisteredError, WorkflowRuntimeError, } from '@workflow/errors'; -import { createWorkflowBaseUrl, withResolvers } from '@workflow/utils'; +import { + createWorkflowBaseUrl, + type PromiseWithResolvers, + withResolvers, +} from '@workflow/utils'; import { parseWorkflowName } from '@workflow/utils/parse-name'; import type { Event, WorkflowRun } from '@workflow/world'; import { SPEC_VERSION_SUPPORTS_COMPRESSION } from '@workflow/world'; @@ -36,7 +40,7 @@ import { WORKFLOW_USE_STEP, } from './symbols.js'; import * as Attribute from './telemetry/semantic-conventions.js'; -import { trace } from './telemetry.js'; +import { applyWorkflowSuspensionToSpan, trace } from './telemetry.js'; import { getWorkflowRunStreamId } from './util.js'; import { createContext } from './vm/index.js'; import { runCachedWorkflowScript } from './vm/script-cache.js'; @@ -133,6 +137,142 @@ async function drainPendingQueueItems( } } +interface WorkflowCompletion { + readonly output: unknown; + readonly resultType: string; +} + +type WorkflowSessionState = + | { + readonly type: 'running'; + readonly interruption: PromiseWithResolvers; + } + | { readonly type: 'suspended'; readonly suspension: WorkflowSuspension } + | { readonly type: 'failed'; readonly error: Error } + | { readonly type: 'replay' } + | { readonly type: 'completed' }; + +function isSameSuspensionBoundary( + previous: WorkflowSuspension, + next: WorkflowSuspension +): boolean { + return ( + previous.stepCount === next.stepCount && + previous.hookCount === next.hookCount && + previous.waitCount === next.waitCount && + previous.attributeCount === next.attributeCount && + previous.hookDisposedCount === next.hookDisposedCount && + previous.abortCount === next.abortCount && + previous.steps.length === next.steps.length && + previous.steps.every((item, index) => item === next.steps[index]) + ); +} + +function updateSuspendedSession( + suspension: WorkflowSuspension, + error: Error +): WorkflowSessionState { + if (!WorkflowSuspension.is(error)) return { type: 'failed', error }; + return isSameSuspensionBoundary(suspension, error) + ? { type: 'suspended', suspension } + : { type: 'replay' }; +} + +export interface WorkflowSession { + readonly workflowRun: WorkflowRun; + readonly argumentCount: number; + resume(events: Event[]): WorkflowSessionResumeResult; +} + +type WorkflowSessionResumeResult = + | { + readonly type: 'resumed'; + readonly execution: Promise; + } + | { readonly type: 'replay' }; + +interface WorkflowSessionOptions { + readonly workflowCode: string; + readonly workflowRun: WorkflowRun; + readonly events: Event[]; + readonly encryptionKey: CryptoKey | undefined; + readonly replayPayloadCache: ReplayPayloadCache; + readonly runReadyBarrier?: Promise; +} + +type WorkflowExecutionRequest = + | ({ readonly type: 'replay' } & WorkflowSessionOptions) + | { + readonly type: 'resume'; + readonly session: WorkflowSession; + readonly events: Event[]; + }; + +export type WorkflowExecutionResult = + | { readonly type: 'replay' } + | { readonly type: 'completed'; readonly output: unknown } + | { + readonly type: 'suspended'; + readonly suspension: WorkflowSuspension; + readonly session: WorkflowSession; + }; + +export async function executeWorkflow( + request: WorkflowExecutionRequest +): Promise { + const workflowRun = + request.type === 'replay' + ? request.workflowRun + : request.session.workflowRun; + const mode = request.type === 'replay' ? 'replay' : 'retained'; + + return trace(`workflow.run ${workflowRun.workflowName}`, async (span) => { + span?.setAttributes({ + ...Attribute.WorkflowName(workflowRun.workflowName), + ...Attribute.WorkflowRunId(workflowRun.runId), + ...Attribute.WorkflowRunStatus(workflowRun.status), + ...Attribute.WorkflowEventsCount(request.events.length), + ...Attribute.WorkflowExecutionMode(mode), + }); + + let session: WorkflowSession; + let execution: Promise; + switch (request.type) { + case 'replay': { + const started = await createWorkflowSession(request); + session = started.session; + execution = started.execution; + break; + } + case 'resume': { + session = request.session; + const resumed = session.resume(request.events); + if (resumed.type === 'replay') return resumed; + execution = resumed.execution; + break; + } + } + + span?.setAttributes({ + ...Attribute.WorkflowArgumentsCount(session.argumentCount), + }); + + try { + const completed = await execution; + span?.setAttributes({ + ...Attribute.WorkflowResultType(completed.resultType), + }); + return { type: 'completed', output: completed.output }; + } catch (error) { + if (WorkflowSuspension.is(error)) { + if (span) applyWorkflowSuspensionToSpan(error, span); + return { type: 'suspended', suspension: error, session }; + } + throw error; + } + }); +} + export async function runWorkflow( workflowCode: string, workflowRun: WorkflowRun, @@ -140,8 +280,7 @@ export async function runWorkflow( encryptionKey: CryptoKey | undefined, /** * Optional per-run cache for replay payload preparation and immutable final - * values. Owned by the inline replay loop so it survives fresh VM contexts - * created by successive iterations of this invocation. + * values. Owned by the inline execution loop for this invocation. */ replayPayloadCache: ReplayPayloadCache = new ReplayPayloadCache( encryptionKey @@ -154,14 +293,36 @@ export async function runWorkflow( */ runReadyBarrier?: Promise ): Promise { - return trace(`workflow.run ${workflowRun.workflowName}`, async (span) => { - span?.setAttributes({ - ...Attribute.WorkflowName(workflowRun.workflowName), - ...Attribute.WorkflowRunId(workflowRun.runId), - ...Attribute.WorkflowRunStatus(workflowRun.status), - ...Attribute.WorkflowEventsCount(events.length), - }); + const result = await executeWorkflow({ + type: 'replay', + workflowCode, + workflowRun, + events, + encryptionKey, + replayPayloadCache, + runReadyBarrier, + }); + if (result.type === 'replay') { + throw new WorkflowRuntimeError( + `Fresh workflow "${workflowRun.runId}" unexpectedly requested replay` + ); + } + if (result.type === 'suspended') throw result.suspension; + return result.output; +} +function createWorkflowSession({ + workflowCode, + workflowRun, + events, + encryptionKey, + replayPayloadCache, + runReadyBarrier, +}: WorkflowSessionOptions): Promise<{ + session: WorkflowSession; + execution: Promise; +}> { + return (async () => { const startedAt = workflowRun.startedAt; if (!startedAt) { throw new WorkflowRuntimeError( @@ -209,7 +370,32 @@ export async function runWorkflow( fixedTimestamp, }); - const workflowDiscontinuation = withResolvers(); + const initialInterruption = withResolvers(); + let state: WorkflowSessionState = { + type: 'running', + interruption: initialInterruption, + }; + + const onWorkflowError = (error: Error): void => { + switch (state.type) { + case 'running': { + const { interruption } = state; + state = WorkflowSuspension.is(error) + ? { type: 'suspended', suspension: error } + : { type: 'failed', error }; + interruption.reject(error); + return; + } + case 'suspended': + state = updateSuspendedSession(state.suspension, error); + return; + case 'failed': + case 'replay': + case 'completed': + return; + } + state satisfies never; + }; const ulid = monotonicFactory(() => vmGlobalThis.Math.random()); const generateNanoid = nanoid.customRandom(nanoid.urlAlphabet, 21, (size) => @@ -226,7 +412,7 @@ export async function runWorkflow( updateTimestamp(+event.createdAt); }, onUnconsumedEvent: (event) => { - workflowDiscontinuation.reject( + onWorkflowError( new ReplayDivergenceError( `Replay could not consume event: eventType=${event.eventType}, correlationId=${event.correlationId}, eventId=${event.eventId}.`, { eventId: event.eventId } @@ -240,7 +426,7 @@ export async function runWorkflow( runId: workflowRun.runId, encryptionKey, globalThis: vmGlobalThis, - onWorkflowError: workflowDiscontinuation.reject, + onWorkflowError, eventsConsumer, // Correlation IDs (step_/wait_/hook_) are derived from `generateUlid`, so // the time prefix fed to `ulid()` MUST be replay-stable across every @@ -873,48 +1059,16 @@ export async function runWorkflow( ); await workflowContext.promiseQueue; - span?.setAttributes({ - ...Attribute.WorkflowArgumentsCount(args.length), - }); - - // Invoke user workflow - try { - const result = await Promise.race([ - workflowFn(...args), - workflowDiscontinuation.promise, - ]); - - const dehydrated = await dehydrateWorkflowReturnValue( - result, - workflowRun.runId, - encryptionKey, - vmGlobalThis, - false, - // Gate payload compression on the run's specVersion: only runs - // marked as possibly containing compressed payloads (spec >= 5) - // get gzip data. - (workflowRun.specVersion ?? 0) >= SPEC_VERSION_SUPPORTS_COMPRESSION - ); + const workflowExecution = (async (): Promise => { + return await workflowFn(...args); + })(); - span?.setAttributes({ - ...Attribute.WorkflowResultType(typeof result), - }); - - await drainPendingQueueItems( - workflowRun.runId, - workflowContext.invocationsQueue, - vmGlobalThis, - workflowRun, - 'completed', - runReadyBarrier - ); - - return dehydrated; - } catch (err) { + const failWorkflow = async (error: unknown): Promise => { + state = { type: 'completed' }; // Control-flow signals are handled by the runtime and do not mean the // workflow has terminally failed. - if (WorkflowSuspension.is(err) || ReplayDivergenceError.is(err)) { - throw err; + if (WorkflowSuspension.is(error) || ReplayDivergenceError.is(error)) { + throw error; } await drainPendingQueueItems( @@ -926,7 +1080,97 @@ export async function runWorkflow( runReadyBarrier ); - throw err; - } - }); + throw error; + }; + + const completeWorkflow = async ( + result: unknown + ): Promise => { + state = { type: 'completed' }; + try { + const output = await dehydrateWorkflowReturnValue( + result, + workflowRun.runId, + encryptionKey, + vmGlobalThis, + false, + // Gate payload compression on the run's specVersion: only runs + // marked as possibly containing compressed payloads (spec >= 5) + // get gzip data. + (workflowRun.specVersion ?? 0) >= SPEC_VERSION_SUPPORTS_COMPRESSION + ); + + await drainPendingQueueItems( + workflowRun.runId, + workflowContext.invocationsQueue, + vmGlobalThis, + workflowRun, + 'completed', + runReadyBarrier + ); + + return { output, resultType: typeof result }; + } catch (error) { + return failWorkflow(error); + } + }; + + const waitForExecution = async ( + interruption: PromiseWithResolvers + ): Promise => { + let result: unknown; + try { + result = await Promise.race([workflowExecution, interruption.promise]); + } catch (error) { + if (state.type === 'suspended' && error === state.suspension) { + throw error; + } + return failWorkflow(error); + } + return completeWorkflow(result); + }; + + const session: WorkflowSession = { + workflowRun, + argumentCount: args.length, + resume(nextEvents) { + switch (state.type) { + case 'suspended': { + const consumedEvents = eventsConsumer.events; + const isStrictExtension = + nextEvents.length > consumedEvents.length && + consumedEvents.every( + (event, index) => event.eventId === nextEvents[index]?.eventId + ); + if (!isStrictExtension) { + state = { type: 'replay' }; + return { type: 'replay' }; + } + const interruption = withResolvers(); + state = { type: 'running', interruption }; + eventsConsumer.append(nextEvents.slice(consumedEvents.length)); + return { + type: 'resumed', + execution: waitForExecution(interruption), + }; + } + case 'failed': + return { type: 'resumed', execution: failWorkflow(state.error) }; + case 'replay': + return { type: 'replay' }; + case 'completed': + case 'running': + throw new WorkflowRuntimeError( + `Cannot resume ${state.type} workflow "${workflowRun.runId}"` + ); + } + state satisfies never; + }, + }; + + return { + session, + execution: waitForExecution(initialInterruption), + }; + })(); }