Skip to content

Commit 04e060a

Browse files
[world] Add WORKFLOW_NODE_HTTP to run the HTTP Worlds on node:http (#3461)
1 parent 707dfe6 commit 04e060a

23 files changed

Lines changed: 1934 additions & 61 deletions
Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,7 @@
1+
---
2+
'@workflow/world-vercel': minor
3+
'@workflow/world-local': minor
4+
'@workflow/world': minor
5+
---
6+
7+
Add `WORKFLOW_NODE_HTTP` to run the Vercel and Local World HTTP layers on Node's built-in `node:http` and `node:https` modules instead of their usual HTTP client library.

‎docs/content/docs/v5/configuration/runtime-tuning.mdx‎

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -214,6 +214,24 @@ For example, a workflow can run a 10-minute inline step even with `WORKFLOW_REPL
214214
- Default: enabled
215215
- On the Vercel World, lets concurrent event-log requests share one HTTP/2 connection instead of one connection per in-flight request.
216216
- Set `0` to take the event-log requests off HTTP/2 entirely, back to one request per HTTP/1.1 connection. Use this as the kill switch if HTTP/2 turns out to be at fault for event delivery problems.
217+
- Has no effect when `WORKFLOW_NODE_HTTP` is enabled, which takes the whole HTTP/2 path away.
218+
219+
### `WORKFLOW_NODE_HTTP`
220+
221+
- Default: disabled
222+
- Makes the Vercel and Local Worlds issue their HTTP requests through Node's built-in `node:http` and `node:https` modules, instead of the HTTP client library those Worlds normally use.
223+
- Set `1` to switch to Node's modules. Read the trade-offs below first: they cost throughput on every deployment, which is why this is opt-in.
224+
- Use it when that library is not an option: a bundler that mangles it, a runtime that does not ship a working copy of it, or a transport-level fault you want to rule out. It is not reached by way of `fetch()` either, so a runtime whose `fetch()` is built on the same library is still covered.
225+
226+
Node's own modules do less than the client they replace, so enabling this drops the per-call-site tuning the Worlds configure:
227+
228+
- Event-log requests lose HTTP/2, so concurrent reads and writes no longer share one connection, and the enlarged HTTP/2 receive windows no longer apply. This is the largest difference, and it slows down replays that read a big event log. It does not apply to event writes on [`WORKFLOW_EVENTS_TRANSPORT=ws`](/docs/configuration/worlds#workflow_events_transport), which take neither transport.
229+
- Requests lose their transport-level retry. Failures still surface to the layers above, which retry event writes and redeliver queue messages, so nothing is silently dropped, but a failure that a same-connection retry would have hidden now costs a full redelivery.
230+
- Stream close loses its retry of retriable server errors. A transient failure at close can leave a stream marked closing until the run expires, where it would previously have resolved on the retry.
231+
232+
Connection pooling, keep-alive, and the request, header, and body deadlines are preserved: pooling and keep-alive are configured on Node's agents, and the deadlines are passed per request — from the Local World's two queue timeouts, and on the Vercel World from the same defaults its HTTP client applies today. Queue sends are the exception in the other direction: that client takes no transport override, so its requests keep using the library either way.
233+
234+
A `dispatcher` passed to `createVercelWorld()` still wins over this variable. The variable chooses which transport the World builds when you have not supplied one.
217235

218236
## Queue namespace
219237

‎docs/content/worlds/v5/local.mdx‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -70,6 +70,10 @@ Maximum time in milliseconds to wait for a local queue handler to begin respondi
7070

7171
Maximum gap in milliseconds between response body chunks from a local queue handler before redelivering the durable message. Default: `30000`. Set to `0` to disable.
7272

73+
### `WORKFLOW_NODE_HTTP`
74+
75+
Whether queue deliveries go out through Node's built-in `node:http` and `node:https` modules instead of the HTTP client library this World normally uses. Default: disabled. Set to `1` to switch to Node's modules. Socket pooling, keep-alive, and both timeouts above apply either way. See [`WORKFLOW_NODE_HTTP`](/docs/configuration/runtime-tuning#workflow_node_http) for what else changes.
76+
7377
### `WORKFLOW_LOCAL_RECOVER_ACTIVE_RUNS`
7478

7579
Whether pending and running runs found in the data directory are re-enqueued when the world starts. Set to `0` or `false` to skip recovery and leave stale runs untouched. Default: `true`

‎packages/world-local/src/queue.test.ts‎

Lines changed: 106 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@ import { createServer, type Server } from 'node:http';
22
import type { AddressInfo } from 'node:net';
33
import { setWorkflowBasePath } from '@workflow/utils';
44
import type { WorkflowInvokePayload } from '@workflow/world';
5-
import { MessageId, ValidQueueName } from '@workflow/world';
5+
import { MessageId, NODE_HTTP_ENV_VAR, ValidQueueName } from '@workflow/world';
66
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
77
import { z } from 'zod/v4';
88
import {
@@ -23,6 +23,19 @@ const workflowPayload: WorkflowInvokePayload = {
2323
stepName: 'test-step',
2424
};
2525

26+
// The suite below covers the queue on undici, including the tests that stub
27+
// the global `fetch`. `WORKFLOW_NODE_HTTP` sends deliveries over `node:http`
28+
// instead, where stubbing `fetch` proves nothing, so pin the flag off rather
29+
// than tracking whichever way its default points. That mode has its own
30+
// describe at the end.
31+
beforeEach(() => {
32+
vi.stubEnv(NODE_HTTP_ENV_VAR, '0');
33+
});
34+
35+
afterEach(() => {
36+
vi.unstubAllEnvs();
37+
});
38+
2639
describe('zod v3/v4 schema compatibility (regression #1587)', () => {
2740
it('ValidQueueName and MessageId from @workflow/world parse correctly in z.object()', () => {
2841
const HeaderParser = z.object({
@@ -615,3 +628,95 @@ describe('queue transport timeouts', () => {
615628
}
616629
});
617630
});
631+
632+
describe('node:http mode', () => {
633+
let server: Server | undefined;
634+
635+
beforeEach(() => {
636+
vi.stubEnv(NODE_HTTP_ENV_VAR, '1');
637+
});
638+
639+
afterEach(async () => {
640+
if (server !== undefined) {
641+
const toClose = server;
642+
server = undefined;
643+
toClose.closeAllConnections();
644+
await new Promise((resolve) => toClose.close(resolve));
645+
}
646+
vi.restoreAllMocks();
647+
vi.unstubAllGlobals();
648+
});
649+
650+
// The equivalence claim the flag makes: a delivery still goes out, still
651+
// carries the VQS headers the handler reads, and still lands on the handler.
652+
// The global `fetch` is stubbed to throw to prove the request left undici
653+
// entirely rather than falling back to the runtime's own pool.
654+
it('delivers without going through fetch', async () => {
655+
const attempts: (string | undefined)[] = [];
656+
server = createServer((request, response) => {
657+
attempts.push(request.headers['x-vqs-message-attempt'] as string);
658+
response.setHeader('content-type', 'application/json');
659+
response.end(JSON.stringify({ ok: true }));
660+
});
661+
await new Promise<void>((resolve) => {
662+
server?.listen(0, '127.0.0.1', resolve);
663+
});
664+
const { port } = server.address() as AddressInfo;
665+
666+
vi.stubGlobal(
667+
'fetch',
668+
vi.fn(() => {
669+
throw new Error('fetch must not be used under WORKFLOW_NODE_HTTP');
670+
})
671+
);
672+
673+
const localQueue = createQueue({ baseUrl: `http://127.0.0.1:${port}` });
674+
try {
675+
await localQueue.queue('__wkf_workflow_test' as any, workflowPayload);
676+
await vi.waitFor(() => expect(attempts.length).toBe(1), {
677+
timeout: 5_000,
678+
});
679+
expect(attempts[0]).toBe('1');
680+
} finally {
681+
await localQueue.close();
682+
}
683+
});
684+
685+
// Redelivery on a transport throw is the queue's own logic, not undici's, so
686+
// it survives the switch. Node's client raises ECONNRESET where undici would
687+
// have raised a TypeError wrapping UND_ERR_SOCKET; the delivery loop keys on
688+
// neither, so both retry the same durable message.
689+
it('still retries a transport throw', async () => {
690+
let calls = 0;
691+
server = createServer((request, response) => {
692+
calls++;
693+
if (calls < 3) {
694+
request.socket.destroy();
695+
return;
696+
}
697+
response.setHeader('content-type', 'application/json');
698+
response.end(JSON.stringify({ ok: true }));
699+
});
700+
await new Promise<void>((resolve) => {
701+
server?.listen(0, '127.0.0.1', resolve);
702+
});
703+
const { port } = server.address() as AddressInfo;
704+
705+
const localQueue = createQueue({ baseUrl: `http://127.0.0.1:${port}` });
706+
try {
707+
await localQueue.queue('__wkf_workflow_test' as any, workflowPayload);
708+
await vi.waitFor(() => expect(calls).toBe(3), { timeout: 20_000 });
709+
} finally {
710+
await localQueue.close();
711+
}
712+
}, 30_000);
713+
714+
// close() owns a socket pool here too, just Node's rather than undici's. It
715+
// still has to settle, still has to be idempotent, and still has to stop
716+
// in-flight deliveries via the abort controller it owns.
717+
it('closes idempotently', async () => {
718+
const localQueue = createQueue({ baseUrl: 'http://localhost:3000' });
719+
await expect(localQueue.close()).resolves.toBeUndefined();
720+
await expect(localQueue.close()).resolves.toBeUndefined();
721+
});
722+
});

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

