import { describe, it } from 'node:test'; import assert from 'node:assert/strict'; import { once } from 'node:events'; import { createServer as createNetServer } from 'node:net'; import { request } from 'node:http'; import { workers, assign } from 'moroutine'; import { serverThreads } from 'moroutine/serve'; import { runSlowServer } from './fixtures/slow-server.ts'; function httpGet(port: number): Promise { return 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('error', reject); req.end(); }); } 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(); server.listen(0); await once(server, 'listening'); const port = (server.address() as any).port; { using run = workers(1); using threads = serverThreads(run.workers, server); const serving = run(threads.map(([w, args]) => assign(w, runSlowServer(200, ...args)))); // 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(); assert.equal(await body, 'done'); await serving; } }); it('drainTimeout force-closes long-hung connections', { timeout: 15_000 }, async () => { const server = createNetServer(); server.listen(0); await once(server, 'listening'); const port = (server.address() as any).port; { using run = workers(1, { shutdownTimeout: 10_000 }); using threads = serverThreads(run.workers, server, { listen: { drainTimeout: 200 }, }); const serving = run(threads.map(([w, args]) => assign(w, runSlowServer(10_000, ...args)))); const t0 = Date.now(); // 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; const elapsed = Date.now() - t0; // Either request aborted or succeeded very quickly (well before 10s delay). assert.ok(elapsed < 2000, `expected fast force-close, got ${elapsed}ms`); assert.equal(result, 'aborted'); await serving; } }); });