Skip to content

Commit 5fc8fb7

Browse files
authored
perf(streams): connect WebSocket after first write (#4104)
## Summary & Motivation Keeps WebSocket setup off the first stream group, then starts the background upgrade when the second HTTP request is dispatched. Later groups continue over HTTP without waiting until the socket is OPEN, when the serialized writer switches transports at a confirmed request boundary. One-group streams never create a socket. An ambiguous HTTP outcome poisons the writer and retires any provisional socket before it can carry a frame. Initial upgrades have a dedicated 10-second background timeout; the existing 250ms bound remains scoped to reconnects after an established socket closes. WS write spans include chunk sequence and count for direct takeover analysis. ## Test Plan Tests cover dispatch overlap, continued HTTP writes while connecting, one-group streams, close during background connection, ambiguous HTTP outcomes, timeout and late OPEN behavior, and trace propagation. The world-vercel suite and typecheck pass locally.
1 parent 01fa7a4 commit 5fc8fb7

7 files changed

Lines changed: 349 additions & 74 deletions

File tree

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
'@workflow/world-vercel': patch
3+
---
4+
5+
Start stream WebSocket connections after the first HTTP group without delaying later HTTP writes.

‎packages/world-vercel/src/http-core.ts‎

Lines changed: 13 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -549,7 +549,13 @@ export interface InstrumentedFetchOptions extends HttpClientSpanOptions {
549549
* Non-2xx bodies consumed by `buildError` are still reported here.
550550
*/
551551
deferTransportSuccessUntilBody?: boolean;
552-
/** Error code used when the request itself fails before a response arrives. */
552+
/**
553+
* Called synchronously after the request promise is created, before awaiting
554+
* its response. This observes local dispatch only; it does not imply that any
555+
* bytes reached the origin. Must not throw.
556+
*/
557+
onRequestDispatched?: () => void;
558+
/** Error code used when the request itself fails before a response arrived. */
553559
transportErrorCode?: 'TRANSPORT' | 'STREAM_ERROR';
554560
}
555561

@@ -585,6 +591,7 @@ export async function instrumentedFetch(
585591
durationAttribute,
586592
onTransportOutcome,
587593
deferTransportSuccessUntilBody = false,
594+
onRequestDispatched,
588595
transportErrorCode = 'TRANSPORT',
589596
} = opts;
590597
const label = logLabel ?? url;
@@ -623,8 +630,8 @@ export async function instrumentedFetch(
623630
span?.setAttributes({
624631
...WorkflowHttpTransport(nodeAgents ? 'node-http' : 'undici'),
625632
});
626-
response = nodeAgents
627-
? await nodeHttpFetch(url, {
633+
const request = nodeAgents
634+
? nodeHttpFetch(url, {
628635
method,
629636
headers,
630637
body,
@@ -637,14 +644,16 @@ export async function instrumentedFetch(
637644
headersTimeoutMs: NODE_HTTP_HEADERS_TIMEOUT_MS,
638645
bodyTimeoutMs: NODE_HTTP_BODY_TIMEOUT_MS,
639646
})
640-
: await fetch(url, {
647+
: fetch(url, {
641648
method,
642649
headers,
643650
body,
644651
signal,
645652
// eslint-disable-next-line @typescript-eslint/no-explicit-any -- undici dispatcher type doesn't match @types/node's RequestInit
646653
dispatcher,
647654
} as any);
655+
onRequestDispatched?.();
656+
response = await request;
648657
} catch (error) {
649658
const elapsed = Date.now() - start;
650659
// Report the raw error, before the timeout mapping below rewraps it: the

‎packages/world-vercel/src/streamer.ts‎

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -211,7 +211,8 @@ export async function writeStreamSessionOverHttp(
211211
name: string,
212212
chunks: (string | Uint8Array)[],
213213
config?: APIConfig,
214-
attributes?: Attributes
214+
attributes?: Attributes,
215+
onRequestDispatched?: () => void
215216
): Promise<void> {
216217
const httpConfig = await getHttpConfig(config);
217218
httpConfig.headers.set('X-Stream-Multi', 'true');
@@ -231,6 +232,7 @@ export async function writeStreamSessionOverHttp(
231232
logLabel: url.pathname,
232233
spanName: 'workflow.stream.write',
233234
durationAttribute: 'workflow.stream.write.chunk_rtt',
235+
onRequestDispatched: offset === 0 ? onRequestDispatched : undefined,
234236
attributes: {
235237
...streamSpanAttributes({
236238
runId,
@@ -282,13 +284,14 @@ export function createStreamer(config?: APIConfig): Streamer {
282284
name,
283285
writerId,
284286
config,
285-
(chunks, attributes) =>
287+
(chunks, attributes, onRequestDispatched) =>
286288
writeStreamSessionOverHttp(
287289
runId,
288290
name,
289291
chunks,
290292
config,
291-
attributes
293+
attributes,
294+
onRequestDispatched
292295
),
293296
() => closeStreamSessionOverHttp(runId, name, config)
294297
);

‎packages/world-vercel/src/trace-propagation.test.ts‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -345,14 +345,16 @@ describe('ws stream transport upgrade trace propagation', () => {
345345
await tracer.startActiveSpan('flow-invocation', async (span) => {
346346
traceId = span.spanContext().traceId;
347347
spanId = span.spanContext().spanId;
348-
createStreamWriteSession(
348+
const session = createStreamWriteSession(
349349
'wrun_1',
350350
'user',
351351
'wrtr_01ARZ3NDEKTSV4RRFFQ69G5FAV',
352352
{ token: 'test-token' },
353-
async () => {},
353+
async (_chunks, _attributes, dispatched) => dispatched?.(),
354354
async () => {}
355355
);
356+
await session.write(0, ['first']);
357+
await session.write(1, ['second']);
356358
await vi.waitFor(() => expect(wsUpgrades).toHaveLength(1));
357359
span.end();
358360
});

‎packages/world-vercel/src/ws-stream-connect.ts‎

Lines changed: 4 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,11 +1,10 @@
11
import type { WebSocket } from 'ws';
22
import { trace } from './telemetry.js';
33

4-
/**
5-
* Initial implementation-level wait for an opted-in stream socket to open.
6-
* Tune from deployment measurements; this is not protocol semantics.
7-
*/
8-
export const STREAM_WS_CONNECT_BUDGET_MS = 250;
4+
/** Background initial upgrade bound; writes continue over HTTP meanwhile. */
5+
export const STREAM_WS_INITIAL_CONNECT_TIMEOUT_MS = 10_000;
6+
/** Bounded wait while reconnecting after an established socket drains/closes. */
7+
export const STREAM_WS_RECONNECT_BUDGET_MS = 250;
98
/** Cleanup-only wait after semantic close success; not protocol semantics. */
109
export const STREAM_WS_CLOSE_BUDGET_MS = 250;
1110

‎packages/world-vercel/src/ws-stream-session.test.ts‎

Lines changed: 156 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -122,6 +122,7 @@ beforeEach(() => {
122122
injectTraceContextIntoHeaders.mockClear();
123123
writeSpans.length = 0;
124124
delete process.env.WORKFLOW_STREAMS_TRANSPORT;
125+
delete process.env.WORKFLOW_REQUEST_TIMEOUT_MS;
125126
});
126127

127128
afterEach(() => {
@@ -131,17 +132,27 @@ afterEach(() => {
131132
});
132133

133134
function makeSession(
134-
config: { token?: string } | undefined = { token: 'token' }
135+
config: { token?: string } | undefined = { token: 'token' },
136+
connectAfterFirstWrite = false
135137
) {
136-
const writeHttp = vi.fn().mockResolvedValue(undefined);
138+
const writeHttp = vi.fn(
139+
async (
140+
_chunks: (string | Uint8Array)[],
141+
_attributes?: Record<string, unknown>,
142+
onRequestDispatched?: () => void
143+
) => {
144+
onRequestDispatched?.();
145+
}
146+
);
137147
const closeHttp = vi.fn().mockResolvedValue(undefined);
138148
const session = createStreamWriteSession(
139149
'wrun_1',
140150
'stream/1',
141151
writerId,
142152
config,
143153
writeHttp,
144-
closeHttp
154+
closeHttp,
155+
connectAfterFirstWrite
145156
);
146157
activeSessions.push(session);
147158
return { session, writeHttp, closeHttp };
@@ -158,6 +169,124 @@ describe('v1 stream WebSocket writer lifecycle', () => {
158169
expect(closeHttp).toHaveBeenCalledTimes(1);
159170
});
160171

172+
it('starts WS during the second HTTP group and keeps writing HTTP until OPEN', async () => {
173+
process.env.WORKFLOW_STREAMS_TRANSPORT = 'ws';
174+
let releaseSecond: (() => void) | undefined;
175+
const secondPending = new Promise<void>((resolve) => {
176+
releaseSecond = resolve;
177+
});
178+
const { session, writeHttp } = makeSession({ token: 'token' }, true);
179+
180+
await session.write(0, ['one']);
181+
expect(sockets).toHaveLength(0);
182+
expect(writeHttp).toHaveBeenNthCalledWith(
183+
1,
184+
['one'],
185+
expect.objectContaining({
186+
'workflow.stream.ws.session_first_write': true,
187+
'workflow.stream.ws.http_group_ordinal': 1,
188+
})
189+
);
190+
191+
writeHttp.mockImplementationOnce(
192+
async (_chunks, _attributes, dispatched) => {
193+
dispatched?.();
194+
await secondPending;
195+
}
196+
);
197+
const second = session.write(1, ['two']);
198+
await vi.waitFor(() => expect(sockets).toHaveLength(1));
199+
expect(writeHttp).toHaveBeenNthCalledWith(
200+
2,
201+
['two'],
202+
expect.objectContaining({
203+
'workflow.stream.ws.connect_after_http_group': 2,
204+
}),
205+
expect.any(Function)
206+
);
207+
expect(sockets[0].sent).toHaveLength(0);
208+
releaseSecond?.();
209+
await second;
210+
211+
const third = session.write(2, ['three']);
212+
await third;
213+
expect(writeHttp.mock.calls[2]?.[0]).toEqual(['three']);
214+
sockets[0].open();
215+
216+
const fourth = session.write(3, ['four']);
217+
await vi.waitFor(() => expect(sockets[0].sent).toHaveLength(1));
218+
expect((await decodeOne(sockets[0].sent[0])).meta).toMatchObject({
219+
type: 'write',
220+
chunkSeq: 3,
221+
numChunks: 1,
222+
});
223+
sockets[0].reply(
224+
encodeFrame({ type: 'write_ack', reqId: 1 }, new Uint8Array())
225+
);
226+
await fourth;
227+
});
228+
229+
it('closes over HTTP without waiting for the background socket', async () => {
230+
process.env.WORKFLOW_STREAMS_TRANSPORT = 'ws';
231+
const { session, closeHttp } = makeSession({ token: 'token' }, true);
232+
await session.write(0, ['one']);
233+
await session.write(1, ['two']);
234+
await vi.waitFor(() => expect(sockets).toHaveLength(1));
235+
236+
await session.close();
237+
expect(closeHttp).toHaveBeenCalledTimes(1);
238+
expect(sockets[0].closed).toContainEqual([1000, 'stream closed over HTTP']);
239+
});
240+
241+
it('retires a provisional socket after an ambiguous second HTTP outcome', async () => {
242+
process.env.WORKFLOW_STREAMS_TRANSPORT = 'ws';
243+
const error = new Error('second HTTP outcome unknown');
244+
const { session, writeHttp } = makeSession({ token: 'token' }, true);
245+
await session.write(0, ['one']);
246+
let rejectSecond: ((error: Error) => void) | undefined;
247+
const secondPending = new Promise<void>((_resolve, reject) => {
248+
rejectSecond = reject;
249+
});
250+
writeHttp.mockImplementationOnce(
251+
async (_chunks, _attributes, dispatched) => {
252+
dispatched?.();
253+
await secondPending;
254+
}
255+
);
256+
257+
const second = session.write(1, ['two']);
258+
await vi.waitFor(() => expect(sockets).toHaveLength(1));
259+
rejectSecond?.(error);
260+
await expect(second).rejects.toBe(error);
261+
await expect(session.write(2, ['three'])).rejects.toBe(error);
262+
expect(sockets[0].sent).toHaveLength(0);
263+
expect(sockets[0].closed).toContainEqual([
264+
1011,
265+
'unknown stream write outcome',
266+
]);
267+
});
268+
269+
it('does not start a socket after an ambiguous first HTTP outcome', async () => {
270+
process.env.WORKFLOW_STREAMS_TRANSPORT = 'ws';
271+
const error = new Error('HTTP outcome unknown');
272+
const { session, writeHttp } = makeSession({ token: 'token' }, true);
273+
writeHttp.mockRejectedValueOnce(error);
274+
275+
await expect(session.write(0, ['one'])).rejects.toBe(error);
276+
await expect(session.write(1, ['two'])).rejects.toBe(error);
277+
expect(sockets).toHaveLength(0);
278+
});
279+
280+
it('closes a one-group stream without starting a socket', async () => {
281+
process.env.WORKFLOW_STREAMS_TRANSPORT = 'ws';
282+
const { session, closeHttp } = makeSession({ token: 'token' }, true);
283+
284+
await session.write(0, ['only']);
285+
await session.close();
286+
expect(closeHttp).toHaveBeenCalledTimes(1);
287+
expect(sockets).toHaveLength(0);
288+
});
289+
161290
it('sends immediately over HTTP while the initial socket connects, then switches to WS', async () => {
162291
process.env.WORKFLOW_STREAMS_TRANSPORT = 'ws';
163292
let releaseHttp: (() => void) | undefined;
@@ -239,27 +368,33 @@ describe('v1 stream WebSocket writer lifecycle', () => {
239368
});
240369
});
241370

242-
it('tombstones to HTTP when the background connect budget expires', async () => {
371+
it('bounds the background connect attempt without delaying HTTP writes', async () => {
243372
vi.useFakeTimers();
244373
process.env.WORKFLOW_STREAMS_TRANSPORT = 'ws';
245-
const { session, writeHttp } = makeSession();
374+
process.env.WORKFLOW_REQUEST_TIMEOUT_MS = '10000';
375+
const { session, writeHttp } = makeSession({ token: 'token' }, true);
246376
const first = session.write(0, ['one']);
247-
await vi.advanceTimersByTimeAsync(0);
248-
expect(writeHttp).toHaveBeenCalledWith(
249-
['one'],
250-
expect.objectContaining({
251-
'workflow.stream.ws.session_first_write': true,
252-
'workflow.stream.ws.connecting_at_write': true,
253-
})
377+
await vi.waitFor(() =>
378+
expect(writeHttp).toHaveBeenCalledWith(
379+
['one'],
380+
expect.objectContaining({
381+
'workflow.stream.ws.session_first_write': true,
382+
'workflow.stream.ws.connecting_at_write': false,
383+
'workflow.stream.ws.connect_deferred_at_write': true,
384+
})
385+
)
254386
);
255387
await first;
388+
expect(sockets).toHaveLength(0);
389+
await session.write(1, ['two']);
390+
await vi.waitFor(() => expect(sockets).toHaveLength(1));
256391

257-
await vi.advanceTimersByTimeAsync(250);
258-
const second = session.write(1, ['two']);
259-
await second;
260-
expect(writeHttp).toHaveBeenCalledTimes(2);
392+
await vi.advanceTimersByTimeAsync(10_000);
393+
await session.write(2, ['three']);
394+
expect(writeHttp).toHaveBeenCalledTimes(3);
261395
expect(writeHttp.mock.calls[0]?.[0]).toEqual(['one']);
262-
expect(writeHttp.mock.calls[1]).toEqual([['two']]);
396+
expect(writeHttp.mock.calls[1]?.[0]).toEqual(['two']);
397+
expect(writeHttp.mock.calls[2]).toEqual([['three']]);
263398
expect(sockets[0].sent).toHaveLength(0);
264399
expect(sockets[0].closed).toContainEqual([1000, 'connect budget expired']);
265400
sockets[0].open();
@@ -766,15 +901,16 @@ describe('v1 stream WebSocket writer lifecycle', () => {
766901
expect(sockets).toHaveLength(0);
767902
});
768903

769-
it('forwards streamer-wrapper disposal after session materialization', async () => {
904+
it('keeps streamer-wrapper disposal socket-free before its first write', async () => {
770905
process.env.WORKFLOW_STREAMS_TRANSPORT = 'ws';
771906
const { createStreamer } = await import('./streamer.js');
772907
const session = createStreamer({
773908
token: 'token',
774909
}).streams.createWriteSession?.('wrun_1', 'stream/1', { writerId });
775-
await vi.waitFor(() => expect(sockets).toHaveLength(1));
910+
await new Promise((resolve) => setTimeout(resolve, 0));
911+
expect(sockets).toHaveLength(0);
776912
await session?.dispose?.();
777-
expect(sockets[0].closed).toContainEqual([1000, 'stream writer disposed']);
913+
expect(sockets).toHaveLength(0);
778914
});
779915

780916
it('disposes transport without sending protocol close', async () => {

0 commit comments

Comments
 (0)