diff --git a/src/implementations/api/index.ts b/src/implementations/api/index.ts index dd38565..f819dc8 100644 --- a/src/implementations/api/index.ts +++ b/src/implementations/api/index.ts @@ -1,15 +1,20 @@ import { GetMetadata } from "@antelopejs/interface-core"; import type { Class } from "@antelopejs/interface-core/decorators"; +import type { RegisteredReadTarget } from "@antelopejs/interface-api/registered-read"; import { type ComputedParameter, ControllerMeta, type CorsConfig, + HTTPResult, type RouteHandler, } from "@antelopejs/interface-api"; +import { assertReadActive } from "../../registered-read-context"; import { getConfig, listenServers, setCorsConfig } from "../../index"; import { type RequestContext, + type RouteCallback, + executeRegisteredRead, registerHandler, unregisterHandler, } from "../../server"; @@ -57,6 +62,7 @@ interface HandlerPlan { const classCacheSymbol = Symbol(); const registeredRoutes = new Map(); +const compiledRoutes = new Map(); const controllerPlans = new WeakMap(); interface RequestContextDev extends RequestContext { @@ -81,16 +87,21 @@ function compileParameter( const { provider } = parameter; const modifiers = [...parameter.modifiers]; if (modifiers.length === 0) { - return (context, controller) => provider.call(controller, context); + return (context, controller) => { + assertReadActive(context); + return provider.call(controller, context); + }; } - return (context, controller) => - applyModifiers( + return (context, controller) => { + assertReadActive(context); + return applyModifiers( provider.call(controller, context), modifiers, context, controller, ); + }; } function applyModifiers( @@ -108,6 +119,7 @@ function applyModifiers( applyModifiers(resolved, modifiers, context, controller, index), ); } + assertReadActive(context); value = modifiers[index].call(controller, context, value); } const then = getThen(value); @@ -193,17 +205,24 @@ function applyComputedProperties( } const pending: Promise[] = []; - for (const property of computedProperties) { - const value = property.resolve(context, controllerInstance); - if (isPromiseLike(value)) { - pending.push( - Promise.resolve(value).then((resolved) => { - controllerInstance[property.key] = resolved; - }), - ); - } else { - controllerInstance[property.key] = value; + try { + for (const property of computedProperties) { + const value = property.resolve(context, controllerInstance); + if (isPromiseLike(value)) { + pending.push( + Promise.resolve(value).then((resolved) => { + assertReadActive(context); + controllerInstance[property.key] = resolved; + }), + ); + } else { + assertReadActive(context); + controllerInstance[property.key] = value; + } } + } catch (error) { + void Promise.allSettled(pending); + throw error; } if (pending.length > 0) { @@ -270,32 +289,52 @@ interface RouteInfo { callbackName: string; } +function resolveParameters( + plan: HandlerPlan, + controllerInstance: UnknownRecord, + context: RequestContextDev, +): unknown[] { + const parameters: unknown[] = []; + try { + for (const resolve of plan.parameters) { + parameters.push(resolve(context, controllerInstance)); + } + } catch (error) { + void Promise.allSettled(parameters); + throw error; + } + return parameters; +} + function invokeCallback( plan: HandlerPlan, controllerInstance: UnknownRecord, context: RequestContextDev, ): unknown { + assertReadActive(context); if (plan.parameters.length === 0) { return plan.callback.call(controllerInstance); } if (plan.parameters.length === 1) { const parameter = plan.parameters[0](context, controllerInstance); if (isPromiseLike(parameter)) { - return Promise.resolve(parameter).then((resolved) => - plan.callback.call(controllerInstance, resolved), - ); + return Promise.resolve(parameter).then((resolved) => { + assertReadActive(context); + return plan.callback.call(controllerInstance, resolved); + }); } + assertReadActive(context); return plan.callback.call(controllerInstance, parameter); } - const parameters = plan.parameters.map((resolve) => - resolve(context, controllerInstance), - ); + const parameters = resolveParameters(plan, controllerInstance, context); if (parameters.some(isPromiseLike)) { - return Promise.all(parameters).then((resolved) => - plan.callback.apply(controllerInstance, resolved), - ); + return Promise.all(parameters).then((resolved) => { + assertReadActive(context); + return plan.callback.apply(controllerInstance, resolved); + }); } + assertReadActive(context); return plan.callback.apply(controllerInstance, parameters); } @@ -317,25 +356,48 @@ function compileHandler(handler: RouteHandler): HandlerPlan { return { callback: handler.callback, controller: compileController(controllerClass, handler.properties), - parameters: handler.parameters.map(compileParameter), + parameters: Array.from(handler.parameters, compileParameter), }; } +const READ_FORBIDDEN = 403; + +export async function ExecuteRegisteredRead( + target: RegisteredReadTarget, + context: RequestContext, +): Promise { + const route = registeredRoutes.get(target.routeId); + const callback = compiledRoutes.get(target.routeId); + if ( + !route || + !callback || + route.method !== "get" || + route.mode !== "handler" + ) { + throw new HTTPResult(READ_FORBIDDEN, "Invalid registered read target"); + } + return executeRegisteredRead(callback, context, target); +} + export const routesProxy = { register: (id: string, handler: RouteHandler): void => { registeredRoutes.set(id, handler); const plan = compileHandler(handler); + const callback = (context: RequestContextDev) => + invokeHandler(plan, context); + compiledRoutes.set(id, callback); registerHandler( `dev/${id}`, handler.mode, handler.method, handler.location, - (context: RequestContextDev) => invokeHandler(plan, context), + callback, handler.priority, ); }, unregister: (id: string): void => { registeredRoutes.delete(id); + compiledRoutes.delete(id); unregisterHandler(`dev/${id}`); }, getRoutes: (): RouteInfo[] => { diff --git a/src/index.ts b/src/index.ts index d75be02..d722abe 100644 --- a/src/index.ts +++ b/src/index.ts @@ -28,6 +28,10 @@ export async function construct(config: Config): Promise { await import("@antelopejs/interface-api"), await import("./implementations/api"), ); + void ImplementInterface( + await import("@antelopejs/interface-api/registered-read"), + await import("./implementations/api"), + ); } export function destroy(): void {} diff --git a/src/registered-read-context.ts b/src/registered-read-context.ts new file mode 100644 index 0000000..5b54718 --- /dev/null +++ b/src/registered-read-context.ts @@ -0,0 +1,176 @@ +import { HTTPResult } from "@antelopejs/interface-api"; +import { IncomingMessage, ServerResponse } from "node:http"; +import type { RegisteredReadTarget } from "@antelopejs/interface-api/registered-read"; + +import type { RequestContext } from "./server"; +import { RegisteredReadSocket } from "./registered-read-socket"; + +const FORBIDDEN = 403; + +const BODY_HEADERS = new Set(["content-length", "transfer-encoding"]); +const HEADER_PAIR_LENGTH = 2; +const readSignals = new WeakMap(); + +export interface RegisteredReadContext extends RequestContext { + signal: AbortSignal; +} + +export function assertReadActive(context: RequestContext): void { + readSignals.get(context)?.throwIfAborted(); +} + +export function isRegisteredRead( + context: RequestContext, +): context is RegisteredReadContext { + return readSignals.has(context); +} + +class ReadRequest extends IncomingMessage { + constructor( + parent: IncomingMessage, + private readonly abort: () => void, + ) { + super(new RegisteredReadSocket(parent.socket, abort)); + } + + override _destroy( + error: Error | null, + callback: (error?: Error | null) => void, + ): void { + if (error || !this.readableEnded) this.abort(); + callback(null); + } +} + +class ReadResponse extends ServerResponse { + constructor( + request: IncomingMessage, + private readonly abort: () => void, + ) { + super(request); + } + + override write(): boolean { + this.abort(); + return false; + } + override end(): this { + this.abort(); + return this; + } + override writeHead(): this { + this.abort(); + return this; + } + override flushHeaders(): void { + this.abort(); + } + override destroy(): this { + this.abort(); + return this; + } + override writeContinue(): void { + this.abort(); + } + override writeProcessing(): void { + this.abort(); + } + override writeEarlyHints(): void { + this.abort(); + } + override setTimeout(): this { + this.abort(); + return this; + } +} + +function readUrl(origin: string, pathname: string): URL { + const url = new URL(pathname, origin); + if ( + !pathname.startsWith("/") || + pathname.startsWith("//") || + url.origin !== origin || + url.pathname !== pathname || + url.search || + url.hash + ) { + throw new HTTPResult(FORBIDDEN, "Invalid registered read path"); + } + return url; +} + +function copyHeaders(parent: IncomingMessage, request: IncomingMessage): void { + request.headers = Object.fromEntries( + Object.entries(parent.headers) + .filter(([name]) => !BODY_HEADERS.has(name.toLowerCase())) + .map(([name, value]) => [ + name, + Array.isArray(value) ? [...value] : value, + ]), + ); + request.headersDistinct = Object.fromEntries( + Object.entries(parent.headersDistinct) + .filter(([name]) => !BODY_HEADERS.has(name.toLowerCase())) + .map(([name, values]) => [name, values && [...values]]), + ); + request.rawHeaders = []; + for ( + let index = 0; + index < parent.rawHeaders.length; + index += HEADER_PAIR_LENGTH + ) { + const name = parent.rawHeaders[index]; + if (!BODY_HEADERS.has(name.toLowerCase())) { + request.rawHeaders.push(name, parent.rawHeaders[index + 1]); + } + } +} + +export function createReadContext( + parent: RequestContext, + target: RegisteredReadTarget, +): RegisteredReadContext { + const url = readUrl(parent.url.origin, target.pathname); + for (const [key, value] of Object.entries(target.query ?? {})) { + url.searchParams.set(key, value); + } + const controller = new AbortController(); + const abort = () => + controller.abort(new HTTPResult(FORBIDDEN, "Registered read interrupted")); + const request = new ReadRequest(parent.rawRequest, abort); + request.method = "GET"; + request.url = `${url.pathname}${url.search}`; + request.httpVersion = parent.rawRequest.httpVersion; + request.httpVersionMajor = parent.rawRequest.httpVersionMajor; + request.httpVersionMinor = parent.rawRequest.httpVersionMinor; + copyHeaders(parent.rawRequest, request); + request.complete = true; + request.push(null); + const context: RegisteredReadContext = { + rawRequest: request, + rawResponse: new ReadResponse(request, abort), + url, + routeParameters: {}, + response: new HTTPResult(), + signal: controller.signal, + }; + readSignals.set(context, controller.signal); + return context; +} + +export async function completeRead( + context: RegisteredReadContext, + run: () => Promise, +): Promise { + let rejectRead: (reason: unknown) => void; + const interrupted = new Promise((_resolve, reject) => { + rejectRead = reject; + }); + const onAbort = () => rejectRead(context.signal.reason); + context.signal.addEventListener("abort", onAbort, { once: true }); + try { + return await Promise.race([run(), interrupted]); + } finally { + context.signal.removeEventListener("abort", onAbort); + } +} diff --git a/src/registered-read-socket.ts b/src/registered-read-socket.ts new file mode 100644 index 0000000..e2f3087 --- /dev/null +++ b/src/registered-read-socket.ts @@ -0,0 +1,58 @@ +import { Socket } from "node:net"; +import type { TLSSocket } from "node:tls"; + +type PeerSocket = Socket & + Partial>; + +const PEER_FIELDS = [ + "remoteAddress", + "remotePort", + "remoteFamily", + "localAddress", + "localPort", + "localFamily", + "encrypted", + "authorized", + "authorizationError", +] as const; + +/** Detached transport retaining peer metadata without exposing the live socket. */ +export class RegisteredReadSocket extends Socket { + constructor( + parent: PeerSocket, + private readonly abort: () => void, + ) { + super(); + for (const field of PEER_FIELDS) { + Object.defineProperty(this, field, { value: parent[field] }); + } + } + + override write(): boolean { + this.abort(); + return false; + } + + override end(): this { + this.abort(); + return this; + } + + override connect(): this { + this.abort(); + return this; + } + + override _destroy( + _error: Error | null, + callback: (error?: Error | null) => void, + ): void { + this.abort(); + callback(); + } + + override setTimeout(): this { + this.abort(); + return this; + } +} diff --git a/src/server.ts b/src/server.ts index d0b1fb7..d2ba026 100644 --- a/src/server.ts +++ b/src/server.ts @@ -2,6 +2,15 @@ import type stream from "node:stream"; import { type WebSocket, WebSocketServer } from "ws"; import { type IncomingMessage, ServerResponse } from "node:http"; import { HandlerPriority, HTTPResult } from "@antelopejs/interface-api"; +import type { RegisteredReadTarget } from "@antelopejs/interface-api/registered-read"; + +import { + assertReadActive, + completeRead, + createReadContext, + isRegisteredRead, + type RegisteredReadContext, +} from "./registered-read-context"; export type RouteCallback = (context: RequestContext) => unknown; interface IdentifiableRouteCallback { @@ -673,6 +682,7 @@ function executePriorityHandlers( handlers.sort((a, b) => a.priority - b.priority); } for (let index = startIndex; index < handlers.length; index += 1) { + assertReadActive(requestContext); const { handler, parameters } = handlers[index]; requestContext.routeParameters = parameters; const result = handler(requestContext); @@ -790,11 +800,75 @@ function executeHandlerAndPostfix( ): Awaitable { const execution = executeHandler(handler, requestContext); return continueExecution(execution, (result) => { + assertReadActive(requestContext); setHandlerResponse(requestContext, result); + if (isRegisteredRead(requestContext)) assertCompletedRead(requestContext); return executePostfix(method, path, requestContext); }); } +const READ_FORBIDDEN = 403; +const SUCCESS_MIN = 200; +const SUCCESS_MAX = 300; + +function assertCompletedRead(context: RegisteredReadContext): void { + const status = context.response.getStatus(); + if ( + context.signal.aborted || + !Number.isInteger(status) || + status < SUCCESS_MIN || + status >= SUCCESS_MAX || + context.response.isStream() || + context.rawResponse.headersSent || + context.rawResponse.writableEnded || + context.rawResponse.destroyed + ) { + throw new HTTPResult(READ_FORBIDDEN, "Registered read did not complete"); + } +} + +export async function executeRegisteredRead( + expected: RouteCallback, + parent: RequestContext, + target: RegisteredReadTarget, +): Promise { + const context = createReadContext(parent, target); + return completeRead(context, () => + runRegisteredRead(expected, target.pathname, context), + ); +} + +async function runRegisteredRead( + expected: RouteCallback, + pathname: string, + context: RegisteredReadContext, +): Promise { + const path = pathname.split("/").filter(Boolean); + const selected = getHandler("get", path, roots.handler, false, pathname); + const callback = + typeof selected === "function" + ? selected + : selected && !Array.isArray(selected) + ? selected.handler + : undefined; + if (callback !== expected) { + throw new HTTPResult(READ_FORBIDDEN, "Registered read target mismatch"); + } + const prefix = await executeMiddleware("prefix", "get", path, context); + assertCompletedRead(context); + if (prefix) { + throw new HTTPResult(READ_FORBIDDEN, "Registered read intercepted"); + } + await executeHandlerAndPostfix( + selected as HandlerResult | RouteCallback, + "get", + path, + context, + ); + assertCompletedRead(context); + return context.response; +} + function executeRequest( method: string, path: string[], diff --git a/src/test/controller-resolution.test.ts b/src/test/controller-resolution.test.ts index b1edb6a..e668bb2 100644 --- a/src/test/controller-resolution.test.ts +++ b/src/test/controller-resolution.test.ts @@ -370,6 +370,29 @@ describe("Controller resolution", () => { assert.equal(thenable.readCount(), 1); }); + it("preserves positional holes in programmatically registered parameters", async () => { + const Controller = createController(); + const location = "/controller-resolution/sparse-parameters"; + const parameters: RouteHandler["parameters"] = []; + parameters[1] = computedParameter(() => "record-a"); + register( + createHandler( + Controller, + (first, second) => { + assert.equal(first, undefined); + assert.equal(second, "record-a"); + return second; + }, + location, + parameters, + ), + ); + assert.deepEqual(await get(port, location), { + status: 200, + body: "record-a", + }); + }); + it("turns provider and modifier failures into request errors", async () => { const ProviderController = createController(); const ModifierController = createController(); diff --git a/src/test/registered-read.test.ts b/src/test/registered-read.test.ts new file mode 100644 index 0000000..1d84646 --- /dev/null +++ b/src/test/registered-read.test.ts @@ -0,0 +1,568 @@ +import { Socket } from "node:net"; +import assert from "node:assert/strict"; +import { IncomingMessage, ServerResponse } from "node:http"; +import { HTTPResult, type RouteHandler } from "@antelopejs/interface-api"; +import { ExecuteRegisteredRead as ExecuteReadInterface } from "@antelopejs/interface-api/registered-read"; + +import type { RequestContext } from "../server"; +import { ExecuteRegisteredRead, routesProxy } from "../implementations/api"; + +const ROOT = "/registered-read"; +const RECORD_PATH = `${ROOT}/records/record-a`; +const FORBIDDEN = 403; +const registrations: string[] = []; +const events: string[] = []; +type TransportOperation = (ctx: RequestContext) => unknown; +const transportOperations: Record = { + requestTimeout: (ctx) => ctx.rawRequest.setTimeout(1), + socketDestroy: (ctx) => ctx.rawRequest.socket.destroy(), + socketEnd: (ctx) => ctx.rawRequest.socket.end(), + socketConnect: (ctx) => ctx.rawRequest.socket.connect(1, "127.0.0.1"), + corkedWrite: (ctx) => { + ctx.rawRequest.socket.cork(); + ctx.rawRequest.socket.write("HTTP/1.1 403 Forbidden\r\n\r\n"); + }, + responseTimeout: (ctx) => + new Promise((resolve) => ctx.rawResponse.setTimeout(1, resolve)), +}; +type InformationalWrite = ( + response: ServerResponse, + callback?: () => void, +) => void; +const informationalWrites: Record = { + continue: (response, callback) => response.writeContinue(callback), + processing: (response, callback) => response.writeProcessing(callback), + hints: (response, callback) => + response.writeEarlyHints( + { link: "; rel=preload; as=style" }, + callback, + ), +}; + +class ReadController { + static location = ROOT; + actor?: string; +} + +function parentContext(): RequestContext { + const request = new IncomingMessage(new Socket()); + request.headers = { + authorization: "Bearer fixture", + "content-length": "999", + }; + request.method = "POST"; + request.url = "/files?bypass=true"; + return { + rawRequest: request, + rawResponse: new ServerResponse(request), + url: new URL("http://localhost/files?bypass=true"), + routeParameters: { id: "untrusted" }, + response: new HTTPResult(), + }; +} + +function register(overrides: Partial = {}): string { + const id = `registered-read-${registrations.length}`; + class FixtureController extends ReadController {} + registrations.push(id); + routesProxy.register(id, { + proto: FixtureController.prototype, + mode: "handler", + method: "get", + location: `${ROOT}/records/:id`, + properties: {}, + parameters: [], + callback: () => ({ id: "record-a" }), + ...overrides, + }); + return id; +} + +function read( + routeId: string, + context = parentContext(), + pathname = RECORD_PATH, +) { + return ExecuteRegisteredRead({ routeId, pathname }, context); +} + +function registerContextPrefix( + callback: (ctx: RequestContext) => unknown, +): void { + register({ + mode: "prefix", + location: ROOT, + parameters: [{ provider: (ctx) => ctx, modifiers: [] }], + callback, + }); +} + +function registerDocumentRead(): string { + return register({ + properties: { + actor: { + provider: (ctx) => { + events.push("property"); + return ctx.rawRequest.headers.authorization; + }, + modifiers: [], + }, + }, + parameters: [ + { + provider: (ctx) => { + events.push("parameter"); + assert.equal(ctx.url.search, ""); + assert.equal(ctx.rawRequest.method, "GET"); + assert.equal(ctx.rawRequest.url, RECORD_PATH); + assert.equal(ctx.rawRequest.headers["content-length"], undefined); + return ctx.routeParameters.id; + }, + modifiers: [], + }, + ], + callback: function (this: ReadController, id: string) { + events.push("handler"); + return { id, actor: this.actor }; + }, + }); +} + +describe("Registered reads", () => { + afterEach(() => { + for (const id of registrations.splice(0)) routesProxy.unregister(id); + events.length = 0; + }); + + it("binds the optional interface to the running API provider", async () => { + const result = await ExecuteReadInterface( + { routeId: register(), pathname: RECORD_PATH }, + parentContext(), + ); + assert.deepEqual(JSON.parse(result.getBody()), { id: "record-a" }); + }); + + it("runs prefixes, properties, parameters and postfixes without mutating the caller", async () => { + register({ + mode: "prefix", + location: ROOT, + callback: () => { + events.push("prefix"); + }, + }); + const routeId = registerDocumentRead(); + register({ + mode: "postfix", + location: ROOT, + callback: () => { + events.push("postfix"); + }, + }); + const parent = parentContext(); + assert.deepEqual(JSON.parse((await read(routeId, parent)).getBody()), { + id: "record-a", + actor: "Bearer fixture", + }); + assert.deepEqual(events, [ + "prefix", + "property", + "parameter", + "handler", + "postfix", + ]); + assert.equal(parent.rawRequest.method, "POST"); + assert.equal(parent.rawRequest.headers["content-length"], "999"); + assert.deepEqual(parent.routeParameters, { id: "untrusted" }); + assert.equal(parent.url.search, "?bypass=true"); + }); + + it("never treats an early successful prefix response as document authorization", async () => { + register({ + mode: "prefix", + location: ROOT, + callback: () => new HTTPResult(200, "intercepted"), + }); + const routeId = register({ + callback: () => { + events.push("handler"); + }, + }); + await assert.rejects(read(routeId), HTTPResult); + assert.deepEqual(events, []); + }); + + it("rejects a prefix that writes to the raw response without returning a value", async () => { + register({ + mode: "prefix", + location: ROOT, + parameters: [{ provider: (ctx) => ctx, modifiers: [] }], + callback: (ctx: RequestContext) => { + ctx.rawResponse.writeHead(200); + }, + }); + const routeId = register({ + callback: () => { + events.push("handler"); + }, + }); + const parent = parentContext(); + await assert.rejects(read(routeId, parent), HTTPResult); + assert.deepEqual(events, []); + assert.equal(parent.rawResponse.headersSent, false); + }); + + it("propagates property and parameter denials before executing the handler", async () => { + const denied = () => { + throw new HTTPResult(FORBIDDEN, "Denied"); + }; + const routeId = register({ + properties: { actor: { provider: denied, modifiers: [] } }, + callback: () => { + events.push("handler"); + }, + }); + await assert.rejects(read(routeId), HTTPResult); + routesProxy.unregister(routeId); + const parameterRoute = register({ + parameters: [{ provider: () => "value", modifiers: [denied] }], + callback: () => { + events.push("handler"); + }, + }); + await assert.rejects(read(parameterRoute), HTTPResult); + assert.deepEqual(events, []); + }); + + it("rejects an awaited raw response end instead of hanging", async () => { + registerContextPrefix( + (ctx) => + new Promise((resolve) => ctx.rawResponse.end("denied", resolve)), + ); + await assert.rejects(read(register()), HTTPResult); + }); + + for (const [name, write] of Object.entries(informationalWrites)) { + for (const awaited of [false, true]) { + it(`rejects ${name} informational output (awaited=${awaited})`, async () => { + registerContextPrefix((ctx) => + awaited + ? new Promise((resolve) => write(ctx.rawResponse, resolve)) + : write(ctx.rawResponse), + ); + const target = register({ + callback: () => { + events.push("handler"); + }, + }); + await assert.rejects(read(target), HTTPResult); + assert.deepEqual(events, []); + }); + } + } + + for (const deferred of [false, true]) { + it(`stops dispatch after a parameter aborts (deferred=${deferred})`, async () => { + const target = register({ + parameters: [ + { + provider: (ctx) => { + ctx.rawResponse.destroy(); + return deferred ? Promise.resolve("id") : "id"; + }, + modifiers: [], + }, + ], + callback: () => { + events.push("handler"); + }, + }); + register({ + mode: "postfix", + location: ROOT, + callback: () => { + events.push("postfix"); + }, + }); + await assert.rejects(read(target), HTTPResult); + await new Promise((resolve) => setImmediate(resolve)); + assert.deepEqual(events, []); + }); + } + + it("does not dispatch later prefixes after an interruption", async () => { + registerContextPrefix((ctx) => { + ctx.rawResponse.destroy(); + }); + registerContextPrefix(() => { + events.push("later prefix"); + }); + await assert.rejects(read(register()), HTTPResult); + assert.deepEqual(events, []); + }); + + for (const deferred of [false, true]) { + it(`does not invoke a computed setter after interruption (deferred=${deferred})`, async () => { + class SetterController { + static location = ROOT; + set actor(_value: unknown) { + events.push("setter"); + } + } + const target = register({ + proto: SetterController.prototype, + properties: { + actor: { + provider: (ctx) => { + ctx.rawResponse.destroy(); + return deferred ? Promise.resolve("actor") : "actor"; + }, + modifiers: [], + }, + }, + callback: () => { + events.push("handler"); + }, + }); + await assert.rejects(read(target), HTTPResult); + await new Promise((resolve) => setImmediate(resolve)); + assert.deepEqual(events, []); + }); + } + + it("rejects raw response destruction before invoking the handler", async () => { + registerContextPrefix((ctx) => { + ctx.rawResponse.destroy(); + }); + const target = register({ + callback: () => { + events.push("handler"); + }, + }); + await assert.rejects(read(target), HTTPResult); + assert.deepEqual(events, []); + }); + + it("rejects child request destruction without destroying the parent socket", async () => { + registerContextPrefix((ctx) => { + ctx.rawRequest.destroy(); + }); + const parent = parentContext(); + const target = register({ + callback: () => { + events.push("handler"); + }, + }); + await assert.rejects(read(target, parent), HTTPResult); + assert.equal(parent.rawRequest.socket.destroyed, false); + assert.equal(parent.rawRequest.destroyed, false); + assert.deepEqual(events, []); + }); + + it("preserves independent credential views and removes framing headers from each", async () => { + const parent = parentContext(); + parent.rawRequest.headers["x-many"] = ["first", "second"]; + parent.rawRequest.rawHeaders = [ + "Authorization", + "Bearer fixture", + "Content-Length", + "999", + ]; + parent.rawRequest.headersDistinct = { + authorization: ["Bearer fixture"], + "content-length": ["999"], + }; + Object.defineProperty(parent.rawRequest.socket, "remoteAddress", { + value: "192.0.2.41", + }); + registerContextPrefix((ctx) => { + assert.notEqual(ctx.rawRequest.socket, parent.rawRequest.socket); + assert.notEqual(ctx.rawRequest.connection, parent.rawRequest.connection); + assert.equal(ctx.rawRequest.socket.remoteAddress, "192.0.2.41"); + assert.equal(ctx.rawRequest.headers.authorization, "Bearer fixture"); + assert.deepEqual(ctx.rawRequest.rawHeaders, [ + "Authorization", + "Bearer fixture", + ]); + assert.deepEqual(ctx.rawRequest.headersDistinct, { + authorization: ["Bearer fixture"], + }); + const many = ctx.rawRequest.headers["x-many"]; + assert.ok(Array.isArray(many)); + many.push("child"); + ctx.rawRequest.headersDistinct.authorization!.push("child"); + ctx.rawRequest.rawHeaders.push("Child", "value"); + }); + await read(register(), parent); + assert.deepEqual(parent.rawRequest.headers["x-many"], ["first", "second"]); + assert.deepEqual(parent.rawRequest.headersDistinct.authorization, [ + "Bearer fixture", + ]); + assert.deepEqual(parent.rawRequest.rawHeaders, [ + "Authorization", + "Bearer fixture", + "Content-Length", + "999", + ]); + }); + + it("rejects non-finite response statuses", async () => { + await assert.rejects( + read( + register({ callback: () => new HTTPResult(Number.NaN, "not success") }), + ), + HTTPResult, + ); + }); + + for (const status of [403, 302, 200]) { + it(`rejects a denied or streaming handler before a successful postfix (${status})`, async () => { + const response = new HTTPResult(status, "denied"); + if (status === 200) response.getWriteStream().end("secret"); + const target = register({ callback: () => response }); + register({ + mode: "postfix", + location: ROOT, + callback: () => { + events.push("postfix"); + return { success: true }; + }, + }); + await assert.rejects(read(target), HTTPResult); + assert.deepEqual(events, []); + }); + } + + it("observes pending properties when a later property interrupts collection", async () => { + const unhandled: unknown[] = []; + const onUnhandled = (reason: unknown) => unhandled.push(reason); + process.on("unhandledRejection", onUnhandled); + try { + const target = register({ + properties: { + first: { provider: () => Promise.resolve("actor"), modifiers: [] }, + second: { + provider: (ctx) => ctx.rawResponse.destroy(), + modifiers: [], + }, + }, + callback: () => events.push("handler"), + }); + await assert.rejects(read(target), HTTPResult); + await new Promise((resolve) => setImmediate(resolve)); + assert.deepEqual(unhandled, []); + assert.deepEqual(events, []); + } finally { + process.off("unhandledRejection", onUnhandled); + } + }); + + for (const [name, operation] of Object.entries(transportOperations)) { + it(`interrupts ${name} without mutating the live connection`, async () => { + const parent = parentContext(); + const timeout = parent.rawRequest.socket.timeout; + registerContextPrefix(operation); + const target = register({ callback: () => events.push("handler") }); + await assert.rejects(read(target, parent), HTTPResult); + assert.equal(parent.rawRequest.socket.timeout, timeout); + assert.equal(parent.rawRequest.socket.destroyed, false); + assert.deepEqual(events, []); + }); + } + + it("uses only server-selected query values in both child URL views", async () => { + const parent = parentContext(); + const id = "record&bypass=true/#"; + const routeId = register({ + location: `${ROOT}/get`, + parameters: [{ provider: (ctx) => ctx, modifiers: [] }], + callback: (ctx: RequestContext) => { + assert.equal(ctx.url.searchParams.get("id"), id); + assert.equal(ctx.url.searchParams.has("bypass"), false); + assert.equal( + ctx.rawRequest.url, + `${ROOT}/get?id=record%26bypass%3Dtrue%2F%23`, + ); + return { id }; + }, + }); + const result = await ExecuteRegisteredRead( + { routeId, pathname: `${ROOT}/get`, query: { id } }, + parent, + ); + assert.deepEqual(JSON.parse(result.getBody()), { id }); + assert.equal(parent.url.search, "?bypass=true"); + }); + + for (const withModifier of [false, true]) { + it(`observes rejecting parameters during interrupted dispatch (modifier=${withModifier})`, async () => { + const unhandled: unknown[] = []; + const onUnhandled = (reason: unknown) => unhandled.push(reason); + process.on("unhandledRejection", onUnhandled); + try { + const target = register({ + parameters: [ + { + provider: (ctx) => { + ctx.rawResponse.destroy(); + return Promise.reject(new Error("parameter denied")); + }, + modifiers: withModifier ? [() => events.push("modifier")] : [], + }, + { provider: () => events.push("later parameter"), modifiers: [] }, + ], + callback: () => events.push("handler"), + }); + await assert.rejects(read(target), HTTPResult); + await new Promise((resolve) => setImmediate(resolve)); + assert.deepEqual(unhandled, []); + assert.deepEqual(events, []); + } finally { + process.off("unhandledRejection", onUnhandled); + } + }); + } + + it("rejects handler and postfix denials", async () => { + const denied = register({ + callback: () => new HTTPResult(FORBIDDEN, "Denied"), + }); + await assert.rejects(read(denied), HTTPResult); + routesProxy.unregister(denied); + const allowed = register(); + register({ + mode: "postfix", + location: ROOT, + callback: () => new HTTPResult(FORBIDDEN, "Denied"), + }); + await assert.rejects(read(allowed), HTTPResult); + }); + + it("rejects a route ID that does not match the selected pathname", async () => { + const routeId = register(); + register({ location: `${ROOT}/other` }); + await assert.rejects( + read(routeId, parentContext(), `${ROOT}/other`), + HTTPResult, + ); + }); + + it("rejects stale registrations and non-GET targets", async () => { + const routeId = register(); + routesProxy.unregister(routeId); + await assert.rejects(read(routeId), HTTPResult); + await assert.rejects(read(register({ method: "post" })), HTTPResult); + }); + + it("rejects origins, query options, fragments and normalized paths", async () => { + const routeId = register(); + for (const path of [ + "https://other.test/read", + "//other.test/read", + `${RECORD_PATH}?bypass=true`, + `${RECORD_PATH}#fragment`, + `${ROOT}/x/../records/record-a`, + ]) { + await assert.rejects(read(routeId, parentContext(), path), HTTPResult); + } + }); +});