Skip to content

Commit 4547e1a

Browse files
authored
feat(streams): add writer session seam (#3832)
## Summary & Motivation Gives one in-memory stream writer a stable identity and its own sequence space, so a transport can preserve chunk ordering across a mid-stream HTTP/WebSocket transition. `Streamer.streams.createWriteSession` is optional — Worlds that don't implement it keep using `write`/`writeMulti`/`close` unchanged. Abort disposes the session rather than closing it, since a producer failure is transport cleanup, not stream completion. ## Test Plan Tests added, plus the full `@workflow/world-vercel` suite and package builds/typecheck pass. Root build/typecheck is blocked locally by a missing Rust toolchain for the unrelated `@workflow/swc-plugin`.
1 parent 51a181a commit 4547e1a

4 files changed

Lines changed: 171 additions & 2 deletions

File tree

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
---
2+
'@workflow/core': patch
3+
'@workflow/world': minor
4+
---
5+
6+
Add an optional stateful stream writer-session seam with stable writer identity and sequence tracking.

‎packages/core/src/serialization.ts‎

Lines changed: 36 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ import {
66
WorkflowRuntimeError,
77
} from '@workflow/errors';
88
import { once } from '@workflow/utils';
9+
import type { StreamWriteSession } from '@workflow/world';
910
import { envNumber } from '@workflow/world/env-config';
1011
import { parse, stringify, unflatten } from 'devalue';
1112
import { monotonicFactory } from 'ulid';
@@ -1310,6 +1311,20 @@ export class WorkflowServerWritableStream extends WritableStream<Uint8Array> {
13101311
}
13111312
};
13121313

