diff --git a/.changeset/server-threads.md b/.changeset/server-threads.md new file mode 100644 index 0000000..a80f568 --- /dev/null +++ b/.changeset/server-threads.md @@ -0,0 +1,22 @@ +--- +'moroutine': minor +--- + +New subpath `moroutine/serve`: fd-fanout server primitive + +Scale an ordinary Node `http.Server` across worker threads by passing raw socket fds from the main thread to workers. Each worker runs its own server; main just accepts connections and routes them. Includes `leastConns()` (default) and `roundRobin()` balance strategies, graceful shutdown via `server.close()`, and composition with any framework that exposes an `http.Server` (express, fastify, etc.). POSIX-only. + +```ts +// worker +export const runServer = mo(import.meta, async (...args: ListenArgs): Promise => { + const srv = createServer((req, res) => { + res.end('hi'); + }); + await listen(srv, ...args); +}); + +// main +await using run = workers(4); +using threads = serverThreads(run.workers, server); +await run(threads.map(([w, args]) => assign(w, runServer(...args)))); +``` diff --git a/README.md b/README.md index f34588a..3b73455 100644 --- a/README.md +++ b/README.md @@ -366,6 +366,58 @@ const pos = shared({ x: int32, y: int32 }); console.log(pos.load()); // { x: 1, y: 1 } ``` +## Server Threads + +Scale a Node HTTP server across worker threads by fanning out socket fds. Main accepts TCP connections; workers each run their own `http.Server` and handle requests. Imported from the `moroutine/serve` subpath. + +```ts +// app-server.ts +import { createServer } from 'node:http'; +import { mo } from 'moroutine'; +import { listen, type ListenArgs } from 'moroutine/serve'; + +export const runServer = mo(import.meta, async (...args: ListenArgs): Promise => { + const srv = createServer((req, res) => { + res.writeHead(200); + res.end('hello'); + }); + await listen(srv, ...args); +}); +``` + +```ts +// main.ts +import { createServer } from 'node:net'; +import { workers, assign } from 'moroutine'; +import { serverThreads } from 'moroutine/serve'; +import { runServer } from './app-server.ts'; + +const server = createServer(); +server.listen(3000); +process.on('SIGINT', () => server.close()); + +{ + await using run = workers(4); + using threads = serverThreads(run.workers, server); + await run(threads.map(([w, args]) => assign(w, runServer(...args)))); +} +``` + +Connection routing defaults to `leastConns()`; pass `{ balance: roundRobin() }` for simple fan-out. TLS termination runs on workers (`https.createServer` / `tls.createServer` inside your moroutine); the main-thread listener stays raw TCP. POSIX-only (Linux, macOS). Frameworks that expose an underlying `http.Server` compose without glue: + +```ts +// express +const srv = createServer(express()); +await listen(srv, ...args); + +// fastify +const app = fastify(); +await app.ready(); +await listen(app.server, ...args); +``` + +Graceful shutdown: call `server.close()` on the main listener; the cascade closes worker channels, drains in-flight requests (up to `drainTimeout`, default 30s), and resolves the moroutines. Combine with `await using run = workers(...)` for full pool teardown. + ## Streaming ### Streaming Moroutines