Skip to content

Commit 7e70d18

Browse files
[core] Add configurable stream flush interval per world (#1533)
1 parent c8dce52 commit 7e70d18

8 files changed

Lines changed: 91 additions & 6 deletions

File tree

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,8 @@
1+
---
2+
"@workflow/world": patch
3+
"@workflow/core": patch
4+
"@workflow/world-local": patch
5+
"@workflow/world-postgres": patch
6+
---
7+
8+
Add `streamFlushIntervalMs` option to `Streamer` interface, optional for worlds to allow overwriting the default of 10ms in low-latency environments.

‎packages/core/src/serialization.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -522,7 +522,7 @@ export class WorkflowServerWritableStream extends WritableStream<Uint8Array> {
522522
for (const w of currentWaiters) w.reject(err);
523523
}
524524
);
525-
}, STREAM_FLUSH_INTERVAL_MS);
525+
}, world.streamFlushIntervalMs ?? STREAM_FLUSH_INTERVAL_MS);
526526
};
527527

528528
super({

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

Lines changed: 49 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@ describe('WorkflowServerWritableStream', () => {
1111
writeToStream: ReturnType<typeof vi.fn>;
1212
writeToStreamMulti: ReturnType<typeof vi.fn>;
1313
closeStream: ReturnType<typeof vi.fn>;
14+
streamFlushIntervalMs?: number;
1415
};
1516

1617
beforeEach(async () => {
@@ -248,4 +249,52 @@ describe('WorkflowServerWritableStream', () => {
248249
);
249250
});
250251
});
252+
253+
describe('streamFlushIntervalMs', () => {
254+
it('should use world.streamFlushIntervalMs when set to 0 (immediate flush)', async () => {
255+
mockWorld.streamFlushIntervalMs = 0;
256+
257+
const stream = new WorkflowServerWritableStream('s', 'run-1');
258+
const writer = stream.getWriter();
259+
260+
// With interval=0, the flush fires on the next microtask tick via setTimeout(fn, 0)
261+
await writer.write(new Uint8Array([1]));
262+
expect(mockWorld.writeToStream).toHaveBeenCalledTimes(1);
263+
264+
await writer.close();
265+
});
266+
267+
it('should fall back to default interval when streamFlushIntervalMs is undefined', async () => {
268+
// mockWorld has no streamFlushIntervalMs set — uses default 10ms
269+
delete mockWorld.streamFlushIntervalMs;
270+
271+
const stream = new WorkflowServerWritableStream('s', 'run-1');
272+
const writer = stream.getWriter();
273+
274+
await writer.write(new Uint8Array([1]));
275+
expect(mockWorld.writeToStream).toHaveBeenCalledTimes(1);
276+
277+
await writer.close();
278+
});
279+
280+
it('should respect a custom non-zero flush interval', async () => {
281+
mockWorld.streamFlushIntervalMs = 50;
282+
283+
const stream = new WorkflowServerWritableStream('s', 'run-1');
284+
const writer = stream.getWriter();
285+
286+
// Start a write — the flush is scheduled 50ms from now
287+
const writePromise = writer.write(new Uint8Array([1]));
288+
289+
// After 10ms (the old default), data should NOT have flushed yet
290+
await new Promise((r) => setTimeout(r, 10));
291+
expect(mockWorld.writeToStream).not.toHaveBeenCalled();
292+
293+
// Wait for the write to complete (will resolve after the 50ms timer fires)
294+
await writePromise;
295+
expect(mockWorld.writeToStream).toHaveBeenCalledTimes(1);
296+
297+
await writer.close();
298+
});
299+
});
251300
});

‎packages/world-local/src/config.ts‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,11 @@ export type Config = {
2222
* `.workflow-data` directory.
2323
*/
2424
tag?: string;
25+
/**
26+
* Override the flush interval (in ms) for buffered stream writes.
27+
* Default is 10ms. Set to 0 for immediate flushing.
28+
*/
29+
streamFlushIntervalMs?: number;
2530
};
2631

2732
export const config = once<Config>(() => {

‎packages/world-local/src/index.ts‎

Lines changed: 7 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -65,11 +65,13 @@ export function createLocalWorld(args?: Partial<Config>): LocalWorld {
6565
const storage = createStorage(mergedConfig.dataDir, tag);
6666
return {
6767
...queue,
68-
...storage,
69-
...instrumentObject(
70-
'world.streams',
71-
createStreamer(mergedConfig.dataDir, tag)
72-
),
68+
...createStorage(mergedConfig.dataDir, tag),
69+
...instrumentObject('world.streams', {
70+
...createStreamer(mergedConfig.dataDir, tag),
71+
...(mergedConfig.streamFlushIntervalMs !== undefined && {
72+
streamFlushIntervalMs: mergedConfig.streamFlushIntervalMs,
73+
}),
74+
}),
7375
async start() {
7476
await initDataDir(mergedConfig.dataDir);
7577
await reenqueueActiveRuns(storage.runs, queue.queue, 'world-local');

‎packages/world-postgres/src/config.ts‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,4 +7,9 @@ type PgConnectionConfig =
77
export type PostgresWorldConfig = PgConnectionConfig & {
88
jobPrefix?: string;
99
queueConcurrency?: number;
10+
/**
11+
* Override the flush interval (in ms) for buffered stream writes.
12+
* Default is 10ms. Set to 0 for immediate flushing.
13+
*/
14+
streamFlushIntervalMs?: number;
1015
};

‎packages/world-postgres/src/index.ts‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -60,6 +60,9 @@ export function createWorld(
6060
...storage,
6161
...streamer,
6262
...queue,
63+
...(config.streamFlushIntervalMs !== undefined && {
64+
streamFlushIntervalMs: config.streamFlushIntervalMs,
65+
}),
6366
async start() {
6467
await queue.start();
6568
await reenqueueActiveRuns(storage.runs, queue.queue, 'world-postgres');

‎packages/world/src/interfaces.ts‎

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,19 @@ import type {
3030
} from './steps.js';
3131

3232
export interface Streamer {
33+
/**
34+
* Override the default flush interval (in milliseconds) for buffered stream writes.
35+
* Chunks are accumulated in a buffer and flushed together on this interval.
36+
*
37+
* The default is 10ms, which is appropriate for HTTP-based backends where
38+
* each flush is a network round-trip. For backends with sub-millisecond writes
39+
* (e.g., Redis, local filesystem), a lower value (or 0 for immediate flushing) reduces
40+
* end-to-end stream latency.
41+
*
42+
* Not supported by all worlds.
43+
*/
44+
streamFlushIntervalMs?: number;
45+
3346
writeToStream(
3447
name: string,
3548
runId: string,

0 commit comments

Comments
 (0)