From f7801db296b0bd0d1a419baecdbe1e873e619c91 Mon Sep 17 00:00:00 2001 From: oliver Date: Wed, 16 Sep 2026 14:39:55 +0800 Subject: [PATCH] Mount /netxops RPC on webServer for DSH 0.1.5 (0.1.39). connection.rpc.handle no longer registers custom channels when webServer is not on the connection ctx, which left the KB settings badge stuck on unconfigured. Co-authored-by: Cursor --- lib/index.js | 349 ++++++++++++++++++++++++++++-------- package.json | 2 +- src/index.ts | 196 +++++++++++--------- src/netx/netxops-web-rpc.ts | 239 ++++++++++++++++++++++++ 4 files changed, 621 insertions(+), 165 deletions(-) create mode 100644 src/netx/netxops-web-rpc.ts diff --git a/lib/index.js b/lib/index.js index ab46a11..c3e2de1 100644 --- a/lib/index.js +++ b/lib/index.js @@ -3043,6 +3043,196 @@ async function sessionsExportHeadResponse(ctx, request) { }); } +// src/netx/netxops-web-rpc.ts +var NETXOPS_RPC_BODY_MAX = 8 * 1024 * 1024; +var ENDPOINT_SEGMENT_PATTERN = /^[A-Za-z0-9_$.-]+$/; +var INVALID_REQUEST_RPC_ID = "invalid-request"; +var LOOPBACK_HOSTNAMES = new Set(["127.0.0.1", "localhost", "::1"]); +function endpointFromPath(channel, pathname) { + if (!pathname.startsWith(`${channel}/`)) + return; + const endpoint = pathname.slice(channel.length + 1); + if (endpoint.split("/").some((seg) => seg === "" || seg === "." || seg === ".." || !ENDPOINT_SEGMENT_PATTERN.test(seg))) { + return; + } + return endpoint; +} +function serverResponseJson(rpcId, result) { + return JSON.stringify({ type: "server-response", rpcId, result }); +} +function isTrustedLoopbackRequest(req) { + const host = req.headers?.host; + if (!host) + return false; + const hostName = host.split(":")[0]; + if (!LOOPBACK_HOSTNAMES.has(hostName ?? "")) + return false; + if (req.headers["sec-fetch-site"] === "cross-site") + return false; + const origin = req.headers.origin; + if (origin === undefined) + return true; + try { + return new URL(origin).host === host; + } catch { + return false; + } +} +function netxopsFetchHandler(channel, handler, log) { + return { + async fetch(request) { + const endpoint = endpointFromPath(channel, new URL(request.url).pathname); + if (request.method !== "POST" || endpoint === undefined) { + return new Response("not found", { status: 404 }); + } + const mediaType = request.headers.get("content-type")?.split(";", 1)[0]?.trim().toLowerCase(); + if (mediaType !== "application/json") { + return new Response("content type must be application/json", { status: 415 }); + } + let body; + try { + body = await request.json(); + } catch { + return new Response("body is not JSON", { status: 400 }); + } + const row = body; + const rpcId = row && typeof row.rpcId === "string" ? row.rpcId : INVALID_REQUEST_RPC_ID; + const method = row && typeof row.method === "string" ? row.method : null; + if (rpcId === INVALID_REQUEST_RPC_ID || method === null) { + return new Response(serverResponseJson(INVALID_REQUEST_RPC_ID, { + ok: false, + error: { code: "bad-request", message: "invalid client-request message", details: { issues: [] } } + }), { status: 200, headers: { "content-type": "application/json" } }); + } + if (method !== endpoint) { + return new Response(serverResponseJson(rpcId, { + ok: false, + error: { + code: "bad-request", + message: `method ${JSON.stringify(method)} does not match endpoint ${JSON.stringify(endpoint)}`, + details: { issues: [] } + } + }), { status: 200, headers: { "content-type": "application/json" } }); + } + try { + const result = await handler(endpoint, row?.payload, request.signal); + return new Response(serverResponseJson(rpcId, result), { + status: 200, + headers: { "content-type": "application/json" } + }); + } catch (error) { + log.error?.("netxops: rpc %s failed: %s", endpoint, error instanceof Error ? error.message : String(error)); + return new Response(`handler failure: ${String(error)}`, { status: 500 }); + } + } + }; +} +async function httpBridge(req, res, fetchHandler, maxBodyBytes) { + const abort = new AbortController; + res.on("close", () => { + if (!res.writableEnded) + abort.abort(); + }); + const declaredLen = req.headers["content-length"]; + if (declaredLen !== undefined && Number(declaredLen) > maxBodyBytes) { + res.writeHead(413, { connection: "close" }); + res.end(); + req.destroy(); + return; + } + const chunks = []; + let received = 0; + for await (const chunk of req) { + const buffer = chunk; + received += buffer.byteLength; + if (received > maxBodyBytes) { + res.writeHead(413, { connection: "close" }); + res.end(); + req.destroy(); + return; + } + chunks.push(buffer); + } + const url = `http://${req.headers.host ?? "127.0.0.1"}${req.url ?? "/"}`; + const request = new Request(url, { + method: req.method ?? "GET", + headers: Object.fromEntries(Object.entries(req.headers).filter(([, value]) => typeof value === "string")), + ...chunks.length > 0 ? { body: Buffer.concat(chunks) } : {}, + signal: abort.signal + }); + const response = await fetchHandler.fetch(request); + const headers = Object.fromEntries(response.headers.entries()); + res.writeHead(response.status, headers); + if (response.body === null) { + res.end(); + return; + } + for await (const chunk of response.body) { + if (!res.write(chunk)) { + await new Promise((resolve3) => { + const done = () => { + res.off("drain", done); + res.off("close", done); + resolve3(); + }; + res.once("drain", done); + res.once("close", done); + }); + } + } + res.end(); +} +function mountNetxopsWebRoute(ctx, channel, handler) { + const webServer = ctx.webServer; + if (!webServer || typeof webServer.register !== "function") + return null; + const connection = ctx.connection; + const fetchHandler = netxopsFetchHandler(channel, handler, ctx.logger); + const route = { + kind: "prefix", + path: channel, + handler: async (req, res) => { + let rejection; + if (typeof connection?.requestRejection === "function") { + try { + rejection = connection.requestRejection(req); + } catch { + rejection = 403; + } + } else if (!isTrustedLoopbackRequest(req)) { + rejection = 403; + } + if (rejection !== undefined) { + res.writeHead(rejection, { "content-type": "text/plain; charset=utf-8" }); + res.end(rejection === 401 ? "unauthorized" : "forbidden"); + return; + } + await httpBridge(req, res, fetchHandler, NETXOPS_RPC_BODY_MAX); + } + }; + const registered = webServer.register(route); + if (typeof registered === "function") { + return () => { + try { + registered(); + } catch {} + }; + } + if (registered && typeof registered.then === "function") { + let done = false; + return () => { + if (done) + return; + done = true; + registered.then((dispose) => { + if (typeof dispose === "function") + dispose(); + }).catch(() => {}); + }; + } + return () => {}; +} + // src/netx/tools.ts import { defineTool as defineTool2 } from "@deepseek-ai/dsh-tools"; @@ -4864,83 +5054,92 @@ function apply(ctx, config = Config({})) { unregisterKbPack?.(); }, "netxops: dispose public skills"); }); + const handleNetxopsRpc = async (endpoint, payload) => { + if (endpoint === "alarm-push.status") { + return { ok: true, value: getAlarmPushStatus() }; + } + if (endpoint === "kb.status") { + return { ok: true, value: getKbContext() }; + } + if (endpoint === "kb.reload") { + const current = source(); + const kb = resolveKbRoot(current.kbRoot ?? ""); + publishKbContext(kb); + applyKbEnv(kb); + if (kb.status === "configured") { + ctx.logger.info("netxops: kb.reload → %s (%s / %s v%s)", kb.realRoot, kb.operatorName, kb.country, kb.version); + } else if (kb.status === "error") { + ctx.logger.warn("netxops: kb.reload error — %s", kb.errorMessage); + } else { + ctx.logger.info("netxops: kb.reload → unconfigured (kbRoot empty)"); + } + return { ok: true, value: kb }; + } + if (endpoint === "kb.resolve") { + const path = extractKbResolvePath(payload); + return { ok: true, value: resolveKbRoot(path) }; + } + if (endpoint === "sessions.export.status") { + const value = await getSessionsExportStatus(ctx); + return { ok: true, value }; + } + if (endpoint === "im-delivery.catalog") { + const fromGet = typeof ctx.get === "function" ? ctx.get("dshIm") : undefined; + const im = fromGet ?? ctx.dshIm; + if (!im || typeof im.listDeliveryCatalog !== "function") { + return { + ok: true, + value: { + available: false, + options: [], + reasonCode: "im_catalog_unavailable", + hint: "dsh-im-ops missing or outdated — install ≥ops.24 for delivery picker" + } + }; + } + try { + const options = await im.listDeliveryCatalog(); + return { + ok: true, + value: { + available: true, + options: Array.isArray(options) ? options : [] + } + }; + } catch (error) { + return { + ok: true, + value: { + available: false, + options: [], + hint: error instanceof Error ? error.message : String(error) + } + }; + } + } + return { ok: false, error: { code: "bad-request", message: "Unknown endpoint." } }; + }; + ctx.inject(["webServer", "connection"], (mountCtx) => { + mountCtx.effect(() => { + const direct = mountNetxopsWebRoute(mountCtx, NETXOPS_RPC_CHANNEL, (endpoint, payload) => handleNetxopsRpc(endpoint, payload)); + if (direct !== null) { + mountCtx.logger.info("netxops: /netxops rpc mounted on webServer"); + return direct; + } + const rpc = mountCtx.connection?.rpc; + if (rpc && typeof rpc.handle === "function") { + mountCtx.logger.info("netxops: /netxops rpc mounted via connection.rpc.handle (legacy)"); + const legacy = rpc.handle(NETXOPS_RPC_CHANNEL, (endpoint, payload) => handleNetxopsRpc(endpoint, payload)); + return () => { + legacy(); + }; + } + mountCtx.logger.warn("netxops: /netxops rpc unavailable — settings card status disabled"); + return () => {}; + }, "netxops: /netxops web rpc"); + }); ctx.inject(["connection"], (connCtx) => { const connection = connCtx.connection; - const rpc = connection?.rpc; - if (!rpc || typeof rpc.handle !== "function") { - connCtx.logger.warn("netxops: connection.rpc.handle unavailable — alarm status UI disabled"); - } else { - connCtx.effect(() => { - const dispose = rpc.handle(NETXOPS_RPC_CHANNEL, async (endpoint, payload) => { - if (endpoint === "alarm-push.status") { - return { ok: true, value: getAlarmPushStatus() }; - } - if (endpoint === "kb.status") { - return { ok: true, value: getKbContext() }; - } - if (endpoint === "kb.reload") { - const current = source(); - const kb = resolveKbRoot(current.kbRoot ?? ""); - publishKbContext(kb); - applyKbEnv(kb); - if (kb.status === "configured") { - connCtx.logger.info("netxops: kb.reload → %s (%s / %s v%s)", kb.realRoot, kb.operatorName, kb.country, kb.version); - } else if (kb.status === "error") { - connCtx.logger.warn("netxops: kb.reload error — %s", kb.errorMessage); - } else { - connCtx.logger.info("netxops: kb.reload → unconfigured (kbRoot empty)"); - } - return { ok: true, value: kb }; - } - if (endpoint === "kb.resolve") { - const path = extractKbResolvePath(payload); - return { ok: true, value: resolveKbRoot(path) }; - } - if (endpoint === "sessions.export.status") { - const value = await getSessionsExportStatus(ctx); - return { ok: true, value }; - } - if (endpoint === "im-delivery.catalog") { - const fromGet = typeof ctx.get === "function" ? ctx.get("dshIm") : undefined; - const im = fromGet ?? ctx.dshIm; - if (!im || typeof im.listDeliveryCatalog !== "function") { - return { - ok: true, - value: { - available: false, - options: [], - reasonCode: "im_catalog_unavailable", - hint: "dsh-im-ops missing or outdated — install ≥ops.24 for delivery picker" - } - }; - } - try { - const options = await im.listDeliveryCatalog(); - return { - ok: true, - value: { - available: true, - options: Array.isArray(options) ? options : [] - } - }; - } catch (error) { - return { - ok: true, - value: { - available: false, - options: [], - hint: error instanceof Error ? error.message : String(error) - } - }; - } - } - return { ok: false, error: { code: "bad-request", message: "Unknown endpoint." } }; - }); - return () => { - dispose(); - }; - }, "netxops: alarm-push status rpc"); - } const fetchApi = connection?.fetch; if (!fetchApi || typeof fetchApi.register !== "function") { connCtx.logger.warn("netxops: connection.fetch.register unavailable — sessions export download disabled"); diff --git a/package.json b/package.json index 5e286e4..b2b53c8 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "dsh-netxops", - "version": "0.1.38", + "version": "0.1.39", "description": "DeepSeek Harness Netx Ops: ops/topology + IM delivery + session export + knowledge-base MANIFEST/_skills + localSkills", "license": "MIT", "type": "module", diff --git a/src/index.ts b/src/index.ts index 6848ab1..e7896ab 100644 --- a/src/index.ts +++ b/src/index.ts @@ -72,6 +72,7 @@ import { sessionsExportHeadResponse, sessionsExportResponse, } from './netx/session-export.ts' +import { mountNetxopsWebRoute } from './netx/netxops-web-rpc.ts' import { registerNetxTools } from './netx/tools.ts' /** Cordis plugin name. */ @@ -595,7 +596,112 @@ export function apply(ctx: Context, config: Config = Config({})): void { }, 'netxops: dispose public skills') }) - // Browser card: alarm-push status + IM catalog (RPC) and all-sessions ZIP (Fetch). + const handleNetxopsRpc = async (endpoint: string, payload?: unknown): Promise => { + if (endpoint === 'alarm-push.status') { + return { ok: true, value: getAlarmPushStatus() } + } + if (endpoint === 'kb.status') { + return { ok: true, value: getKbContext() } + } + if (endpoint === 'kb.reload') { + const current = source() + const kb = resolveKbRoot(current.kbRoot ?? '') + publishKbContext(kb) + applyKbEnv(kb) + if (kb.status === 'configured') { + ctx.logger.info( + 'netxops: kb.reload → %s (%s / %s v%s)', + kb.realRoot, + kb.operatorName, + kb.country, + kb.version, + ) + } else if (kb.status === 'error') { + ctx.logger.warn('netxops: kb.reload error — %s', kb.errorMessage) + } else { + ctx.logger.info('netxops: kb.reload → unconfigured (kbRoot empty)') + } + return { ok: true, value: kb } + } + if (endpoint === 'kb.resolve') { + const path = extractKbResolvePath(payload) + return { ok: true, value: resolveKbRoot(path) } + } + if (endpoint === 'sessions.export.status') { + const value = await getSessionsExportStatus(ctx) + return { ok: true, value } + } + if (endpoint === 'im-delivery.catalog') { + type DshImCatalog = { listDeliveryCatalog?: () => Promise } + const fromGet = typeof (ctx as { get?: (name: string) => unknown }).get === 'function' + ? (ctx as { get: (name: string) => unknown }).get('dshIm') as DshImCatalog | undefined + : undefined + const im = fromGet ?? (ctx as { dshIm?: DshImCatalog }).dshIm + if (!im || typeof im.listDeliveryCatalog !== 'function') { + return { + ok: true, + value: { + available: false, + options: [], + reasonCode: 'im_catalog_unavailable', + hint: 'dsh-im-ops missing or outdated — install ≥ops.24 for delivery picker', + }, + } + } + try { + const options = await im.listDeliveryCatalog() + return { + ok: true, + value: { + available: true, + options: Array.isArray(options) ? options : [], + }, + } + } catch (error) { + return { + ok: true, + value: { + available: false, + options: [], + hint: error instanceof Error ? error.message : String(error), + }, + } + } + } + return { ok: false, error: { code: 'bad-request', message: 'Unknown endpoint.' } } + } + + // DSH ≥0.1.5: mount /netxops on this plugin's webServer inject (connection.rpc.handle no longer registers). + ctx.inject(['webServer', 'connection'], (mountCtx) => { + mountCtx.effect(() => { + const direct = mountNetxopsWebRoute( + mountCtx, + NETXOPS_RPC_CHANNEL, + (endpoint, payload) => handleNetxopsRpc(endpoint, payload), + ) + if (direct !== null) { + mountCtx.logger.info('netxops: /netxops rpc mounted on webServer') + return direct + } + type RpcHandle = ( + channel: string, + handler: (endpoint: string, payload?: unknown) => Promise, + ) => (() => void) | Promise + const rpc = (mountCtx.connection as { rpc?: { handle?: RpcHandle } } | undefined)?.rpc + if (rpc && typeof rpc.handle === 'function') { + mountCtx.logger.info('netxops: /netxops rpc mounted via connection.rpc.handle (legacy)') + const legacy = rpc.handle( + NETXOPS_RPC_CHANNEL, + (endpoint, payload) => handleNetxopsRpc(endpoint, payload), + ) + return () => { void legacy() } + } + mountCtx.logger.warn('netxops: /netxops rpc unavailable — settings card status disabled') + return () => {} + }, 'netxops: /netxops web rpc') + }) + + // Browser card: sessions export download via shared /api Fetch route. ctx.inject(['connection'], (connCtx) => { const connection = connCtx.connection as { rpc?: { handle?: ( @@ -610,94 +716,6 @@ export function apply(ctx: Context, config: Config = Config({})): void { }) => () => Promise } } | undefined - const rpc = connection?.rpc - if (!rpc || typeof rpc.handle !== 'function') { - connCtx.logger.warn('netxops: connection.rpc.handle unavailable — alarm status UI disabled') - } else { - connCtx.effect(() => { - const dispose = rpc.handle( - NETXOPS_RPC_CHANNEL, - async (endpoint: string, payload?: unknown) => { - if (endpoint === 'alarm-push.status') { - return { ok: true, value: getAlarmPushStatus() } - } - if (endpoint === 'kb.status') { - return { ok: true, value: getKbContext() } - } - if (endpoint === 'kb.reload') { - // Re-read live settings and publish — used by the settings card after - // browse/save so the badge does not stick on a stale unconfigured snapshot. - const current = source() - const kb = resolveKbRoot(current.kbRoot ?? '') - publishKbContext(kb) - applyKbEnv(kb) - if (kb.status === 'configured') { - connCtx.logger.info( - 'netxops: kb.reload → %s (%s / %s v%s)', - kb.realRoot, - kb.operatorName, - kb.country, - kb.version, - ) - } else if (kb.status === 'error') { - connCtx.logger.warn('netxops: kb.reload error — %s', kb.errorMessage) - } else { - connCtx.logger.info('netxops: kb.reload → unconfigured (kbRoot empty)') - } - return { ok: true, value: kb } - } - if (endpoint === 'kb.resolve') { - const path = extractKbResolvePath(payload) - return { ok: true, value: resolveKbRoot(path) } - } - if (endpoint === 'sessions.export.status') { - const value = await getSessionsExportStatus(ctx) - return { ok: true, value } - } - if (endpoint === 'im-delivery.catalog') { - type DshImCatalog = { listDeliveryCatalog?: () => Promise } - const fromGet = typeof (ctx as { get?: (name: string) => unknown }).get === 'function' - ? (ctx as { get: (name: string) => unknown }).get('dshIm') as DshImCatalog | undefined - : undefined - const im = fromGet ?? (ctx as { dshIm?: DshImCatalog }).dshIm - if (!im || typeof im.listDeliveryCatalog !== 'function') { - return { - ok: true, - value: { - available: false, - options: [], - reasonCode: 'im_catalog_unavailable', - hint: 'dsh-im-ops missing or outdated — install ≥ops.24 for delivery picker', - }, - } - } - try { - const options = await im.listDeliveryCatalog() - return { - ok: true, - value: { - available: true, - options: Array.isArray(options) ? options : [], - }, - } - } catch (error) { - return { - ok: true, - value: { - available: false, - options: [], - hint: error instanceof Error ? error.message : String(error), - }, - } - } - } - return { ok: false, error: { code: 'bad-request', message: 'Unknown endpoint.' } } - }, - ) - return () => { void dispose() } - }, 'netxops: alarm-push status rpc') - } - const fetchApi = connection?.fetch if (!fetchApi || typeof fetchApi.register !== 'function') { connCtx.logger.warn('netxops: connection.fetch.register unavailable — sessions export download disabled') diff --git a/src/netx/netxops-web-rpc.ts b/src/netx/netxops-web-rpc.ts new file mode 100644 index 0000000..b478607 --- /dev/null +++ b/src/netx/netxops-web-rpc.ts @@ -0,0 +1,239 @@ +/** + * Mount `/netxops/*` on the plugin's own `webServer` inject (DSH ≥0.1.5). + * + * `connection.rpc.handle('/netxops')` registers via the connection plugin ctx, + * which no longer injects `webServer` — the route never mounts and the settings + * card sees HTTP 405 → "未配置". Same pattern as dsh-pocket `web-rpc.js`. + */ + +import type { Context } from '@deepseek-ai/cordis' +import type { IncomingMessage, ServerResponse } from 'node:http' + +/** Max buffered JSON body for control-plane RPC (status / kb.resolve paths). */ +const NETXOPS_RPC_BODY_MAX = 8 * 1024 * 1024 + +const ENDPOINT_SEGMENT_PATTERN = /^[A-Za-z0-9_$.-]+$/ +const INVALID_REQUEST_RPC_ID = 'invalid-request' +const LOOPBACK_HOSTNAMES = new Set(['127.0.0.1', 'localhost', '::1']) + +export type NetxopsRpcHandler = ( + endpoint: string, + payload: unknown, + signal: AbortSignal, +) => Promise + +function endpointFromPath(channel: string, pathname: string): string | undefined { + if (!pathname.startsWith(`${channel}/`)) return undefined + const endpoint = pathname.slice(channel.length + 1) + if (endpoint.split('/').some((seg) => + seg === '' || seg === '.' || seg === '..' || !ENDPOINT_SEGMENT_PATTERN.test(seg))) { + return undefined + } + return endpoint +} + +function serverResponseJson(rpcId: string, result: unknown): string { + return JSON.stringify({ type: 'server-response', rpcId, result }) +} + +function isTrustedLoopbackRequest(req: IncomingMessage): boolean { + const host = req.headers?.host + if (!host) return false + const hostName = host.split(':')[0] + if (!LOOPBACK_HOSTNAMES.has(hostName ?? '')) return false + if (req.headers['sec-fetch-site'] === 'cross-site') return false + const origin = req.headers.origin + if (origin === undefined) return true + try { + return new URL(origin).host === host + } catch { + return false + } +} + +function netxopsFetchHandler( + channel: string, + handler: NetxopsRpcHandler, + log: Pick, +) { + return { + async fetch(request: Request): Promise { + const endpoint = endpointFromPath(channel, new URL(request.url).pathname) + if (request.method !== 'POST' || endpoint === undefined) { + return new Response('not found', { status: 404 }) + } + const mediaType = request.headers.get('content-type')?.split(';', 1)[0]?.trim().toLowerCase() + if (mediaType !== 'application/json') { + return new Response('content type must be application/json', { status: 415 }) + } + let body: unknown + try { + body = await request.json() + } catch { + return new Response('body is not JSON', { status: 400 }) + } + const row = body as { rpcId?: unknown; method?: unknown; payload?: unknown } | null + const rpcId = row && typeof row.rpcId === 'string' ? row.rpcId : INVALID_REQUEST_RPC_ID + const method = row && typeof row.method === 'string' ? row.method : null + if (rpcId === INVALID_REQUEST_RPC_ID || method === null) { + return new Response( + serverResponseJson(INVALID_REQUEST_RPC_ID, { + ok: false, + error: { code: 'bad-request', message: 'invalid client-request message', details: { issues: [] } }, + }), + { status: 200, headers: { 'content-type': 'application/json' } }, + ) + } + if (method !== endpoint) { + return new Response( + serverResponseJson(rpcId, { + ok: false, + error: { + code: 'bad-request', + message: `method ${JSON.stringify(method)} does not match endpoint ${JSON.stringify(endpoint)}`, + details: { issues: [] }, + }, + }), + { status: 200, headers: { 'content-type': 'application/json' } }, + ) + } + try { + const result = await handler(endpoint, row?.payload, request.signal) + return new Response(serverResponseJson(rpcId, result), { + status: 200, + headers: { 'content-type': 'application/json' }, + }) + } catch (error) { + log.error?.('netxops: rpc %s failed: %s', endpoint, error instanceof Error ? error.message : String(error)) + return new Response(`handler failure: ${String(error)}`, { status: 500 }) + } + }, + } +} + +async function httpBridge( + req: IncomingMessage, + res: ServerResponse, + fetchHandler: { fetch: (request: Request) => Promise }, + maxBodyBytes: number, +): Promise { + const abort = new AbortController() + res.on('close', () => { + if (!res.writableEnded) abort.abort() + }) + + const declaredLen = req.headers['content-length'] + if (declaredLen !== undefined && Number(declaredLen) > maxBodyBytes) { + res.writeHead(413, { connection: 'close' }) + res.end() + req.destroy() + return + } + const chunks: Buffer[] = [] + let received = 0 + for await (const chunk of req) { + const buffer = chunk as Buffer + received += buffer.byteLength + if (received > maxBodyBytes) { + res.writeHead(413, { connection: 'close' }) + res.end() + req.destroy() + return + } + chunks.push(buffer) + } + + const url = `http://${req.headers.host ?? '127.0.0.1'}${req.url ?? '/'}` + const request = new Request(url, { + method: req.method ?? 'GET', + headers: Object.fromEntries( + Object.entries(req.headers).filter(([, value]) => typeof value === 'string') as [string, string][], + ), + ...chunks.length > 0 ? { body: Buffer.concat(chunks) } : {}, + signal: abort.signal, + }) + + const response = await fetchHandler.fetch(request) + const headers = Object.fromEntries(response.headers.entries()) + res.writeHead(response.status, headers) + if (response.body === null) { + res.end() + return + } + for await (const chunk of response.body) { + if (!res.write(chunk)) { + await new Promise((resolve) => { + const done = (): void => { + res.off('drain', done) + res.off('close', done) + resolve() + } + res.once('drain', done) + res.once('close', done) + }) + } + } + res.end() +} + +type ConnectionRejection = { requestRejection?: (req: IncomingMessage) => number | undefined } + +/** + * Register prefix route `channel` on `ctx.webServer` when available. + * @returns disposer, or null when `webServer` is not on this inject ctx. + */ +export function mountNetxopsWebRoute( + ctx: Context & { webServer?: { register: (route: unknown) => unknown }; connection?: ConnectionRejection }, + channel: string, + handler: NetxopsRpcHandler, +): (() => void) | null { + const webServer = ctx.webServer + if (!webServer || typeof webServer.register !== 'function') return null + + const connection = ctx.connection + const fetchHandler = netxopsFetchHandler(channel, handler, ctx.logger) + const route = { + kind: 'prefix' as const, + path: channel, + handler: async (req: IncomingMessage, res: ServerResponse) => { + let rejection: number | undefined + if (typeof connection?.requestRejection === 'function') { + try { + rejection = connection.requestRejection(req) + } catch { + rejection = 403 + } + } else if (!isTrustedLoopbackRequest(req)) { + rejection = 403 + } + if (rejection !== undefined) { + res.writeHead(rejection, { 'content-type': 'text/plain; charset=utf-8' }) + res.end(rejection === 401 ? 'unauthorized' : 'forbidden') + return + } + await httpBridge(req, res, fetchHandler, NETXOPS_RPC_BODY_MAX) + }, + } + + const registered = webServer.register(route) + if (typeof registered === 'function') { + return () => { + try { + registered() + } catch { + // already disposed + } + } + } + if (registered && typeof (registered as Promise).then === 'function') { + let done = false + return () => { + if (done) return + done = true + void (registered as Promise<(() => void) | void>).then((dispose) => { + if (typeof dispose === 'function') dispose() + }).catch(() => {}) + } + } + return () => {} +}