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/events-streamed-body-h1.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
'@workflow/world-vercel': patch
---

Send event requests with bodies too large to re-buffer over HTTP/1.1 so they no longer stall behind in-flight HTTP/2 streams, bound events requests at 60s instead of 300s, and retry event-log reads whose HTTP/2 stream the peer reset.
7 changes: 7 additions & 0 deletions .changeset/undici-7-30-h2-running-slots.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
---
'@workflow/world-vercel': patch
'@workflow/world-local': patch
'@workflow/world-postgres': patch
---

Upgrade undici to 7.30.0, which stops a failed HTTP/2 stream from leaving a phantom in-flight request on its connection.
235 changes: 234 additions & 1 deletion packages/world-vercel/src/http-client.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,9 @@
import { createServer as createHttpServer } from 'node:http';
import { createSecureServer, type Http2SecureServer } from 'node:http2';
import {
createSecureServer,
type Http2SecureServer,
constants as http2Constants,
} from 'node:http2';
import { type AddressInfo, connect, createServer, type Server } from 'node:net';
import type { TLSSocket } from 'node:tls';
import { NODE_HTTP_ENV_VAR } from '@workflow/world';
Expand All @@ -25,6 +29,7 @@ import {
EVENTS_AGENT_OPTIONS,
EVENTS_AGENT_OPTIONS_NO_H2,
EVENTS_RECYCLE_AFTER_CONSECUTIVE_FAILURES,
EVENTS_REQUEST_TIMEOUT_MS,
getDispatcher,
getEventsDispatcher,
getNodeHttpAgents,
Expand All @@ -33,6 +38,7 @@ import {
getQueueRequestTimeoutMs,
getStreamCloseDispatcher,
getStreamDispatcher,
h2MultiplexInterceptor,
isRecyclableTransportError,
NODE_HTTP_BODY_TIMEOUT_MS,
NODE_HTTP_HEADERS_TIMEOUT_MS,
Expand Down Expand Up @@ -145,6 +151,17 @@ describe('agent transport', () => {
// fault, so it has to leave nothing of H2 behind: `allowH2: true` with the
// interceptor skipped still keeps a wedged session (see
// EVENTS_AGENT_OPTIONS_NO_H2).
// undici's 300s defaults outlast the runtime's 240s replay budget, so a
// silent stream would cost the whole delivery before it even failed.
it('arms events deadlines well inside the replay budget', () => {
const REPLAY_BUDGET_MS = 240_000;
for (const options of [EVENTS_AGENT_OPTIONS, EVENTS_AGENT_OPTIONS_NO_H2]) {
expect(options.headersTimeout).toBe(EVENTS_REQUEST_TIMEOUT_MS);
expect(options.bodyTimeout).toBe(EVENTS_REQUEST_TIMEOUT_MS);
}
expect(EVENTS_REQUEST_TIMEOUT_MS).toBeLessThan(REPLAY_BUDGET_MS / 2);
});

it('gives the kill switch an events agent with no HTTP/2 at all', () => {
expect(EVENTS_AGENT_OPTIONS_NO_H2.allowH2).toBe(false);
expect(EVENTS_AGENT_OPTIONS_NO_H2.pipelining).toBe(1);
Expand Down Expand Up @@ -476,6 +493,222 @@ describe('HTTP/2 multiplexing (events vs stream-write agents)', () => {
});
});

// undici dispatches a streamed request body on HTTP/2 only once the connection
// has no stream in flight, and stops its whole queue behind it until then. The
// multiplexing interceptor re-buffers small bodies to avoid that, so the
// production shape that still hit it was a body too large to re-buffer (a
// multi-MiB run input or event batch) on a connection that never drained.
describe('events dispatcher on a busy HTTP/2 connection', () => {
const LARGE_BODY_BYTES = 2 * 1024 * 1024;
const LOOPBACK = { connect: { rejectUnauthorized: false } };

let server: Http2SecureServer;
let origin: string;
let seen: Array<{ path: string; httpVersion: string; bytes: number }>;
let releaseHold: (() => void) | undefined;
let holdArrived: Promise<void>;
let resetAttempts: number;

beforeAll(async () => {
// allowHTTP1 so one origin serves both the h2 stream and the h1 fallback.
server = createSecureServer({
key: TEST_KEY,
cert: TEST_CERT,
allowHTTP1: true,
});
server.on('sessionError', () => undefined);
server.on('clientError', () => undefined);
server.on('request', (req, res) => {
req.stream?.on('error', () => undefined);
const path = String(req.url);
let bytes = 0;
req.on('data', (chunk: Buffer) => {
bytes += chunk.length;
});
req.on('end', () => {
seen.push({ path, httpVersion: req.httpVersion, bytes });
if (path === '/hold') {
holdArrivedResolve();
releaseHold = () => {
res.writeHead(200);
res.end('held');
};
return;
}
// The edge's per-stream reset on a busy connection. `/reset-once`
// answers the retry; `/reset-always` never does.
if (
(path === '/reset-once' && ++resetAttempts === 1) ||
path === '/reset-always'
) {
if (path === '/reset-always') resetAttempts++;
req.stream.close(http2Constants.NGHTTP2_ENHANCE_YOUR_CALM);
return;
}
res.writeHead(200);
res.end(path);
});
});
await new Promise<void>((resolve) => {
server.listen(0, '127.0.0.1', resolve);
});
origin = `https://127.0.0.1:${(server.address() as AddressInfo).port}`;
});

afterAll(async () => {
await new Promise<void>((resolve) => {
server.close(() => resolve());
});
});

let holdArrivedResolve: () => void;
beforeEach(() => {
seen = [];
releaseHold = undefined;
resetAttempts = 0;
holdArrived = new Promise<void>((resolve) => {
holdArrivedResolve = resolve;
});
});

/**
* Opens a GET that the origin holds, so the H2 connection never drains.
* Resolves once the origin has it, to a wrapper (not the response promise
* itself, which an async return would await).
*/
async function holdStreamOpen(dispatcher: unknown) {
const held = fetch(`${origin}/hold`, {
dispatcher,
// eslint-disable-next-line @typescript-eslint/no-explicit-any -- undici dispatcher type doesn't match @types/node's RequestInit
} as any).then((r) => r.text());
await holdArrived;
return { held };
}

function postLarge(dispatcher: unknown) {
return fetch(`${origin}/large`, {
method: 'POST',
body: new Uint8Array(LARGE_BODY_BYTES),
dispatcher,
// eslint-disable-next-line @typescript-eslint/no-explicit-any -- undici dispatcher type doesn't match @types/node's RequestInit
} as any).then((r) => r.text());
}

/** Resolves to 'settled' or 'pending' after `ms`. */
function stateAfter(promise: Promise<unknown>, ms: number) {
return Promise.race([
promise.then(() => 'settled' as const),
new Promise<'pending'>((resolve) => {
setTimeout(() => resolve('pending'), ms);
}),
]);
}

// Establishes the premise with the interceptor alone (no fallback): the
// large POST cannot start while the held stream is in flight.
it('queues behind an in-flight H2 stream without the fallback (the failure this fixes)', async () => {
const agent = new Agent({ ...EVENTS_AGENT_OPTIONS, ...LOOPBACK });
const dispatcher = agent.compose(h2MultiplexInterceptor);
try {
const { held } = await holdStreamOpen(dispatcher);
const large = postLarge(dispatcher);
expect(await stateAfter(large, 1_000)).toBe('pending');
expect(seen.map((s) => s.path)).toEqual(['/hold']);

releaseHold?.();
expect(await held).toBe('held');
expect(await large).toBe('/large');
} finally {
await agent.close();
}
});

it('sends it over HTTP/1.1 so it does not wait for the H2 connection to drain', async () => {
const agent = createEventsDispatcher(LOOPBACK);
try {
const { held } = await holdStreamOpen(agent);
const large = postLarge(agent);
expect(await stateAfter(large, 5_000)).toBe('settled');
expect(await large).toBe('/large');
expect(seen).toContainEqual({
path: '/large',
httpVersion: '1.1',
bytes: LARGE_BODY_BYTES,
});

releaseHold?.();
expect(await held).toBe('held');
expect(seen).toContainEqual({
path: '/hold',
httpVersion: '2.0',
bytes: 0,
});
} finally {
await agent.close();
}
});

it('retries an event-log read whose stream the peer reset', async () => {
const agent = createEventsDispatcher(LOOPBACK);
try {
const response = await fetch(`${origin}/reset-once`, {
dispatcher: agent,
// eslint-disable-next-line @typescript-eslint/no-explicit-any -- undici dispatcher type doesn't match @types/node's RequestInit
} as any);
expect(await response.text()).toBe('/reset-once');
expect(resetAttempts).toBe(2);
} finally {
await agent.close();
}
});

// undici before 7.30.0 retired the wrong request when an H2 stream failed out
// of order, leaving a phantom "running" slot on the connection for good
// (nodejs/undici#5410, #5569). A streamed body then never dispatched on that
// connection again, since it waits for zero running streams. Pins the undici
// version this package depends on.
it('leaves no phantom running stream after a peer reset', async () => {
const agent = new Agent({ ...EVENTS_AGENT_OPTIONS, ...LOOPBACK });
try {
await (
await fetch(`${origin}/warm`, {
dispatcher: agent,
// eslint-disable-next-line @typescript-eslint/no-explicit-any -- undici dispatcher type doesn't match @types/node's RequestInit
} as any)
).text();
await expect(
fetch(`${origin}/reset-always`, {
dispatcher: agent,
// eslint-disable-next-line @typescript-eslint/no-explicit-any -- undici dispatcher type doesn't match @types/node's RequestInit
} as any)
).rejects.toThrow();
await new Promise((resolve) => setTimeout(resolve, 100));
expect(agent.stats[origin]?.running).toBe(0);
} finally {
await agent.close();
}
});

// An event write may already have committed when its stream is reset, so it
// must surface the error rather than be replayed.
it('does not retry an event write whose stream the peer reset', async () => {
const agent = createEventsDispatcher(LOOPBACK);
try {
await expect(
fetch(`${origin}/reset-always`, {
method: 'POST',
body: 'x',
dispatcher: agent,
// eslint-disable-next-line @typescript-eslint/no-explicit-any -- undici dispatcher type doesn't match @types/node's RequestInit
} as any)
).rejects.toThrow();
expect(resetAttempts).toBe(1);
} finally {
await agent.close();
}
});
});

// The transport fault this whole mechanism exists for: an HTTP/2 session whose
// TCP connection stays established while no bytes cross it. undici keeps such a
// session in service — on a stream timeout it deliberately does not destroy the
Expand Down
Loading
Loading