Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
110 changes: 86 additions & 24 deletions src/implementations/api/index.ts
Original file line number Diff line number Diff line change
@@ -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";

Check failure on line 3 in src/implementations/api/index.ts

View workflow job for this annotation

GitHub Actions / checks

Cannot find module '@antelopejs/interface-api/registered-read' or its corresponding type declarations.
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";
Expand Down Expand Up @@ -57,6 +62,7 @@

const classCacheSymbol = Symbol();
const registeredRoutes = new Map<string, RouteHandler>();
const compiledRoutes = new Map<string, RouteCallback>();
const controllerPlans = new WeakMap<ControllerClass, ControllerPlan>();

interface RequestContextDev extends RequestContext {
Expand All @@ -81,16 +87,21 @@
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(
Expand All @@ -108,6 +119,7 @@
applyModifiers(resolved, modifiers, context, controller, index),
);
}
assertReadActive(context);
value = modifiers[index].call(controller, context, value);
}
const then = getThen(value);
Expand Down Expand Up @@ -193,17 +205,24 @@
}

const pending: Promise<void>[] = [];
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) {
Expand Down Expand Up @@ -270,32 +289,52 @@
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);
}

Expand All @@ -317,25 +356,48 @@
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<HTTPResult> {
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[] => {
Expand Down
4 changes: 4 additions & 0 deletions src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,10 @@
await import("@antelopejs/interface-api"),
await import("./implementations/api"),
);
void ImplementInterface(
await import("@antelopejs/interface-api/registered-read"),

Check failure on line 32 in src/index.ts

View workflow job for this annotation

GitHub Actions / checks

Cannot find module '@antelopejs/interface-api/registered-read' or its corresponding type declarations.
await import("./implementations/api"),
);
}

export function destroy(): void {}
Expand Down
176 changes: 176 additions & 0 deletions src/registered-read-context.ts
Original file line number Diff line number Diff line change
@@ -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";

Check failure on line 3 in src/registered-read-context.ts

View workflow job for this annotation

GitHub Actions / checks

Cannot find module '@antelopejs/interface-api/registered-read' or its corresponding type declarations.

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<RequestContext, AbortSignal>();

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);

Check failure on line 135 in src/registered-read-context.ts

View workflow job for this annotation

GitHub Actions / checks

Argument of type 'unknown' is not assignable to parameter of type 'string'.
}
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<HTTPResult>,
): Promise<HTTPResult> {
let rejectRead: (reason: unknown) => void;
const interrupted = new Promise<never>((_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);
}
}
Loading
Loading