1314+
// One stable identity and sequence space per in-memory sink. Start session
1315+
// construction as soon as the run-ready barrier permits, so transports may
1316+
// negotiate eagerly without allowing a write to overtake run creation.
1317+
// This eager chain cannot strand a new rejection: ensureRunReady absorbs its
1318+
// ordering-only failure, and every path that can observe worldPromise also
1319+
// awaits this session promise before it writes, closes, or disposes.
1320+
const writerId = `wrtr_${defaultUlid()}` as const;
1321+
const writeSessionPromise: Promise<StreamWriteSession | undefined> =
1322+
ensureRunReady().then(async () => {
1323+
const world = await worldPromise;
1324+
return world.streams.createWriteSession?.(runId, name, { writerId });
1325+
});
1326+
let nextChunkSeq = 0;
1327+
13131328
// ------------------------------------------------------------------
13141329
// Group-commit buffering.
13151330
//
@@ -1427,15 +1442,24 @@ export class WorkflowServerWritableStream extends WritableStream<Uint8Array> {
14271442
): Promise<void> => {
14281443
await ensureRunReady();
14291444
const world = await worldPromise;
1445+
const session = await writeSessionPromise;
14301446
const dispatchAt = Date.now();
1431-
if (typeof world.streams.writeMulti === 'function' && group.length > 1) {
1447+
if (session) {
1448+
await session.write(nextChunkSeq, group);
1449+
} else if (
1450+
typeof world.streams.writeMulti === 'function' &&
1451+
group.length > 1
1452+
) {
14321453
await world.streams.writeMulti(runId, name, group);
14331454
} else {
14341455
// Fall back to sequential writes
14351456
for (const chunk of group) {
14361457
await world.streams.write(runId, name, chunk);
14371458
}
14381459
}
1460+
// `inFlight` admits only one dispatch loop, so no second group can read
1461+
// this sequence space until the current group has advanced it.
1462+
nextChunkSeq += group.length;
14391463
if (groupT0 !== undefined) {
14401464
recordStreamWriteFlush(
14411465
groupT0,
@@ -1635,8 +1659,13 @@ export class WorkflowServerWritableStream extends WritableStream<Uint8Array> {
16351659
await ensureRunReady();
16361660

16371661
const world = await worldPromise;
1662+
const session = await writeSessionPromise;
16381663
const closeStart = Date.now();
1639-
await world.streams.close(runId, name);
1664+
if (session) {
1665+
await session.close();
1666+
} else {
1667+
await world.streams.close(runId, name);
1668+
}
16401669
recordStreamClose(closeStart, runId, name);
16411670
},
16421671
async abort(reason) {
@@ -1671,6 +1700,11 @@ export class WorkflowServerWritableStream extends WritableStream<Uint8Array> {
16711700
// Reject blocked writers and drain waiters so nothing leaks or
16721701
// hangs on a promise whose timer was just cleared.
16731702
rejectWaiters(sinkError);
1703+
// Abort is transport cleanup, not semantic stream completion. A
1704+
// stateful World releases its socket without sending close; stateless
1705+
// Worlds keep the existing no-op behavior.
1706+
const session = await writeSessionPromise.catch(() => undefined);
1707+
await session?.dispose?.();
16741708
},
16751709
});
16761710

‎packages/core/src/writable-stream.test.ts‎

Lines changed: 99 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -31,6 +31,7 @@ async function waitFor(
3131

3232
describe('WorkflowServerWritableStream', () => {
3333
let mockStreams: {
34+
createWriteSession?: ReturnType<typeof vi.fn>;
3435
write: ReturnType<typeof vi.fn>;
3536
writeMulti: ReturnType<typeof vi.fn>;
3637
close: ReturnType<typeof vi.fn>;
@@ -516,6 +517,82 @@ describe('WorkflowServerWritableStream', () => {
516517
});
517518
});
518519

520+
describe('stateful writer sessions', () => {
521+
it('uses one stable writer id and writer-local sequence across groups', async () => {
522+
const session = {
523+
write: vi.fn().mockResolvedValue(undefined),
524+
close: vi.fn().mockResolvedValue(undefined),
525+
};
526+
mockStreams.createWriteSession = vi.fn(() => session);
527+
528+
const stream = new WorkflowServerWritableStream('run-123', 'test-stream');
529+
const writer = stream.getWriter();
530+
await writer.write(new Uint8Array([1]));
531+
await waitFor(() => expect(session.write).toHaveBeenCalledTimes(1));
532+
await writer.write(new Uint8Array([2]));
533+
await waitFor(() => expect(session.write).toHaveBeenCalledTimes(2));
534+
await writer.close();
535+
536+
expect(mockStreams.createWriteSession).toHaveBeenCalledTimes(1);
537+
expect(mockStreams.createWriteSession).toHaveBeenCalledWith(
538+
'run-123',
539+
'test-stream',
540+
{
541+
writerId: expect.stringMatching(
542+
/^wrtr_[0123456789ABCDEFGHJKMNPQRSTVWXYZ]{26}$/
543+
),
544+
}
545+
);
546+
expect(session.write).toHaveBeenNthCalledWith(1, 0, [
547+
new Uint8Array([1]),
548+
]);
549+
expect(session.write).toHaveBeenNthCalledWith(2, 1, [
550+
new Uint8Array([2]),
551+
]);
552+
expect(session.close).toHaveBeenCalledTimes(1);
553+
expect(mockStreams.write).not.toHaveBeenCalled();
554+
expect(mockStreams.writeMulti).not.toHaveBeenCalled();
555+
expect(mockStreams.close).not.toHaveBeenCalled();
556+
});
557+
558+
it('creates a new writer id for each in-memory sink lifetime', async () => {
559+
const writerIds: string[] = [];
560+
mockStreams.createWriteSession = vi.fn((_runId, _name, options) => {
561+
writerIds.push(options.writerId);
562+
return {
563+
write: vi.fn().mockResolvedValue(undefined),
564+
close: vi.fn().mockResolvedValue(undefined),
565+
};
566+
});
567+
568+
new WorkflowServerWritableStream('run-123', 'test-stream');
569+
new WorkflowServerWritableStream('run-123', 'test-stream');
570+
await waitFor(() => expect(writerIds).toHaveLength(2));
571+
expect(new Set(writerIds).size).toBe(2);
572+
});
573+
});
574+
575+
describe('abort cleanup', () => {
576+
it('drains accepted chunks then disposes without semantic close', async () => {
577+
const session = {
578+
write: vi.fn().mockResolvedValue(undefined),
579+
close: vi.fn().mockResolvedValue(undefined),
580+
dispose: vi.fn().mockResolvedValue(undefined),
581+
};
582+
mockStreams.createWriteSession = vi.fn(() => session);
583+
584+
const stream = new WorkflowServerWritableStream('run-123', 'test-stream');
585+
const writer = stream.getWriter();
586+
await writer.write(new Uint8Array([1]));
587+
await writer.abort(new Error('producer failed'));
588+
589+
expect(session.write).toHaveBeenCalledWith(0, [new Uint8Array([1])]);
590+
expect(session.dispose).toHaveBeenCalledTimes(1);
591+
expect(session.close).not.toHaveBeenCalled();
592+
expect(mockStreams.close).not.toHaveBeenCalled();
593+
});
594+
});
595+
519596
describe('error handling', () => {
520597
it('surfaces a dispatch failure at the durability barrier (close) and poisons later writes', async () => {
521598
mockStreams.write.mockRejectedValueOnce(new Error('write error'));
@@ -549,6 +626,28 @@ describe('WorkflowServerWritableStream', () => {
549626
expect(mockStreams.write).toHaveBeenCalledTimes(1);
550627
});
551628

629+
it('keeps a session failure sticky without advancing its sequence', async () => {
630+
const session = {
631+
write: vi.fn().mockRejectedValueOnce(new Error('session failed')),
632+
close: vi.fn().mockResolvedValue(undefined),
633+
};
634+
mockStreams.createWriteSession = vi.fn(() => session);
635+
636+
const writer = new WorkflowServerWritableStream(
637+
'run-123',
638+
'test-stream'
639+
).getWriter();
640+
await writer.write(new Uint8Array([1]));
641+
await waitFor(() => expect(session.write).toHaveBeenCalledTimes(1));
642+
await expect(writer.write(new Uint8Array([2]))).rejects.toThrow(
643+
'session failed'
644+
);
645+
await expect(writer.close()).rejects.toThrow();
646+
expect(session.write).toHaveBeenCalledTimes(1);
647+
expect(session.write).toHaveBeenCalledWith(0, [new Uint8Array([1])]);
648+
expect(session.close).not.toHaveBeenCalled();
649+
});
650+
552651
it('should propagate close errors', async () => {
553652
mockStreams.close.mockRejectedValueOnce(new Error('close error'));
554653

‎packages/world/src/interfaces.ts‎

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,25 @@ import type {
4040
StepWithoutData,
4141
} from './steps.js';
4242

43+
export interface StreamWriteSession {
44+
/**
45+
* Write one ordered group from this in-memory writer lifetime.
46+
* `chunkSeq` is writer-local and identifies the first chunk in `chunks`.
47+
*/
48+
write(chunkSeq: number, chunks: (string | Uint8Array)[]): Promise<void>;
49+
50+
/** Close this writer lifetime after all prior writes are durable. */
51+
close(): Promise<void>;
52+
53+
/** Release transport resources without semantically closing the stream. */
54+
dispose?(): Promise<void> | void;
55+
}
56+
57+
export interface CreateStreamWriteSessionOptions {
58+
/** Stable observational id for this in-memory writer lifetime. */
59+
writerId: `wrtr_${string}`;
60+
}
61+
4362
export interface Streamer {
4463
/**
4564
* Number of milliseconds a stream waits for additional chunks to arrive
@@ -57,6 +76,17 @@ export interface Streamer {
5776
streamFlushIntervalMs?: number;
5877

5978
streams: {
79+
/**
80+
* Optionally create a stateful writer session. Core creates at most one
81+
* session per in-memory WritableStream and otherwise uses the stateless
82+
* write/writeMulti/close methods below unchanged.
83+
*/
84+
createWriteSession?(
85+
runId: string,
86+
name: string,
87+
options: CreateStreamWriteSessionOptions
88+
): StreamWriteSession;
89+
6090
write(
6191
runId: string,
6292
name: string,

0 commit comments

Comments
 (0)