Lines changed: 42 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -2,12 +2,18 @@ import { setTimeout } from 'node:timers/promises';
22
import type { Transport } from '@vercel/queue';
33
import { createWorkflowUrl } from '@workflow/utils';
44
import {
5+
isNodeHttpEnabled,
56
MessageId,
67
parseQueueName,
78
type Queue,
89
type QueuePrefix,
910
ValidQueueName,
1011
} from '@workflow/world';
12+
import {
13+
createNodeHttpAgents,
14+
destroyNodeHttpAgents,
15+
nodeHttpFetch,
16+
} from '@workflow/world/node-http.js';
1117
import { Sema } from 'async-sema';
1218
import { monotonicFactory } from 'ulid';
1319
import { Agent } from 'undici';
@@ -78,6 +84,11 @@ function envTimeoutMs(name: string, fallback: number): number {
7884
* Bounds a stalled delivery below the queue handler's retry horizon. A
7985
* transport timeout is retried by the delivery loop with the same durable
8086
* message; `0` remains available for applications that need unbounded calls.
87+
*
88+
* Both transports honour every field. Over undici this is the `Agent`'s own
89+
* configuration; over `node:http` (`WORKFLOW_NODE_HTTP`) the two timeouts are
90+
* passed per request and `connections` / `keepAliveTimeout` size the socket
91+
* pool, so the two env vars tune a delivery identically either way.
8192
*/
8293
export function getQueueAgentOptions() {
8394
return {
@@ -128,7 +139,17 @@ function isDetachedArrayBufferQueueError(error: unknown): boolean {
128139
}
129140

130141
export function createQueue(config: Partial<Config>): LocalQueue {
131-
const httpAgent = new Agent(getQueueAgentOptions());
142+
// Exactly one of these is built, and close() shuts down whichever it is.
143+
// Resolved once per queue rather than per delivery so a single queue never
144+
// mixes transports mid-flight.
145+
const agentOptions = getQueueAgentOptions();
146+
const nodeHttpAgents = isNodeHttpEnabled()
147+
? createNodeHttpAgents({
148+
maxSockets: agentOptions.connections,
149+
keepAliveMs: agentOptions.keepAliveTimeout,
150+
})
151+
: undefined;
152+
const httpAgent = nodeHttpAgents ? undefined : new Agent(agentOptions);
132153
const transport = new TypedJsonTransport();
133154
const generateId = monotonicFactory();
134155
const semaphore = new Sema(WORKFLOW_LOCAL_QUEUE_CONCURRENCY);
@@ -231,17 +252,24 @@ export function createQueue(config: Partial<Config>): LocalQueue {
231252
response = await directHandler(req);
232253
} else {
233254
const baseUrl = await resolveBaseUrl(config);
234-
// eslint-disable-next-line @typescript-eslint/no-explicit-any -- undici v7 dispatcher types don't match @types/node's RequestInit
235-
response = await fetch(
236-
createWorkflowUrl(baseUrl, { type: 'flow' }),
237-
{
238-
method: 'POST',
239-
duplex: 'half',
240-
dispatcher: httpAgent,
241-
headers,
242-
body,
243-
} as any
244-
);
255+
const url = createWorkflowUrl(baseUrl, { type: 'flow' });
256+
response = nodeHttpAgents
257+
? await nodeHttpFetch(url, {
258+
method: 'POST',
259+
headers: new Headers(headers),
260+
body,
261+
agents: nodeHttpAgents,
262+
headersTimeoutMs: agentOptions.headersTimeout,
263+
bodyTimeoutMs: agentOptions.bodyTimeout,
264+
})
265+
: // eslint-disable-next-line @typescript-eslint/no-explicit-any -- undici v7 dispatcher types don't match @types/node's RequestInit
266+
await fetch(url, {
267+
method: 'POST',
268+
duplex: 'half',
269+
dispatcher: httpAgent,
270+
headers,
271+
body,
272+
} as any);
245273
}
246274
delivery++;
247275
text = await response.text();
@@ -433,7 +461,8 @@ export function createQueue(config: Partial<Config>): LocalQueue {
433461
// may close the queue more than once.
434462
if (closeSignal.aborted) return;
435463
closeController.abort();
436-
await httpAgent.close();
464+
if (nodeHttpAgents) destroyNodeHttpAgents(nodeHttpAgents);
465+
await httpAgent?.close();
437466
},
438467
};
439468
}

‎packages/world-vercel/src/events-v4.test.ts‎

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -8,9 +8,10 @@ import {
88
WorkflowWorldError,
99
} from '@workflow/errors';
1010
import type { AnyEventRequest } from '@workflow/world';
11+
import { NODE_HTTP_ENV_VAR } from '@workflow/world';
1112
import { decode, encode } from 'cbor-x';
1213
import { MockAgent } from 'undici';
13-
import { afterEach, describe, expect, it, vi } from 'vitest';
14+
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
1415
import { splitEventDataForV4 } from './events.js';
1516
import {
1617
createWorkflowRunEventV4,
@@ -1506,7 +1507,15 @@ describe('v4 POST frame meta forwards every field the splitter produces', () =>
15061507
* `fetchV4` the recycler is never told anything and the pool lives forever.
15071508
*/
15081509
describe('v4 transport reports failures to the events recycler', () => {
1510+
// There is only an undici pool to retire while the adapter owns one:
1511+
// `WORKFLOW_NODE_HTTP` takes the request off undici and makes
1512+
// getEventsDispatcher return `undefined`, so pin the flag off here.
1513+
beforeEach(() => {
1514+
vi.stubEnv(NODE_HTTP_ENV_VAR, '0');
1515+
});
1516+
15091517
afterEach(() => {
1518+
vi.unstubAllEnvs();
15101519
vi.restoreAllMocks();
15111520
});
15121521

0 commit comments

Comments
 (0)