From bd21d01a78fd0a9441b26ff22b307f3f95d6521a Mon Sep 17 00:00:00 2001 From: OpenCode Date: Sun, 21 Jun 2026 16:57:12 -0700 Subject: [PATCH] feat: add worker thread adapter runtime --- package.json | 12 +-- src/index.ts | 11 +++ src/server.ts | 152 +++++++++++++++++++++++++++++ src/types.ts | 40 ++++++++ src/worker-route.ts | 227 ++++++++++++++++++++++++++++++++++++++++++++ src/worker.ts | 7 ++ 6 files changed, 443 insertions(+), 6 deletions(-) create mode 100644 src/index.ts create mode 100644 src/server.ts create mode 100644 src/types.ts create mode 100644 src/worker-route.ts create mode 100644 src/worker.ts diff --git a/package.json b/package.json index ce0e263..897c629 100644 --- a/package.json +++ b/package.json @@ -3,19 +3,19 @@ "version": "0.0.1", "description": "Node.js Worker Threads runtime adapter for Hono", "type": "module", - "main": "dist/index.js", - "types": "dist/index.d.ts", + "main": "dist/index.mjs", + "types": "dist/index.d.mts", "files": [ "dist" ], "exports": { ".": { - "types": "./dist/index.d.ts", - "import": "./dist/index.js" + "types": "./dist/index.d.mts", + "import": "./dist/index.mjs" }, "./worker": { - "types": "./dist/worker.d.ts", - "import": "./dist/worker.js" + "types": "./dist/worker.d.mts", + "import": "./dist/worker.mjs" } }, "scripts": { diff --git a/src/index.ts b/src/index.ts new file mode 100644 index 0000000..b2d4f19 --- /dev/null +++ b/src/index.ts @@ -0,0 +1,11 @@ +export { createAdaptorServer, getRequestListener, serve } from './server.js' +export { workerRoute } from './worker-route.js' +export type { + FetchCallback, + HttpBindings, + ServeOptions, + ServerType, + WorkerRouteContext, + WorkerRouteHandler, + WorkerRouteOptions, +} from './types.js' diff --git a/src/server.ts b/src/server.ts new file mode 100644 index 0000000..4899fa6 --- /dev/null +++ b/src/server.ts @@ -0,0 +1,152 @@ +import { createServer as createHttpServer } from 'node:http' +import type { AddressInfo } from 'node:net' +import { Readable } from 'node:stream' +import type { ReadableStream as NodeReadableStream } from 'node:stream/web' +import type { IncomingMessage, ServerResponse } from 'node:http' +import type { FetchCallback, HttpBindings, ServeOptions, ServerType } from './types.js' + +type RequestInitWithDuplex = RequestInit & { duplex?: 'half' } + +const hasRequestBody = (incoming: IncomingMessage): boolean => + incoming.method !== 'GET' && incoming.method !== 'HEAD' + +const headersFromIncoming = (incoming: IncomingMessage): Headers => { + const headers = new Headers() + + for (let i = 0; i < incoming.rawHeaders.length; i += 2) { + headers.append(incoming.rawHeaders[i], incoming.rawHeaders[i + 1]) + } + + return headers +} + +const requestUrl = (incoming: IncomingMessage, hostname?: string): string => { + const incomingUrl = incoming.url || '/' + + if (incomingUrl.startsWith('http://') || incomingUrl.startsWith('https://')) { + return incomingUrl + } + + const host = incoming.headers.host || hostname + if (!host) { + throw new Error('Missing host header') + } + + const protocol = incoming.socket && 'encrypted' in incoming.socket ? 'https' : 'http' + return new URL(incomingUrl, `${protocol}://${host}`).href +} + +const toRequest = ( + incoming: IncomingMessage, + outgoing: ServerResponse, + hostname?: string +): Request => { + const controller = new AbortController() + + outgoing.on('close', () => { + if (!outgoing.writableFinished) { + controller.abort() + } + }) + + const init: RequestInitWithDuplex = { + method: incoming.method, + headers: headersFromIncoming(incoming), + signal: controller.signal, + } + + if (hasRequestBody(incoming)) { + init.body = Readable.toWeb(incoming) as ReadableStream + init.duplex = 'half' + } + + return new Request(requestUrl(incoming, hostname), init) +} + +const outgoingHeaders = (headers: Headers): Record => { + const output: Record = {} + const getSetCookie = (headers as Headers & { getSetCookie?: () => string[] }).getSetCookie + const setCookie = getSetCookie?.call(headers) + + for (const [key, value] of headers) { + if (key !== 'set-cookie') { + output[key] = value + } + } + + if (setCookie && setCookie.length > 0) { + output['set-cookie'] = setCookie + } + + return output +} + +const writeResponse = async (response: Response, outgoing: ServerResponse): Promise => { + outgoing.writeHead(response.status, outgoingHeaders(response.headers)) + + if (!response.body) { + outgoing.end() + return + } + + await new Promise((resolve, reject) => { + const body = Readable.fromWeb(response.body as unknown as NodeReadableStream) + body.on('error', reject) + outgoing.on('error', reject) + outgoing.on('finish', resolve) + body.pipe(outgoing) + }) +} + +const errorResponse = (status: number): Response => + new Response(null, { + status, + }) + +export const getRequestListener = ( + fetchCallback: FetchCallback, + options: { hostname?: string } = {} +) => { + return async (incoming: IncomingMessage, outgoing: ServerResponse): Promise => { + let request: Request + + try { + request = toRequest(incoming, outgoing, options.hostname) + } catch { + await writeResponse(errorResponse(400), outgoing) + return + } + + let response: Response + try { + response = await fetchCallback(request, { incoming, outgoing } satisfies HttpBindings) + } catch { + response = errorResponse(500) + } + + try { + await writeResponse(response, outgoing) + } catch { + if (!outgoing.headersSent) { + outgoing.writeHead(500) + } + outgoing.end() + } + } +} + +export const createAdaptorServer = (options: ServeOptions): ServerType => { + const createServer = options.createServer || createHttpServer + return createServer(options.serverOptions || {}, getRequestListener(options.fetch, options)) +} + +export const serve = ( + options: ServeOptions, + listeningListener?: (info: AddressInfo) => void +): ServerType => { + const server = createAdaptorServer(options) + server.listen(options.port ?? 3000, options.hostname, () => { + listeningListener?.(server.address() as AddressInfo) + }) + return server +} diff --git a/src/types.ts b/src/types.ts new file mode 100644 index 0000000..4be0b94 --- /dev/null +++ b/src/types.ts @@ -0,0 +1,40 @@ +import type { IncomingMessage, Server, ServerOptions, ServerResponse, createServer } from 'node:http' +import type { Context, Env, MiddlewareHandler } from 'hono' + +export type HttpBindings = { + incoming: IncomingMessage + outgoing: ServerResponse +} + +export type FetchCallback = (request: Request, env: HttpBindings) => Response | Promise + +export type ServerType = Server + +export type ServeOptions = { + fetch: FetchCallback + hostname?: string + port?: number + serverOptions?: ServerOptions + createServer?: typeof createServer +} + +export type WorkerRouteContext = { + params: Record + env: Bindings +} + +export type WorkerRouteHandler = ( + request: Request, + context: WorkerRouteContext +) => Response | Promise + +export type WorkerRouteOptions = { + exportName?: string + env?: (context: Context) => Bindings + timeoutMs?: number +} + +export type WorkerRoute = ( + moduleSpecifier: string | URL, + options?: WorkerRouteOptions +) => MiddlewareHandler diff --git a/src/worker-route.ts b/src/worker-route.ts new file mode 100644 index 0000000..af4d5a5 --- /dev/null +++ b/src/worker-route.ts @@ -0,0 +1,227 @@ +import { Worker, type TransferListItem } from 'node:worker_threads' +import type { Context, Env, MiddlewareHandler } from 'hono' +import type { WorkerRouteOptions, WorkerRouteContext } from './types.js' + +type WorkerSuccessMessage = { + type: 'response' + status: number + statusText: string + headers: [string, string][] + body: ReadableStream | null +} + +type WorkerErrorMessage = { + type: 'error' + error: { + name?: string + message?: string + stack?: string + } +} + +type WorkerMessage = WorkerSuccessMessage | WorkerErrorMessage + +const workerSource = ` +import { parentPort } from 'node:worker_threads'; + +const serializeHeaders = (headers) => { + const entries = []; + const setCookie = headers.getSetCookie?.(); + + for (const [key, value] of headers) { + if (key !== 'set-cookie') { + entries.push([key, value]); + } + } + + if (setCookie) { + for (const value of setCookie) { + entries.push(['set-cookie', value]); + } + } + + return entries; +}; + +const serializeError = (error) => ({ + name: error?.name, + message: error?.message || String(error), + stack: error?.stack, +}); + +parentPort.on('message', async (message) => { + try { + const mod = await import(message.moduleSpecifier); + const handler = mod[message.exportName]; + + if (typeof handler !== 'function') { + throw new TypeError('Worker route module does not export a handler function.'); + } + + const requestInit = { + method: message.method, + headers: message.headers, + }; + + if (message.body) { + requestInit.body = message.body; + requestInit.duplex = 'half'; + } + + const request = new Request(message.url, requestInit); + const response = await handler(request, message.context); + + if (!(response instanceof Response)) { + throw new TypeError('Worker route handler must return a Response.'); + } + + const responseMessage = { + type: 'response', + status: response.status, + statusText: response.statusText, + headers: serializeHeaders(response.headers), + body: response.body, + }; + + const transferList = response.body ? [response.body] : []; + parentPort.postMessage(responseMessage, transferList); + } catch (error) { + parentPort.postMessage({ type: 'error', error: serializeError(error) }); + } +}); +` + +const workerUrl = `data:text/javascript;base64,${Buffer.from(workerSource).toString('base64')}` + +const requestHeaders = (request: Request): [string, string][] => { + const entries: [string, string][] = [] + for (const [key, value] of request.headers) { + entries.push([key, value]) + } + return entries +} + +const toError = (message: WorkerErrorMessage): Error => { + const error = new Error(message.error.message || 'Worker route failed') + error.name = message.error.name || 'Error' + error.stack = message.error.stack + return error +} + +const workerResponse = (message: WorkerSuccessMessage): Response => { + return new Response(message.body, { + status: message.status, + statusText: message.statusText, + headers: message.headers, + }) +} + +const runWorkerRoute = async ( + moduleSpecifier: string, + exportName: string, + request: Request, + context: WorkerRouteContext, + timeoutMs?: number +): Promise => { + const worker = new Worker(new URL(workerUrl)) + const body = request.body + const transferList: TransferListItem[] = body ? [body as unknown as TransferListItem] : [] + + let timeout: NodeJS.Timeout | undefined + + try { + return await new Promise((resolve, reject) => { + const cleanup = () => { + if (timeout) { + clearTimeout(timeout) + } + request.signal.removeEventListener('abort', abort) + worker.off('message', onMessage) + worker.off('error', onError) + worker.off('exit', onExit) + } + + const abort = () => { + worker.terminate().catch(() => {}) + cleanup() + reject(request.signal.reason || new DOMException('The request was aborted.', 'AbortError')) + } + + const onMessage = (message: WorkerMessage) => { + cleanup() + if (message.type === 'error') { + reject(toError(message)) + } else { + resolve(workerResponse(message)) + } + } + + const onError = (error: Error) => { + cleanup() + reject(error) + } + + const onExit = (code: number) => { + cleanup() + if (code !== 0) { + reject(new Error(`Worker route exited with code ${code}.`)) + } + } + + if (timeoutMs !== undefined) { + timeout = setTimeout(() => { + worker.terminate().catch(() => {}) + cleanup() + reject(new Error(`Worker route timed out after ${timeoutMs}ms.`)) + }, timeoutMs) + } + + request.signal.addEventListener('abort', abort, { once: true }) + worker.on('message', onMessage) + worker.on('error', onError) + worker.on('exit', onExit) + worker.postMessage( + { + moduleSpecifier, + exportName, + method: request.method, + url: request.url, + headers: requestHeaders(request), + body, + context, + }, + transferList + ) + }) + } finally { + worker.terminate().catch(() => {}) + } +} + +const routeContext = ( + context: Context, + options: WorkerRouteOptions +): WorkerRouteContext => { + return { + params: context.req.param(), + env: options.env ? options.env(context) : (context.env as Bindings), + } +} + +export const workerRoute = ( + moduleSpecifier: string | URL, + options: WorkerRouteOptions = {} +): MiddlewareHandler => { + const moduleUrl = String(moduleSpecifier) + const exportName = options.exportName || 'default' + + return async (context) => { + return runWorkerRoute( + moduleUrl, + exportName, + context.req.raw, + routeContext(context, options), + options.timeoutMs + ) + } +} diff --git a/src/worker.ts b/src/worker.ts new file mode 100644 index 0000000..0210f17 --- /dev/null +++ b/src/worker.ts @@ -0,0 +1,7 @@ +import type { WorkerRouteHandler } from './types.js' + +export type { WorkerRouteContext, WorkerRouteHandler } from './types.js' + +export const defineWorkerRoute = ( + handler: WorkerRouteHandler +): WorkerRouteHandler => handler -- 2.51.2