diff --git a/src/serve/server-threads.ts b/src/serve/server-threads.ts index 0d2e05f..a1f1b01 100644 --- a/src/serve/server-threads.ts +++ b/src/serve/server-threads.ts @@ -1,8 +1,4 @@ -import { closeSync, openSync } from 'node:fs'; -import { platform } from 'node:process'; import type { Server } from 'node:net'; - -const fdPath = (fd: number) => (platform === 'linux' ? `/proc/self/fd/${fd}` : `/dev/fd/${fd}`); import { channel } from '../channel.ts'; import { shared } from '../shared/shared.ts'; import { int32atomic } from '../shared/descriptors.ts'; @@ -49,20 +45,22 @@ export function serverThreads( const selected = balance.select(threads); const idx = workers.indexOf(selected.worker); - // Dup the fd via /dev/fd so the new descriptor is independent of libuv's - // event loop tracking. This lets socket.destroy() properly close the - // original fd (freeing libuv state) while the dup'd fd survives for the - // worker to open as a fresh Socket. - const fd = openSync(fdPath(origFd), 'r+'); + // Detach the socket from libuv without closing the fd. readStop() + // deregisters the read watcher; no-oping close() prevents destroy() + // from closing the underlying fd. The worker will open it fresh. + const handle = socket._handle; + handle.readStop(); + handle.close = (cb: any) => { + if (cb) cb(); + }; + socket._handle = null; socket.destroy(); counters.elements[idx].add(1); try { - pushChannels[idx].send(fd); + pushChannels[idx].send(origFd); } catch { - // send threw because the channel is closed — undo the increment and close the fd. counters.elements[idx].sub(1); - closeSync(fd); } }; diff --git a/test/serve/listen.test.ts b/test/serve/listen.test.ts index 8ae7077..7755b6f 100644 --- a/test/serve/listen.test.ts +++ b/test/serve/listen.test.ts @@ -1,11 +1,7 @@ import { describe, it } from 'node:test'; import assert from 'node:assert/strict'; import { once } from 'node:events'; -import { openSync } from 'node:fs'; -import { platform } from 'node:process'; import { createServer as createNetServer, connect } from 'node:net'; - -const fdPath = (fd: number) => (platform === 'linux' ? `/proc/self/fd/${fd}` : `/dev/fd/${fd}`); import { createServer } from 'node:http'; import { int32atomic } from 'moroutine'; import { Int32Atomic } from '../../src/shared/int32-atomic.ts'; @@ -13,8 +9,8 @@ import { pushChannel } from '../../src/serve/push-channel.ts'; import { listen } from 'moroutine/serve'; // Helper: accept a real TCP connection on a scratch listener and yield its fd. -// We dup the fd via /dev/fd so that the new fd is not registered in libuv's -// event loop — this lets new Socket({ fd }) succeed in the same thread. +// Detaches the socket from libuv (readStop + no-op close) so the fd can be +// reused by listen() in the same thread without conflict. async function acquireLocalFd(): Promise<{ fd: number; close: () => void }> { const srv = createNetServer(); srv.listen(0); @@ -22,9 +18,12 @@ async function acquireLocalFd(): Promise<{ fd: number; close: () => void }> { const port = (srv.address() as any).port; const client = connect(port); const [peer] = (await once(srv, 'connection')) as [any]; - const origFd: number = peer._handle.fd; - const fd = openSync(fdPath(origFd), 'r+'); // dup — new fd, not tracked by libuv - peer.pause(); + const handle = peer._handle; + const fd: number = handle.fd; + handle.readStop(); + handle.close = (cb: any) => { + if (cb) cb(); + }; peer._handle = null; peer.destroy(); return { diff --git a/test/serve/spike.test.ts b/test/serve/spike.test.ts index 3c424af..2ce55ce 100644 --- a/test/serve/spike.test.ts +++ b/test/serve/spike.test.ts @@ -1,11 +1,7 @@ import { describe, it } from 'node:test'; import assert from 'node:assert/strict'; import { once } from 'node:events'; -import { openSync } from 'node:fs'; -import { platform } from 'node:process'; import { createServer, connect } from 'node:net'; - -const fdPath = (fd: number) => (platform === 'linux' ? `/proc/self/fd/${fd}` : `/dev/fd/${fd}`); import { Worker } from 'node:worker_threads'; describe('fd-passing spike', () => { @@ -25,9 +21,15 @@ describe('fd-passing spike', () => { // Main TCP server const server = createServer(); - server.on('connection', (socket) => { - // Dup the fd via /dev/fd so the new descriptor is independent of libuv. - const fd = openSync(fdPath((socket as any)._handle.fd), 'r+'); + server.on('connection', (socket: any) => { + const handle = socket._handle; + const fd: number = handle.fd; + // Deregister libuv's read interest; no-op close to preserve the fd. + handle.readStop(); + handle.close = (cb: any) => { + if (cb) cb(); + }; + socket._handle = null; socket.destroy(); worker.postMessage({ fd }); });