From da14615835283e38a7a19a48d6a3c89469bda7d6 Mon Sep 17 00:00:00 2001 From: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> Date: Thu, 13 Aug 2026 20:15:08 -0700 Subject: [PATCH 1/9] Prepare replay payloads during event streaming --- .../prepare-streamed-replay-payloads.md | 5 + .../core/src/replay-payload-cache.test.ts | 95 ++++- packages/core/src/replay-payload-cache.ts | 385 +++++++++++------- packages/core/src/runtime.ts | 135 ++++-- packages/core/src/runtime/helpers.test.ts | 12 + packages/core/src/runtime/helpers.ts | 10 +- .../runtime/quickjs-partial-preload.test.ts | 1 + .../core/src/step-delivery-ordering.test.ts | 2 +- packages/core/src/step.ts | 65 +-- packages/core/src/workflow/hook.ts | 59 +-- 10 files changed, 532 insertions(+), 237 deletions(-) create mode 100644 .changeset/prepare-streamed-replay-payloads.md diff --git a/.changeset/prepare-streamed-replay-payloads.md b/.changeset/prepare-streamed-replay-payloads.md new file mode 100644 index 0000000000..9678641541 --- /dev/null +++ b/.changeset/prepare-streamed-replay-payloads.md @@ -0,0 +1,5 @@ +--- +"@workflow/core": patch +--- + +Prepare replay payloads as event frames arrive and reuse immutable primitive values across fresh workflow VMs. diff --git a/packages/core/src/replay-payload-cache.test.ts b/packages/core/src/replay-payload-cache.test.ts index f09ec0f4ca..55adb4b196 100644 --- a/packages/core/src/replay-payload-cache.test.ts +++ b/packages/core/src/replay-payload-cache.test.ts @@ -111,6 +111,47 @@ describe('ReplayPayloadCache', () => { allSettled.mockRestore(); }); + it('prepares streamed events synchronously inside the decoder callback', async () => { + const payload = new Uint8Array([1]); + const order: string[] = []; + const preparer = vi.fn((value) => { + order.push('prepare'); + return { data: value }; + }); + const cache = new ReplayPayloadCache(undefined, preparer); + const [event] = makeEvents([payload]); + + const preparation = cache.observeEvent(event, () => order.push('start')); + expect(preparation).toBeDefined(); + expect(preparer).toHaveBeenCalledOnce(); + expect(order).toEqual(['start', 'prepare']); + + // Re-observing a cache hit does not move the preparation-span boundary. + cache.observeEvent(event, () => order.push('cached-start')); + expect(order).toEqual(['start', 'prepare']); + + await expect(preparation).resolves.toEqual({ data: payload }); + }); + + it('prepares queued stream events as soon as the run key resolves', async () => { + const payload = new Uint8Array([1]); + const preparer = vi.fn((value) => ({ data: value })); + let resolveKey!: (key: undefined) => void; + const key = new Promise((resolve) => { + resolveKey = resolve; + }); + const cache = new ReplayPayloadCache(key, preparer); + const [event] = makeEvents([payload]); + + const preparation = cache.observeEvent(event); + expect(preparation).toBeDefined(); + expect(preparer).not.toHaveBeenCalled(); + + resolveKey(undefined); + await expect(preparation).resolves.toEqual({ data: payload }); + expect(preparer).toHaveBeenCalledOnce(); + }); + it('caches real decrypt/decompress output but revives fresh objects', async () => { const key = await importKey(new Uint8Array(32).fill(7)); const serialized = await dehydrateStepReturnValue( @@ -126,6 +167,10 @@ describe('ReplayPayloadCache', () => { const preparer = vi.fn(prepareReplayPayload); const cache = new ReplayPayloadCache(key, preparer); + const directPreparation = prepareReplayPayload(serialized, key); + expect(directPreparation).not.toBeInstanceOf(Promise); + await directPreparation; + const prepared = await cache.prepareEventPayload( 'evnt_encrypted', 'result', @@ -196,12 +241,34 @@ describe('ReplayPayloadCache', () => { const cache = new ReplayPayloadCache(undefined); const hydrate = vi.fn().mockResolvedValue(value); - expect(await cache.getStepResult('evnt_result', hydrate)).toBe(value); - expect(await cache.getStepResult('evnt_result', hydrate)).toBe(value); + expect( + await cache.getPrimitiveValue('evnt_result', 'result', hydrate) + ).toBe(value); + expect( + await cache.getPrimitiveValue('evnt_result', 'result', hydrate) + ).toBe(value); expect(hydrate).toHaveBeenCalledOnce(); } }); + it('isolates primitive values by event payload field', async () => { + const cache = new ReplayPayloadCache(undefined); + const result = vi.fn().mockResolvedValue('result'); + const error = vi.fn().mockResolvedValue('error'); + + await expect( + cache.getPrimitiveValue('evnt_shared', 'result', result) + ).resolves.toBe('result'); + await expect( + cache.getPrimitiveValue('evnt_shared', 'error', error) + ).resolves.toBe('error'); + await expect( + cache.getPrimitiveValue('evnt_shared', 'result', result) + ).resolves.toBe('result'); + expect(result).toHaveBeenCalledOnce(); + expect(error).toHaveBeenCalledOnce(); + }); + it('rehydrates mutable and oversized step results', async () => { const oversized = 'x'.repeat(4097); for (const value of [{ count: 0 }, oversized]) { @@ -212,8 +279,16 @@ describe('ReplayPayloadCache', () => { typeof value === 'object' ? { ...value } : value ); - const first = await cache.getStepResult('evnt_result', hydrate); - const second = await cache.getStepResult('evnt_result', hydrate); + const first = await cache.getPrimitiveValue( + 'evnt_result', + 'result', + hydrate + ); + const second = await cache.getPrimitiveValue( + 'evnt_result', + 'result', + hydrate + ); expect(hydrate).toHaveBeenCalledTimes(2); if (typeof value === 'object') expect(second).not.toBe(first); } @@ -226,12 +301,12 @@ describe('ReplayPayloadCache', () => { .mockRejectedValueOnce(new Error('boom')) .mockResolvedValueOnce('ok'); - await expect(cache.getStepResult('evnt_result', hydrate)).rejects.toThrow( - 'boom' - ); - await expect(cache.getStepResult('evnt_result', hydrate)).resolves.toBe( - 'ok' - ); + await expect( + cache.getPrimitiveValue('evnt_result', 'result', hydrate) + ).rejects.toThrow('boom'); + await expect( + cache.getPrimitiveValue('evnt_result', 'result', hydrate) + ).resolves.toBe('ok'); expect(hydrate).toHaveBeenCalledTimes(2); }); }); diff --git a/packages/core/src/replay-payload-cache.ts b/packages/core/src/replay-payload-cache.ts index 084c6fd3c9..f9c6baba18 100644 --- a/packages/core/src/replay-payload-cache.ts +++ b/packages/core/src/replay-payload-cache.ts @@ -5,204 +5,311 @@ import { prepareReplayPayload, } from './serialization/replay.js'; -const MAX_MEMOIZED_PRIMITIVE_LENGTH = 4096; type ReplayPayloadField = 'result' | 'error' | 'payload'; -function isMemoizablePrimitive(value: unknown): boolean { - if (value === null) return true; - const type = typeof value; - if (type === 'object' || type === 'function') return false; - if (type === 'string') { - return (value as string).length <= MAX_MEMOIZED_PRIMITIVE_LENGTH; - } - if (type === 'bigint') { - return (value as bigint).toString().length <= MAX_MEMOIZED_PRIMITIVE_LENGTH; - } - return true; +type KeyState = + | { state: 'pending'; promise: Promise } + | { state: 'ready'; value: DecryptionKey | undefined } + | { state: 'failed'; error: unknown }; + +type Preparation = + | { state: 'waiting'; value: Uint8Array } + | { state: 'ready'; value: PreparedReplayPayload } + | { state: 'pending'; promise: Promise } + | { state: 'failed'; error: unknown }; + +function isPrimitive(value: unknown): boolean { + return value === null || !['object', 'function'].includes(typeof value); } /** * Invocation-scoped cache for replay payload hydration. * - * A workflow invocation may replay the same event log through several fresh - * VMs. This cache keeps the VM-independent decrypt/decompress result across - * those replays. Deserialization still runs against each VM's globals so every - * replay receives fresh object graphs and correctly revived Workflow objects. - * - * Successful prepared plaintext remains resident for the invocation lifetime. - * Its memory cost is the sum of decrypted and decompressed payload sizes, but - * it never crosses workflow runs or queue deliveries. + * The cache retains VM-independent decrypt/decompress output across fresh VMs. + * Deserialization still runs against each VM's globals so object graphs and + * Workflow objects remain realm-local. Primitive final values are safe to + * share and skip that repeated deserialization entirely. */ export class ReplayPayloadCache { - private readonly preparedPayloads = new Map< - string, + private readonly preparations = new Map(); + private readonly pendingPreparations = new Set< Promise >(); - private readonly primitiveStepResults = new Map(); + private readonly primitiveValues = new Map(); private nextUnscannedEventIndex = 0; - constructor( - private readonly encryptionKey: DecryptionKey | undefined, + private constructor( + private key: KeyState, private readonly preparer: typeof prepareReplayPayload = prepareReplayPayload - ) {} + ) { + if (key.state === 'pending') { + void key.promise.then( + (value) => this.resolveKey(value), + (error) => this.rejectKey(error) + ); + } + } + + static unencrypted( + preparer: typeof prepareReplayPayload = prepareReplayPayload + ): ReplayPayloadCache { + return new ReplayPayloadCache( + { state: 'ready', value: undefined }, + preparer + ); + } + + static withKey( + key: DecryptionKey, + preparer: typeof prepareReplayPayload = prepareReplayPayload + ): ReplayPayloadCache { + return new ReplayPayloadCache({ state: 'ready', value: key }, preparer); + } + + static waitingForKey( + key: Promise, + preparer: typeof prepareReplayPayload = prepareReplayPayload + ): ReplayPayloadCache { + return new ReplayPayloadCache({ state: 'pending', promise: key }, preparer); + } + + /** Start preparing an event as soon as its frame has been decoded. */ + observeEvent(event: Event, onPreparationStart?: () => void): void { + switch (event.eventType) { + case 'run_created': + this.start( + this.workflowInputKey(event.runId), + event.eventData.input, + onPreparationStart + ); + break; + case 'run_started': + this.start( + this.workflowInputKey(event.runId), + event.eventData?.input, + onPreparationStart + ); + break; + case 'step_completed': + this.start( + this.eventPayloadKey(event.eventId, 'result'), + event.eventData?.result, + onPreparationStart + ); + break; + case 'step_failed': + this.start( + this.eventPayloadKey(event.eventId, 'error'), + event.eventData?.error, + onPreparationStart + ); + break; + case 'hook_received': + this.start( + this.eventPayloadKey(event.eventId, 'payload'), + event.eventData?.payload, + onPreparationStart + ); + } + } /** - * Start every missing binary preparation before workflow execution. Failures - * are intentionally retained: the ordered event consumer must observe the - * original rejection before that entry becomes retryable. + * Start every preparation not already observed from the event stream. + * Returns a Promise only when a sealed or portable codec is still running. */ - async prewarm(workflowRun: WorkflowRun, events: Event[]): Promise { - const preparations: Promise[] = []; - const start = (cacheKey: string, value: unknown): void => { - // Legacy flattened values may be mutated by devalue's unflatten and are - // therefore prepared only by their eventual consumer, never cached. - if (!(value instanceof Uint8Array)) return; - - // Each replay scans the full event log, so awaiting cached promises here - // would add O(N^2) promise reactions over an N-step invocation. Only wait - // for preparations first discovered by this prewarm pass. - if (this.preparedPayloads.has(cacheKey)) return; - preparations.push(this.ensurePreparation(cacheKey, value)); - }; - - start(this.workflowInputKey(workflowRun.runId), workflowRun.input); - // This cache is scoped to one invocation. Incremental loads and write - // response deltas only ever append, so the scanned length locates the - // events added since the previous replay. A reload that can insert events - // BELOW that length — a stale-snapshot restart replacing the log with a - // corrected one — must call `resetScan()` first, or the inserted events are - // never scanned. Prepared entries stay valid across that: they are keyed by - // event id, not by position. + prewarm(workflowRun: WorkflowRun, events: Event[]): void | Promise { + this.start( + this.workflowInputKey(workflowRun.runId), + workflowRun.input + ); for ( let index = this.nextUnscannedEventIndex; index < events.length; index++ ) { - const event = events[index]; - switch (event.eventType) { - case 'step_completed': - start( - this.eventPayloadKey(event.eventId, 'result'), - event.eventData?.result - ); - break; - case 'step_failed': - start( - this.eventPayloadKey(event.eventId, 'error'), - event.eventData?.error - ); - break; - case 'hook_received': - start( - this.eventPayloadKey(event.eventId, 'payload'), - event.eventData?.payload - ); - break; - } + this.observeEvent(events[index]); } this.nextUnscannedEventIndex = events.length; - - // Prewarming is speculative and must not fail replay before the matching - // event is consumed. allSettled also attaches rejection handlers eagerly. - await Promise.allSettled(preparations); + return this.waitForPending(); } - /** - * Forget how much of the event log has been scanned, so the next - * {@link prewarm} walks it from the start again. - * - * Required before a replay whose event log was reloaded rather than extended: - * a corrected log inserts the events the previous load was missing, which - * shifts every later position, so a positional resume would skip exactly the - * events the reload was for. Already-prepared payloads are kept — they are - * keyed by event id, so re-scanning re-observes them for free. - */ + /** A corrected reload may insert events before the previous scan position. */ resetScan(): void { this.nextUnscannedEventIndex = 0; } - /** Return the workflow input after shared host-side preparation. */ prepareWorkflowInput( workflowRun: WorkflowRun - ): Promise { - return this.consumePreparation( + ): PreparedReplayPayload | Promise { + return this.consume( this.workflowInputKey(workflowRun.runId), workflowRun.input ); } - /** - * Return an event payload after shared host-side preparation. A rejected - * preparation is evicted only after this ordered consumer requests it, so a - * later replay can retry without hiding the original failure. - */ prepareEventPayload( eventId: string, field: ReplayPayloadField, value: unknown - ): Promise { - return this.consumePreparation(this.eventPayloadKey(eventId, field), value); + ): PreparedReplayPayload | Promise { + return this.consume(this.eventPayloadKey(eventId, field), value); } - /** - * Reuse final step values only when sharing them across VMs is unobservable. - * Objects and large strings/bigints always run `hydrate` again, producing a - * fresh VM-specific value from the separately cached prepared payload. - */ - async getStepResult( + getEventValue( eventId: string, - hydrate: () => Promise - ): Promise { - if (this.primitiveStepResults.has(eventId)) { - return this.primitiveStepResults.get(eventId); + field: ReplayPayloadField, + serializedValue: unknown, + hydrate: (prepared: PreparedReplayPayload) => unknown | Promise + ): unknown | Promise { + const cacheKey = this.eventPayloadKey(eventId, field); + if (this.primitiveValues.has(cacheKey)) { + return this.primitiveValues.get(cacheKey); } - const value = await hydrate(); - if (isMemoizablePrimitive(value)) { - this.primitiveStepResults.set(eventId, value); - } + const prepared = this.prepareEventPayload( + eventId, + field, + serializedValue + ); + const hydrateAndCache = (payload: PreparedReplayPayload) => { + const hydrated = hydrate(payload); + return hydrated instanceof Promise + ? hydrated.then((value) => this.cachePrimitive(cacheKey, value)) + : this.cachePrimitive(cacheKey, hydrated); + }; + return prepared instanceof Promise + ? prepared.then(hydrateAndCache) + : hydrateAndCache(prepared); + } + + private cachePrimitive(cacheKey: string, value: unknown): unknown { + if (isPrimitive(value)) this.primitiveValues.set(cacheKey, value); return value; } - /** - * Consumer-facing lookup. Binary payloads share preparation; legacy values - * bypass the cache because their flattened representation may be mutated. - */ - private consumePreparation( + private start( cacheKey: string, - value: unknown - ): Promise { - if (!(value instanceof Uint8Array)) { - return Promise.resolve({ legacy: value }); + value: unknown, + onPreparationStart?: () => void + ): void { + if (!(value instanceof Uint8Array) || this.preparations.has(cacheKey)) { + return; + } + + onPreparationStart?.(); + switch (this.key.state) { + case 'pending': + this.preparations.set(cacheKey, { state: 'waiting', value }); + break; + case 'ready': + this.prepare(cacheKey, value, this.key.value); + break; + case 'failed': + this.preparations.set(cacheKey, { + state: 'failed', + error: this.key.error, + }); + break; } + } - const preparation = this.ensurePreparation(cacheKey, value); - void preparation.catch(() => { - if (this.preparedPayloads.get(cacheKey) === preparation) { - this.preparedPayloads.delete(cacheKey); + private prepare( + cacheKey: string, + value: Uint8Array, + key: DecryptionKey | undefined + ): void { + try { + const result = this.preparer(value, key); + if (!(result instanceof Promise)) { + this.preparations.set(cacheKey, { state: 'ready', value: result }); + return; } - }); - return preparation; + + this.preparations.set(cacheKey, { state: 'pending', promise: result }); + this.pendingPreparations.add(result); + void result.then( + (prepared) => { + this.pendingPreparations.delete(result); + const current = this.preparations.get(cacheKey); + if (current?.state === 'pending' && current.promise === result) { + this.preparations.set(cacheKey, { + state: 'ready', + value: prepared, + }); + } + }, + (error) => { + this.pendingPreparations.delete(result); + const current = this.preparations.get(cacheKey); + if (current?.state === 'pending' && current.promise === result) { + this.preparations.set(cacheKey, { state: 'failed', error }); + } + } + ); + } catch (error) { + this.preparations.set(cacheKey, { state: 'failed', error }); + } } - /** Start preparation once and share the exact in-flight promise. */ - private ensurePreparation( + private consume( cacheKey: string, - value: Uint8Array - ): Promise { - const cached = this.preparedPayloads.get(cacheKey); - if (cached) return cached; + value: unknown + ): PreparedReplayPayload | Promise { + if (!(value instanceof Uint8Array)) return { legacy: value }; + + this.start(cacheKey, value); + const preparation = this.preparations.get(cacheKey); + if (!preparation) { + throw new Error(`Replay payload preparation was not started: ${cacheKey}`); + } + + switch (preparation.state) { + case 'ready': + return preparation.value; + case 'pending': + return preparation.promise; + case 'failed': + this.preparations.delete(cacheKey); + throw preparation.error; + case 'waiting': + if (this.key.state !== 'pending') { + throw new Error(`Replay payload key was not resolved: ${cacheKey}`); + } + return this.key.promise.then(() => this.consume(cacheKey, value)); + } + } + + private waitForPending(): void | Promise { + if ( + this.key.state === 'pending' && + [...this.preparations.values()].some( + (preparation) => preparation.state === 'waiting' + ) + ) { + return this.key.promise.then(() => this.waitForPending()); + } + if (this.pendingPreparations.size === 0) return; + return Promise.allSettled([...this.pendingPreparations]).then(() => {}); + } - const preparation = this.runPreparation(value); - this.preparedPayloads.set(cacheKey, preparation); - return preparation; + private resolveKey(value: DecryptionKey | undefined): void { + if (this.key.state !== 'pending') return; + this.key = { state: 'ready', value }; + for (const [cacheKey, preparation] of this.preparations) { + if (preparation.state === 'waiting') { + this.prepare(cacheKey, preparation.value, value); + } + } } - /** Normalize synchronous and asynchronous preparers to one promise contract. */ - private async runPreparation( - value: Uint8Array - ): Promise { - return this.preparer(value, this.encryptionKey); + private rejectKey(error: unknown): void { + if (this.key.state !== 'pending') return; + this.key = { state: 'failed', error }; + for (const [cacheKey, preparation] of this.preparations) { + if (preparation.state === 'waiting') { + this.preparations.set(cacheKey, { state: 'failed', error }); + } + } } private workflowInputKey(runId: string): string { diff --git a/packages/core/src/runtime.ts b/packages/core/src/runtime.ts index 0c9bc472d2..c997becbb7 100644 --- a/packages/core/src/runtime.ts +++ b/packages/core/src/runtime.ts @@ -14,7 +14,7 @@ import { WorkflowRuntimeError, WorkflowWorldError, } from '@workflow/errors'; -import { once, setWorkflowBasePath } from '@workflow/utils'; +import { once, setWorkflowBasePath, withResolvers } from '@workflow/utils'; import { parseWorkflowName, workflowDisplayName, @@ -540,6 +540,32 @@ function appendEventLog(log: LoadedEventLog, appended: LoadedEventLog): void { log.cursor = appended.cursor ?? log.cursor; } +function replayEventDeploymentId(event: Event): string | undefined { + if (event.eventType === 'run_created' || event.eventType === 'run_started') { + return event.eventData?.deploymentId; + } + return undefined; +} + +function createReplayEventObserver({ + runId, + cache, + resolveKey, +}: { + runId: string; + cache: ReplayPayloadCache; + resolveKey: ( + runOrId: WorkflowRun | string, + context?: Record + ) => void; +}): (event: Event) => void { + return (event) => { + const deploymentId = replayEventDeploymentId(event); + if (deploymentId) resolveKey(runId, { deploymentId }); + cache.observeEvent(event); + }; +} + /** * The whole retention predicate: keep the session only for a pure step * boundary (every suspension item is a step — any other item type, present @@ -1013,6 +1039,43 @@ export function workflowEntrypoint( // describe the next load exactly. let eventLog: ReplayEventLog = { type: 'loadAll' }; + // Resolve the run-scoped key as soon as the deployment id is + // known. On a normal replay that is the streamed run_created + // frame; resilient start and turbo already carry it in the + // queue payload. This starts key resolution and payload + // preparation before the remainder of the event log arrives, + // without guessing the key for a cross-deployment run. + const { + promise: replayKeySource, + resolve: resolveReplayKeySource, + } = withResolvers<{ + runOrId: WorkflowRun | string; + context?: Record; + }>(); + const resolveReplayKey = ( + runOrId: WorkflowRun | string, + context?: Record + ): void => { + resolveReplayKeySource({ runOrId, context }); + }; + const encryptionKeyPromise = replayKeySource.then( + ({ runOrId, context }) => + memoizeEncryptionKey(world, runOrId, context)() + ); + const replayPayloadCache = new ReplayPayloadCache( + encryptionKeyPromise + ); + const observeReplayEvent = createReplayEventObserver({ + runId, + cache: replayPayloadCache, + resolveKey: resolveReplayKey, + }); + if (runInput?.deploymentId) { + resolveReplayKey(runId, { + deploymentId: runInput.deploymentId, + }); + } + // Shared state: set by either the background step path // or the run_started setup below. let workflowRun: WorkflowRun | undefined; @@ -1983,6 +2046,7 @@ export function workflowEntrypoint( resumeId: hookResumeInput.resumeId, resumePayloadDigest: hookResumeInput.payloadDigest, preloadEvents: true, + onEvent: observeReplayEvent, } ); hookEnsured = true; @@ -2257,6 +2321,7 @@ export function workflowEntrypoint( }); const result = await createEvent(runStartedEvent, { requestId, + onEvent: observeReplayEvent, }); workflowRun = result.run; maxEventsLimit = clampMaxEvents(result.maxEvents); @@ -2541,31 +2606,18 @@ export function workflowEntrypoint( // do we fall back to reloading the complete log. if (eventLog.type !== 'loadAll' && ensuredEvent) { insertEventByEventId(eventLog.events, ensuredEvent); + observeReplayEvent(ensuredEvent); } else { eventLog = { type: 'loadAll' }; } } // end else (re-ensure needed) } - // Resolve the encryption key for this run's deployment. - // Used eagerly here since both workflow execution (input - // hydration / hook payload decryption) and the run_failed - // dehydrate path below need it. Memoized accessor: first - // call triggers the actual fetch / HKDF derivation, - // subsequent calls await the cached promise. - const getEncryptionKey = memoizeEncryptionKey( - world, - workflowRun - ); - const encryptionKey = await getEncryptionKey(); - - // Invocation-scoped cache of VM-independent prepared payloads - // and immutable final values. It survives the fresh workflow - // VM created by each inline replay, but never crosses runs or - // queue deliveries. - const replayPayloadCache = new ReplayPayloadCache( - encryptionKey - ); + // Worlds that do not implement streamed observation still + // resolve from the materialized run. This is also the final + // cross-deployment-safe source of truth. + resolveReplayKey(workflowRun); + const encryptionKey = await encryptionKeyPromise; // The live VM parked at the previous boundary, when the // retention decision kept it. null → this iteration cold- @@ -2660,7 +2712,11 @@ export function workflowEntrypoint( if (eventLog.type === 'loadAfter') { appendEventLog( eventLog, - await loadWorkflowRunEvents(runId, eventLog.cursor) + await loadWorkflowRunEvents( + runId, + eventLog.cursor, + observeReplayEvent + ) ); eventLog = { ...eventLog, type: 'ready' }; } @@ -2716,7 +2772,8 @@ export function workflowEntrypoint( runId, eventLog.type === 'loadAfter' ? eventLog.cursor - : undefined + : undefined, + observeReplayEvent ); if (eventLog.type === 'loadAfter') { appendEventLog(eventLog, page); @@ -2836,7 +2893,8 @@ export function workflowEntrypoint( if (eventLog.cursor) { const page = await loadWorkflowRunEvents( runId, - eventLog.cursor + eventLog.cursor, + observeReplayEvent ); const completedWaitIdsAfterCursor = new Set( page.events @@ -2854,13 +2912,21 @@ export function workflowEntrypoint( appendEventLog(eventLog, page); } else { eventLog = { - ...(await loadWorkflowRunEvents(runId)), + ...(await loadWorkflowRunEvents( + runId, + undefined, + observeReplayEvent + )), type: 'ready', }; } } else { eventLog = { - ...(await loadWorkflowRunEvents(runId)), + ...(await loadWorkflowRunEvents( + runId, + undefined, + observeReplayEvent + )), type: 'ready', }; } @@ -2943,15 +3009,20 @@ export function workflowEntrypoint( if (resumeTracking) { resumeTracking.replayStartedAtMs ??= replayStart; } - // Start every missing decrypt/decompress operation up - // front (already-prepared payloads are skipped). Web - // Crypto work overlaps VM setup on the replay path and - // the appended events' consumption on the resume path; - // consumers still deserialize and resolve in event order. + // Finish scheduling every missing decrypt/decompress + // operation (stream-observed payloads are already in + // flight). Preparation overlaps VM setup on replay and + // appended-event consumption on resume; consumers still + // deserialize and resolve in event order. + const replayEvents = eventLog.events; const payloadPrewarm = replayPayloadCache.prewarm( workflowRun, - eventLog.events + replayEvents ); + // Consumers await their own prepared payloads in event + // order. Do not delay a suspension on speculative work + // for payloads this replay never touched. + void payloadPrewarm.catch(() => {}); let workflowResult: WorkflowResumeResult = retainedSession ? await resumeWorkflow(retainedSession, eventLog.events) : { type: 'replay' }; @@ -2971,8 +3042,6 @@ export function workflowEntrypoint( worldCapabilities: world.capabilities, }); } - await payloadPrewarm; - if (workflowResult.type === 'suspended') { // Park the live session; the suspension catch below // makes the one retention decision — keep it for the diff --git a/packages/core/src/runtime/helpers.test.ts b/packages/core/src/runtime/helpers.test.ts index 0ce474ef07..0854f5d381 100644 --- a/packages/core/src/runtime/helpers.test.ts +++ b/packages/core/src/runtime/helpers.test.ts @@ -971,6 +971,18 @@ describe('memoizeEncryptionKey', () => { expect(spy).toHaveBeenCalledTimes(1); }); + it('passes deployment context when resolving before the run is materialized', async () => { + const spy = vi.fn().mockResolvedValue(MATERIAL); + const getKey = memoizeEncryptionKey(worldWithKey(spy), 'wrun_1', { + deploymentId: 'dpl_streamed', + }); + + await getKey(); + expect(spy).toHaveBeenCalledWith('wrun_1', { + deploymentId: 'dpl_streamed', + }); + }); + it('resolves undefined when encryption is not configured', async () => { const getKey = memoizeEncryptionKey(worldWithKey(undefined), 'wrun_1'); await expect(getKey()).resolves.toBeUndefined(); diff --git a/packages/core/src/runtime/helpers.ts b/packages/core/src/runtime/helpers.ts index f24944478f..c82cdee697 100644 --- a/packages/core/src/runtime/helpers.ts +++ b/packages/core/src/runtime/helpers.ts @@ -8,6 +8,7 @@ import type { CreateEventRequest, Event, EventResult, + EventStreamObserver, HealthCheckPayload, ValidQueueName, WorkflowRun, @@ -594,7 +595,8 @@ function shouldRetryWithoutEventCursor( */ export async function loadWorkflowRunEvents( runId: string, - afterCursor?: string + afterCursor?: string, + onEvent?: EventStreamObserver ): Promise { const incremental = afterCursor !== undefined; return trace( @@ -630,6 +632,7 @@ export async function loadWorkflowRunEvents( sortOrder: 'asc', cursor: requestedCursor ?? undefined, }, + onEvent, }); } catch (error) { if ( @@ -1246,7 +1249,8 @@ export function getQueueOverhead(message: { requestedAt?: Date }) { */ export function memoizeEncryptionKey( world: World, - runOrId: WorkflowRun | string + runOrId: WorkflowRun | string, + context?: Record ): () => Promise { let cached: Promise | undefined; return () => { @@ -1257,7 +1261,7 @@ export function memoizeEncryptionKey( // here so TypeScript picks the right overload for each shape. const rawKey = typeof runOrId === 'string' - ? await world.getEncryptionKeyForRun?.(runOrId) + ? await world.getEncryptionKeyForRun?.(runOrId, context) : await world.getEncryptionKeyForRun?.(runOrId); // Resolve the *full* capability, not just the symmetric key: a run // reading its own event log may encounter sealed (`encp`) payloads diff --git a/packages/core/src/runtime/quickjs-partial-preload.test.ts b/packages/core/src/runtime/quickjs-partial-preload.test.ts index d9d3edf528..33f7b1df80 100644 --- a/packages/core/src/runtime/quickjs-partial-preload.test.ts +++ b/packages/core/src/runtime/quickjs-partial-preload.test.ts @@ -116,6 +116,7 @@ describe('QuickJS partial run_started preload', () => { expect(listEvents).toHaveBeenCalledWith({ runId, pagination: { sortOrder: 'asc', cursor: preloadCursor }, + onEvent: expect.any(Function), }); expect(runWorkflowWithQuickJS).toHaveBeenCalledWith( expect.objectContaining({ diff --git a/packages/core/src/step-delivery-ordering.test.ts b/packages/core/src/step-delivery-ordering.test.ts index 300cb6dba4..bbe81626a1 100644 --- a/packages/core/src/step-delivery-ordering.test.ts +++ b/packages/core/src/step-delivery-ordering.test.ts @@ -41,7 +41,7 @@ import { createSleep } from './workflow/sleep.js'; * - `wait_completed` resolves through a detached chain with a fixed, small * microtask-hop count (`workflow/sleep.ts`). A `step_completed` instead * resolves inside a serial `ctx.promiseQueue` slot that first hydrates the - * payload via `ReplayPayloadCache.getStepResult(...)`. That hop count is not + * payload via `ReplayPayloadCache.getPrimitiveValue(...)`. That hop count is not * fixed: the first hydration pays async decrypt/deserialize, while a later * replay sharing the same `ReplayPayloadCache` hits the * `primitiveStepResults` memo for small primitive results and resolves in diff --git a/packages/core/src/step.ts b/packages/core/src/step.ts index 1fcc4c11df..77e078613f 100644 --- a/packages/core/src/step.ts +++ b/packages/core/src/step.ts @@ -190,18 +190,25 @@ export function createUseStep(ctx: WorkflowOrchestratorContext) { ctx.pendingDeliveries++; ctx.promiseQueue = ctx.promiseQueue.then(async () => { try { - const prepared = await ctx.replayPayloadCache.prepareEventPayload( + rejection = await ctx.replayPayloadCache.getPrimitiveValue( event.eventId, 'error', - event.eventData.error - ); - rejection = await hydrateStepError( - event.eventData.error, - ctx.runId, - ctx.encryptionKey, - ctx.globalThis, - {}, - prepared + async () => { + const prepared = + await ctx.replayPayloadCache.prepareEventPayload( + event.eventId, + 'error', + event.eventData.error + ); + return hydrateStepError( + event.eventData.error, + ctx.runId, + ctx.encryptionKey, + ctx.globalThis, + {}, + prepared + ); + } ); } catch (hydrateErr) { // If hydration fails for any reason, fall back to a generic @@ -301,25 +308,27 @@ export function createUseStep(ctx: WorkflowOrchestratorContext) { ctx.pendingDeliveries++; ctx.promiseQueue = ctx.promiseQueue.then(async () => { try { - const hydratedResult = await ctx.replayPayloadCache.getStepResult( - completedEventId, - async () => { - const prepared = - await ctx.replayPayloadCache.prepareEventPayload( - completedEventId, - 'result', - serializedResult + const hydratedResult = + await ctx.replayPayloadCache.getPrimitiveValue( + completedEventId, + 'result', + async () => { + const prepared = + await ctx.replayPayloadCache.prepareEventPayload( + completedEventId, + 'result', + serializedResult + ); + return await hydrateStepReturnValue( + serializedResult, + ctx.runId, + ctx.encryptionKey, + ctx.globalThis, + {}, + prepared ); - return await hydrateStepReturnValue( - serializedResult, - ctx.runId, - ctx.encryptionKey, - ctx.globalThis, - {}, - prepared - ); - } - ); + } + ); outcome = { ok: true, value: hydratedResult as Result }; } catch (error) { outcome = { ok: false, error }; diff --git a/packages/core/src/workflow/hook.ts b/packages/core/src/workflow/hook.ts index bcb714d89d..aed3f44a02 100644 --- a/packages/core/src/workflow/hook.ts +++ b/packages/core/src/workflow/hook.ts @@ -352,19 +352,25 @@ export function createCreateHook(ctx: WorkflowOrchestratorContext) { | { ok: false; error: unknown }; ctx.promiseQueue = ctx.promiseQueue.then(async () => { try { - const prepared = - await ctx.replayPayloadCache.prepareEventPayload( - event.eventId, - 'payload', - event.eventData.payload - ); - const payload = await hydrateStepReturnValue( - event.eventData.payload, - ctx.runId, - ctx.encryptionKey, - ctx.globalThis, - {}, - prepared + const payload = await ctx.replayPayloadCache.getPrimitiveValue( + event.eventId, + 'payload', + async () => { + const prepared = + await ctx.replayPayloadCache.prepareEventPayload( + event.eventId, + 'payload', + event.eventData.payload + ); + return hydrateStepReturnValue( + event.eventData.payload, + ctx.runId, + ctx.encryptionKey, + ctx.globalThis, + {}, + prepared + ); + } ); hydrateOutcome = { ok: true, value: payload as T }; } catch (error) { @@ -425,18 +431,25 @@ export function createCreateHook(ctx: WorkflowOrchestratorContext) { ctx.pendingDeliveries++; ctx.promiseQueue = ctx.promiseQueue.then(async () => { try { - const prepared = await ctx.replayPayloadCache.prepareEventPayload( + const payload = await ctx.replayPayloadCache.getPrimitiveValue( event.eventId, 'payload', - event.eventData.payload - ); - const payload = await hydrateStepReturnValue( - event.eventData.payload, - ctx.runId, - ctx.encryptionKey, - ctx.globalThis, - {}, - prepared + async () => { + const prepared = + await ctx.replayPayloadCache.prepareEventPayload( + event.eventId, + 'payload', + event.eventData.payload + ); + return hydrateStepReturnValue( + event.eventData.payload, + ctx.runId, + ctx.encryptionKey, + ctx.globalThis, + {}, + prepared + ); + } ); outcome = { ok: true, value: payload as T }; } catch (error) { From c485c65dbf4062d8dbb0f6ded935941fc791b73e Mon Sep 17 00:00:00 2001 From: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> Date: Thu, 13 Aug 2026 23:53:47 -0700 Subject: [PATCH 2/9] [core] Centralize replay event hydration --- packages/core/src/step.ts | 51 +++++++++++------------------- packages/core/src/workflow/hook.ts | 32 ++++++------------- 2 files changed, 29 insertions(+), 54 deletions(-) diff --git a/packages/core/src/step.ts b/packages/core/src/step.ts index 77e078613f..9e43d2740b 100644 --- a/packages/core/src/step.ts +++ b/packages/core/src/step.ts @@ -190,25 +190,19 @@ export function createUseStep(ctx: WorkflowOrchestratorContext) { ctx.pendingDeliveries++; ctx.promiseQueue = ctx.promiseQueue.then(async () => { try { - rejection = await ctx.replayPayloadCache.getPrimitiveValue( + rejection = await ctx.replayPayloadCache.getEventValue( event.eventId, 'error', - async () => { - const prepared = - await ctx.replayPayloadCache.prepareEventPayload( - event.eventId, - 'error', - event.eventData.error - ); - return hydrateStepError( + event.eventData.error, + (prepared) => + hydrateStepError( event.eventData.error, ctx.runId, ctx.encryptionKey, ctx.globalThis, {}, prepared - ); - } + ) ); } catch (hydrateErr) { // If hydration fails for any reason, fall back to a generic @@ -308,27 +302,20 @@ export function createUseStep(ctx: WorkflowOrchestratorContext) { ctx.pendingDeliveries++; ctx.promiseQueue = ctx.promiseQueue.then(async () => { try { - const hydratedResult = - await ctx.replayPayloadCache.getPrimitiveValue( - completedEventId, - 'result', - async () => { - const prepared = - await ctx.replayPayloadCache.prepareEventPayload( - completedEventId, - 'result', - serializedResult - ); - return await hydrateStepReturnValue( - serializedResult, - ctx.runId, - ctx.encryptionKey, - ctx.globalThis, - {}, - prepared - ); - } - ); + const hydratedResult = await ctx.replayPayloadCache.getEventValue( + completedEventId, + 'result', + serializedResult, + (prepared) => + hydrateStepReturnValue( + serializedResult, + ctx.runId, + ctx.encryptionKey, + ctx.globalThis, + {}, + prepared + ) + ); outcome = { ok: true, value: hydratedResult as Result }; } catch (error) { outcome = { ok: false, error }; diff --git a/packages/core/src/workflow/hook.ts b/packages/core/src/workflow/hook.ts index aed3f44a02..53d5e4ced2 100644 --- a/packages/core/src/workflow/hook.ts +++ b/packages/core/src/workflow/hook.ts @@ -352,25 +352,19 @@ export function createCreateHook(ctx: WorkflowOrchestratorContext) { | { ok: false; error: unknown }; ctx.promiseQueue = ctx.promiseQueue.then(async () => { try { - const payload = await ctx.replayPayloadCache.getPrimitiveValue( + const payload = await ctx.replayPayloadCache.getEventValue( event.eventId, 'payload', - async () => { - const prepared = - await ctx.replayPayloadCache.prepareEventPayload( - event.eventId, - 'payload', - event.eventData.payload - ); - return hydrateStepReturnValue( + event.eventData.payload, + (prepared) => + hydrateStepReturnValue( event.eventData.payload, ctx.runId, ctx.encryptionKey, ctx.globalThis, {}, prepared - ); - } + ) ); hydrateOutcome = { ok: true, value: payload as T }; } catch (error) { @@ -431,25 +425,19 @@ export function createCreateHook(ctx: WorkflowOrchestratorContext) { ctx.pendingDeliveries++; ctx.promiseQueue = ctx.promiseQueue.then(async () => { try { - const payload = await ctx.replayPayloadCache.getPrimitiveValue( + const payload = await ctx.replayPayloadCache.getEventValue( event.eventId, 'payload', - async () => { - const prepared = - await ctx.replayPayloadCache.prepareEventPayload( - event.eventId, - 'payload', - event.eventData.payload - ); - return hydrateStepReturnValue( + event.eventData.payload, + (prepared) => + hydrateStepReturnValue( event.eventData.payload, ctx.runId, ctx.encryptionKey, ctx.globalThis, {}, prepared - ); - } + ) ); outcome = { ok: true, value: payload as T }; } catch (error) { From 537159cc776daafb88a067f9fd358f682f2a71f5 Mon Sep 17 00:00:00 2001 From: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> Date: Fri, 14 Aug 2026 13:56:04 -0700 Subject: [PATCH 3/9] [core] Simplify replay payload caching --- packages/core/src/abort-consistency.test.ts | 2 +- packages/core/src/abort-controller.test.ts | 2 +- .../core/src/abort-replay-ordering.test.ts | 2 +- .../async-deserialization-ordering.test.ts | 2 +- .../src/delivery-barrier-coverage.test.ts | 2 +- .../core/src/replay-payload-cache.test.ts | 132 ++++++++++++------ packages/core/src/replay-payload-cache.ts | 97 +++++-------- packages/core/src/runtime.ts | 14 +- .../core/src/step-delivery-hop-count.test.ts | 8 +- .../core/src/step-delivery-ordering.test.ts | 10 +- .../src/step-hydration-memoization.test.ts | 8 +- packages/core/src/step.test.ts | 2 +- .../src/test-support/orchestrator-context.ts | 2 +- packages/core/src/workflow/hook.test.ts | 2 +- packages/core/src/workflow/sleep.test.ts | 2 +- 15 files changed, 151 insertions(+), 136 deletions(-) diff --git a/packages/core/src/abort-consistency.test.ts b/packages/core/src/abort-consistency.test.ts index ec93e93227..c1f2f5a4b9 100644 --- a/packages/core/src/abort-consistency.test.ts +++ b/packages/core/src/abort-consistency.test.ts @@ -37,7 +37,7 @@ function setupWorkflowContext(events: Event[]): WorkflowOrchestratorContext { return { runId: 'wrun_test', encryptionKey: undefined, - replayPayloadCache: new ReplayPayloadCache(undefined), + replayPayloadCache: new ReplayPayloadCache(), globalThis: context.globalThis, eventsConsumer: new EventsConsumer(events, { // Fake context: no deliveries are modeled, so the gate is a no-op here. diff --git a/packages/core/src/abort-controller.test.ts b/packages/core/src/abort-controller.test.ts index 3adf469366..96c4f8bc04 100644 --- a/packages/core/src/abort-controller.test.ts +++ b/packages/core/src/abort-controller.test.ts @@ -35,7 +35,7 @@ function setupWorkflowContext( return { runId: 'wrun_test', encryptionKey: undefined, - replayPayloadCache: new ReplayPayloadCache(undefined), + replayPayloadCache: new ReplayPayloadCache(), globalThis: context.globalThis, eventsConsumer: new EventsConsumer(events, { // Fake context: no deliveries are modeled, so the gate is a no-op here. diff --git a/packages/core/src/abort-replay-ordering.test.ts b/packages/core/src/abort-replay-ordering.test.ts index 8961a64925..06296ee362 100644 --- a/packages/core/src/abort-replay-ordering.test.ts +++ b/packages/core/src/abort-replay-ordering.test.ts @@ -70,7 +70,7 @@ function setupWorkflowContext(events: Event[]): WorkflowOrchestratorContext { return { runId: 'wrun_test', encryptionKey: undefined, - replayPayloadCache: new ReplayPayloadCache(undefined), + replayPayloadCache: new ReplayPayloadCache(), globalThis: context.globalThis, eventsConsumer: new EventsConsumer(events, { // Fake context: no deliveries are modeled, so the gate is a no-op here. diff --git a/packages/core/src/async-deserialization-ordering.test.ts b/packages/core/src/async-deserialization-ordering.test.ts index 41d4684563..b23346e3e7 100644 --- a/packages/core/src/async-deserialization-ordering.test.ts +++ b/packages/core/src/async-deserialization-ordering.test.ts @@ -50,7 +50,7 @@ function setupWorkflowContext(events: Event[]): WorkflowOrchestratorContext { return { runId: 'wrun_test', encryptionKey: undefined, - replayPayloadCache: new ReplayPayloadCache(undefined), + replayPayloadCache: new ReplayPayloadCache(), globalThis: context.globalThis, eventsConsumer: new EventsConsumer(events, { // Fake context: no deliveries are modeled, so the gate is a no-op here. diff --git a/packages/core/src/delivery-barrier-coverage.test.ts b/packages/core/src/delivery-barrier-coverage.test.ts index ee113b8ff0..27b29d2b5d 100644 --- a/packages/core/src/delivery-barrier-coverage.test.ts +++ b/packages/core/src/delivery-barrier-coverage.test.ts @@ -86,7 +86,7 @@ function setupWorkflowContext(events: Event[]): WorkflowOrchestratorContext { suspensionGeneration: 0, runId: 'wrun_test', encryptionKey: undefined, - replayPayloadCache: new ReplayPayloadCache(undefined), + replayPayloadCache: new ReplayPayloadCache(), globalThis: context.globalThis, eventsConsumer: new EventsConsumer(events, { // Fake context: no deliveries are modeled, so the gate is a no-op here. diff --git a/packages/core/src/replay-payload-cache.test.ts b/packages/core/src/replay-payload-cache.test.ts index 55adb4b196..c0aaff97c5 100644 --- a/packages/core/src/replay-payload-cache.test.ts +++ b/packages/core/src/replay-payload-cache.test.ts @@ -54,7 +54,7 @@ function makeEvents(payloads: unknown[]): Event[] { } describe('ReplayPayloadCache', () => { - it('deduplicates preparation and accepts a synchronous preparer', async () => { + it('deduplicates synchronous preparation without creating a promise', () => { const payload = new Uint8Array([1]); const preparer = vi.fn((value) => value); const cache = new ReplayPayloadCache(undefined, preparer); @@ -63,10 +63,24 @@ describe('ReplayPayloadCache', () => { const second = cache.prepareEventPayload('evnt_one', 'result', payload); expect(first).toBe(second); - await expect(first).resolves.toBe(payload); + expect(first).toEqual(payload); expect(preparer).toHaveBeenCalledOnce(); }); + it('hydrates and memoizes a primitive without creating a promise', () => { + const payload = new Uint8Array([1]); + const hydrate = vi.fn(() => 42); + const cache = new ReplayPayloadCache(undefined, (value) => value); + + expect(cache.getEventValue('evnt_one', 'result', payload, hydrate)).toBe( + 42 + ); + expect(cache.getEventValue('evnt_one', 'result', payload, hydrate)).toBe( + 42 + ); + expect(hydrate).toHaveBeenCalledOnce(); + }); + it('keeps a failed prewarm until its consumer observes it, then retries', async () => { const payload = new Uint8Array([1]); const run = makeRun(payload); @@ -76,13 +90,12 @@ describe('ReplayPayloadCache', () => { .mockReturnValueOnce(payload); const cache = new ReplayPayloadCache(undefined, preparer); - await cache.prewarm(run, []); - await expect(cache.prepareWorkflowInput(run)).rejects.toThrow( - 'decrypt failed' - ); + cache.prewarm(run, []); + await Promise.resolve(); + expect(() => cache.prepareWorkflowInput(run)).toThrow('decrypt failed'); expect(preparer).toHaveBeenCalledOnce(); - await expect(cache.prepareWorkflowInput(run)).resolves.toBe(payload); + expect(cache.prepareWorkflowInput(run)).toEqual(payload); expect(preparer).toHaveBeenCalledTimes(2); }); @@ -99,30 +112,52 @@ describe('ReplayPayloadCache', () => { const run = makeRun(payloads[0]); const events = makeEvents(payloads.slice(1)); - const warming = cache.prewarm(run, events); + cache.prewarm(run, events); expect(preparer).toHaveBeenCalledTimes(4); for (const resolve of resolvers.reverse()) resolve(); - await warming; - - const allSettled = vi.spyOn(Promise, 'allSettled'); - await cache.prewarm(run, events); + await Promise.all([ + cache.prepareWorkflowInput(run), + ...events.map((event) => { + switch (event.eventType) { + case 'step_completed': + return cache.prepareEventPayload( + event.eventId, + 'result', + event.eventData?.result + ); + case 'step_failed': + return cache.prepareEventPayload( + event.eventId, + 'error', + event.eventData?.error + ); + case 'hook_received': + return cache.prepareEventPayload( + event.eventId, + 'payload', + event.eventData?.payload + ); + default: + throw new Error(`Unexpected event: ${event.eventType}`); + } + }), + ]); + + cache.prewarm(run, events); expect(preparer).toHaveBeenCalledTimes(4); - expect(allSettled).toHaveBeenLastCalledWith([]); - allSettled.mockRestore(); }); it('prepares streamed events synchronously inside the decoder callback', async () => { const payload = new Uint8Array([1]); const order: string[] = []; - const preparer = vi.fn((value) => { + const preparer = vi.fn((value) => { order.push('prepare'); - return { data: value }; + return value; }); const cache = new ReplayPayloadCache(undefined, preparer); const [event] = makeEvents([payload]); - const preparation = cache.observeEvent(event, () => order.push('start')); - expect(preparation).toBeDefined(); + cache.observeEvent(event, () => order.push('start')); expect(preparer).toHaveBeenCalledOnce(); expect(order).toEqual(['start', 'prepare']); @@ -130,25 +165,31 @@ describe('ReplayPayloadCache', () => { cache.observeEvent(event, () => order.push('cached-start')); expect(order).toEqual(['start', 'prepare']); - await expect(preparation).resolves.toEqual({ data: payload }); + expect(cache.prepareEventPayload(event.eventId, 'result', payload)).toEqual( + payload + ); }); it('prepares queued stream events as soon as the run key resolves', async () => { const payload = new Uint8Array([1]); - const preparer = vi.fn((value) => ({ data: value })); + const preparer = vi.fn((value) => value); let resolveKey!: (key: undefined) => void; const key = new Promise((resolve) => { resolveKey = resolve; }); - const cache = new ReplayPayloadCache(key, preparer); + const cache = ReplayPayloadCache.waitingForKey(key, preparer); const [event] = makeEvents([payload]); - const preparation = cache.observeEvent(event); - expect(preparation).toBeDefined(); + cache.observeEvent(event); expect(preparer).not.toHaveBeenCalled(); + const preparation = cache.prepareEventPayload( + event.eventId, + 'result', + payload + ); resolveKey(undefined); - await expect(preparation).resolves.toEqual({ data: payload }); + await expect(preparation).resolves.toEqual(payload); expect(preparer).toHaveBeenCalledOnce(); }); @@ -238,74 +279,81 @@ describe('ReplayPayloadCache', () => { it('memoizes primitive step results, including undefined', async () => { for (const value of [0, false, '', null, undefined]) { - const cache = new ReplayPayloadCache(undefined); + const cache = new ReplayPayloadCache(); const hydrate = vi.fn().mockResolvedValue(value); expect( - await cache.getPrimitiveValue('evnt_result', 'result', hydrate) + await cache.getEventValue('evnt_result', 'result', undefined, hydrate) ).toBe(value); expect( - await cache.getPrimitiveValue('evnt_result', 'result', hydrate) + await cache.getEventValue('evnt_result', 'result', undefined, hydrate) ).toBe(value); expect(hydrate).toHaveBeenCalledOnce(); } }); it('isolates primitive values by event payload field', async () => { - const cache = new ReplayPayloadCache(undefined); + const cache = new ReplayPayloadCache(); const result = vi.fn().mockResolvedValue('result'); const error = vi.fn().mockResolvedValue('error'); await expect( - cache.getPrimitiveValue('evnt_shared', 'result', result) + cache.getEventValue('evnt_shared', 'result', undefined, result) ).resolves.toBe('result'); await expect( - cache.getPrimitiveValue('evnt_shared', 'error', error) + cache.getEventValue('evnt_shared', 'error', undefined, error) ).resolves.toBe('error'); - await expect( - cache.getPrimitiveValue('evnt_shared', 'result', result) - ).resolves.toBe('result'); + expect( + cache.getEventValue('evnt_shared', 'result', undefined, result) + ).toBe('result'); expect(result).toHaveBeenCalledOnce(); expect(error).toHaveBeenCalledOnce(); }); - it('rehydrates mutable and oversized step results', async () => { + it('rehydrates mutable results and memoizes primitives of any size', async () => { const oversized = 'x'.repeat(4097); for (const value of [{ count: 0 }, oversized]) { - const cache = new ReplayPayloadCache(undefined); + const cache = new ReplayPayloadCache(); const hydrate = vi .fn() .mockImplementation(async () => typeof value === 'object' ? { ...value } : value ); - const first = await cache.getPrimitiveValue( + const first = await cache.getEventValue( 'evnt_result', 'result', + undefined, hydrate ); - const second = await cache.getPrimitiveValue( + const second = await cache.getEventValue( 'evnt_result', 'result', + undefined, hydrate ); - expect(hydrate).toHaveBeenCalledTimes(2); - if (typeof value === 'object') expect(second).not.toBe(first); + if (typeof value === 'object') { + expect(hydrate).toHaveBeenCalledTimes(2); + expect(second).not.toBe(first); + } else { + expect(hydrate).toHaveBeenCalledOnce(); + expect(second).toBe(first); + } } }); it('does not memoize failed step hydration', async () => { - const cache = new ReplayPayloadCache(undefined); + const cache = new ReplayPayloadCache(); const hydrate = vi .fn() .mockRejectedValueOnce(new Error('boom')) .mockResolvedValueOnce('ok'); await expect( - cache.getPrimitiveValue('evnt_result', 'result', hydrate) + cache.getEventValue('evnt_result', 'result', undefined, hydrate) ).rejects.toThrow('boom'); await expect( - cache.getPrimitiveValue('evnt_result', 'result', hydrate) + cache.getEventValue('evnt_result', 'result', undefined, hydrate) ).resolves.toBe('ok'); expect(hydrate).toHaveBeenCalledTimes(2); }); diff --git a/packages/core/src/replay-payload-cache.ts b/packages/core/src/replay-payload-cache.ts index f9c6baba18..f288cc62e2 100644 --- a/packages/core/src/replay-payload-cache.ts +++ b/packages/core/src/replay-payload-cache.ts @@ -18,8 +18,12 @@ type Preparation = | { state: 'pending'; promise: Promise } | { state: 'failed'; error: unknown }; -function isPrimitive(value: unknown): boolean { - return value === null || !['object', 'function'].includes(typeof value); +function isCacheablePrimitive(value: unknown): boolean { + const type = typeof value; + return ( + value === null || + (type !== 'object' && type !== 'function' && type !== 'symbol') + ); } /** @@ -31,46 +35,29 @@ function isPrimitive(value: unknown): boolean { * share and skip that repeated deserialization entirely. */ export class ReplayPayloadCache { + private key: KeyState; private readonly preparations = new Map(); - private readonly pendingPreparations = new Set< - Promise - >(); private readonly primitiveValues = new Map(); private nextUnscannedEventIndex = 0; - private constructor( - private key: KeyState, + constructor( + key?: DecryptionKey, private readonly preparer: typeof prepareReplayPayload = prepareReplayPayload ) { - if (key.state === 'pending') { - void key.promise.then( - (value) => this.resolveKey(value), - (error) => this.rejectKey(error) - ); - } - } - - static unencrypted( - preparer: typeof prepareReplayPayload = prepareReplayPayload - ): ReplayPayloadCache { - return new ReplayPayloadCache( - { state: 'ready', value: undefined }, - preparer - ); - } - - static withKey( - key: DecryptionKey, - preparer: typeof prepareReplayPayload = prepareReplayPayload - ): ReplayPayloadCache { - return new ReplayPayloadCache({ state: 'ready', value: key }, preparer); + this.key = { state: 'ready', value: key }; } static waitingForKey( key: Promise, preparer: typeof prepareReplayPayload = prepareReplayPayload ): ReplayPayloadCache { - return new ReplayPayloadCache({ state: 'pending', promise: key }, preparer); + const cache = new ReplayPayloadCache(undefined, preparer); + cache.key = { state: 'pending', promise: key }; + void key.then( + (value) => cache.resolveKey(value), + (error) => cache.rejectKey(error) + ); + return cache; } /** Start preparing an event as soon as its frame has been decoded. */ @@ -115,13 +102,10 @@ export class ReplayPayloadCache { /** * Start every preparation not already observed from the event stream. - * Returns a Promise only when a sealed or portable codec is still running. + * Consumers await the few codecs that cannot complete synchronously. */ - prewarm(workflowRun: WorkflowRun, events: Event[]): void | Promise { - this.start( - this.workflowInputKey(workflowRun.runId), - workflowRun.input - ); + prewarm(workflowRun: WorkflowRun, events: Event[]): void { + this.start(this.workflowInputKey(workflowRun.runId), workflowRun.input); for ( let index = this.nextUnscannedEventIndex; index < events.length; @@ -130,7 +114,6 @@ export class ReplayPayloadCache { this.observeEvent(events[index]); } this.nextUnscannedEventIndex = events.length; - return this.waitForPending(); } /** A corrected reload may insert events before the previous scan position. */ @@ -166,11 +149,7 @@ export class ReplayPayloadCache { return this.primitiveValues.get(cacheKey); } - const prepared = this.prepareEventPayload( - eventId, - field, - serializedValue - ); + const prepared = this.prepareEventPayload(eventId, field, serializedValue); const hydrateAndCache = (payload: PreparedReplayPayload) => { const hydrated = hydrate(payload); return hydrated instanceof Promise @@ -183,7 +162,7 @@ export class ReplayPayloadCache { } private cachePrimitive(cacheKey: string, value: unknown): unknown { - if (isPrimitive(value)) this.primitiveValues.set(cacheKey, value); + if (isCacheablePrimitive(value)) this.primitiveValues.set(cacheKey, value); return value; } @@ -226,10 +205,8 @@ export class ReplayPayloadCache { } this.preparations.set(cacheKey, { state: 'pending', promise: result }); - this.pendingPreparations.add(result); void result.then( (prepared) => { - this.pendingPreparations.delete(result); const current = this.preparations.get(cacheKey); if (current?.state === 'pending' && current.promise === result) { this.preparations.set(cacheKey, { @@ -239,7 +216,6 @@ export class ReplayPayloadCache { } }, (error) => { - this.pendingPreparations.delete(result); const current = this.preparations.get(cacheKey); if (current?.state === 'pending' && current.promise === result) { this.preparations.set(cacheKey, { state: 'failed', error }); @@ -260,14 +236,26 @@ export class ReplayPayloadCache { this.start(cacheKey, value); const preparation = this.preparations.get(cacheKey); if (!preparation) { - throw new Error(`Replay payload preparation was not started: ${cacheKey}`); + throw new Error( + `Replay payload preparation was not started: ${cacheKey}` + ); } switch (preparation.state) { case 'ready': return preparation.value; case 'pending': - return preparation.promise; + return preparation.promise.catch((error) => { + const current = this.preparations.get(cacheKey); + if ( + current?.state === 'failed' || + (current?.state === 'pending' && + current.promise === preparation.promise) + ) { + this.preparations.delete(cacheKey); + } + throw error; + }); case 'failed': this.preparations.delete(cacheKey); throw preparation.error; @@ -279,19 +267,6 @@ export class ReplayPayloadCache { } } - private waitForPending(): void | Promise { - if ( - this.key.state === 'pending' && - [...this.preparations.values()].some( - (preparation) => preparation.state === 'waiting' - ) - ) { - return this.key.promise.then(() => this.waitForPending()); - } - if (this.pendingPreparations.size === 0) return; - return Promise.allSettled([...this.pendingPreparations]).then(() => {}); - } - private resolveKey(value: DecryptionKey | undefined): void { if (this.key.state !== 'pending') return; this.key = { state: 'ready', value }; diff --git a/packages/core/src/runtime.ts b/packages/core/src/runtime.ts index c997becbb7..c7d3e61dec 100644 --- a/packages/core/src/runtime.ts +++ b/packages/core/src/runtime.ts @@ -1062,9 +1062,8 @@ export function workflowEntrypoint( ({ runOrId, context }) => memoizeEncryptionKey(world, runOrId, context)() ); - const replayPayloadCache = new ReplayPayloadCache( - encryptionKeyPromise - ); + const replayPayloadCache = + ReplayPayloadCache.waitingForKey(encryptionKeyPromise); const observeReplayEvent = createReplayEventObserver({ runId, cache: replayPayloadCache, @@ -3015,14 +3014,7 @@ export function workflowEntrypoint( // appended-event consumption on resume; consumers still // deserialize and resolve in event order. const replayEvents = eventLog.events; - const payloadPrewarm = replayPayloadCache.prewarm( - workflowRun, - replayEvents - ); - // Consumers await their own prepared payloads in event - // order. Do not delay a suspension on speculative work - // for payloads this replay never touched. - void payloadPrewarm.catch(() => {}); + replayPayloadCache.prewarm(workflowRun, replayEvents); let workflowResult: WorkflowResumeResult = retainedSession ? await resumeWorkflow(retainedSession, eventLog.events) : { type: 'replay' }; diff --git a/packages/core/src/step-delivery-hop-count.test.ts b/packages/core/src/step-delivery-hop-count.test.ts index 7653f85d4e..65f628206f 100644 --- a/packages/core/src/step-delivery-hop-count.test.ts +++ b/packages/core/src/step-delivery-hop-count.test.ts @@ -44,7 +44,7 @@ import { createSleep } from './workflow/sleep.js'; function setupWorkflowContext( events: Event[], - replayPayloadCache: ReplayPayloadCache = new ReplayPayloadCache(undefined) + replayPayloadCache: ReplayPayloadCache = new ReplayPayloadCache() ): WorkflowOrchestratorContext { const context = createContext({ seed: 'test', @@ -229,7 +229,7 @@ describe('step delivery ordering is independent of consumer hop count: hook payl const spy = await slowHydration(); try { const events = await buildEventLog(); - const cache = new ReplayPayloadCache(undefined); + const cache = new ReplayPayloadCache(); const c1 = setupWorkflowContext(events, cache); const r1 = await runWithDiscontinuation(c1, body(c1, extraHops)); if (!WorkflowSuspension.is(r1.error)) { @@ -346,7 +346,7 @@ describe('step delivery ordering is independent of consumer hop count: wait comp const spy = await slowHydration(); try { const events = await buildWaitEventLog(); - const cache = new ReplayPayloadCache(undefined); + const cache = new ReplayPayloadCache(); const c1 = setupWorkflowContext(events, cache); const r1 = await runWithDiscontinuation(c1, waitBody(c1, extraHops)); if (!WorkflowSuspension.is(r1.error)) @@ -472,7 +472,7 @@ describe('step delivery ordering is independent of consumer hop count: step fail const spy = await slowHydration(); try { const events = await buildFailedEventLog(); - const cache = new ReplayPayloadCache(undefined); + const cache = new ReplayPayloadCache(); const c1 = setupWorkflowContext(events, cache); const r1 = await runWithDiscontinuation(c1, failedBody(c1, extraHops)); if (!WorkflowSuspension.is(r1.error)) { diff --git a/packages/core/src/step-delivery-ordering.test.ts b/packages/core/src/step-delivery-ordering.test.ts index bbe81626a1..18df1c7b77 100644 --- a/packages/core/src/step-delivery-ordering.test.ts +++ b/packages/core/src/step-delivery-ordering.test.ts @@ -41,7 +41,7 @@ import { createSleep } from './workflow/sleep.js'; * - `wait_completed` resolves through a detached chain with a fixed, small * microtask-hop count (`workflow/sleep.ts`). A `step_completed` instead * resolves inside a serial `ctx.promiseQueue` slot that first hydrates the - * payload via `ReplayPayloadCache.getPrimitiveValue(...)`. That hop count is not + * payload via `ReplayPayloadCache.getEventValue(...)`. That hop count is not * fixed: the first hydration pays async decrypt/deserialize, while a later * replay sharing the same `ReplayPayloadCache` hits the * `primitiveStepResults` memo for small primitive results and resolves in @@ -91,7 +91,7 @@ import { createSleep } from './workflow/sleep.js'; */ function setupWorkflowContext( events: Event[], - replayPayloadCache: ReplayPayloadCache = new ReplayPayloadCache(undefined) + replayPayloadCache: ReplayPayloadCache = new ReplayPayloadCache() ): WorkflowOrchestratorContext { const context = createContext({ seed: 'test', @@ -340,7 +340,7 @@ describe('step result delivery ordering across replays', () => { // One cache for both replays: production shares a single // `ReplayPayloadCache` across every replay of one queue delivery. - const sharedCache = new ReplayPayloadCache(undefined); + const sharedCache = new ReplayPayloadCache(); const firstCtx = setupWorkflowContext(events, sharedCache); const first = await runWithDiscontinuation( @@ -532,7 +532,7 @@ describe('step result delivery ordering across replays', () => { const hydration = delayHydration(); spy = await hydration.install(); const events = await buildEventLog(); - const sharedCache = new ReplayPayloadCache(undefined); + const sharedCache = new ReplayPayloadCache(); const firstCtx = setupWorkflowContext(events, sharedCache); const first = await runWithDiscontinuation( @@ -578,7 +578,7 @@ describe('step result delivery ordering across replays', () => { const hydration = delayHydration(); spy = await hydration.install(); const events = await buildEventLog(); - const sharedCache = new ReplayPayloadCache(undefined); + const sharedCache = new ReplayPayloadCache(); for (const replay of [1, 2]) { const ctx = setupWorkflowContext(events, sharedCache); diff --git a/packages/core/src/step-hydration-memoization.test.ts b/packages/core/src/step-hydration-memoization.test.ts index e8aabc1d0d..ede78c99b2 100644 --- a/packages/core/src/step-hydration-memoization.test.ts +++ b/packages/core/src/step-hydration-memoization.test.ts @@ -25,7 +25,7 @@ import { createContext } from './vm/index.js'; // the inline loop threads one cache across replay iterations. function setupWorkflowContext( events: Event[], - replayPayloadCache = new ReplayPayloadCache(undefined) + replayPayloadCache = new ReplayPayloadCache() ): WorkflowOrchestratorContext { const context = createContext({ seed: 'test', @@ -90,7 +90,7 @@ describe('step hydration memoization through the step consumer', () => { it('skips re-hydration of primitive step results on a second replay sharing the cache', async () => { const events = await makeStepEvents(); - const cache = new ReplayPayloadCache(undefined); + const cache = new ReplayPayloadCache(); const serialization = await import('./serialization.js'); const hydrateSpy = vi.spyOn(serialization, 'hydrateStepReturnValue'); @@ -122,7 +122,7 @@ describe('step hydration memoization through the step consumer', () => { it('preserves event-log resolution order on cache hits even with variable timing', async () => { const events = await makeStepEvents(); - const cache = new ReplayPayloadCache(undefined); + const cache = new ReplayPayloadCache(); // Replay 1: populate the cache (no timing games needed). const ctx1 = setupWorkflowContext(events, cache); @@ -168,7 +168,7 @@ describe('step hydration memoization through the step consumer', () => { createdAt: new Date(), }, ]; - const cache = new ReplayPayloadCache(undefined); + const cache = new ReplayPayloadCache(); // Replay 1: hydrate the object, then mutate it (as workflow code might). const ctx1 = setupWorkflowContext(events, cache); diff --git a/packages/core/src/step.test.ts b/packages/core/src/step.test.ts index 574e441110..cd759f9004 100644 --- a/packages/core/src/step.test.ts +++ b/packages/core/src/step.test.ts @@ -54,7 +54,7 @@ function setupWorkflowContext(events: Event[]): WorkflowOrchestratorContext { return { runId: 'wrun_test', encryptionKey: undefined, - replayPayloadCache: new ReplayPayloadCache(undefined), + replayPayloadCache: new ReplayPayloadCache(), globalThis: context.globalThis, eventsConsumer: new EventsConsumer(events, { // Fake context: no deliveries are modeled, so the gate is a no-op here. diff --git a/packages/core/src/test-support/orchestrator-context.ts b/packages/core/src/test-support/orchestrator-context.ts index 1ae7a78204..a9aa11f3e9 100644 --- a/packages/core/src/test-support/orchestrator-context.ts +++ b/packages/core/src/test-support/orchestrator-context.ts @@ -34,7 +34,7 @@ export function setupWorkflowContext( suspensionGeneration: 0, runId: 'wrun_test', encryptionKey: undefined, - replayPayloadCache: new ReplayPayloadCache(undefined), + replayPayloadCache: new ReplayPayloadCache(), globalThis: context.globalThis, eventsConsumer: new EventsConsumer(events, { // Fake context: no deliveries are modeled, so the gate is a no-op here. diff --git a/packages/core/src/workflow/hook.test.ts b/packages/core/src/workflow/hook.test.ts index c6c59727db..0ac3f008ae 100644 --- a/packages/core/src/workflow/hook.test.ts +++ b/packages/core/src/workflow/hook.test.ts @@ -39,7 +39,7 @@ function setupWorkflowContext(events: Event[]): WorkflowOrchestratorContext { runId: 'wrun_test', encryptionKey: undefined, worldCapabilities: { hookRetention: { active: true } }, - replayPayloadCache: new ReplayPayloadCache(undefined), + replayPayloadCache: new ReplayPayloadCache(), globalThis: context.globalThis, eventsConsumer: new EventsConsumer(events, { // Fake context: no deliveries are modeled, so the gate is a no-op here. diff --git a/packages/core/src/workflow/sleep.test.ts b/packages/core/src/workflow/sleep.test.ts index 7780bc05f1..139b185653 100644 --- a/packages/core/src/workflow/sleep.test.ts +++ b/packages/core/src/workflow/sleep.test.ts @@ -23,7 +23,7 @@ function setupWorkflowContext(events: Event[]): WorkflowOrchestratorContext { suspensionGeneration: 0, runId: 'wrun_test', encryptionKey: undefined, - replayPayloadCache: new ReplayPayloadCache(undefined), + replayPayloadCache: new ReplayPayloadCache(), globalThis: context.globalThis, // ctx.onWorkflowError is accessed via closure — it's defined below on the same object eventsConsumer: new EventsConsumer(events, { From d527e73006aab4e494425da1896613382b4f8713 Mon Sep 17 00:00:00 2001 From: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> Date: Fri, 14 Aug 2026 14:05:41 -0700 Subject: [PATCH 4/9] [core] Name replay preparation states explicitly --- packages/core/src/replay-payload-cache.ts | 69 ++++++++++++----------- 1 file changed, 37 insertions(+), 32 deletions(-) diff --git a/packages/core/src/replay-payload-cache.ts b/packages/core/src/replay-payload-cache.ts index f288cc62e2..af6c102294 100644 --- a/packages/core/src/replay-payload-cache.ts +++ b/packages/core/src/replay-payload-cache.ts @@ -8,14 +8,14 @@ import { type ReplayPayloadField = 'result' | 'error' | 'payload'; type KeyState = - | { state: 'pending'; promise: Promise } + | { state: 'loading'; promise: Promise } | { state: 'ready'; value: DecryptionKey | undefined } | { state: 'failed'; error: unknown }; type Preparation = - | { state: 'waiting'; value: Uint8Array } + | { state: 'waitingForKey'; value: Uint8Array } | { state: 'ready'; value: PreparedReplayPayload } - | { state: 'pending'; promise: Promise } + | { state: 'preparing'; promise: Promise } | { state: 'failed'; error: unknown }; function isCacheablePrimitive(value: unknown): boolean { @@ -52,7 +52,7 @@ export class ReplayPayloadCache { preparer: typeof prepareReplayPayload = prepareReplayPayload ): ReplayPayloadCache { const cache = new ReplayPayloadCache(undefined, preparer); - cache.key = { state: 'pending', promise: key }; + cache.key = { state: 'loading', promise: key }; void key.then( (value) => cache.resolveKey(value), (error) => cache.rejectKey(error) @@ -64,35 +64,35 @@ export class ReplayPayloadCache { observeEvent(event: Event, onPreparationStart?: () => void): void { switch (event.eventType) { case 'run_created': - this.start( + this.startPreparation( this.workflowInputKey(event.runId), event.eventData.input, onPreparationStart ); break; case 'run_started': - this.start( + this.startPreparation( this.workflowInputKey(event.runId), event.eventData?.input, onPreparationStart ); break; case 'step_completed': - this.start( + this.startPreparation( this.eventPayloadKey(event.eventId, 'result'), event.eventData?.result, onPreparationStart ); break; case 'step_failed': - this.start( + this.startPreparation( this.eventPayloadKey(event.eventId, 'error'), event.eventData?.error, onPreparationStart ); break; case 'hook_received': - this.start( + this.startPreparation( this.eventPayloadKey(event.eventId, 'payload'), event.eventData?.payload, onPreparationStart @@ -105,7 +105,10 @@ export class ReplayPayloadCache { * Consumers await the few codecs that cannot complete synchronously. */ prewarm(workflowRun: WorkflowRun, events: Event[]): void { - this.start(this.workflowInputKey(workflowRun.runId), workflowRun.input); + this.startPreparation( + this.workflowInputKey(workflowRun.runId), + workflowRun.input + ); for ( let index = this.nextUnscannedEventIndex; index < events.length; @@ -124,7 +127,7 @@ export class ReplayPayloadCache { prepareWorkflowInput( workflowRun: WorkflowRun ): PreparedReplayPayload | Promise { - return this.consume( + return this.getPreparedPayload( this.workflowInputKey(workflowRun.runId), workflowRun.input ); @@ -135,7 +138,7 @@ export class ReplayPayloadCache { field: ReplayPayloadField, value: unknown ): PreparedReplayPayload | Promise { - return this.consume(this.eventPayloadKey(eventId, field), value); + return this.getPreparedPayload(this.eventPayloadKey(eventId, field), value); } getEventValue( @@ -166,7 +169,7 @@ export class ReplayPayloadCache { return value; } - private start( + private startPreparation( cacheKey: string, value: unknown, onPreparationStart?: () => void @@ -177,11 +180,11 @@ export class ReplayPayloadCache { onPreparationStart?.(); switch (this.key.state) { - case 'pending': - this.preparations.set(cacheKey, { state: 'waiting', value }); + case 'loading': + this.preparations.set(cacheKey, { state: 'waitingForKey', value }); break; case 'ready': - this.prepare(cacheKey, value, this.key.value); + this.runPreparation(cacheKey, value, this.key.value); break; case 'failed': this.preparations.set(cacheKey, { @@ -192,7 +195,7 @@ export class ReplayPayloadCache { } } - private prepare( + private runPreparation( cacheKey: string, value: Uint8Array, key: DecryptionKey | undefined @@ -204,11 +207,11 @@ export class ReplayPayloadCache { return; } - this.preparations.set(cacheKey, { state: 'pending', promise: result }); + this.preparations.set(cacheKey, { state: 'preparing', promise: result }); void result.then( (prepared) => { const current = this.preparations.get(cacheKey); - if (current?.state === 'pending' && current.promise === result) { + if (current?.state === 'preparing' && current.promise === result) { this.preparations.set(cacheKey, { state: 'ready', value: prepared, @@ -217,7 +220,7 @@ export class ReplayPayloadCache { }, (error) => { const current = this.preparations.get(cacheKey); - if (current?.state === 'pending' && current.promise === result) { + if (current?.state === 'preparing' && current.promise === result) { this.preparations.set(cacheKey, { state: 'failed', error }); } } @@ -227,13 +230,13 @@ export class ReplayPayloadCache { } } - private consume( + private getPreparedPayload( cacheKey: string, value: unknown ): PreparedReplayPayload | Promise { if (!(value instanceof Uint8Array)) return { legacy: value }; - this.start(cacheKey, value); + this.startPreparation(cacheKey, value); const preparation = this.preparations.get(cacheKey); if (!preparation) { throw new Error( @@ -244,12 +247,12 @@ export class ReplayPayloadCache { switch (preparation.state) { case 'ready': return preparation.value; - case 'pending': + case 'preparing': return preparation.promise.catch((error) => { const current = this.preparations.get(cacheKey); if ( current?.state === 'failed' || - (current?.state === 'pending' && + (current?.state === 'preparing' && current.promise === preparation.promise) ) { this.preparations.delete(cacheKey); @@ -259,29 +262,31 @@ export class ReplayPayloadCache { case 'failed': this.preparations.delete(cacheKey); throw preparation.error; - case 'waiting': - if (this.key.state !== 'pending') { + case 'waitingForKey': + if (this.key.state !== 'loading') { throw new Error(`Replay payload key was not resolved: ${cacheKey}`); } - return this.key.promise.then(() => this.consume(cacheKey, value)); + return this.key.promise.then(() => + this.getPreparedPayload(cacheKey, value) + ); } } private resolveKey(value: DecryptionKey | undefined): void { - if (this.key.state !== 'pending') return; + if (this.key.state !== 'loading') return; this.key = { state: 'ready', value }; for (const [cacheKey, preparation] of this.preparations) { - if (preparation.state === 'waiting') { - this.prepare(cacheKey, preparation.value, value); + if (preparation.state === 'waitingForKey') { + this.runPreparation(cacheKey, preparation.value, value); } } } private rejectKey(error: unknown): void { - if (this.key.state !== 'pending') return; + if (this.key.state !== 'loading') return; this.key = { state: 'failed', error }; for (const [cacheKey, preparation] of this.preparations) { - if (preparation.state === 'waiting') { + if (preparation.state === 'waitingForKey') { this.preparations.set(cacheKey, { state: 'failed', error }); } } From 87e4649ed9320f82e73c406d10cd4b8a5981f3f3 Mon Sep 17 00:00:00 2001 From: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> Date: Fri, 14 Aug 2026 14:16:06 -0700 Subject: [PATCH 5/9] [core] Key replay payloads by event --- .../core/src/replay-payload-cache.test.ts | 61 +++++++------------ packages/core/src/replay-payload-cache.ts | 38 ++++-------- packages/core/src/step.ts | 2 - .../core/src/workflow/abort-controller.ts | 1 - packages/core/src/workflow/hook.ts | 2 - 5 files changed, 35 insertions(+), 69 deletions(-) diff --git a/packages/core/src/replay-payload-cache.test.ts b/packages/core/src/replay-payload-cache.test.ts index c0aaff97c5..e8d48e58ca 100644 --- a/packages/core/src/replay-payload-cache.test.ts +++ b/packages/core/src/replay-payload-cache.test.ts @@ -59,8 +59,8 @@ describe('ReplayPayloadCache', () => { const preparer = vi.fn((value) => value); const cache = new ReplayPayloadCache(undefined, preparer); - const first = cache.prepareEventPayload('evnt_one', 'result', payload); - const second = cache.prepareEventPayload('evnt_one', 'result', payload); + const first = cache.prepareEventPayload('evnt_one', payload); + const second = cache.prepareEventPayload('evnt_one', payload); expect(first).toBe(second); expect(first).toEqual(payload); @@ -72,12 +72,8 @@ describe('ReplayPayloadCache', () => { const hydrate = vi.fn(() => 42); const cache = new ReplayPayloadCache(undefined, (value) => value); - expect(cache.getEventValue('evnt_one', 'result', payload, hydrate)).toBe( - 42 - ); - expect(cache.getEventValue('evnt_one', 'result', payload, hydrate)).toBe( - 42 - ); + expect(cache.getEventValue('evnt_one', payload, hydrate)).toBe(42); + expect(cache.getEventValue('evnt_one', payload, hydrate)).toBe(42); expect(hydrate).toHaveBeenCalledOnce(); }); @@ -122,19 +118,16 @@ describe('ReplayPayloadCache', () => { case 'step_completed': return cache.prepareEventPayload( event.eventId, - 'result', event.eventData?.result ); case 'step_failed': return cache.prepareEventPayload( event.eventId, - 'error', event.eventData?.error ); case 'hook_received': return cache.prepareEventPayload( event.eventId, - 'payload', event.eventData?.payload ); default: @@ -165,9 +158,7 @@ describe('ReplayPayloadCache', () => { cache.observeEvent(event, () => order.push('cached-start')); expect(order).toEqual(['start', 'prepare']); - expect(cache.prepareEventPayload(event.eventId, 'result', payload)).toEqual( - payload - ); + expect(cache.prepareEventPayload(event.eventId, payload)).toEqual(payload); }); it('prepares queued stream events as soon as the run key resolves', async () => { @@ -182,11 +173,7 @@ describe('ReplayPayloadCache', () => { cache.observeEvent(event); expect(preparer).not.toHaveBeenCalled(); - const preparation = cache.prepareEventPayload( - event.eventId, - 'result', - payload - ); + const preparation = cache.prepareEventPayload(event.eventId, payload); resolveKey(undefined); await expect(preparation).resolves.toEqual(payload); @@ -214,12 +201,10 @@ describe('ReplayPayloadCache', () => { const prepared = await cache.prepareEventPayload( 'evnt_encrypted', - 'result', serialized ); const samePrepared = await cache.prepareEventPayload( 'evnt_encrypted', - 'result', serialized ); const first = deserializePreparedReplayPayload(prepared) as { @@ -267,8 +252,8 @@ describe('ReplayPayloadCache', () => { const preparer = vi.fn((value) => value); const cache = new ReplayPayloadCache(undefined, preparer); - await cache.prepareEventPayload('evnt_legacy', 'result', legacy); - await cache.prepareEventPayload('evnt_legacy', 'result', legacy); + await cache.prepareEventPayload('evnt_legacy', legacy); + await cache.prepareEventPayload('evnt_legacy', legacy); expect(preparer).not.toHaveBeenCalled(); const events = makeEvents([legacy, legacy, legacy]); @@ -282,30 +267,30 @@ describe('ReplayPayloadCache', () => { const cache = new ReplayPayloadCache(); const hydrate = vi.fn().mockResolvedValue(value); - expect( - await cache.getEventValue('evnt_result', 'result', undefined, hydrate) - ).toBe(value); - expect( - await cache.getEventValue('evnt_result', 'result', undefined, hydrate) - ).toBe(value); + expect(await cache.getEventValue('evnt_result', undefined, hydrate)).toBe( + value + ); + expect(await cache.getEventValue('evnt_result', undefined, hydrate)).toBe( + value + ); expect(hydrate).toHaveBeenCalledOnce(); } }); - it('isolates primitive values by event payload field', async () => { + it('isolates primitive values by event id', async () => { const cache = new ReplayPayloadCache(); const result = vi.fn().mockResolvedValue('result'); const error = vi.fn().mockResolvedValue('error'); await expect( - cache.getEventValue('evnt_shared', 'result', undefined, result) + cache.getEventValue('evnt_result', undefined, result) ).resolves.toBe('result'); await expect( - cache.getEventValue('evnt_shared', 'error', undefined, error) + cache.getEventValue('evnt_error', undefined, error) ).resolves.toBe('error'); - expect( - cache.getEventValue('evnt_shared', 'result', undefined, result) - ).toBe('result'); + expect(cache.getEventValue('evnt_result', undefined, result)).toBe( + 'result' + ); expect(result).toHaveBeenCalledOnce(); expect(error).toHaveBeenCalledOnce(); }); @@ -322,13 +307,11 @@ describe('ReplayPayloadCache', () => { const first = await cache.getEventValue( 'evnt_result', - 'result', undefined, hydrate ); const second = await cache.getEventValue( 'evnt_result', - 'result', undefined, hydrate ); @@ -350,10 +333,10 @@ describe('ReplayPayloadCache', () => { .mockResolvedValueOnce('ok'); await expect( - cache.getEventValue('evnt_result', 'result', undefined, hydrate) + cache.getEventValue('evnt_result', undefined, hydrate) ).rejects.toThrow('boom'); await expect( - cache.getEventValue('evnt_result', 'result', undefined, hydrate) + cache.getEventValue('evnt_result', undefined, hydrate) ).resolves.toBe('ok'); expect(hydrate).toHaveBeenCalledTimes(2); }); diff --git a/packages/core/src/replay-payload-cache.ts b/packages/core/src/replay-payload-cache.ts index af6c102294..fe29a7f7e9 100644 --- a/packages/core/src/replay-payload-cache.ts +++ b/packages/core/src/replay-payload-cache.ts @@ -5,7 +5,7 @@ import { prepareReplayPayload, } from './serialization/replay.js'; -type ReplayPayloadField = 'result' | 'error' | 'payload'; +const WORKFLOW_INPUT_CACHE_KEY = 'workflow-input'; type KeyState = | { state: 'loading'; promise: Promise } @@ -65,35 +65,35 @@ export class ReplayPayloadCache { switch (event.eventType) { case 'run_created': this.startPreparation( - this.workflowInputKey(event.runId), + WORKFLOW_INPUT_CACHE_KEY, event.eventData.input, onPreparationStart ); break; case 'run_started': this.startPreparation( - this.workflowInputKey(event.runId), + WORKFLOW_INPUT_CACHE_KEY, event.eventData?.input, onPreparationStart ); break; case 'step_completed': this.startPreparation( - this.eventPayloadKey(event.eventId, 'result'), + this.eventPayloadKey(event.eventId), event.eventData?.result, onPreparationStart ); break; case 'step_failed': this.startPreparation( - this.eventPayloadKey(event.eventId, 'error'), + this.eventPayloadKey(event.eventId), event.eventData?.error, onPreparationStart ); break; case 'hook_received': this.startPreparation( - this.eventPayloadKey(event.eventId, 'payload'), + this.eventPayloadKey(event.eventId), event.eventData?.payload, onPreparationStart ); @@ -105,10 +105,7 @@ export class ReplayPayloadCache { * Consumers await the few codecs that cannot complete synchronously. */ prewarm(workflowRun: WorkflowRun, events: Event[]): void { - this.startPreparation( - this.workflowInputKey(workflowRun.runId), - workflowRun.input - ); + this.startPreparation(WORKFLOW_INPUT_CACHE_KEY, workflowRun.input); for ( let index = this.nextUnscannedEventIndex; index < events.length; @@ -127,32 +124,27 @@ export class ReplayPayloadCache { prepareWorkflowInput( workflowRun: WorkflowRun ): PreparedReplayPayload | Promise { - return this.getPreparedPayload( - this.workflowInputKey(workflowRun.runId), - workflowRun.input - ); + return this.getPreparedPayload(WORKFLOW_INPUT_CACHE_KEY, workflowRun.input); } prepareEventPayload( eventId: string, - field: ReplayPayloadField, value: unknown ): PreparedReplayPayload | Promise { - return this.getPreparedPayload(this.eventPayloadKey(eventId, field), value); + return this.getPreparedPayload(this.eventPayloadKey(eventId), value); } getEventValue( eventId: string, - field: ReplayPayloadField, serializedValue: unknown, hydrate: (prepared: PreparedReplayPayload) => unknown | Promise ): unknown | Promise { - const cacheKey = this.eventPayloadKey(eventId, field); + const cacheKey = this.eventPayloadKey(eventId); if (this.primitiveValues.has(cacheKey)) { return this.primitiveValues.get(cacheKey); } - const prepared = this.prepareEventPayload(eventId, field, serializedValue); + const prepared = this.prepareEventPayload(eventId, serializedValue); const hydrateAndCache = (payload: PreparedReplayPayload) => { const hydrated = hydrate(payload); return hydrated instanceof Promise @@ -292,11 +284,7 @@ export class ReplayPayloadCache { } } - private workflowInputKey(runId: string): string { - return `run:${runId}:input`; - } - - private eventPayloadKey(eventId: string, field: ReplayPayloadField): string { - return `event:${eventId}:${field}`; + private eventPayloadKey(eventId: string): string { + return `event:${eventId}`; } } diff --git a/packages/core/src/step.ts b/packages/core/src/step.ts index 9e43d2740b..78a9379cee 100644 --- a/packages/core/src/step.ts +++ b/packages/core/src/step.ts @@ -192,7 +192,6 @@ export function createUseStep(ctx: WorkflowOrchestratorContext) { try { rejection = await ctx.replayPayloadCache.getEventValue( event.eventId, - 'error', event.eventData.error, (prepared) => hydrateStepError( @@ -304,7 +303,6 @@ export function createUseStep(ctx: WorkflowOrchestratorContext) { try { const hydratedResult = await ctx.replayPayloadCache.getEventValue( completedEventId, - 'result', serializedResult, (prepared) => hydrateStepReturnValue( diff --git a/packages/core/src/workflow/abort-controller.ts b/packages/core/src/workflow/abort-controller.ts index 4d5d9a08af..02b4f39bba 100644 --- a/packages/core/src/workflow/abort-controller.ts +++ b/packages/core/src/workflow/abort-controller.ts @@ -233,7 +233,6 @@ export function createCreateAbortController(ctx: WorkflowOrchestratorContext) { const prepared = await ctx.replayPayloadCache.prepareEventPayload( event.eventId, - 'payload', rawPayload ); const hydrated = (await hydrateStepReturnValue( diff --git a/packages/core/src/workflow/hook.ts b/packages/core/src/workflow/hook.ts index 53d5e4ced2..16e14ccedc 100644 --- a/packages/core/src/workflow/hook.ts +++ b/packages/core/src/workflow/hook.ts @@ -354,7 +354,6 @@ export function createCreateHook(ctx: WorkflowOrchestratorContext) { try { const payload = await ctx.replayPayloadCache.getEventValue( event.eventId, - 'payload', event.eventData.payload, (prepared) => hydrateStepReturnValue( @@ -427,7 +426,6 @@ export function createCreateHook(ctx: WorkflowOrchestratorContext) { try { const payload = await ctx.replayPayloadCache.getEventValue( event.eventId, - 'payload', event.eventData.payload, (prepared) => hydrateStepReturnValue( From 62ed28959a1f56faad058327f8220c1972614a8a Mon Sep 17 00:00:00 2001 From: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> Date: Fri, 14 Aug 2026 15:42:26 -0700 Subject: [PATCH 6/9] [core] Collapse replay cache state machinery --- packages/core/src/replay-payload-cache.ts | 130 ++++++---------------- 1 file changed, 32 insertions(+), 98 deletions(-) diff --git a/packages/core/src/replay-payload-cache.ts b/packages/core/src/replay-payload-cache.ts index fe29a7f7e9..445f484a52 100644 --- a/packages/core/src/replay-payload-cache.ts +++ b/packages/core/src/replay-payload-cache.ts @@ -7,17 +7,6 @@ import { const WORKFLOW_INPUT_CACHE_KEY = 'workflow-input'; -type KeyState = - | { state: 'loading'; promise: Promise } - | { state: 'ready'; value: DecryptionKey | undefined } - | { state: 'failed'; error: unknown }; - -type Preparation = - | { state: 'waitingForKey'; value: Uint8Array } - | { state: 'ready'; value: PreparedReplayPayload } - | { state: 'preparing'; promise: Promise } - | { state: 'failed'; error: unknown }; - function isCacheablePrimitive(value: unknown): boolean { const type = typeof value; return ( @@ -35,27 +24,31 @@ function isCacheablePrimitive(value: unknown): boolean { * share and skip that repeated deserialization entirely. */ export class ReplayPayloadCache { - private key: KeyState; - private readonly preparations = new Map(); + private readonly preparations = new Map< + string, + Uint8Array | Promise | { readonly error: unknown } + >(); private readonly primitiveValues = new Map(); private nextUnscannedEventIndex = 0; + private encryptionKeyPromise?: Promise; constructor( - key?: DecryptionKey, + private encryptionKey?: DecryptionKey, private readonly preparer: typeof prepareReplayPayload = prepareReplayPayload - ) { - this.key = { state: 'ready', value: key }; - } + ) {} static waitingForKey( key: Promise, preparer: typeof prepareReplayPayload = prepareReplayPayload ): ReplayPayloadCache { const cache = new ReplayPayloadCache(undefined, preparer); - cache.key = { state: 'loading', promise: key }; + cache.encryptionKeyPromise = key; void key.then( - (value) => cache.resolveKey(value), - (error) => cache.rejectKey(error) + (value) => { + cache.encryptionKey = value; + cache.encryptionKeyPromise = undefined; + }, + () => {} ); return cache; } @@ -171,54 +164,32 @@ export class ReplayPayloadCache { } onPreparationStart?.(); - switch (this.key.state) { - case 'loading': - this.preparations.set(cacheKey, { state: 'waitingForKey', value }); - break; - case 'ready': - this.runPreparation(cacheKey, value, this.key.value); - break; - case 'failed': - this.preparations.set(cacheKey, { - state: 'failed', - error: this.key.error, - }); - break; - } + this.runPreparation(cacheKey, value); } - private runPreparation( - cacheKey: string, - value: Uint8Array, - key: DecryptionKey | undefined - ): void { + private runPreparation(cacheKey: string, value: Uint8Array): void { try { - const result = this.preparer(value, key); + const result = this.encryptionKeyPromise + ? this.encryptionKeyPromise.then((key) => this.preparer(value, key)) + : this.preparer(value, this.encryptionKey); if (!(result instanceof Promise)) { - this.preparations.set(cacheKey, { state: 'ready', value: result }); + this.preparations.set(cacheKey, result); return; } - this.preparations.set(cacheKey, { state: 'preparing', promise: result }); + this.preparations.set(cacheKey, result); void result.then( (prepared) => { const current = this.preparations.get(cacheKey); - if (current?.state === 'preparing' && current.promise === result) { - this.preparations.set(cacheKey, { - state: 'ready', - value: prepared, - }); - } + if (current === result) this.preparations.set(cacheKey, prepared); }, (error) => { const current = this.preparations.get(cacheKey); - if (current?.state === 'preparing' && current.promise === result) { - this.preparations.set(cacheKey, { state: 'failed', error }); - } + if (current === result) this.preparations.set(cacheKey, { error }); } ); } catch (error) { - this.preparations.set(cacheKey, { state: 'failed', error }); + this.preparations.set(cacheKey, { error }); } } @@ -229,59 +200,22 @@ export class ReplayPayloadCache { if (!(value instanceof Uint8Array)) return { legacy: value }; this.startPreparation(cacheKey, value); - const preparation = this.preparations.get(cacheKey); - if (!preparation) { + const prepared = this.preparations.get(cacheKey); + if (!prepared) { throw new Error( `Replay payload preparation was not started: ${cacheKey}` ); } - switch (preparation.state) { - case 'ready': - return preparation.value; - case 'preparing': - return preparation.promise.catch((error) => { - const current = this.preparations.get(cacheKey); - if ( - current?.state === 'failed' || - (current?.state === 'preparing' && - current.promise === preparation.promise) - ) { - this.preparations.delete(cacheKey); - } - throw error; - }); - case 'failed': + if (prepared instanceof Uint8Array) return prepared; + if (prepared instanceof Promise) { + return prepared.catch((error) => { this.preparations.delete(cacheKey); - throw preparation.error; - case 'waitingForKey': - if (this.key.state !== 'loading') { - throw new Error(`Replay payload key was not resolved: ${cacheKey}`); - } - return this.key.promise.then(() => - this.getPreparedPayload(cacheKey, value) - ); - } - } - - private resolveKey(value: DecryptionKey | undefined): void { - if (this.key.state !== 'loading') return; - this.key = { state: 'ready', value }; - for (const [cacheKey, preparation] of this.preparations) { - if (preparation.state === 'waitingForKey') { - this.runPreparation(cacheKey, preparation.value, value); - } - } - } - - private rejectKey(error: unknown): void { - if (this.key.state !== 'loading') return; - this.key = { state: 'failed', error }; - for (const [cacheKey, preparation] of this.preparations) { - if (preparation.state === 'waitingForKey') { - this.preparations.set(cacheKey, { state: 'failed', error }); - } + throw error; + }); } + this.preparations.delete(cacheKey); + throw prepared.error; } private eventPayloadKey(eventId: string): string { From 8d01f3027bf8cbebab7cc0b980a48f51b24c1bd3 Mon Sep 17 00:00:00 2001 From: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> Date: Fri, 14 Aug 2026 15:51:18 -0700 Subject: [PATCH 7/9] [core] Pass replay event loading options directly --- packages/core/src/runtime.ts | 47 +++++++++++------------ packages/core/src/runtime/helpers.test.ts | 28 +++++++++----- packages/core/src/runtime/helpers.ts | 17 ++++---- 3 files changed, 52 insertions(+), 40 deletions(-) diff --git a/packages/core/src/runtime.ts b/packages/core/src/runtime.ts index c7d3e61dec..6f58607cf1 100644 --- a/packages/core/src/runtime.ts +++ b/packages/core/src/runtime.ts @@ -1696,7 +1696,7 @@ export function workflowEntrypoint( getStepFunction(incomingStepName)?.maxRetries ?? DEFAULT_STEP_MAX_RETRIES; if (metadata.attempt > bgMaxRetries + 1) { - const loaded = await loadWorkflowRunEvents(runId); + const loaded = await loadWorkflowRunEvents({ runId }); bgAuthoritativeAttempt = countStepStartedEvents( loaded.events, @@ -1810,7 +1810,7 @@ export function workflowEntrypoint( // Load events to check if all parallel steps are done. // Use cursor-based loading so the main loop can continue // incrementally from here. - const loaded = await loadWorkflowRunEvents(runId); + const loaded = await loadWorkflowRunEvents({ runId }); eventLog = nextEventLogLoad(loaded); // Check for pending steps: any step_created without @@ -2711,11 +2711,11 @@ export function workflowEntrypoint( if (eventLog.type === 'loadAfter') { appendEventLog( eventLog, - await loadWorkflowRunEvents( + await loadWorkflowRunEvents({ runId, - eventLog.cursor, - observeReplayEvent - ) + afterCursor: eventLog.cursor, + onEvent: observeReplayEvent, + }) ); eventLog = { ...eventLog, type: 'ready' }; } @@ -2767,13 +2767,14 @@ export function workflowEntrypoint( } if (eventLog.type !== 'ready') { - const page = await loadWorkflowRunEvents( + const page = await loadWorkflowRunEvents({ runId, - eventLog.type === 'loadAfter' - ? eventLog.cursor - : undefined, - observeReplayEvent - ); + afterCursor: + eventLog.type === 'loadAfter' + ? eventLog.cursor + : undefined, + onEvent: observeReplayEvent, + }); if (eventLog.type === 'loadAfter') { appendEventLog(eventLog, page); eventLog = { ...eventLog, type: 'ready' }; @@ -2890,11 +2891,11 @@ export function workflowEntrypoint( // not include the wait completion this handler just // attempted. if (eventLog.cursor) { - const page = await loadWorkflowRunEvents( + const page = await loadWorkflowRunEvents({ runId, - eventLog.cursor, - observeReplayEvent - ); + afterCursor: eventLog.cursor, + onEvent: observeReplayEvent, + }); const completedWaitIdsAfterCursor = new Set( page.events .filter((e) => e.eventType === 'wait_completed') @@ -2911,21 +2912,19 @@ export function workflowEntrypoint( appendEventLog(eventLog, page); } else { eventLog = { - ...(await loadWorkflowRunEvents( + ...(await loadWorkflowRunEvents({ runId, - undefined, - observeReplayEvent - )), + onEvent: observeReplayEvent, + })), type: 'ready', }; } } else { eventLog = { - ...(await loadWorkflowRunEvents( + ...(await loadWorkflowRunEvents({ runId, - undefined, - observeReplayEvent - )), + onEvent: observeReplayEvent, + })), type: 'ready', }; } diff --git a/packages/core/src/runtime/helpers.test.ts b/packages/core/src/runtime/helpers.test.ts index 0854f5d381..cc397963bd 100644 --- a/packages/core/src/runtime/helpers.test.ts +++ b/packages/core/src/runtime/helpers.test.ts @@ -414,7 +414,7 @@ describe('loadWorkflowRunEvents', () => { hasMore: false, }); - const result = await loadWorkflowRunEvents('wrun_test'); + const result = await loadWorkflowRunEvents({ runId: 'wrun_test' }); expect(result.events).toHaveLength(2); expect(result.cursor).toBe('eid:evnt_b'); @@ -447,7 +447,7 @@ describe('loadWorkflowRunEvents', () => { hasMore: false, }); - const result = await loadWorkflowRunEvents('wrun_test'); + const result = await loadWorkflowRunEvents({ runId: 'wrun_test' }); expect(result.events).toHaveLength(2); expect(result.cursor).toBe('eid:evnt_b'); @@ -461,7 +461,7 @@ describe('loadWorkflowRunEvents', () => { hasMore: false, }); - const result = await loadWorkflowRunEvents('wrun_test'); + const result = await loadWorkflowRunEvents({ runId: 'wrun_test' }); expect(result.events).toHaveLength(0); expect(result.cursor).toBeNull(); @@ -484,7 +484,7 @@ describe('loadWorkflowRunEvents', () => { hasMore: false, }); - const result = await loadWorkflowRunEvents('wrun_test'); + const result = await loadWorkflowRunEvents({ runId: 'wrun_test' }); expect(result.events.map((e) => e.eventId)).toEqual([ 'evnt_a', @@ -501,7 +501,10 @@ describe('loadWorkflowRunEvents', () => { hasMore: false, }); - const result = await loadWorkflowRunEvents('wrun_test', 'eid:evnt_z'); + const result = await loadWorkflowRunEvents({ + runId: 'wrun_test', + afterCursor: 'eid:evnt_z', + }); expect(result.events).toHaveLength(0); // Preserving the input cursor avoids the runtime treating "no new events @@ -521,7 +524,7 @@ describe('loadWorkflowRunEvents', () => { hasMore: false, }); - const result = await loadWorkflowRunEvents('wrun_test'); + const result = await loadWorkflowRunEvents({ runId: 'wrun_test' }); expect(result.events.map((event) => event.eventId)).toEqual([ 'evnt_a', @@ -540,7 +543,10 @@ describe('loadWorkflowRunEvents', () => { hasMore: false, }); - const result = await loadWorkflowRunEvents('wrun_test', 'opaque-cursor'); + const result = await loadWorkflowRunEvents({ + runId: 'wrun_test', + afterCursor: 'opaque-cursor', + }); expect(result.events.map((event) => event.eventId)).toEqual([ 'evnt_a', @@ -568,7 +574,9 @@ describe('loadWorkflowRunEvents', () => { hasMore: true, }); - await expect(loadWorkflowRunEvents('wrun_test')).rejects.toMatchObject({ + await expect( + loadWorkflowRunEvents({ runId: 'wrun_test' }) + ).rejects.toMatchObject({ code: 'WORLD_CONTRACT_ERROR', }); expect(eventsListMock).toHaveBeenCalledTimes(2); @@ -581,7 +589,9 @@ describe('loadWorkflowRunEvents', () => { hasMore: true, }); - await expect(loadWorkflowRunEvents('wrun_test')).rejects.toMatchObject({ + await expect( + loadWorkflowRunEvents({ runId: 'wrun_test' }) + ).rejects.toMatchObject({ code: 'WORLD_CONTRACT_ERROR', }); expect(eventsListMock).toHaveBeenCalledTimes(1); diff --git a/packages/core/src/runtime/helpers.ts b/packages/core/src/runtime/helpers.ts index c82cdee697..3a2fe82b1b 100644 --- a/packages/core/src/runtime/helpers.ts +++ b/packages/core/src/runtime/helpers.ts @@ -8,7 +8,6 @@ import type { CreateEventRequest, Event, EventResult, - EventStreamObserver, HealthCheckPayload, ValidQueueName, WorkflowRun, @@ -593,11 +592,15 @@ function shouldRetryWithoutEventCursor( * The returned cursor can be passed back in on a subsequent call for * incremental loading. */ -export async function loadWorkflowRunEvents( - runId: string, - afterCursor?: string, - onEvent?: EventStreamObserver -): Promise { +export async function loadWorkflowRunEvents({ + runId, + afterCursor, + onEvent, +}: { + runId: string; + afterCursor?: string; + onEvent?: (event: Event) => void; +}): Promise { const incremental = afterCursor !== undefined; return trace( incremental ? 'workflow.loadNewEvents' : 'workflow.loadEvents', @@ -984,7 +987,7 @@ export async function settleEventSlotGap( await new Promise((resolve) => setTimeout(resolve, SLOT_GAP_RECHECK_BASE_DELAY_MS * 2 ** attempt) ); - log = await loadWorkflowRunEvents(runId); + log = await loadWorkflowRunEvents({ runId }); gap = findEventSlotGap(log.events); } return { log, gap }; From eda627a9bdfb7ed3b17ac5c20e7e00ba86f65698 Mon Sep 17 00:00:00 2001 From: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> Date: Fri, 14 Aug 2026 20:29:12 -0700 Subject: [PATCH 8/9] Simplify streamed replay payload preparation --- .../core/src/replay-payload-cache.test.ts | 104 ++++------ packages/core/src/replay-payload-cache.ts | 189 ++++++------------ packages/core/src/runtime.ts | 122 +++++------ packages/core/src/runtime/helpers.ts | 33 +-- packages/core/src/workflow.ts | 2 +- .../core/src/workflow/abort-controller.ts | 22 +- 6 files changed, 194 insertions(+), 278 deletions(-) diff --git a/packages/core/src/replay-payload-cache.test.ts b/packages/core/src/replay-payload-cache.test.ts index e8d48e58ca..dc2f36eede 100644 --- a/packages/core/src/replay-payload-cache.test.ts +++ b/packages/core/src/replay-payload-cache.test.ts @@ -57,14 +57,16 @@ describe('ReplayPayloadCache', () => { it('deduplicates synchronous preparation without creating a promise', () => { const payload = new Uint8Array([1]); const preparer = vi.fn((value) => value); + const hydrate = vi.fn((prepared: unknown) => prepared); const cache = new ReplayPayloadCache(undefined, preparer); - const first = cache.prepareEventPayload('evnt_one', payload); - const second = cache.prepareEventPayload('evnt_one', payload); + const first = cache.getEventValue('evnt_one', payload, hydrate); + const second = cache.getEventValue('evnt_one', payload, hydrate); expect(first).toBe(second); expect(first).toEqual(payload); expect(preparer).toHaveBeenCalledOnce(); + expect(hydrate).toHaveBeenCalledTimes(2); }); it('hydrates and memoizes a primitive without creating a promise', () => { @@ -86,12 +88,12 @@ describe('ReplayPayloadCache', () => { .mockReturnValueOnce(payload); const cache = new ReplayPayloadCache(undefined, preparer); - cache.prewarm(run, []); + cache.prepareAll(run, []); await Promise.resolve(); - expect(() => cache.prepareWorkflowInput(run)).toThrow('decrypt failed'); + await expect(cache.getWorkflowInput(run)).rejects.toThrow('decrypt failed'); expect(preparer).toHaveBeenCalledOnce(); - expect(cache.prepareWorkflowInput(run)).toEqual(payload); + expect(cache.getWorkflowInput(run)).toEqual(payload); expect(preparer).toHaveBeenCalledTimes(2); }); @@ -108,27 +110,30 @@ describe('ReplayPayloadCache', () => { const run = makeRun(payloads[0]); const events = makeEvents(payloads.slice(1)); - cache.prewarm(run, events); + cache.prepareAll(run, events); expect(preparer).toHaveBeenCalledTimes(4); for (const resolve of resolvers.reverse()) resolve(); await Promise.all([ - cache.prepareWorkflowInput(run), + cache.getWorkflowInput(run), ...events.map((event) => { switch (event.eventType) { case 'step_completed': - return cache.prepareEventPayload( + return cache.getEventValue( event.eventId, - event.eventData?.result + event.eventData?.result, + (prepared) => prepared ); case 'step_failed': - return cache.prepareEventPayload( + return cache.getEventValue( event.eventId, - event.eventData?.error + event.eventData?.error, + (prepared) => prepared ); case 'hook_received': - return cache.prepareEventPayload( + return cache.getEventValue( event.eventId, - event.eventData?.payload + event.eventData?.payload, + (prepared) => prepared ); default: throw new Error(`Unexpected event: ${event.eventType}`); @@ -136,7 +141,7 @@ describe('ReplayPayloadCache', () => { }), ]); - cache.prewarm(run, events); + cache.prepareAll(run, events); expect(preparer).toHaveBeenCalledTimes(4); }); @@ -150,34 +155,16 @@ describe('ReplayPayloadCache', () => { const cache = new ReplayPayloadCache(undefined, preparer); const [event] = makeEvents([payload]); - cache.observeEvent(event, () => order.push('start')); + cache.prepareEvent(event); expect(preparer).toHaveBeenCalledOnce(); - expect(order).toEqual(['start', 'prepare']); + expect(order).toEqual(['prepare']); - // Re-observing a cache hit does not move the preparation-span boundary. - cache.observeEvent(event, () => order.push('cached-start')); - expect(order).toEqual(['start', 'prepare']); + cache.prepareEvent(event); + expect(order).toEqual(['prepare']); - expect(cache.prepareEventPayload(event.eventId, payload)).toEqual(payload); - }); - - it('prepares queued stream events as soon as the run key resolves', async () => { - const payload = new Uint8Array([1]); - const preparer = vi.fn((value) => value); - let resolveKey!: (key: undefined) => void; - const key = new Promise((resolve) => { - resolveKey = resolve; - }); - const cache = ReplayPayloadCache.waitingForKey(key, preparer); - const [event] = makeEvents([payload]); - - cache.observeEvent(event); - expect(preparer).not.toHaveBeenCalled(); - const preparation = cache.prepareEventPayload(event.eventId, payload); - - resolveKey(undefined); - await expect(preparation).resolves.toEqual(payload); - expect(preparer).toHaveBeenCalledOnce(); + expect( + cache.getEventValue(event.eventId, payload, (prepared) => prepared) + ).toEqual(payload); }); it('caches real decrypt/decompress output but revives fresh objects', async () => { @@ -199,13 +186,15 @@ describe('ReplayPayloadCache', () => { expect(directPreparation).not.toBeInstanceOf(Promise); await directPreparation; - const prepared = await cache.prepareEventPayload( + const prepared = await cache.getEventValue( 'evnt_encrypted', - serialized + serialized, + (value) => value ); - const samePrepared = await cache.prepareEventPayload( + const samePrepared = await cache.getEventValue( 'evnt_encrypted', - serialized + serialized, + (value) => value ); const first = deserializePreparedReplayPayload(prepared) as { count: number; @@ -220,45 +209,36 @@ describe('ReplayPayloadCache', () => { expect(second.count).toBe(0); }); - it('rescans a log whose missing events were filled in below the scanned prefix', async () => { - // A stale-snapshot (412) restart replaces the log with a corrected one, so - // the events it was missing appear BELOW the length already scanned and - // shift every later position. Resuming from that length skips exactly the - // events the reload was for, which is what `resetScan` exists to prevent. + it('finds events inserted below a previously prepared prefix', () => { + // A stale-snapshot restart can replace the log with a corrected one whose + // missing events appear below the old tail. Full scans are cheap because + // event-id cache hits do no payload work. const payloads = [0, 1, 2].map((value) => new Uint8Array([value])); const preparer = vi.fn((value) => value); const cache = new ReplayPayloadCache(undefined, preparer); const run = makeRun(undefined); const [first, missing, second] = makeEvents(payloads); - await cache.prewarm(run, [first, second]); - expect(preparer).toHaveBeenCalledTimes(2); - - // Positional resume: `missing` sits inside the scanned prefix, so it is - // skipped and its payload is only prepared on demand. - await cache.prewarm(run, [first, missing, second]); + cache.prepareAll(run, [first, second]); expect(preparer).toHaveBeenCalledTimes(2); - cache.resetScan(); - await cache.prewarm(run, [first, missing, second]); - // Only the inserted event is new: the other two are keyed by event id and - // stay prepared across the rescan. + cache.prepareAll(run, [first, missing, second]); expect(preparer).toHaveBeenCalledTimes(3); expect(preparer).toHaveBeenLastCalledWith(payloads[1], undefined); }); - it('bypasses legacy values and ignores missing event data during prewarm', async () => { + it('bypasses legacy values and ignores missing event data during preparation', async () => { const legacy = [0, { value: 1 }]; const preparer = vi.fn((value) => value); const cache = new ReplayPayloadCache(undefined, preparer); - await cache.prepareEventPayload('evnt_legacy', legacy); - await cache.prepareEventPayload('evnt_legacy', legacy); + await cache.getEventValue('evnt_legacy', legacy, (prepared) => prepared); + await cache.getEventValue('evnt_legacy', legacy, (prepared) => prepared); expect(preparer).not.toHaveBeenCalled(); const events = makeEvents([legacy, legacy, legacy]); events[2] = { ...events[2], eventData: undefined } as unknown as Event; - await cache.prewarm(makeRun(legacy), events); + cache.prepareAll(makeRun(legacy), events); expect(preparer).not.toHaveBeenCalled(); }); diff --git a/packages/core/src/replay-payload-cache.ts b/packages/core/src/replay-payload-cache.ts index 445f484a52..c5d91622b7 100644 --- a/packages/core/src/replay-payload-cache.ts +++ b/packages/core/src/replay-payload-cache.ts @@ -5,7 +5,9 @@ import { prepareReplayPayload, } from './serialization/replay.js'; -const WORKFLOW_INPUT_CACHE_KEY = 'workflow-input'; +const WORKFLOW_INPUT = Symbol('workflow-input'); +type ReplayPayloadKey = string | typeof WORKFLOW_INPUT; +type CachedPreparation = Uint8Array | Promise; function isCacheablePrimitive(value: unknown): boolean { const type = typeof value; @@ -22,109 +24,54 @@ function isCacheablePrimitive(value: unknown): boolean { * Deserialization still runs against each VM's globals so object graphs and * Workflow objects remain realm-local. Primitive final values are safe to * share and skip that repeated deserialization entirely. + * + * Key lookup is deliberately outside this class. The runtime creates the + * cache once the run's key has resolved, then feeds it decoded events. Most + * Node preparation is synchronous; only codecs that are inherently async + * leave a Promise in the cache. */ export class ReplayPayloadCache { private readonly preparations = new Map< - string, - Uint8Array | Promise | { readonly error: unknown } + ReplayPayloadKey, + CachedPreparation >(); private readonly primitiveValues = new Map(); - private nextUnscannedEventIndex = 0; - private encryptionKeyPromise?: Promise; constructor( - private encryptionKey?: DecryptionKey, + private readonly encryptionKey?: DecryptionKey, private readonly preparer: typeof prepareReplayPayload = prepareReplayPayload ) {} - static waitingForKey( - key: Promise, - preparer: typeof prepareReplayPayload = prepareReplayPayload - ): ReplayPayloadCache { - const cache = new ReplayPayloadCache(undefined, preparer); - cache.encryptionKeyPromise = key; - void key.then( - (value) => { - cache.encryptionKey = value; - cache.encryptionKeyPromise = undefined; - }, - () => {} - ); - return cache; - } - - /** Start preparing an event as soon as its frame has been decoded. */ - observeEvent(event: Event, onPreparationStart?: () => void): void { + /** Prepare a payload as soon as its event frame has been decoded. */ + prepareEvent(event: Event): void { switch (event.eventType) { case 'run_created': - this.startPreparation( - WORKFLOW_INPUT_CACHE_KEY, - event.eventData.input, - onPreparationStart - ); + this.cachePayload(WORKFLOW_INPUT, event.eventData.input); break; case 'run_started': - this.startPreparation( - WORKFLOW_INPUT_CACHE_KEY, - event.eventData?.input, - onPreparationStart - ); + this.cachePayload(WORKFLOW_INPUT, event.eventData?.input); break; case 'step_completed': - this.startPreparation( - this.eventPayloadKey(event.eventId), - event.eventData?.result, - onPreparationStart - ); + this.cachePayload(event.eventId, event.eventData?.result); break; case 'step_failed': - this.startPreparation( - this.eventPayloadKey(event.eventId), - event.eventData?.error, - onPreparationStart - ); + this.cachePayload(event.eventId, event.eventData?.error); break; case 'hook_received': - this.startPreparation( - this.eventPayloadKey(event.eventId), - event.eventData?.payload, - onPreparationStart - ); - } - } - - /** - * Start every preparation not already observed from the event stream. - * Consumers await the few codecs that cannot complete synchronously. - */ - prewarm(workflowRun: WorkflowRun, events: Event[]): void { - this.startPreparation(WORKFLOW_INPUT_CACHE_KEY, workflowRun.input); - for ( - let index = this.nextUnscannedEventIndex; - index < events.length; - index++ - ) { - this.observeEvent(events[index]); + this.cachePayload(event.eventId, event.eventData?.payload); } - this.nextUnscannedEventIndex = events.length; } - /** A corrected reload may insert events before the previous scan position. */ - resetScan(): void { - this.nextUnscannedEventIndex = 0; + /** Prepare every payload not already seen through the event stream. */ + prepareAll(workflowRun: WorkflowRun, events: Event[]): void { + this.cachePayload(WORKFLOW_INPUT, workflowRun.input); + for (const event of events) this.prepareEvent(event); } - prepareWorkflowInput( + getWorkflowInput( workflowRun: WorkflowRun ): PreparedReplayPayload | Promise { - return this.getPreparedPayload(WORKFLOW_INPUT_CACHE_KEY, workflowRun.input); - } - - prepareEventPayload( - eventId: string, - value: unknown - ): PreparedReplayPayload | Promise { - return this.getPreparedPayload(this.eventPayloadKey(eventId), value); + return this.getPayload(WORKFLOW_INPUT, workflowRun.input); } getEventValue( @@ -132,93 +79,75 @@ export class ReplayPayloadCache { serializedValue: unknown, hydrate: (prepared: PreparedReplayPayload) => unknown | Promise ): unknown | Promise { - const cacheKey = this.eventPayloadKey(eventId); - if (this.primitiveValues.has(cacheKey)) { - return this.primitiveValues.get(cacheKey); + if (this.primitiveValues.has(eventId)) { + return this.primitiveValues.get(eventId); } - const prepared = this.prepareEventPayload(eventId, serializedValue); + const prepared = this.getPayload(eventId, serializedValue); const hydrateAndCache = (payload: PreparedReplayPayload) => { const hydrated = hydrate(payload); return hydrated instanceof Promise - ? hydrated.then((value) => this.cachePrimitive(cacheKey, value)) - : this.cachePrimitive(cacheKey, hydrated); + ? hydrated.then((value) => this.cachePrimitive(eventId, value)) + : this.cachePrimitive(eventId, hydrated); }; return prepared instanceof Promise ? prepared.then(hydrateAndCache) : hydrateAndCache(prepared); } - private cachePrimitive(cacheKey: string, value: unknown): unknown { - if (isCacheablePrimitive(value)) this.primitiveValues.set(cacheKey, value); + private cachePrimitive(eventId: string, value: unknown): unknown { + if (isCacheablePrimitive(value)) { + this.primitiveValues.set(eventId, value); + } return value; } - private startPreparation( - cacheKey: string, - value: unknown, - onPreparationStart?: () => void - ): void { + private cachePayload(cacheKey: ReplayPayloadKey, value: unknown): void { if (!(value instanceof Uint8Array) || this.preparations.has(cacheKey)) { return; } - onPreparationStart?.(); - this.runPreparation(cacheKey, value); - } - - private runPreparation(cacheKey: string, value: Uint8Array): void { + let preparation: CachedPreparation; try { - const result = this.encryptionKeyPromise - ? this.encryptionKeyPromise.then((key) => this.preparer(value, key)) - : this.preparer(value, this.encryptionKey); - if (!(result instanceof Promise)) { - this.preparations.set(cacheKey, result); - return; - } + preparation = this.preparer(value, this.encryptionKey); + } catch (error) { + // Preparation is speculative. Preserve a synchronous failure for the + // ordered consumer without failing event loading or creating an + // unhandled rejection. + preparation = Promise.reject(error); + } + this.preparations.set(cacheKey, preparation); - this.preparations.set(cacheKey, result); - void result.then( + if (preparation instanceof Promise) { + void preparation.then( (prepared) => { - const current = this.preparations.get(cacheKey); - if (current === result) this.preparations.set(cacheKey, prepared); + if (this.preparations.get(cacheKey) === preparation) { + this.preparations.set(cacheKey, prepared); + } }, - (error) => { - const current = this.preparations.get(cacheKey); - if (current === result) this.preparations.set(cacheKey, { error }); - } + () => {} ); - } catch (error) { - this.preparations.set(cacheKey, { error }); } } - private getPreparedPayload( - cacheKey: string, + private getPayload( + cacheKey: ReplayPayloadKey, value: unknown ): PreparedReplayPayload | Promise { if (!(value instanceof Uint8Array)) return { legacy: value }; - this.startPreparation(cacheKey, value); + this.cachePayload(cacheKey, value); const prepared = this.preparations.get(cacheKey); if (!prepared) { - throw new Error( - `Replay payload preparation was not started: ${cacheKey}` - ); + throw new Error('Replay payload preparation was not cached'); } - if (prepared instanceof Uint8Array) return prepared; - if (prepared instanceof Promise) { - return prepared.catch((error) => { + if (!(prepared instanceof Promise)) return prepared; + return prepared.catch((error) => { + if (this.preparations.get(cacheKey) === prepared) { this.preparations.delete(cacheKey); - throw error; - }); - } - this.preparations.delete(cacheKey); - throw prepared.error; - } - - private eventPayloadKey(eventId: string): string { - return `event:${eventId}`; + } + throw error; + }); } } diff --git a/packages/core/src/runtime.ts b/packages/core/src/runtime.ts index 6f58607cf1..645e1d722b 100644 --- a/packages/core/src/runtime.ts +++ b/packages/core/src/runtime.ts @@ -14,7 +14,7 @@ import { WorkflowRuntimeError, WorkflowWorldError, } from '@workflow/errors'; -import { once, setWorkflowBasePath, withResolvers } from '@workflow/utils'; +import { once, setWorkflowBasePath } from '@workflow/utils'; import { parseWorkflowName, workflowDisplayName, @@ -81,6 +81,7 @@ import { parseHealthCheckPayload, preconditionEventDelta, queueMessage, + resolveRunEncryptionKey, type SlotSnapshotParams, settleEventSlotGap, slotSnapshotParams, @@ -118,6 +119,7 @@ import { type WorldHandlers, } from './runtime/world.js'; import { dehydrateRunError } from './serialization.js'; +import type { DecryptionKey } from './serialization/encryption.js'; import { remapErrorStack } from './source-map.js'; import * as Attribute from './telemetry/semantic-conventions.js'; import { @@ -547,25 +549,6 @@ function replayEventDeploymentId(event: Event): string | undefined { return undefined; } -function createReplayEventObserver({ - runId, - cache, - resolveKey, -}: { - runId: string; - cache: ReplayPayloadCache; - resolveKey: ( - runOrId: WorkflowRun | string, - context?: Record - ) => void; -}): (event: Event) => void { - return (event) => { - const deploymentId = replayEventDeploymentId(event); - if (deploymentId) resolveKey(runId, { deploymentId }); - cache.observeEvent(event); - }; -} - /** * The whole retention predicate: keep the session only for a pure step * boundary (every suspension item is a step — any other item type, present @@ -1045,32 +1028,55 @@ export function workflowEntrypoint( // queue payload. This starts key resolution and payload // preparation before the remainder of the event log arrives, // without guessing the key for a cross-deployment run. - const { - promise: replayKeySource, - resolve: resolveReplayKeySource, - } = withResolvers<{ - runOrId: WorkflowRun | string; - context?: Record; - }>(); - const resolveReplayKey = ( + let replayEncryptionKey: + | Promise + | undefined; + let streamedPayloadCache: ReplayPayloadCache | undefined; + const eventsWaitingForKey: Event[] = []; + const activatePayloadCache = ( + key: DecryptionKey | undefined + ): ReplayPayloadCache => { + if (!streamedPayloadCache) { + streamedPayloadCache = new ReplayPayloadCache(key); + for (const event of eventsWaitingForKey) { + streamedPayloadCache.prepareEvent(event); + } + eventsWaitingForKey.length = 0; + } + return streamedPayloadCache; + }; + const getReplayEncryptionKey = ( runOrId: WorkflowRun | string, context?: Record - ): void => { - resolveReplayKeySource({ runOrId, context }); + ): Promise => { + if (!replayEncryptionKey) { + replayEncryptionKey = resolveRunEncryptionKey( + world, + runOrId, + context + ); + void replayEncryptionKey.then( + activatePayloadCache, + () => {} + ); + } + return replayEncryptionKey; + }; + const onReplayEvent = (event: Event): void => { + const deploymentId = replayEventDeploymentId(event); + if (deploymentId) { + void getReplayEncryptionKey(runId, { + deploymentId, + }); + } + if (streamedPayloadCache) { + streamedPayloadCache.prepareEvent(event); + } else { + eventsWaitingForKey.push(event); + } }; - const encryptionKeyPromise = replayKeySource.then( - ({ runOrId, context }) => - memoizeEncryptionKey(world, runOrId, context)() - ); - const replayPayloadCache = - ReplayPayloadCache.waitingForKey(encryptionKeyPromise); - const observeReplayEvent = createReplayEventObserver({ - runId, - cache: replayPayloadCache, - resolveKey: resolveReplayKey, - }); if (runInput?.deploymentId) { - resolveReplayKey(runId, { + void getReplayEncryptionKey(runId, { deploymentId: runInput.deploymentId, }); } @@ -1395,10 +1401,6 @@ export function workflowEntrypoint( // incremental load starts above the hole and never // returns it. eventLog = { type: 'loadAll' }; - // The corrected log inserts the missing events BELOW the - // length already scanned for payload prewarming, shifting - // every later position. Only a full rescan sees them. - replayPayloadCache.resetScan(); } runtimeLogger.warn( 'Event creation rejected as stale; restarting replay in-process', @@ -2045,7 +2047,7 @@ export function workflowEntrypoint( resumeId: hookResumeInput.resumeId, resumePayloadDigest: hookResumeInput.payloadDigest, preloadEvents: true, - onEvent: observeReplayEvent, + onEvent: onReplayEvent, } ); hookEnsured = true; @@ -2320,7 +2322,7 @@ export function workflowEntrypoint( }); const result = await createEvent(runStartedEvent, { requestId, - onEvent: observeReplayEvent, + onEvent: onReplayEvent, }); workflowRun = result.run; maxEventsLimit = clampMaxEvents(result.maxEvents); @@ -2605,7 +2607,7 @@ export function workflowEntrypoint( // do we fall back to reloading the complete log. if (eventLog.type !== 'loadAll' && ensuredEvent) { insertEventByEventId(eventLog.events, ensuredEvent); - observeReplayEvent(ensuredEvent); + onReplayEvent(ensuredEvent); } else { eventLog = { type: 'loadAll' }; } @@ -2615,8 +2617,10 @@ export function workflowEntrypoint( // Worlds that do not implement streamed observation still // resolve from the materialized run. This is also the final // cross-deployment-safe source of truth. - resolveReplayKey(workflowRun); - const encryptionKey = await encryptionKeyPromise; + const encryptionKey = + await getReplayEncryptionKey(workflowRun); + const replayPayloadCache = + activatePayloadCache(encryptionKey); // The live VM parked at the previous boundary, when the // retention decision kept it. null → this iteration cold- @@ -2714,7 +2718,7 @@ export function workflowEntrypoint( await loadWorkflowRunEvents({ runId, afterCursor: eventLog.cursor, - onEvent: observeReplayEvent, + onEvent: onReplayEvent, }) ); eventLog = { ...eventLog, type: 'ready' }; @@ -2773,7 +2777,7 @@ export function workflowEntrypoint( eventLog.type === 'loadAfter' ? eventLog.cursor : undefined, - onEvent: observeReplayEvent, + onEvent: onReplayEvent, }); if (eventLog.type === 'loadAfter') { appendEventLog(eventLog, page); @@ -2894,7 +2898,7 @@ export function workflowEntrypoint( const page = await loadWorkflowRunEvents({ runId, afterCursor: eventLog.cursor, - onEvent: observeReplayEvent, + onEvent: onReplayEvent, }); const completedWaitIdsAfterCursor = new Set( page.events @@ -2914,7 +2918,7 @@ export function workflowEntrypoint( eventLog = { ...(await loadWorkflowRunEvents({ runId, - onEvent: observeReplayEvent, + onEvent: onReplayEvent, })), type: 'ready', }; @@ -2923,7 +2927,7 @@ export function workflowEntrypoint( eventLog = { ...(await loadWorkflowRunEvents({ runId, - onEvent: observeReplayEvent, + onEvent: onReplayEvent, })), type: 'ready', }; @@ -3013,7 +3017,7 @@ export function workflowEntrypoint( // appended-event consumption on resume; consumers still // deserialize and resolve in event order. const replayEvents = eventLog.events; - replayPayloadCache.prewarm(workflowRun, replayEvents); + replayPayloadCache.prepareAll(workflowRun, replayEvents); let workflowResult: WorkflowResumeResult = retainedSession ? await resumeWorkflow(retainedSession, eventLog.events) : { type: 'replay' }; @@ -3290,12 +3294,10 @@ export function workflowEntrypoint( } if (suspensionResult.reportedEventCount > 0) { // Bump-and-report merged events BELOW the tail and - // re-sorted the array to slot order, shifting every - // position the prewarm scan had already recorded. + // re-sorted the array to slot order. // The cursor is deliberately left alone: the report // is a lower bound on what was skipped, so the next // incremental read still has to cover the same range. - replayPayloadCache.resetScan(); } // Open hooks/waits in the log as loaded for this diff --git a/packages/core/src/runtime/helpers.ts b/packages/core/src/runtime/helpers.ts index 3a2fe82b1b..f774ed0c6b 100644 --- a/packages/core/src/runtime/helpers.ts +++ b/packages/core/src/runtime/helpers.ts @@ -1250,6 +1250,24 @@ export function getQueueOverhead(message: { requestedAt?: Date }) { * outer try/catch to log and surface the issue; the queue's redelivery * semantics will retry the key fetch on the next attempt. */ +export async function resolveRunEncryptionKey( + world: World, + runOrId: WorkflowRun | string, + context?: Record +): Promise { + // The `getEncryptionKeyForRun` overload set takes either a `WorkflowRun` or + // a `runId: string` (with optional context). Branch here so TypeScript picks + // the right overload for each shape. + const rawKey = + typeof runOrId === 'string' + ? await world.getEncryptionKeyForRun?.(runOrId, context) + : await world.getEncryptionKeyForRun?.(runOrId); + // Resolve the *full* capability, not just the symmetric key: a run reading + // its own event log may encounter sealed (`encp`) payloads that another run + // wrote to it, and opening those needs the run's X25519 scalar as well. + return rawKey ? await deriveRunPayloadKeys(rawKey) : undefined; +} + export function memoizeEncryptionKey( world: World, runOrId: WorkflowRun | string, @@ -1258,20 +1276,7 @@ export function memoizeEncryptionKey( let cached: Promise | undefined; return () => { if (!cached) { - cached = (async () => { - // The `getEncryptionKeyForRun` overload set takes either a - // `WorkflowRun` or a `runId: string` (with optional context). Branch - // here so TypeScript picks the right overload for each shape. - const rawKey = - typeof runOrId === 'string' - ? await world.getEncryptionKeyForRun?.(runOrId, context) - : await world.getEncryptionKeyForRun?.(runOrId); - // Resolve the *full* capability, not just the symmetric key: a run - // reading its own event log may encounter sealed (`encp`) payloads - // that another run wrote to it, and opening those needs the run's - // X25519 scalar as well. - return rawKey ? await deriveRunPayloadKeys(rawKey) : undefined; - })(); + cached = resolveRunEncryptionKey(world, runOrId, context); } return cached; }; diff --git a/packages/core/src/workflow.ts b/packages/core/src/workflow.ts index 3bcd3f08fa..de9a248cf1 100644 --- a/packages/core/src/workflow.ts +++ b/packages/core/src/workflow.ts @@ -1085,7 +1085,7 @@ async function createWorkflowSession({ // workflow function subscribing its first step callbacks. let args: unknown[] = []; workflowContext.promiseQueue = workflowContext.promiseQueue.then(async () => { - const prepared = await replayPayloadCache.prepareWorkflowInput(workflowRun); + const prepared = await replayPayloadCache.getWorkflowInput(workflowRun); args = await hydrateWorkflowArguments( workflowRun.input, workflowRun.runId, diff --git a/packages/core/src/workflow/abort-controller.ts b/packages/core/src/workflow/abort-controller.ts index 02b4f39bba..6b7cacd6cb 100644 --- a/packages/core/src/workflow/abort-controller.ts +++ b/packages/core/src/workflow/abort-controller.ts @@ -230,18 +230,18 @@ export function createCreateAbortController(ctx: WorkflowOrchestratorContext) { try { if (rawPayload !== undefined) { try { - const prepared = - await ctx.replayPayloadCache.prepareEventPayload( - event.eventId, - rawPayload - ); - const hydrated = (await hydrateStepReturnValue( + const hydrated = (await ctx.replayPayloadCache.getEventValue( + event.eventId, rawPayload, - ctx.runId, - ctx.encryptionKey, - ctx.globalThis, - {}, - prepared + (prepared) => + hydrateStepReturnValue( + rawPayload, + ctx.runId, + ctx.encryptionKey, + ctx.globalThis, + {}, + prepared + ) )) as { reason?: unknown } | undefined; if ( hydrated && From 2c3379f69ef20c827f287fc945b9320907a971bc Mon Sep 17 00:00:00 2001 From: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> Date: Fri, 14 Aug 2026 20:32:23 -0700 Subject: [PATCH 9/9] Sort replay runtime imports --- packages/core/src/runtime.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/core/src/runtime.ts b/packages/core/src/runtime.ts index 645e1d722b..fe8e1e00e7 100644 --- a/packages/core/src/runtime.ts +++ b/packages/core/src/runtime.ts @@ -118,8 +118,8 @@ import { getWorldHandlers, type WorldHandlers, } from './runtime/world.js'; -import { dehydrateRunError } from './serialization.js'; import type { DecryptionKey } from './serialization/encryption.js'; +import { dehydrateRunError } from './serialization.js'; import { remapErrorStack } from './source-map.js'; import * as Attribute from './telemetry/semantic-conventions.js'; import {