diff --git a/src/serve/listen.ts b/src/serve/listen.ts index 1ae4c49..c64dd51 100644 --- a/src/serve/listen.ts +++ b/src/serve/listen.ts @@ -1,21 +1,19 @@ import { Socket } from 'node:net'; import type { Server } from 'node:net'; -import type { Tuple } from '../shared/tuple.ts'; import type { Int32Atomic } from '../shared/int32-atomic.ts'; import type { ListenOptions } from './types.ts'; /** - * Drain fds from `ch` and emit each as a `'connection'` on `server`. Maintains - * the per-worker counter (decrements on each socket's 'close' event). Resolves - * after the channel ends and either: (a) all active sockets have closed, or - * (b) `opts.drainTimeout` elapses — in which case any remaining sockets are + * Drain fds from the iterable and emit each as a `'connection'` on `server`. + * Decrements `counter` on each socket's 'close' event. Resolves after the + * iterable ends and either: (a) all active sockets have closed, or + * (b) `opts.drainTimeout` elapses — in which case remaining sockets are * force-destroyed. */ export async function listen( server: Server, fds: AsyncIterable, - counters: Tuple, - slot: number, + counter: Int32Atomic, opts: Required, ): Promise { const active = new Set(); @@ -35,7 +33,7 @@ export async function listen( active.add(socket); socket.once('close', () => { active.delete(socket); - counters.elements[slot].sub(1); + counter.sub(1); tryResolve(); }); server.emit('connection', socket); diff --git a/src/serve/server-threads.ts b/src/serve/server-threads.ts index de07803..3b1e486 100644 --- a/src/serve/server-threads.ts +++ b/src/serve/server-threads.ts @@ -76,7 +76,7 @@ export function serverThreads( server.once('close', onClose); // Build the indexable/iterable/disposable result. - const entries = workers.map((w, i) => [w, [channels[i], counters, i, listenOpts] as ListenArgs] as const); + const entries = workers.map((w, i) => [w, [channels[i], counters.elements[i], listenOpts] as ListenArgs] as const); const result: readonly (readonly [WorkerHandle, ListenArgs])[] = entries; Object.defineProperty(result, Symbol.dispose, { value: dispose, enumerable: false }); return result as unknown as ServerThreads; diff --git a/src/serve/types.ts b/src/serve/types.ts index b074efa..84a1ea3 100644 --- a/src/serve/types.ts +++ b/src/serve/types.ts @@ -1,4 +1,3 @@ -import type { Tuple } from '../shared/tuple.ts'; import type { Int32Atomic } from '../shared/int32-atomic.ts'; import type { WorkerHandle } from '../runner.ts'; @@ -12,12 +11,7 @@ export interface ListenOptions { * Opaque tuple passed from `serverThreads` to each worker's moroutine. * The shape may evolve; users only ever spread it into `listen()`. */ -export type ListenArgs = readonly [ - fds: AsyncIterable, - counters: Tuple, - slot: number, - opts: Required, -]; +export type ListenArgs = readonly [fds: AsyncIterable, counter: Int32Atomic, opts: Required]; /** Connection-routing strategy. Picks a worker index given a counter snapshot. */ export interface Balance { diff --git a/test/serve/listen.test.ts b/test/serve/listen.test.ts index ae525e1..2143693 100644 --- a/test/serve/listen.test.ts +++ b/test/serve/listen.test.ts @@ -4,7 +4,8 @@ import { once } from 'node:events'; import { openSync } from 'node:fs'; import { createServer as createNetServer, connect } from 'node:net'; import { createServer } from 'node:http'; -import { shared, int32atomic } from 'moroutine'; +import { int32atomic } from 'moroutine'; +import { Int32Atomic } from '../../src/shared/int32-atomic.ts'; import { pushChannel } from '../../src/serve/push-channel.ts'; import { listen } from 'moroutine/serve'; @@ -35,11 +36,11 @@ async function acquireLocalFd(): Promise<{ fd: number; close: () => void }> { describe('listen()', () => { it('emits received fds as connections on the user server', { timeout: 5000 }, async () => { const ch = pushChannel({ highWaterMark: 4 }); - const counters = shared([int32atomic]); + const counter = new Int32Atomic(); const http = createServer((req, res) => { res.end('ok'); }); - const drained = listen(http, ch, counters, 0, { drainTimeout: 5_000 }); + const drained = listen(http, ch, counter, { drainTimeout: 5_000 }); const { fd, close: closeSrc } = await acquireLocalFd(); ch.send(fd); @@ -53,19 +54,19 @@ describe('listen()', () => { it('decrements counter on socket close', { timeout: 5000 }, async () => { const ch = pushChannel({ highWaterMark: 4 }); - const counters = shared([int32atomic]); - counters.elements[0].store(1); // simulate main having incremented + const counter = new Int32Atomic(); + counter.store(1); // simulate main having incremented const http = createServer((req, res) => { res.end('ok'); }); - const drained = listen(http, ch, counters, 0, { drainTimeout: 5_000 }); + const drained = listen(http, ch, counter, { drainTimeout: 5_000 }); const { fd, close: closeSrc } = await acquireLocalFd(); ch.send(fd); await new Promise((r) => setTimeout(r, 50)); closeSrc(); // destroy client so socket closes → counter decrements await new Promise((r) => setTimeout(r, 100)); - assert.equal(counters.elements[0].load(), 0); + assert.equal(counter.load(), 0); ch.close(); await drained; });