import { createServer, request as httpRequest, type IncomingMessage, type Server, type ServerResponse, } from "node:http"; import { request as httpsRequest } from "node:https"; import type { Socket } from "node:net"; import { Readable, type Duplex } from "node:stream"; import { pipeline } from "node:stream/promises"; import type { Redis } from "ioredis"; import type { Registry as MetricsRegistry } from "prom-client"; import { Agent, request as undiciRequest } from "undici"; import { createPinboard } from "../server.js"; import type { Registry } from "../registry.js"; import { json } from "../http.js"; import { authenticateRegistry } from "../registry-api.js"; import type { HostRule } from "../auth.js"; import { isConnectError, isReservedPinboardHeader, REFUSAL_HEADER, type Transport, } from "../proxy.js"; import { RedisKVMarker } from "../keys.js"; import { noopMetrics } from "../metrics.js"; import { IoredisKV } from "./redis.js"; import { createMetrics } from "./metrics.js"; const REASONS: Record = { 400: "Bad Request", 403: "Forbidden", 404: "Not Found", 503: "Service Unavailable", 508: "Loop Detected", }; export const toRequest = ( req: IncomingMessage, signal?: AbortSignal, ): Request => { const headers = new Headers(); for (let i = 0; i < req.rawHeaders.length; i += 2) { headers.append(req.rawHeaders[i]!, req.rawHeaders[i + 1]!); } const method = req.method ?? "GET"; const hasBody = method !== "GET" && method !== "HEAD" && (req.headers["content-length"] !== undefined || req.headers["transfer-encoding"] !== undefined); return new Request(new URL(req.url ?? "/", "http://pinboard.invalid"), { method, headers, // toWeb instead of the raw IncomingMessage: Bun's undici cannot consume a // Node stream body, and web streams work on both runtimes. ...(hasBody && { body: Readable.toWeb(req), duplex: "half" }), ...(signal !== undefined && { signal }), } as RequestInit); }; export const writeResponse = async ( res: ServerResponse, response: Response, ): Promise => { const headers: Record = {}; response.headers.forEach((value, name) => { if (name !== "set-cookie") headers[name] = value; }); const cookies = response.headers.getSetCookie(); if (cookies.length > 0) headers["set-cookie"] = cookies; res.writeHead(response.status, headers); if (response.body === null) { res.end(); return; } try { await pipeline( Readable.fromWeb(response.body as Parameters[0]), res, ); } catch { res.destroy(); await response.body.cancel().catch(() => {}); } }; type UpgradeContext = { socket: Duplex; head: Buffer }; export namespace createNodeTransport { export type Options = { /** Time allowed for a backend to produce response headers (or the WS 101), default 30s. */ responseHeaderTimeoutMs?: number | undefined; }; } export type NodeTransport = Transport.Options & { bindUpgrade(request: Request, ctx: UpgradeContext): void; completed(response: Response): boolean; close(): Promise; }; export const createNodeTransport = ( options?: createNodeTransport.Options, ): NodeTransport => { const responseHeaderTimeoutMs = options?.responseHeaderTimeoutMs ?? 30_000; // Body idle: backends heartbeat SSE streams every 15s (`: ping`), so 60s of // body silence means a dead peer, not a quiet stream. const agent = new Agent({ keepAliveTimeout: 30_000, connectTimeout: 5_000, headersTimeout: responseHeaderTimeoutMs, bodyTimeout: 60_000, }); const upgradeSockets = new Set(); const contexts = new WeakMap(); const markers = new WeakSet(); const upgrade: Transport.Upgrade = (req, outbound, hooks) => { const ctx = contexts.get(req); if (ctx === undefined) { throw new Error("upgrade requires a bound node socket"); } const { socket, head } = ctx; const headers: Record = {}; outbound.headers.forEach((value, name) => { headers[name] = value; }); headers["connection"] = "Upgrade"; headers["upgrade"] = "websocket"; const requestUpgrade = new URL(outbound.url).protocol === "https:" ? httpsRequest : httpRequest; // Mirror the undici agent's phases: 5s to connect, then the configured // response-header window for the 101. const upstream = requestUpgrade(outbound.url, { method: req.method, headers, timeout: 5_000, }); upstream.on("socket", (upstreamSocket) => { const widen = (): void => { upstream.setTimeout(responseHeaderTimeoutMs); }; if (upstreamSocket.connecting) upstreamSocket.once("connect", widen); else widen(); }); // A Response cannot carry status 101; the marker tells the adapter the // socket was handled in place. const marker = new Response(null, { status: 200 }); markers.add(marker); return new Promise((resolve) => { let backendSocket: Duplex | null = null; const onSocketClose = (): void => { upgradeSockets.delete(socket); upstream.destroy(); backendSocket?.destroy(); hooks.onClose(); resolve(marker); }; const onSocketError = (): void => { hooks.onFailure("aborted"); socket.destroy(); }; upgradeSockets.add(socket); socket.on("close", onSocketClose); socket.on("error", onSocketError); upstream.on("timeout", () => upstream.destroy( upstream.socket?.connecting === true ? Object.assign(new Error("connect timeout"), { code: "UND_ERR_CONNECT_TIMEOUT", }) : new Error("upstream timeout"), ), ); upstream.on("error", (err) => { hooks.onFailure(isConnectError(err) ? "unreachable" : "failed"); if (socket.writable) { socket.end("HTTP/1.1 502 Bad Gateway\r\nconnection: close\r\n\r\n"); } else { socket.destroy(); } resolve(marker); }); upstream.on("response", (res) => { upstream.destroy(); const refusal = res.headers[REFUSAL_HEADER]; if (typeof refusal === "string") { // A refusal rides back to the adapter for the NACK re-place; the // client socket stays open for the retried upgrade. socket.off("close", onSocketClose); socket.off("error", onSocketError); upgradeSockets.delete(socket); hooks.onClose(); resolve( new Response(null, { status: res.statusCode ?? 503, headers: { [REFUSAL_HEADER]: refusal }, }), ); return; } // The backend refused the upgrade: relay the status line and close. socket.end( `HTTP/1.1 ${res.statusCode} ${res.statusMessage ?? ""}\r\nconnection: close\r\n\r\n`, ); resolve(marker); }); upstream.on("upgrade", (res, upstreamSocket, upstreamHead) => { backendSocket = upstreamSocket; const teardown = (): void => { socket.destroy(); upstreamSocket.destroy(); }; socket.off("error", onSocketError); socket.on("error", teardown); upstreamSocket.setTimeout(0); // TCP keepalive is the tunnel's only dead-peer detection: a silently // vanished peer (power loss, partition) never emits close/error. (upstreamSocket as Socket).setKeepAlive?.(true, 30_000); (socket as Socket).setKeepAlive?.(true, 30_000); const lines = [ `HTTP/1.1 101 ${res.statusMessage ?? "Switching Protocols"}`, ]; for (let i = 0; i < res.rawHeaders.length; i += 2) { const name = res.rawHeaders[i]!; if ( isReservedPinboardHeader(name) || name.toLowerCase() === REFUSAL_HEADER ) { continue; } lines.push(`${name}: ${res.rawHeaders[i + 1]}`); } socket.write(lines.join("\r\n") + "\r\n\r\n"); if (upstreamHead.length > 0) socket.write(upstreamHead); if (head.length > 0) upstreamSocket.write(head); // WebSocket close is frame-level: TCP half-close from either side // (server sockets are allowHalfOpen) tears the whole tunnel down. upstreamSocket.on("error", () => { hooks.onFailure("failed"); teardown(); }); upstreamSocket.on("close", teardown); upstreamSocket.on("end", teardown); socket.on("end", teardown); socket.pipe(upstreamSocket); upstreamSocket.pipe(socket); resolve(marker); }); upstream.end(); }); }; return { // undici request over undici fetch: fetch injects user-agent/accept-* // headers and transparently decompresses, breaking byte-identical proxying. fetch: async (request) => { const upstream = await undiciRequest(request.url, { method: request.method as "GET", headers: Object.fromEntries(request.headers), // Web streams are async iterables at runtime (hence the cast); Bun's // undici cannot consume a Node stream body. ...(request.body !== null && { body: request.body as unknown as Readable, }), signal: request.signal, dispatcher: agent, }); const headers = new Headers(); for (const [name, value] of Object.entries(upstream.headers)) { if (value === undefined) continue; for (const entry of Array.isArray(value) ? value : [value]) { headers.append(name, entry); } } const bodyless = upstream.statusCode === 204 || upstream.statusCode === 205 || upstream.statusCode === 304; if (bodyless) void upstream.body.dump().catch(() => {}); return new Response( bodyless ? null : (Readable.toWeb(upstream.body) as ReadableStream), { status: upstream.statusCode, headers }, ); }, upgrade, bindUpgrade: (request, ctx) => { contexts.set(request, ctx); }, completed: (response) => markers.has(response), close: async () => { for (const socket of upgradeSockets) socket.destroy(); // Bun's builtin undici Agent has no close(). await agent.close?.(); }, }; }; export namespace servePinboard { export type Options = Omit< createPinboard.Options, "redis" | "transport" | "metrics" > & createNodeTransport.Options & { /** An ioredis instance (owned and quit by `close()`), or any RedisKV. */ redis: Redis | createPinboard.RedisKV; /** prom-client Registry backing GET /metrics, served to pinboard:operator credentials. Absent = metrics off. */ metricsRegistry?: MetricsRegistry | undefined; }; export type Instance = { server: Server; registry: Registry; fetch: createPinboard.Instance["fetch"]; sweep(): Promise; handleRequest(req: IncomingMessage, res: ServerResponse): void; handleUpgrade(req: IncomingMessage, socket: Duplex, head: Buffer): void; close(): Promise; }; } export const servePinboard = ( options: servePinboard.Options, ): servePinboard.Instance => { const { redis: rawRedis, metricsRegistry, responseHeaderTimeoutMs, ...core } = options; const redis = RedisKVMarker in rawRedis ? rawRedis : new IoredisKV(rawRedis); const ownsRedis = redis !== rawRedis; const metrics = metricsRegistry === undefined ? noopMetrics : createMetrics(metricsRegistry); const transport = createNodeTransport({ responseHeaderTimeoutMs }); const pinboard = createPinboard({ ...core, redis, metrics, transport }); pinboard.start(); // Metrics are fleet-global and cannot be filtered per caller, so scrapes // require an operator whose policy is itself unrestricted. const anyHost = (rules: HostRule[]): boolean => rules.length === 0 || rules.some((rule) => rule.kind === "any"); const metricsResponse = async ( req: IncomingMessage, registry: MetricsRegistry, ): Promise => { const auth = await authenticateRegistry( core.passport, core.hostnameAware ?? false, toRequest(req), "pinboard:operator", ); if (auth instanceof Response) return auth; if ( !auth.namespaces.includes("*") || !anyHost(auth.hostnames) || !anyHost(auth.envs) ) { return json(403, { error: "forbidden" }); } return new Response(await registry.metrics(), { headers: { "content-type": registry.contentType }, }); }; const handleRequest = (req: IncomingMessage, res: ServerResponse): void => { if ( metricsRegistry !== undefined && req.method === "GET" && new URL(req.url ?? "/", "http://pinboard.invalid").pathname === "/metrics" ) { void metricsResponse(req, metricsRegistry) .catch(() => json(500, { error: "metrics collection failed" })) .then((response) => writeResponse(res, response)) .catch(() => res.destroy()); return; } const abort = new AbortController(); res.on("close", () => { if (!res.writableEnded) abort.abort(); }); void pinboard .fetch(toRequest(req, abort.signal), { address: req.socket.remoteAddress, }) .then((response) => writeResponse(res, response)) .catch(() => res.destroy()); }; const handleUpgrade = ( req: IncomingMessage, socket: Duplex, head: Buffer, ): void => { const request = toRequest(req); transport.bindUpgrade(request, { socket, head }); void pinboard .fetch(request, { address: req.socket.remoteAddress }) .then((response) => { if (transport.completed(response)) return; if (response.status === 500) { socket.destroy(); return; } const reason = REASONS[response.status] ?? ""; const retryAfter = response.headers.get("retry-after"); socket.end( `HTTP/1.1 ${response.status} ${reason}\r\n` + (retryAfter === null ? "" : `retry-after: ${retryAfter}\r\n`) + "connection: close\r\n\r\n", ); }) .catch(() => socket.destroy()); }; const server = createServer(handleRequest); server.on("upgrade", handleUpgrade); let closePromise: Promise | undefined; return { server, registry: pinboard.registry, fetch: pinboard.fetch, sweep: pinboard.sweep, handleRequest, handleUpgrade, close() { if (closePromise !== undefined) return closePromise; closePromise = (async () => { const errors: unknown[] = []; const attempt = async ( operation: () => Promise, ): Promise => { try { await operation(); } catch (error) { errors.push(error); } }; let serverClose: Promise = Promise.resolve(); if (server.listening) { serverClose = new Promise((resolve, reject) => { try { server.close((error) => (error ? reject(error) : resolve())); server.closeIdleConnections(); } catch (error) { reject(error); } }); } await attempt(() => transport.close()); await attempt(() => serverClose); await attempt(() => pinboard.close()); if (ownsRedis) await attempt(() => redis.quit()); if (errors.length > 0) throw errors[0]; })(); return closePromise; }, }; };