Skip to content

Commit 49da6c5

Browse files
authored
feat(core): support passing parent WritableStream to child workflow via start() (#2059)
* test(e2e): cover WritableStream passed as start() argument Adds an e2e workflow + test where a parent workflow gets a WritableStream via getWritable(), forwards it through start() to a child workflow, and the child step writes raw bytes to it. Asserts the external reader on the parent's stream observes the exact bytes the child wrote. * fix(core): avoid double-framing when WritableStream is forwarded via start() When a workflow's getWritable() handle is passed across start() to a child workflow, the parent step's reviver wraps it in a serialize transform that pipes into a workflow server stream. Until now, getExternalReducers.WritableStream then installed a second serialize transform on top of that — so every chunk the child step wrote got devalue-framed twice but only deframed once on the reader side, and external consumers saw the inner frame instead of the original bytes. Fix: tag every user-visible writable that's already backed by a workflow server stream with its (runId, name). When the external reducer recognizes those tags during dehydration, it bridges bytes straight from the new child-side server stream to the original server stream instead of piping through the user's writable. That leaves the producer-side serialize transform (installed once by the child's step reviver) as the only framing layer in the chain. * fix(core): forward (runId, name) when a tagged WritableStream crosses start() Replaces the previous in-process bridge with first-class writable forwarding at the descriptor level. When a parent workflow's getWritable() handle is passed as an argument to a child workflow, the dehydrated descriptor now carries the original (runId, name). The child run's step-side reviver opens the writable against the parent's server stream directly and resolves the parent run's encryption key (encrypt-only) via getEncryptionKeyForRun. This removes the architectural limitation that the bridge could only stay alive for the duration of the parent step process — on Vercel that capped forwarding at ~15 minutes regardless of the child run's lifetime, dropping any writes the child made after the parent step process exited. importKey() now accepts a usages parameter, defaulting to ['encrypt', 'decrypt']. The cross-run forwarding path imports with ['encrypt'] only so a compromised child run cannot decrypt any existing data on the parent's stream — only contribute new writes. * test: rename writable-forwarded workflows and cover step-context getWritable() Addresses PR review: - Rename writableForwardedToChildChildWorkflow → writableForwardedChildWorkflow (drops the duplicated 'Child' segment). - Split writableForwardedToChildWorkflow into two variants covered by a test.each: writableForwardedFromWorkflowWorkflow (workflow-context getWritable, the original test) and writableForwardedFromStepWorkflow (step-context getWritable passed directly into start() from the same step that called getWritable()). - Terser changeset description.
1 parent 1d3959e commit 49da6c5

9 files changed

Lines changed: 343 additions & 17 deletions

File tree

‎.changeset/tired-pigs-hug.md‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
---
2+
"@workflow/core": minor
3+
"workflow": minor
4+
---
5+
6+
A `WritableStream` from a workflow's `getWritable()` can now be passed as an argument to a child workflow via `start()`; the child's writes land on the parent run's stream directly for the full lifetime of the child run.

‎packages/core/e2e/e2e.test.ts‎

Lines changed: 44 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -927,6 +927,50 @@ describe('e2e', () => {
927927
expect(await run.returnValue).toEqual('done');
928928
});
929929

