diff --git a/bun.lock b/bun.lock index 98b8fcd9..59a150b1 100644 --- a/bun.lock +++ b/bun.lock @@ -26,6 +26,7 @@ "@wrnexus/authz": "workspace:*", "@wrnexus/captcha": "workspace:*", "@wrnexus/core": "workspace:*", + "@wrnexus/rpc": "workspace:*", "@wrnexus/validation": "workspace:*", }, "devDependencies": { @@ -245,6 +246,7 @@ "@wrnexus/pubsub": "workspace:*", "@wrnexus/pwa": "workspace:*", "@wrnexus/router": "workspace:*", + "@wrnexus/rpc": "workspace:*", "@wrnexus/security": "workspace:*", "@wrnexus/ssr": "workspace:*", "@wrnexus/store": "workspace:*", @@ -450,6 +452,7 @@ "dependencies": { "@wrnexus/authz": "workspace:*", "@wrnexus/core": "workspace:*", + "@wrnexus/helpers": "workspace:*", "@wrnexus/jwt": "workspace:*", "@wrnexus/validation": "workspace:*", }, diff --git a/docs/plans/2026-08-05-inter-app-comms-design.md b/docs/plans/2026-08-05-inter-app-comms-design.md index 3597d22d..3445c5e2 100644 --- a/docs/plans/2026-08-05-inter-app-comms-design.md +++ b/docs/plans/2026-08-05-inter-app-comms-design.md @@ -1,7 +1,7 @@ # Inter-app communication design (`@wrnexus/rpc`) Date: 2026-08-05 -Status: approved, not yet implemented +Status: phase 1 implemented (2026-08-05); phases 2–4 remain deferred Depends on: the permissions system (`@wrnexus/authz`), merged 2026-08-05 ## Problem diff --git a/docs/plans/2026-08-05-inter-app-comms-implementation.md b/docs/plans/2026-08-05-inter-app-comms-implementation.md index f460a49d..46245b9a 100644 --- a/docs/plans/2026-08-05-inter-app-comms-implementation.md +++ b/docs/plans/2026-08-05-inter-app-comms-implementation.md @@ -8,6 +8,8 @@ **Tech Stack:** TypeScript, Bun (`bun:test`), `@wrnexus/validation` (input schemas), `@wrnexus/jwt` (identity token, HS256), `@wrnexus/authz` (permission checks), `@wrnexus/core` (Context/Middleware types only). +**Status:** Implemented 2026-08-05. The deferred phases at the end of this document remain out of scope. + ## Global Constraints - Every `@wrnexus/*` package is version `0.8.4`. Do not change versions. diff --git a/docs/public-api-0.8.json b/docs/public-api-0.8.json index 67398c01..cc8fdd3f 100644 --- a/docs/public-api-0.8.json +++ b/docs/public-api-0.8.json @@ -2372,6 +2372,53 @@ "sortRoutes" ] }, + "@wrnexus/rpc": { + ".": [ + "AnyProcedures", + "CallOptions", + "ExportOptions", + "HandlerContext", + "HttpTransportOptions", + "ImplementOptions", + "ImportOptions", + "InProcessHandler", + "InferInput", + "InferProcedureInput", + "InferProcedureOutput", + "InputSchema", + "ProcedureBuilder", + "ProcedureDef", + "RPC_ERROR_CODES", + "RPC_IDENTITY_HEADER", + "RPC_INTERNAL_HEADER", + "RPC_PATH_PREFIX", + "RpcErrorCode", + "RpcTarget", + "ServiceClient", + "ServiceClientOptions", + "ServiceContract", + "ServiceError", + "ServiceHandlers", + "ServiceImplementation", + "ServiceResult", + "SubjectContext", + "ToResultOptions", + "Transport", + "defineService", + "exportSubjectContext", + "failure", + "httpTransport", + "implement", + "importSubjectContext", + "inProcessTransport", + "isRetryableStatus", + "procedure", + "rpcPath", + "rpcSecret", + "serviceClient", + "success" + ] + }, "@wrnexus/security": { ".": [ "RequestHardeningOptions", diff --git a/examples/auth-showcase/app/services/greeter.ts b/examples/auth-showcase/app/services/greeter.ts new file mode 100644 index 00000000..61cddb0e --- /dev/null +++ b/examples/auth-showcase/app/services/greeter.ts @@ -0,0 +1,23 @@ +import { defineService, implement, procedure } from "@wrnexus/rpc"; +import { v } from "@wrnexus/validation"; + +export const greeter = defineService({ + name: "greeter", + procedures: { + greet: procedure + .input(v.object({ name: v.string() })) + .output<{ message: string; subject: string }>() + .build(), + }, +}); + +export default implement( + greeter, + { + greet: async ({ name }, ctx) => ({ + message: `Hello, ${name}`, + subject: ctx.subject?.subjectId ?? "anonymous", + }), + }, + { selfApp: "auth-showcase" }, +); diff --git a/examples/auth-showcase/package.json b/examples/auth-showcase/package.json index 2c455272..618376ee 100644 --- a/examples/auth-showcase/package.json +++ b/examples/auth-showcase/package.json @@ -13,6 +13,7 @@ "dependencies": { "@wrnexus/auth": "workspace:*", "@wrnexus/authz": "workspace:*", + "@wrnexus/rpc": "workspace:*", "@wrnexus/captcha": "workspace:*", "@wrnexus/core": "workspace:*", "@wrnexus/validation": "workspace:*" diff --git a/packages/dev-server/package.json b/packages/dev-server/package.json index da910339..0969b7ff 100644 --- a/packages/dev-server/package.json +++ b/packages/dev-server/package.json @@ -9,6 +9,7 @@ }, "dependencies": { "@wrnexus/authz": "workspace:*", + "@wrnexus/rpc": "workspace:*", "@wrnexus/core": "workspace:*", "@wrnexus/dev-toolbar": "workspace:*", "@wrnexus/router": "workspace:*", diff --git a/packages/dev-server/src/gateway.ts b/packages/dev-server/src/gateway.ts index e5cf7da8..761d0dc0 100644 --- a/packages/dev-server/src/gateway.ts +++ b/packages/dev-server/src/gateway.ts @@ -15,6 +15,9 @@ import { dirname, join, resolve } from "node:path"; import { fileURLToPath } from "node:url"; import { RESTART_EXIT_CODE } from "./restart.ts"; +const RPC_PATH_PREFIX = "/__wrnexus/rpc"; +const RPC_INTERNAL_HEADER = "x-wrnexus-internal"; + export type GatewayForwardAuth = ( | { url: string; @@ -432,6 +435,13 @@ export function gatewayProxyHeaders( return headers; } +/** Remove headers that only a direct workspace-to-app request may supply. */ +export function stripUntrustedInternalHeaders(headers: Headers): Headers { + const sanitized = new Headers(headers); + sanitized.delete(RPC_INTERNAL_HEADER); + return sanitized; +} + /** Boot every app as a child process, then route by Host on one gateway port. */ export async function startGateway(opts: GatewayOptions): Promise { const port = opts.port ?? 3000; @@ -615,6 +625,10 @@ export async function startGateway(opts: GatewayOptions): Promise req.headers.has(name)) + ); +} + +function json(body: unknown, status = 200): Response { + return Response.json(body, { status, headers: { "cache-control": "private, no-store" } }); +} + +export async function handleRpcRequest( + req: Request, + url: URL, + services: Map, +): Promise { + if (!isRpcPath(url.pathname)) return null; + if (!isInternalCaller(req)) return new Response("Not found", { status: 404 }); + if (req.method !== "POST") return new Response("Method not allowed", { status: 405 }); + const segments = url.pathname.split("/"); + const service = segments[3] ? services.get(segments[3]) : undefined; + const procedure = segments[4]; + if (!service || !procedure || segments.length !== 5) { + return json({ ok: false, code: "RPC_UNKNOWN", message: "Unknown procedure", retryable: false }); + } + let payload: unknown; + try { + payload = await req.json(); + } catch { + return json({ ok: false, code: "RPC_INVALID", message: "Invalid input", retryable: false }); + } + return json( + await service.invoke(procedure, payload, req.headers.get(RPC_IDENTITY_HEADER) ?? undefined), + ); +} diff --git a/packages/dev-server/src/runtime.ts b/packages/dev-server/src/runtime.ts index e8534bc2..332f9969 100644 --- a/packages/dev-server/src/runtime.ts +++ b/packages/dev-server/src/runtime.ts @@ -89,6 +89,8 @@ import { type ResolvedI18n, } from "@wrnexus/i18n"; import { runMiddleware } from "./pipeline.ts"; +import { handleRpcRequest } from "./rpc-dispatch.ts"; +import type { ServiceImplementation } from "@wrnexus/rpc"; import type { HmrHub } from "./hmr.ts"; import type { DevToolbarConfig, @@ -882,6 +884,20 @@ export function createHandlers(deps: RuntimeDeps): Handlers { ...(await getMiddleware()), ]; const maxBodyBytes = deps.maxBodyBytes ?? 10 * 1024 * 1024; // 10 MB default + let servicesPromise: Promise> | undefined; + const loadServices = () => + (servicesPromise ??= (async () => { + const services = new Map(); + for (const entry of router.services) { + const imported = await loadModule(entry.file); + const implementation = imported.default as ServiceImplementation | undefined; + if (!implementation || typeof implementation.invoke !== "function") { + throw new Error(`RPC service ${entry.file} must default-export implement(...)`); + } + services.set(entry.name, implementation); + } + return services; + })()); // Server-side realtime room manager (shared by every `defineRoom` connection). const realtime = createRealtimeRegistry(); @@ -948,6 +964,9 @@ export function createHandlers(deps: RuntimeDeps): Handlers { const secure = (res: Response): Response => withSecurityHeaders(req, res, mode, runtimeSecurity, nonce); + const rpcResponse = await handleRpcRequest(req, url, await loadServices()); + if (rpcResponse) return secure(rpcResponse); + const preflight = createCorsPreflightResponse(req, deps.security); if (preflight) return secure(preflight); diff --git a/packages/dev-server/test/gateway.test.ts b/packages/dev-server/test/gateway.test.ts index e553ef39..cbde0a22 100644 --- a/packages/dev-server/test/gateway.test.ts +++ b/packages/dev-server/test/gateway.test.ts @@ -4,6 +4,7 @@ import { forwardAuthFailure, forwardAuthHeaders, gatewayProxyHeaders, + stripUntrustedInternalHeaders, gatewayRestartDelay, internalError, stripInternalError, @@ -34,6 +35,16 @@ test("gateway disables compression for its internal proxy hop", () => { expect(headers.get("x-forwarded-for")).toBe("127.0.0.1"); }); +test("gateway proxy headers do not preserve the RPC internal marker", () => { + const request = new Request("http://localhost:3000/path", { + headers: { "x-wrnexus-internal": "1" }, + }); + const headers = stripUntrustedInternalHeaders( + gatewayProxyHeaders(request, new URL(request.url), "127.0.0.1", true), + ); + expect(headers.has("x-wrnexus-internal")).toBe(false); +}); + test("forward auth preserves intentional verifier redirects", () => { const redirected = forwardAuthFailure( new Response(null, { status: 302, headers: { location: "/login?returnTo=%2Fadmin" } }), diff --git a/packages/dev-server/test/observability-runtime.test.ts b/packages/dev-server/test/observability-runtime.test.ts index 154b48f1..febd89c1 100644 --- a/packages/dev-server/test/observability-runtime.test.ts +++ b/packages/dev-server/test/observability-runtime.test.ts @@ -14,6 +14,7 @@ function runtime(health: HealthRegistry, trustProxy = false) { stores: [], schemas: [], authz: [], + services: [], matchPage: () => null, matchApi: () => null, matchRealtime: () => null, diff --git a/packages/dev-server/test/rpc-endpoint.test.ts b/packages/dev-server/test/rpc-endpoint.test.ts new file mode 100644 index 00000000..272ee439 --- /dev/null +++ b/packages/dev-server/test/rpc-endpoint.test.ts @@ -0,0 +1,51 @@ +import { describe, expect, test } from "bun:test"; +import { defineService, implement, procedure } from "@wrnexus/rpc"; +import { v } from "@wrnexus/validation"; +import { handleRpcRequest, isInternalCaller, isRpcPath } from "../src/rpc-dispatch.ts"; + +const demo = defineService({ + name: "demo", + procedures: { + add: procedure + .input(v.object({ a: v.number() })) + .output<{ a: number }>() + .build(), + }, +}); +const services = new Map([ + ["demo", implement(demo, { add: async ({ a }) => ({ a }) }, { selfApp: "demo-app" })], +]); + +function request(path: string, headers: Record = {}) { + return new Request(`http://demo.test${path}`, { + method: "POST", + headers: { "content-type": "application/json", ...headers }, + body: JSON.stringify({ a: 2 }), + }); +} + +describe("RPC endpoint", () => { + test("only matches its reserved prefix", () => { + expect(isRpcPath("/__wrnexus/rpc/demo/add")).toBe(true); + expect(isRpcPath("/__wrnexus/rpcx/demo/add")).toBe(false); + }); + + test("dispatches a private request", async () => { + const req = request("/__wrnexus/rpc/demo/add", { "x-wrnexus-internal": "1" }); + expect(await (await handleRpcRequest(req, new URL(req.url), services))!.json()).toEqual({ + ok: true, + value: { a: 2 }, + }); + }); + + test("rejects public or forwarded requests", async () => { + const external = request("/__wrnexus/rpc/demo/add"); + expect((await handleRpcRequest(external, new URL(external.url), services))!.status).toBe(404); + const forwarded = request("/__wrnexus/rpc/demo/add", { + "x-wrnexus-internal": "1", + "x-forwarded-for": "203.0.113.1", + }); + expect(isInternalCaller(forwarded)).toBe(false); + expect((await handleRpcRequest(forwarded, new URL(forwarded.url), services))!.status).toBe(404); + }); +}); diff --git a/packages/router/src/index.ts b/packages/router/src/index.ts index 9060e178..1fcd3911 100644 --- a/packages/router/src/index.ts +++ b/packages/router/src/index.ts @@ -61,6 +61,8 @@ export interface Router { schemas: ComponentRef[]; /** Authorization declarations (`app/authz/.ts`) merged into the catalog. */ authz: ComponentRef[]; + /** Service implementations (`app/services/.ts`) mounted for inter-app calls. */ + services: ComponentRef[]; matchPage(pathname: string): RouteMatch | null; matchApi(pathname: string): RouteMatch | null; matchRealtime(pathname: string): RouteMatch | null; @@ -302,6 +304,17 @@ export function buildRouter(appDir: string, opts: RouterOptions = {}): Router { authz.push({ name, file: f.file }); } + const services: ComponentRef[] = []; + for (const f of scanDir(join(appDir, "services"), [".js"])) { + if (!/\.(ts|js)$/.test(f.file) || /[.]gen[.](ts|js)$/.test(f.file)) continue; + const name = basename(f.file).replace(/\.(ts|js)$/, ""); + if (!isSafeIslandName(name)) { + console.warn(`[wrnexus] skipping service with unsafe name: ${name}`); + continue; + } + services.push({ name, file: f.file }); + } + return { pages, api, @@ -312,6 +325,7 @@ export function buildRouter(appDir: string, opts: RouterOptions = {}): Router { stores, schemas, authz, + services, matchPage: (p) => matchRoute(pages, p), matchApi: (p) => matchRoute(api, p), matchRealtime: (p) => matchRoute(realtime, p), diff --git a/packages/router/test/services-discovery.test.ts b/packages/router/test/services-discovery.test.ts new file mode 100644 index 00000000..bf9f660f --- /dev/null +++ b/packages/router/test/services-discovery.test.ts @@ -0,0 +1,21 @@ +import { describe, expect, test } from "bun:test"; +import { mkdirSync, mkdtempSync, writeFileSync } from "node:fs"; +import { join } from "node:path"; +import { buildRouter } from "../src/index.ts"; + +const root = join(import.meta.dir, ".tmp-services"); +mkdirSync(root, { recursive: true }); + +describe("service discovery", () => { + test("discovers source services and skips generated files", () => { + const base = mkdtempSync(join(root, "app-")); + const services = join(base, "app", "services"); + mkdirSync(services, { recursive: true }); + mkdirSync(join(base, "app", "pages"), { recursive: true }); + writeFileSync(join(services, "billing.ts"), "export default {};"); + writeFileSync(join(services, "types.gen.ts"), "export type T = string;"); + expect(buildRouter(join(base, "app")).services.map((service) => service.name)).toEqual([ + "billing", + ]); + }); +}); diff --git a/packages/rpc/README.md b/packages/rpc/README.md new file mode 100644 index 00000000..9cf468d5 --- /dev/null +++ b/packages/rpc/README.md @@ -0,0 +1,44 @@ +# `@wrnexus/rpc` + +Define a service contract in a shared workspace package, then import that same contract from the caller and callee. + +```ts +import { + defineService, + implement, + inProcessTransport, + procedure, + serviceClient, +} from "@wrnexus/rpc"; +import { v } from "@wrnexus/validation"; + +const greeter = defineService({ + name: "greeter", + procedures: { + greet: procedure + .input(v.object({ name: v.string() })) + .output<{ message: string }>() + .build(), + }, +}); + +const service = implement( + greeter, + { greet: async ({ name }) => ({ message: `Hello, ${name}` }) }, + { selfApp: "greeter" }, +); + +const client = serviceClient(greeter, { + app: "greeter", + transport: inProcessTransport({ + "greeter/greet": (input, identity) => service.invoke("greet", input, identity), + }), +}); +await client.greet({ name: "Ada" }); +``` + +Service files default-export `implement(...)` from `app/services`. The development server mounts them under the private `/__wrnexus/rpc` prefix. + +Pass `{ as: ctx }` to `serviceClient` to propagate the subject. The signed token contains only subject and tenant identifiers; permissions are always checked by the callee. Set `WRNEXUS_RPC_SECRET` in every app, use at least 32 characters, and never reuse the session secret. + +Calls time out by default. Retrying is intentionally deferred; when introduced, only procedures marked `.idempotent()` may be retried. diff --git a/packages/rpc/package.json b/packages/rpc/package.json index b1832c35..fef417e6 100644 --- a/packages/rpc/package.json +++ b/packages/rpc/package.json @@ -22,6 +22,7 @@ "@wrnexus/authz": "workspace:*", "@wrnexus/core": "workspace:*", "@wrnexus/jwt": "workspace:*", + "@wrnexus/helpers": "workspace:*", "@wrnexus/validation": "workspace:*" }, "devDependencies": { diff --git a/packages/rpc/src/client.ts b/packages/rpc/src/client.ts new file mode 100644 index 00000000..457b2c5d --- /dev/null +++ b/packages/rpc/src/client.ts @@ -0,0 +1,57 @@ +import type { Context } from "@wrnexus/core"; +import { RPC_ERROR_CODES, ServiceError } from "./errors.ts"; +import { exportSubjectContext } from "./identity.ts"; +import type { Transport } from "./transport.ts"; +import type { + AnyProcedures, + InferProcedureInput, + InferProcedureOutput, + ServiceContract, +} from "./types.ts"; + +export interface ServiceClientOptions { + app?: string; + transport: Transport; + as?: Context; + timeoutMs?: number; +} + +export type ServiceClient = { + [K in keyof Procedures]: ( + input: InferProcedureInput, + ) => Promise>; +}; + +const DEFAULT_TIMEOUT_MS = 10_000; + +export function serviceClient( + contract: ServiceContract, + options: ServiceClientOptions, +): ServiceClient { + const app = options.app ?? contract.name; + const timeoutMs = options.timeoutMs ?? DEFAULT_TIMEOUT_MS; + return new Proxy({} as ServiceClient, { + get(_target, property) { + if (typeof property !== "string") return undefined; + return async (input: unknown) => { + if (!Object.hasOwn(contract.procedures, property)) { + throw new ServiceError(RPC_ERROR_CODES.unknown, "Unknown procedure"); + } + const identity = options.as ? await exportSubjectContext(options.as, app) : undefined; + const controller = new AbortController(); + const timer = setTimeout(() => controller.abort(), timeoutMs); + try { + const result = await options.transport.call( + { app, service: contract.name, procedure: property }, + input, + { signal: controller.signal, ...(identity ? { identity } : {}) }, + ); + if (result.ok) return result.value; + throw new ServiceError(result.code, result.message); + } finally { + clearTimeout(timer); + } + }; + }, + }); +} diff --git a/packages/rpc/src/http.ts b/packages/rpc/src/http.ts new file mode 100644 index 00000000..d4f72bf7 --- /dev/null +++ b/packages/rpc/src/http.ts @@ -0,0 +1,73 @@ +import { appOrigin } from "@wrnexus/helpers"; +import { RPC_ERROR_CODES, failure, isRetryableStatus } from "./errors.ts"; +import { RPC_IDENTITY_HEADER } from "./identity.ts"; +import type { CallOptions, RpcTarget, Transport } from "./transport.ts"; +import type { ServiceResult } from "./types.ts"; + +export const RPC_PATH_PREFIX = "/__wrnexus/rpc"; +export const RPC_INTERNAL_HEADER = "x-wrnexus-internal"; + +export function rpcPath(service: string, procedure: string): string { + return `${RPC_PATH_PREFIX}/${service}/${procedure}`; +} + +export interface HttpTransportOptions { + resolveOrigin?: (app: string) => string; + fetch?: typeof fetch; +} + +function isServiceResult(value: unknown): value is ServiceResult { + if (!value || typeof value !== "object" || !("ok" in value)) return false; + const result = value as Record; + return ( + result.ok === true || + (result.ok === false && + typeof result.code === "string" && + typeof result.message === "string" && + typeof result.retryable === "boolean") + ); +} + +export function httpTransport(options: HttpTransportOptions = {}): Transport { + const resolveOrigin = options.resolveOrigin ?? appOrigin; + const doFetch = options.fetch ?? fetch; + return { + async call(target: RpcTarget, payload: unknown, callOptions: CallOptions) { + let response: Response; + try { + const headers: Record = { + "content-type": "application/json", + [RPC_INTERNAL_HEADER]: "1", + }; + if (callOptions.identity) headers[RPC_IDENTITY_HEADER] = callOptions.identity; + response = await doFetch( + `${resolveOrigin(target.app)}${rpcPath(target.service, target.procedure)}`, + { + method: "POST", + headers, + body: JSON.stringify(payload ?? {}), + signal: callOptions.signal, + }, + ); + } catch { + return failure(RPC_ERROR_CODES.transport, "Service unreachable"); + } + if (!response.ok) { + return { + ok: false, + code: RPC_ERROR_CODES.transport, + message: `Service returned ${response.status}`, + retryable: isRetryableStatus(response.status), + }; + } + try { + const result: unknown = await response.json(); + return isServiceResult(result) + ? result + : failure(RPC_ERROR_CODES.malformed, "Malformed service response"); + } catch { + return failure(RPC_ERROR_CODES.malformed, "Malformed service response"); + } + }, + }; +} diff --git a/packages/rpc/src/index.ts b/packages/rpc/src/index.ts index c6e2591e..b4a854b3 100644 --- a/packages/rpc/src/index.ts +++ b/packages/rpc/src/index.ts @@ -30,3 +30,17 @@ export { rpcSecret, } from "./identity.ts"; export type { ExportOptions, ImportOptions, SubjectContext } from "./identity.ts"; + +export { inProcessTransport } from "./transport.ts"; +export type { CallOptions, InProcessHandler, RpcTarget, Transport } from "./transport.ts"; +export { implement } from "./server.ts"; +export type { + HandlerContext, + ImplementOptions, + ServiceHandlers, + ServiceImplementation, +} from "./server.ts"; +export { serviceClient } from "./client.ts"; +export type { ServiceClient, ServiceClientOptions } from "./client.ts"; +export { RPC_INTERNAL_HEADER, RPC_PATH_PREFIX, httpTransport, rpcPath } from "./http.ts"; +export type { HttpTransportOptions } from "./http.ts"; diff --git a/packages/rpc/src/server.ts b/packages/rpc/src/server.ts new file mode 100644 index 00000000..e806537e --- /dev/null +++ b/packages/rpc/src/server.ts @@ -0,0 +1,81 @@ +import { RPC_ERROR_CODES, failure, success } from "./errors.ts"; +import { importSubjectContext, type SubjectContext } from "./identity.ts"; +import type { + AnyProcedures, + InferProcedureInput, + InferProcedureOutput, + ServiceContract, + ServiceResult, +} from "./types.ts"; + +export interface HandlerContext { + subject?: SubjectContext; +} + +export type ServiceHandlers = { + [K in keyof Procedures]: ( + input: InferProcedureInput, + ctx: HandlerContext, + ) => Promise> | InferProcedureOutput; +}; + +export interface ImplementOptions { + selfApp: string; + checkPermission?: (permission: string, subject?: SubjectContext) => Promise | boolean; +} + +export interface ServiceImplementation { + contract: ServiceContract; + invoke(procedure: string, payload: unknown, identity?: string): Promise; +} + +export function implement( + contract: ServiceContract, + handlers: ServiceHandlers, + options: ImplementOptions, +): ServiceImplementation { + return { + contract, + async invoke(procedureName, payload, identity) { + const definition = contract.procedures[procedureName as keyof Procedures]; + const handler = handlers[procedureName as keyof Procedures]; + if (!definition || !handler) return failure(RPC_ERROR_CODES.unknown, "Unknown procedure"); + + let subject: SubjectContext | undefined; + if (identity !== undefined) { + try { + subject = await importSubjectContext(identity, options.selfApp); + } catch { + return failure(RPC_ERROR_CODES.identity, "Invalid identity"); + } + } + + if (definition.permission) { + if (!options.checkPermission) return failure(RPC_ERROR_CODES.denied, "Forbidden"); + try { + if (!(await options.checkPermission(definition.permission, subject))) { + return failure(RPC_ERROR_CODES.denied, "Forbidden"); + } + } catch { + return failure(RPC_ERROR_CODES.denied, "Forbidden"); + } + } + + let input: unknown = payload; + if (definition.input) { + const parsed = definition.input.parse(payload as Record); + if (!parsed.ok) return failure(RPC_ERROR_CODES.invalid, "Invalid input"); + input = parsed.value; + } + + try { + const value = await (handler as (value: unknown, ctx: HandlerContext) => unknown)(input, { + subject, + }); + return success(value); + } catch { + return failure(RPC_ERROR_CODES.handler, "Internal error"); + } + }, + }; +} diff --git a/packages/rpc/src/transport.ts b/packages/rpc/src/transport.ts new file mode 100644 index 00000000..19fab891 --- /dev/null +++ b/packages/rpc/src/transport.ts @@ -0,0 +1,38 @@ +import { RPC_ERROR_CODES, failure } from "./errors.ts"; +import type { ServiceResult } from "./types.ts"; + +export interface RpcTarget { + app: string; + service: string; + procedure: string; +} + +export interface CallOptions { + signal?: AbortSignal; + identity?: string; +} + +export interface Transport { + call(target: RpcTarget, payload: unknown, options: CallOptions): Promise; +} + +export type InProcessHandler = ( + payload: unknown, + identity?: string, +) => Promise | ServiceResult; + +/** Direct transport for tests and local integration harnesses. */ +export function inProcessTransport(handlers: Record): Transport { + return { + async call(target, payload, options) { + if (options.signal?.aborted) return failure(RPC_ERROR_CODES.transport, "Call aborted"); + const handler = handlers[`${target.service}/${target.procedure}`]; + if (!handler) return failure(RPC_ERROR_CODES.unknown, "Unknown procedure"); + try { + return await handler(payload, options.identity); + } catch { + return failure(RPC_ERROR_CODES.handler, "Internal error"); + } + }, + }; +} diff --git a/packages/rpc/test/http.test.ts b/packages/rpc/test/http.test.ts new file mode 100644 index 00000000..856e4ba5 --- /dev/null +++ b/packages/rpc/test/http.test.ts @@ -0,0 +1,45 @@ +import { describe, expect, test } from "bun:test"; +import { RPC_ERROR_CODES } from "../src/errors.ts"; +import { RPC_IDENTITY_HEADER } from "../src/identity.ts"; +import { httpTransport, rpcPath } from "../src/http.ts"; + +const target = { app: "billing", service: "billing", procedure: "createInvoice" }; + +function transportWith(handler: (request: Request) => Response | Promise) { + return httpTransport({ + resolveOrigin: () => "http://billing.test", + fetch: (async (input: RequestInfo | URL, init?: RequestInit) => + handler(new Request(input, init))) as typeof fetch, + }); +} + +describe("httpTransport", () => { + test("posts to the private endpoint with identity", async () => { + const transport = transportWith(async (request) => { + expect(request.url).toBe("http://billing.test/__wrnexus/rpc/billing/createInvoice"); + expect(request.headers.get(RPC_IDENTITY_HEADER)).toBe("token"); + expect(request.headers.get("x-wrnexus-internal")).toBe("1"); + expect(await request.json()).toEqual({ amountCents: 5 }); + return Response.json({ ok: true, value: { invoiceId: "inv_1" } }); + }); + expect(await transport.call(target, { amountCents: 5 }, { identity: "token" })).toEqual({ + ok: true, + value: { invoiceId: "inv_1" }, + }); + }); + + test("classifies unavailable and malformed responses safely", async () => { + const unavailable = transportWith(() => new Response("busy", { status: 503 })); + const failed = await unavailable.call(target, {}, {}); + expect(failed).toMatchObject({ code: RPC_ERROR_CODES.transport, retryable: true }); + const malformed = transportWith(() => Response.json({ hello: "world" })); + expect(await malformed.call(target, {}, {})).toMatchObject({ + code: RPC_ERROR_CODES.malformed, + retryable: false, + }); + }); + + test("uses the stable reserved path", () => { + expect(rpcPath("billing", "createInvoice")).toBe("/__wrnexus/rpc/billing/createInvoice"); + }); +}); diff --git a/packages/rpc/test/integration.test.ts b/packages/rpc/test/integration.test.ts new file mode 100644 index 00000000..f858befe --- /dev/null +++ b/packages/rpc/test/integration.test.ts @@ -0,0 +1,74 @@ +import { afterEach, describe, expect, test } from "bun:test"; +import type { Context } from "@wrnexus/core"; +import { v } from "@wrnexus/validation"; +import { serviceClient } from "../src/client.ts"; +import { defineService, procedure } from "../src/contract.ts"; +import { RPC_ERROR_CODES } from "../src/errors.ts"; +import { implement } from "../src/server.ts"; +import { inProcessTransport } from "../src/transport.ts"; + +const originalRpcSecret = process.env.WRNEXUS_RPC_SECRET; +const originalAppName = process.env.WRNEXUS_APP_NAME; +afterEach(() => { + if (originalRpcSecret === undefined) delete process.env.WRNEXUS_RPC_SECRET; + else process.env.WRNEXUS_RPC_SECRET = originalRpcSecret; + if (originalAppName === undefined) delete process.env.WRNEXUS_APP_NAME; + else process.env.WRNEXUS_APP_NAME = originalAppName; +}); + +const billing = defineService({ + name: "billing", + procedures: { + createInvoice: procedure + .input(v.object({ amountCents: v.number() })) + .output<{ invoiceId: string; forSubject: string }>() + .permission("invoice:create") + .build(), + }, +}); + +function wire(allowed: boolean) { + const service = implement( + billing, + { + createInvoice: async ({ amountCents }, ctx) => ({ + invoiceId: `inv_${amountCents}`, + forSubject: ctx.subject?.subjectId ?? "anon", + }), + }, + { selfApp: "billing", checkPermission: async () => allowed }, + ); + return inProcessTransport({ + "billing/createInvoice": (payload, identity) => + service.invoke("createInvoice", payload, identity), + }); +} + +describe("RPC integration", () => { + test("propagates identity and validates the callee permission", async () => { + process.env.WRNEXUS_RPC_SECRET = "test-rpc-secret-at-least-32-chars-long"; + process.env.WRNEXUS_APP_NAME = "web"; + const client = serviceClient(billing, { + app: "billing", + as: { user: { id: "u1" }, tenant: { id: "acme" }, locals: {} } as unknown as Context, + transport: wire(true), + }); + expect(await client.createInvoice({ amountCents: 250 })).toEqual({ + invoiceId: "inv_250", + forSubject: "u1", + }); + }); + + test("denies when the callee permission check refuses", async () => { + process.env.WRNEXUS_RPC_SECRET = "test-rpc-secret-at-least-32-chars-long"; + process.env.WRNEXUS_APP_NAME = "web"; + const client = serviceClient(billing, { + app: "billing", + as: { user: { id: "u1" }, locals: {} } as unknown as Context, + transport: wire(false), + }); + await expect(client.createInvoice({ amountCents: 1 })).rejects.toMatchObject({ + code: RPC_ERROR_CODES.denied, + }); + }); +});