diff --git a/test/channel-fanout.test.ts b/test/channel-fanout.test.ts index 1c77969..136d637 100644 --- a/test/channel-fanout.test.ts +++ b/test/channel-fanout.test.ts @@ -67,9 +67,9 @@ describe('channel fan-out', () => { lengths.reduce((a, b) => a + b, 0), 4000, ); - // Skip-based RR doesn't guarantee even splits under heterogeneous drain - // rates, but no consumer should receive zero items. + // Skip-based RR guarantees each consumer receives at least its initial + // high-water-mark fill (default 16) before faster consumers pull ahead. const min = Math.min(...lengths); - assert.ok(min > 0, `expected each consumer > 0 items, got ${lengths.join(',')}`); + assert.ok(min >= 16, `expected each consumer ≥ 16 items, got ${lengths.join(',')}`); }); }); diff --git a/test/serve/shutdown.test.ts b/test/serve/shutdown.test.ts index 47183b9..dd024d9 100644 --- a/test/serve/shutdown.test.ts +++ b/test/serve/shutdown.test.ts @@ -19,6 +19,44 @@ function httpGet(port: number): Promise { }); } +function httpGetWithSignal(port: number): { body: Promise; responded: Promise } { + let onResponded: () => void; + const responded = new Promise((r) => { + onResponded = r; + }); + const body = new Promise((resolve, reject) => { + const req = request({ host: '127.0.0.1', port, path: '/' }, (res) => { + onResponded!(); + const chunks: Buffer[] = []; + res.on('data', (c) => chunks.push(Buffer.isBuffer(c) ? c : Buffer.from(c))); + res.on('end', () => resolve(Buffer.concat(chunks).toString('utf8'))); + }); + req.on('error', reject); + req.end(); + }); + return { body, responded }; +} + +function httpGetWithConnect(port: number): { body: Promise; connected: Promise } { + let onConnected: () => void; + const connected = new Promise((r) => { + onConnected = r; + }); + const body = new Promise((resolve, reject) => { + const req = request({ host: '127.0.0.1', port, path: '/' }, (res) => { + const chunks: Buffer[] = []; + res.on('data', (c) => chunks.push(Buffer.isBuffer(c) ? c : Buffer.from(c))); + res.on('end', () => resolve(Buffer.concat(chunks).toString('utf8'))); + }); + req.on('socket', (sock) => { + sock.on('connect', () => onConnected!()); + }); + req.on('error', (e) => reject(e)); + req.end(); + }); + return { body: body.catch(() => 'aborted'), connected }; +} + describe('graceful shutdown', () => { it('in-flight request completes before pool exits', { timeout: 15_000 }, async () => { const server = createNetServer(); @@ -31,13 +69,13 @@ describe('graceful shutdown', () => { using threads = serverThreads(run.workers, server); const serving = run(threads.map(([w, args]) => assign(w, runSlowServer(200, ...args)))); - const responsePromise = httpGet(port); - // Give the worker a moment to accept the connection. - await new Promise((r) => setTimeout(r, 50)); + // Wait for the worker to actually begin handling: send the request and + // wait for the response to start (headers received = worker is mid-flight). + const { body, responded } = httpGetWithSignal(port); + await responded; server.close(); - const body = await responsePromise; - assert.equal(body, 'done'); + assert.equal(await body, 'done'); await serving; } }); @@ -56,8 +94,15 @@ describe('graceful shutdown', () => { const serving = run(threads.map(([w, args]) => assign(w, runSlowServer(10_000, ...args)))); const t0 = Date.now(); - const responsePromise = httpGet(port).catch(() => 'aborted'); - await new Promise((r) => setTimeout(r, 50)); + // Use httpGetWithSignal — the slow server won't respond for 10s, but + // once the worker accepts, it sends headers (which fires `responded`). + // Actually, slow-server only responds after delay, so headers won't + // arrive until 10s. Instead, wait for the TCP socket to connect, then + // give a small grace for the worker to accept the fd. + const { body: responsePromise, connected } = httpGetWithConnect(port); + await connected; + // Small grace for the fd to transit to the worker and be opened. + await new Promise((r) => setTimeout(r, 100)); server.close(); const result = await responsePromise;