930+
// A WritableStream passed as a workflow argument to start() should
931+
// land raw bytes on the parent's output stream when the child step
932+
// writes to it. Covered for both:
933+
// - `writableForwardedFromWorkflowWorkflow`: parent calls
934+
// `getWritable()` in workflow context (fake handle revived in the
935+
// intermediary step).
936+
// - `writableForwardedFromStepWorkflow`: parent calls `getWritable()`
937+
// in step context (real `serialize.writable` passed straight to
938+
// `start()`).
939+
test.each([
940+
'writableForwardedFromWorkflowWorkflow',
941+
'writableForwardedFromStepWorkflow',
942+
] as const)('%s', { timeout: 120_000 }, async (workflowName) => {
943+
const payload = `hello-from-child-${Date.now()}\n`;
944+
const run = await start(await e2e(workflowName), [payload]);
945+
946+
const reader = run.getReadable().getReader();
947+
// `fatal: true` makes the decoder throw on any invalid UTF-8
948+
// sequence, so a successful decode is itself a round-trip
949+
// assertion that the bytes survived intact.
950+
const decoder = new TextDecoder('utf-8', { fatal: true });
951+
952+
// The child step performs exactly one write of `payload` as
953+
// UTF-8 bytes, so we should receive a single chunk containing
954+
// exactly those bytes before the stream closes.
955+
const { value, done } = await reader.read();
956+
expect(done).toBeFalsy();
957+
assert(value);
958+
assert(value instanceof Uint8Array);
959+
960+
const expectedBytes = new TextEncoder().encode(payload);
961+
expect(value.byteLength).toBe(expectedBytes.byteLength);
962+
expect(decoder.decode(value)).toBe(payload);
963+
964+
// Default stream should close cleanly after the parent closes its
965+
// writable.
966+
expect((await reader.read()).done).toBe(true);
967+
968+
const returnValue = await run.returnValue;
969+
expect(returnValue).toMatchObject({
970+
childRunId: expect.stringMatching(/^wrun_/),
971+
});
972+
});
973+
930974
test('fetchWorkflow', { timeout: 60_000 }, async () => {
931975
const run = await start(await e2e('fetchWorkflow'), []);
932976
const returnValue = await run.returnValue;

‎packages/core/src/encryption.ts‎

Lines changed: 20 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -33,19 +33,34 @@ const KEY_LENGTH = 32; // bytes (AES-256)
3333
* Callers should call this once per run (after `getEncryptionKeyForRun()`)
3434
* and pass the resulting `CryptoKey` to all subsequent encrypt/decrypt calls.
3535
*
36+
* Pass `usages: ['encrypt']` (or `['decrypt']`) for cross-run scenarios
37+
* where the caller should not be able to perform the inverse operation
38+
* with the key — for example a child workflow writing into a parent
39+
* run's forwarded WritableStream only needs to encrypt, never decrypt.
40+
*
3641
* @param raw - Raw 32-byte AES-256 key (from World.getEncryptionKeyForRun)
42+
* @param usages - Key usages. Defaults to `['encrypt', 'decrypt']`.
3743
* @returns CryptoKey ready for AES-GCM operations
3844
*/
39-
export async function importKey(raw: Uint8Array) {
45+
export async function importKey(
46+
raw: Uint8Array,
47+
usages: ReadonlyArray<'encrypt' | 'decrypt'> = ['encrypt', 'decrypt']
48+
) {
4049
if (raw.byteLength !== KEY_LENGTH) {
4150
throw new WorkflowRuntimeError(
4251
`Encryption key must be exactly ${KEY_LENGTH} bytes, got ${raw.byteLength}`
4352
);
4453
}
45-
return globalThis.crypto.subtle.importKey('raw', raw, 'AES-GCM', false, [
46-
'encrypt',
47-
'decrypt',
48-
]);
54+
return globalThis.crypto.subtle.importKey(
55+
'raw',
56+
raw,
57+
'AES-GCM',
58+
false,
59+
// `KeyUsage` is a DOM-lib type that's not in scope under `es2022`.
60+
// The `ReadonlyArray<'encrypt' | 'decrypt'>` parameter type matches
61+
// a strict subset of `KeyUsage[]`, so this cast is sound.
62+
usages as ('encrypt' | 'decrypt')[]
63+
);
4964
}
5065

5166
/**

‎packages/core/src/serialization.test.ts‎

Lines changed: 39 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,7 @@ import {
4141
ABORT_STREAM_NAME,
4242
STABLE_ULID,
4343
STREAM_NAME_SYMBOL,
44+
STREAM_SERVER_RUN_ID_SYMBOL,
4445
} from './symbols.js';
4546
import { createContext } from './vm/index.js';
4647

@@ -498,6 +499,44 @@ describe('workflow arguments', () => {
498499
expect(streamName).toMatch(/^strm_[0-9A-Z]{26}$/);
499500
});
500501

502+
// When a user writable is already backed by a workflow server
503+
// stream (because it was hydrated by a step-side reviver or created
504+
// via step-context `getWritable()`), forwarding it across a
505+
// `start()` boundary must emit the original `(runId, name)` in the
506+
// dehydrated descriptor and MUST NOT install any pipe through the
507+
// user's writable. The child run's step-side reviver then opens a
508+
// server writable against the original `(runId, name)` directly,
509+
// so writes survive for the full lifetime of the child run — not
510+
// just for the dehydrating step's process.
511+
it('forwards original (runId, name) for a tagged WritableStream', async () => {
512+
const userWritable = new WritableStream();
513+
Object.defineProperty(userWritable, STREAM_NAME_SYMBOL, {
514+
value: 'strm_parentstreamname',
515+
writable: false,
516+
});
517+
Object.defineProperty(userWritable, STREAM_SERVER_RUN_ID_SYMBOL, {
518+
value: 'wrun_parent',
519+
writable: false,
520+
});
521+
522+
expect(userWritable.locked).toBe(false);
523+
const serialized = await dehydrateWorkflowArguments(
524+
userWritable,
525+
'wrun_child',
526+
noEncryptionKey,
527+
[]
528+
);
529+
// If the reducer had piped through the user's writable, the lock
530+
// would be acquired here.
531+
expect(userWritable.locked).toBe(false);
532+
// The dehydrated descriptor should carry both the original name
533+
// and the original runId so the child's reviver can open the
534+
// writable against the parent's server stream directly.
535+
const text = new TextDecoder().decode(serialized as Uint8Array);
536+
expect(text).toContain('strm_parentstreamname');
537+
expect(text).toContain('wrun_parent');
538+
});
539+
501540
it('should work with ReadableStream', async () => {
502541
const stream = new ReadableStream();
503542
const serialized = await dehydrateWorkflowArguments(

‎packages/core/src/serialization.ts‎

Lines changed: 108 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ import {
55
decrypt as aesGcmDecrypt,
66
encrypt as aesGcmEncrypt,
77
type CryptoKey,
8+
importKey,
89
} from './encryption.js';
910
import {
1011
createFlushableState,
@@ -62,6 +63,7 @@ import {
6263
BODY_INIT_SYMBOL,
6364
STABLE_ULID,
6465
STREAM_NAME_SYMBOL,
66+
STREAM_SERVER_RUN_ID_SYMBOL,
6567
STREAM_TYPE_SYMBOL,
6668
WEBHOOK_RESPONSE_WRITABLE,
6769
} from './symbols.js';
@@ -775,9 +777,26 @@ export function getExternalReducers(
775777
WritableStream: (value) => {
776778
if (!(value instanceof global.WritableStream)) return false;
777779

780+
// Fast path: when the writable is already backed by a workflow
781+
// server stream (e.g. it came from a step-context `getWritable()`
782+
// or was hydrated from a workflow input by `getStepRevivers`),
783+
// forward its underlying `(runId, name)` to the receiving run.
784+
// The receiving run's step-side reviver opens a server writable
785+
// against the original `(runId, name)` and resolves that run's
786+
// encryption key directly, so writes land on the original stream
787+
// for the full lifetime of the receiving run — no in-process
788+
// bridge tied to the dehydrating step's lifetime.
789+
const existingName = (value as any)[STREAM_NAME_SYMBOL];
790+
const existingRunId = (value as any)[STREAM_SERVER_RUN_ID_SYMBOL];
791+
if (
792+
typeof existingName === 'string' &&
793+
typeof existingRunId === 'string'
794+
) {
795+
return { name: existingName, runId: existingRunId };
796+
}
797+
778798
const streamId = ((global as any)[STABLE_ULID] || defaultUlid)();
779799
const name = `strm_${streamId}`;
780-
781800
const readable = new WorkflowServerReadableStream(runId, name);
782801
ops.push(readable.pipeTo(value));
783802

@@ -861,7 +880,13 @@ export function getWorkflowReducers(
861880
if (!name) {
862881
throw new WorkflowRuntimeError('WritableStream `name` is not set');
863882
}
864-
return { name };
883+
const s: SerializableSpecial['WritableStream'] = { name };
884+
// When the handle was forwarded from another run (parent → child
885+
// via `start()`), preserve the foreign runId so the step-side
886+
// reviver opens the writable against the original stream.
887+
const foreignRunId = value[STREAM_SERVER_RUN_ID_SYMBOL];
888+
if (typeof foreignRunId === 'string') s.runId = foreignRunId;
889+
return s;
865890
},
866891

867892
// AbortController/AbortSignal in workflow context — just read symbols (handles).
@@ -961,6 +986,7 @@ function getStepReducers(
961986
if (!(value instanceof global.WritableStream)) return false;
962987

963988
let name = value[STREAM_NAME_SYMBOL];
989+
const foreignRunId = (value as any)[STREAM_SERVER_RUN_ID_SYMBOL];
964990
if (!name) {
965991
const streamId = ((global as any)[STABLE_ULID] || defaultUlid)();
966992
name = `strm_${streamId}`;
@@ -976,7 +1002,9 @@ function getStepReducers(
9761002
);
9771003
}
9781004

979-
return { name };
1005+
const s: SerializableSpecial['WritableStream'] = { name };
1006+
if (typeof foreignRunId === 'string') s.runId = foreignRunId;
1007+
return s;
9801008
},
9811009

9821010
AbortController: (value) => {
@@ -1414,12 +1442,25 @@ export function getExternalRevivers(
14141442
}
14151443
},
14161444
WritableStream: (value) => {
1445+
// Same handling as `getStepRevivers.WritableStream` — see comments
1446+
// there for the cross-run case (writable carries `runId` from
1447+
// parent → child forwarding via `start()`).
1448+
const targetRunId = typeof value.runId === 'string' ? value.runId : runId;
1449+
const targetKey: EncryptionKeyParam =
1450+
targetRunId === runId
1451+
? cryptoKey
1452+
: (async () => {
1453+
const world = await getWorldLazy();
1454+
const rawKey = await world.getEncryptionKeyForRun?.(targetRunId);
1455+
return rawKey ? await importKey(rawKey, ['encrypt']) : undefined;
1456+
})();
1457+
14171458
const serialize = getSerializeStream(
1418-
getExternalReducers(global, ops, runId, cryptoKey),
1419-
cryptoKey
1459+
getExternalReducers(global, ops, targetRunId, targetKey),
1460+
targetKey
14201461
);
14211462
const serverWritable = new WorkflowServerWritableStream(
1422-
runId,
1463+
targetRunId,
14231464
value.name
14241465
);
14251466

@@ -1435,6 +1476,15 @@ export function getExternalRevivers(
14351476
// Start polling to detect when user releases lock
14361477
pollWritableLock(serialize.writable, state);
14371478

1479+
Object.defineProperty(serialize.writable, STREAM_NAME_SYMBOL, {
1480+
value: value.name,
1481+
writable: false,
1482+
});
1483+
Object.defineProperty(serialize.writable, STREAM_SERVER_RUN_ID_SYMBOL, {
1484+
value: targetRunId,
1485+
writable: false,
1486+
});
1487+
14381488
return serialize.writable;
14391489
},
14401490

@@ -1519,12 +1569,22 @@ export function getWorkflowRevivers(
15191569
});
15201570
},
15211571
WritableStream: (value) => {
1522-
return Object.create(global.WritableStream.prototype, {
1572+
const descriptor: PropertyDescriptorMap = {
15231573
[STREAM_NAME_SYMBOL]: {
15241574
value: value.name,
15251575
writable: false,
15261576
},
1527-
});
1577+
};
1578+
// Preserve the foreign runId, if present, so that when the
1579+
// handle is later passed to a step the workflow reducer can
1580+
// forward it through to the step reviver.
1581+
if (typeof value.runId === 'string') {
1582+
descriptor[STREAM_SERVER_RUN_ID_SYMBOL] = {
1583+
value: value.runId,
1584+
writable: false,
1585+
};
1586+
}
1587+
return Object.create(global.WritableStream.prototype, descriptor);
15281588
},
15291589

15301590
// AbortController/AbortSignal revived inside the workflow VM. Use the
@@ -1742,12 +1802,34 @@ function getStepRevivers(
17421802
}
17431803
},
17441804
WritableStream: (value) => {
1805+
// Same-run case: the writable belongs to the current run. Use the
1806+
// local cryptoKey and write to the local runId's server stream.
1807+
//
1808+
// Cross-run case (parent → child via `start()`): the descriptor
1809+
// carries the original `runId` and `name`. Open a server writable
1810+
// against the original `(runId, name)` and resolve THAT run's key
1811+
// for encryption. The resolution is async but doesn't need to
1812+
// block reviver return — `getSerializeStream` accepts the
1813+
// `Promise<CryptoKey | undefined>` directly and awaits it lazily
1814+
// on the first chunk written. The key is imported encrypt-only
1815+
// so the receiving run can never decrypt anything else on the
1816+
// owning run's stream — it can only contribute new writes.
1817+
const targetRunId = typeof value.runId === 'string' ? value.runId : runId;
1818+
const targetKey: EncryptionKeyParam =
1819+
targetRunId === runId
1820+
? cryptoKey
1821+
: (async () => {
1822+
const world = await getWorldLazy();
1823+
const rawKey = await world.getEncryptionKeyForRun?.(targetRunId);
1824+
return rawKey ? await importKey(rawKey, ['encrypt']) : undefined;
1825+
})();
1826+
17451827
const serialize = getSerializeStream(
1746-
getStepReducers(global, ops, runId, cryptoKey),
1747-
cryptoKey
1828+
getStepReducers(global, ops, targetRunId, targetKey),
1829+
targetKey
17481830
);
17491831
const serverWritable = new WorkflowServerWritableStream(
1750-
runId,
1832+
targetRunId,
17511833
value.name
17521834
);
17531835

@@ -1763,6 +1845,21 @@ function getStepRevivers(
17631845
// Start polling to detect when user releases lock
17641846
pollWritableLock(serialize.writable, state);
17651847

1848+
// Record the underlying `(runId, name)` so downstream reducers can
1849+
// recognize that this writable is already backed by a workflow
1850+
// server stream. When forwarded across `start()` again — e.g.
1851+
// the child passes this writable on to a grandchild — the
1852+
// external reducer needs both to emit the original `runId` in
1853+
// the descriptor.
1854+
Object.defineProperty(serialize.writable, STREAM_NAME_SYMBOL, {
1855+
value: value.name,
1856+
writable: false,
1857+
});
1858+
Object.defineProperty(serialize.writable, STREAM_SERVER_RUN_ID_SYMBOL, {
1859+
value: targetRunId,
1860+
writable: false,
1861+
});
1862+
17661863
return serialize.writable;
17671864
},
17681865

‎packages/core/src/serialization/types.ts‎

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -152,7 +152,16 @@ export interface SerializableSpecial {
152152
cause?: unknown;
153153
errors: unknown[];
154154
};
155-
WritableStream: { name: string };
155+
WritableStream: {
156+
name: string;
157+
/**
158+
* The runId of the workflow run that owns the underlying server
159+
* stream. Present only when the writable was forwarded across a
160+
* `start()` boundary (parent → child). When omitted, the writable
161+
* belongs to the receiving run (the normal in-run case).
162+
*/
163+
runId?: string;
164+
};
156165
AbortController: {
157166
streamName: string;
158167
hookToken: string;

0 commit comments

Comments
 (0)