Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/step-stream-drain-barrier.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
'@workflow/core': patch
---

Wait for writes queued by released step stream writers to reach durable storage before recording step completion.
5 changes: 5 additions & 0 deletions docs/content/docs/v5/configuration/runtime-tuning.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -314,6 +314,11 @@ These variables are primarily for tests, debugging, or unusual deployments.
- Group-commit window for the *leading* chunk of an idle stream. `0` sends it at once; a positive value holds it up to that many milliseconds to collect a group. This opt-in setting trades first-chunk latency for larger batches and can benefit slow-but-steady producers. Chunks arriving while a request is already in flight always coalesce into the next group regardless of this setting.
- Also available as `streamFlushIntervalMs` on Worlds that expose it (the env var, when set, takes precedence over the World option).

### `WORKFLOW_STEP_STREAM_DRAIN_TIMEOUT_MS`

- Default: `30000` (30 seconds)
- Maximum time a step waits for writes queued by a released workflow stream writer to be acknowledged by the Workflow server before recording `step_completed`. A timeout fails the step instead of exposing a completion while released-writer data remains client-side. A writer intentionally kept locked does not trigger this durability wait.

### `WORKFLOW_STREAM_MAX_INFLIGHT_CHUNKS`

- Default: `1000`
Expand Down
126 changes: 126 additions & 0 deletions packages/core/src/flushable-stream.test.ts
Original file line number Diff line number Diff line change
@@ -1,10 +1,12 @@
import { afterEach, describe, expect, it } from 'vitest';
import {
createFlushableState,
drainFlushableSnapshot,
flushablePipe,
LOCK_POLL_INTERVAL_MS,
pollReadableLock,
pollWritableLock,
trackFlushableWritable,
} from './flushable-stream.js';
import { STREAM_DRAIN_SYMBOL } from './symbols.js';

Expand Down Expand Up @@ -446,6 +448,130 @@ describe('flushablePipe drain barrier (group-commit sinks)', () => {
delete process.env.WORKFLOW_STREAM_MAX_BYTES_PER_BATCH;
});

it('waits for a produced frame to reach the sink before draining', async () => {
let releaseWrite!: () => void;
const writeGate = new Promise<void>((resolve) => {
releaseWrite = resolve;
});
let drained = false;
const sink = new WritableStream<Uint8Array>({
async write() {
await writeGate;
},
});
Object.defineProperty(sink, STREAM_DRAIN_SYMBOL, {
value: async () => {
drained = true;
},
});
const transform = new TransformStream<Uint8Array, Uint8Array>();
const state = createFlushableState();
const pipe = flushablePipe(transform.readable, sink, state).catch(() => {});
const writable = trackFlushableWritable(transform.writable, state);
const writer = writable.getWriter();

await writer.write(new Uint8Array([1]));
const snapshot = drainFlushableSnapshot(state);
await tick();
expect(drained).toBe(false);

releaseWrite();
await snapshot;
expect(drained).toBe(true);

await writer.close();
await pipe;
});

it('propagates an idle downstream failure through the tracked writable', async () => {
const downstreamError = new Error('downstream failed while producer idle');
const transform = new TransformStream<Uint8Array, Uint8Array>();
const state = createFlushableState();
const sink = new WritableStream<Uint8Array>({
write() {
throw downstreamError;
},
});
const pipe = flushablePipe(transform.readable, sink, state).catch(() => {});
const writable = trackFlushableWritable(transform.writable, state);
let sourceCancelled = false;
const source = new ReadableStream<Uint8Array>({
start(controller) {
controller.enqueue(new Uint8Array([1]));
},
cancel() {
sourceCancelled = true;
},
});

await expect(source.pipeTo(writable)).rejects.toThrow(
'downstream failed while producer idle'
);
expect(sourceCancelled).toBe(true);
await pipe;
});

it('drains an accepted prefix before rejecting its snapshot', async () => {
let releaseDrain!: () => void;
const drainGate = new Promise<void>((resolve) => {
releaseDrain = resolve;
});
const { sink } = makeDrainSink(() => drainGate);
const { source, controller } = makeControlledSource();
const state = createFlushableState();
state.producedFrames = 2;
const pipe = flushablePipe(source, sink, state).catch(() => {});

controller().enqueue(new Uint8Array([1]));
await tick();
controller().error(new Error('second frame failed'));
await tick();

let settled = false;
const snapshot = drainFlushableSnapshot(state).finally(() => {
settled = true;
});
await tick();
expect(settled).toBe(false);

releaseDrain();
await expect(snapshot).rejects.toThrow('second frame failed');
await pipe;
});

it('rejects snapshot waiters with the upstream pipe error', async () => {
const { sink } = makeDrainSink(async () => {});
const source = new ReadableStream<Uint8Array>({
start(controller) {
controller.error(new Error('producer failed before acceptance'));
},
});
const state = createFlushableState();
state.producedFrames = 1;
const pipe = flushablePipe(source, sink, state).catch(() => {});

const snapshot = drainFlushableSnapshot(state);

await expect(snapshot).rejects.toThrow('producer failed before acceptance');
await pipe;
});

it('reports a pipe error to a snapshot registered after failure', async () => {
const { sink } = makeDrainSink(async () => {});
const source = new ReadableStream<Uint8Array>({
start(controller) {
controller.error(new Error('producer already failed'));
},
});
const state = createFlushableState();
state.producedFrames = 1;
await flushablePipe(source, sink, state).catch(() => {});

await expect(drainFlushableSnapshot(state)).rejects.toThrow(
'producer already failed'
);
});

it('adopts the sink drain barrier onto the flushable state', async () => {
const { sink } = makeDrainSink(async () => {});
const { source, controller } = makeControlledSource();
Expand Down
Loading
Loading