| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638 |
- /**
- * Live TypeRT Remote dispatch over Cordis Services and registered providers.
- * Transport, request correlation, and response envelopes belong to Connection.
- * @module @deepseek-ai/dsh-api-gateway
- */
- import { Context, Service, symbols } from '@deepseek-ai/cordis'
- import type { ConnectionRpcHandler } from '@deepseek-ai/dsh-client-connection'
- import {
- remoteMethods,
- TypeRTLookupFailure,
- type InvocationDescriptor,
- type InvocationParameterDescriptor,
- type TypeRTCodec,
- type TypeRTGatewayBinding,
- } from '@deepseek-ai/dsh-type-meta'
- import type {
- InvokeRemoteRequest,
- TypertGateway,
- TypertGatewayErrorCode,
- } from './types.ts'
- export type {
- InvokeRemoteRequest,
- TypertGateway,
- TypertGatewayErrorCode,
- } from './types.ts'
- interface GatewayErrorOptions {
- readonly cause?: unknown
- readonly field?: string
- }
- interface ResolvedBinding {
- readonly binding: TypeRTGatewayBinding
- readonly original: object
- }
- type ConnectionRpcResult = Awaited<ReturnType<ConnectionRpcHandler>>
- type ConnectionRpcError = Extract<ConnectionRpcResult, { readonly ok: false }>['error']
- const NEVER_ABORTED_SIGNAL = new AbortController().signal
- /** Dispatch failure produced outside the invoked business method. */
- export class TypertGatewayError extends Error {
- /** Machine-readable failure category. */
- readonly code: TypertGatewayErrorCode
- /** Canonical `<namespace>/<method>` endpoint. */
- readonly endpoint: string
- /** Affected wire field when the failure is field-specific. */
- readonly field: string | undefined
- /**
- * Construct a Gateway failure without embedding boundary values in its message.
- * @param code - stable failure category.
- * @param endpoint - canonical Remote endpoint.
- * @param message - correction-oriented diagnostic without sensitive values.
- * @param options - optional field and contained cause.
- */
- constructor(
- code: TypertGatewayErrorCode,
- endpoint: string,
- message: string,
- options: GatewayErrorOptions = {},
- ) {
- super(`typert gateway: ${endpoint}: ${message}`, options.cause === undefined ? undefined : { cause: options.cause })
- this.name = 'TypertGatewayError'
- this.code = code
- this.endpoint = endpoint
- this.field = options.field
- }
- }
- /**
- * Resolve strict generated definitions or conservative SRC markers against
- * current Cordis Services and TypeRT providers.
- * @typert service typertGateway
- */
- export class TypertGatewayService extends Service implements TypertGateway {
- static inject = ['typert']
- private srcClaims: ReadonlySet<string> | undefined
- /**
- * Register the Gateway against the active TypeRT registry.
- * @param ctx - owning Host Context with TypeRT registry access.
- */
- constructor(ctx: Context) {
- super(ctx, 'typertGateway')
- ctx.on('internal/service', () => {
- this.srcClaims = undefined
- })
- ctx.inject(['connection'], (connectionCtx) => {
- connectionCtx.connection.rpc.intercept(
- '/api',
- endpoint => this.claimsEndpoint(endpoint),
- (endpoint, payload, signal) => this.dispatchRpc(endpoint, payload, signal),
- { authority: 'trusted-host' },
- )
- })
- }
- private claimsEndpoint(endpoint: string): boolean {
- const segments = endpoint.split('/')
- if (segments.length !== 2 || segments[0] === '' || segments[1] === '') return false
- if (this.ctx.typert.local.get(endpoint) !== undefined || this.ctx.typert.local.hasSeen(endpoint)) return true
- this.srcClaims ??= this.collectSrcClaims()
- return this.srcClaims.has(endpoint)
- }
- private collectSrcClaims(): ReadonlySet<string> {
- const claims = new Set<string>()
- for (const [serviceKey, definition] of Object.entries(this.ctx.reflect.props)) {
- if (definition.type !== 'service') continue
- const receiver = this.ctx.get(serviceKey) as unknown
- if (!isObject(receiver)) continue
- const original = originalOf(receiver)
- const binding = Reflect.get(original, 'typertGateway') as unknown
- if (!isObject(binding) || typeof Reflect.get(binding, 'namespace') !== 'string') continue
- const namespace = Reflect.get(binding, 'namespace') as string
- for (const candidate of remoteMethods(original)) {
- claims.add(endpointOf(namespace, candidate.exportName ?? candidate.method))
- }
- }
- return claims
- }
- /**
- * Invoke one live Remote method through strict generated reflection or SRC markers.
- * @param request - decoded endpoint and exact named wire arguments.
- * @returns the validated business result.
- * @throws {@link TypertGatewayError} for dispatch, provider, or boundary failures; lookup-policy and business errors retain identity.
- */
- async invoke(request: InvokeRemoteRequest): Promise<unknown> {
- const endpoint = endpointOf(request.namespace, request.method)
- const descriptor = this.resolveDescriptor(request.namespace, request.method, endpoint)
- assertExactArguments(request.args, descriptor, endpoint)
- const receiverContext = await this.resolveReceiverContext(descriptor, request.args, endpoint)
- const receiver = receiverContext.get(descriptor.service) as unknown
- if (!isObject(receiver)) {
- throw new TypertGatewayError(
- 'service-unavailable',
- endpoint,
- `active Service ${JSON.stringify(descriptor.service)} is unavailable`,
- )
- }
- validateBinding(receiver, descriptor.service, descriptor.namespace, endpoint)
- const args = await Promise.all(descriptor.parameters.map(parameter =>
- this.resolveParameter(parameter, request.args, endpoint)))
- if (descriptor.cancellation !== undefined) args.push(request.signal ?? NEVER_ABORTED_SIGNAL)
- const implementation = descriptor.implementation ?? descriptor.method
- const method = Reflect.get(receiver, implementation) as unknown
- if (typeof method !== 'function') {
- throw new TypertGatewayError(
- 'method-unavailable',
- endpoint,
- `active Service ${JSON.stringify(descriptor.service)} has no callable method ${JSON.stringify(implementation)}`,
- )
- }
- const result = await Reflect.apply(method, receiver, args) as unknown
- return decode(descriptor.result, result, 'result-invalid', endpoint, 'result')
- }
- private async dispatchRpc(
- endpoint: string,
- payload: unknown,
- signal: AbortSignal,
- ): Promise<ConnectionRpcResult> {
- return this.invokeRpc(endpoint, payload, signal)
- }
- private async invokeRpc(endpoint: string, payload: unknown, signal: AbortSignal): Promise<ConnectionRpcResult> {
- try {
- const segments = endpoint.split('/')
- if (segments.length !== 2 || segments[0] === '' || segments[1] === '') {
- throw new Error(`invalid Remote endpoint ${JSON.stringify(endpoint)}`)
- }
- const [namespace, method] = segments as [string, string]
- if (!isObject(payload)
- || !isPlainObject(payload)
- || Reflect.ownKeys(payload).length !== 1
- || !Object.hasOwn(payload, 'args')
- || !isObject(payload.args)
- || !isPlainObject(payload.args)) {
- throw new Error('Remote payload must contain exactly one plain-object args field')
- }
- const value = await this.invoke({
- namespace,
- method,
- args: payload.args,
- signal,
- })
- return { ok: true, value }
- } catch (error) {
- return rpcFailure(error)
- }
- }
- private resolveDescriptor(namespace: string, method: string, endpoint: string): InvocationDescriptor {
- const strict = this.ctx.typert.local.get(endpoint)
- if (strict !== undefined) return strict
- if (this.ctx.typert.local.hasSeen(endpoint)) {
- throw new TypertGatewayError(
- 'definition-unavailable',
- endpoint,
- 'its strict definition was withdrawn and SRC fallback is forbidden',
- )
- }
- return this.resolveSrcDescriptor(namespace, method, endpoint)
- }
- private resolveSrcDescriptor(namespace: string, method: string, endpoint: string): InvocationDescriptor {
- const candidates: InvocationDescriptor[] = []
- for (const [serviceKey, definition] of Object.entries(this.ctx.reflect.props)) {
- if (definition.type !== 'service') continue
- const receiver = this.ctx.get(serviceKey) as unknown
- if (!isObject(receiver)) continue
- const original = originalOf(receiver)
- const value = Reflect.get(original, 'typertGateway') as unknown
- if (value === undefined) continue
- const binding = readBinding(value, original, serviceKey, endpoint)
- if (binding.namespace !== namespace) continue
- const marker = remoteMethods(original).find(candidate => (candidate.exportName ?? candidate.method) === method)
- if (marker === undefined) continue
- candidates.push(this.srcDescriptor(binding, marker, method, endpoint))
- }
- if (candidates.length === 0) {
- throw new TypertGatewayError('invocation-unavailable', endpoint, 'no active Remote method exports this endpoint')
- }
- if (candidates.length > 1) {
- throw new TypertGatewayError(
- 'ambiguous-endpoint',
- endpoint,
- `multiple active Services export this endpoint: ${candidates.map(candidate => candidate.service).sort().join(', ')}`,
- )
- }
- return candidates[0] as InvocationDescriptor
- }
- private srcDescriptor(
- binding: TypeRTGatewayBinding,
- marker: ReturnType<typeof remoteMethods>[number],
- method: string,
- endpoint: string,
- ): InvocationDescriptor {
- const names = methodParameterNames(binding.service, marker.method, endpoint)
- const signalIndex = names.indexOf('signal')
- if (signalIndex >= 0 && signalIndex !== names.length - 1) {
- throw new TypertGatewayError(
- 'signature-invalid',
- endpoint,
- 'SRC cancellation parameter signal must be the final parameter',
- { field: 'signal' },
- )
- }
- const cancellation = signalIndex >= 0
- ? { parameter: 'signal' as const }
- : undefined
- const businessNames = cancellation === undefined ? names : names.slice(0, -1)
- const parameters: InvocationParameterDescriptor[] = []
- const wires = new Set<string>()
- for (const name of businessNames) {
- const matches = this.ctx.typert.lookups.definitions()
- .filter(definition => definition.parameter === name)
- if (matches.length > 1) {
- throw new TypertGatewayError(
- 'signature-invalid',
- endpoint,
- `parameter ${JSON.stringify(name)} matches multiple lookup providers`,
- { field: name },
- )
- }
- const match = matches[0]
- const parameter: InvocationParameterDescriptor = match === undefined
- ? { name, wire: name, source: 'json', codec: { mode: 'src-json' } }
- : {
- name,
- wire: match.wire,
- source: 'lookup',
- lookup: match.key,
- codec: { mode: 'src-json' },
- }
- if (wires.has(parameter.wire)) {
- throw new TypertGatewayError(
- 'signature-invalid',
- endpoint,
- `multiple parameters use wire field ${JSON.stringify(parameter.wire)}`,
- { field: parameter.wire },
- )
- }
- wires.add(parameter.wire)
- parameters.push(parameter)
- }
- let receiver: InvocationDescriptor['invocation'] = { kind: 'direct' }
- if (marker.invocation.kind === 'context') {
- const provider = this.ctx.typert.contexts.getHost(marker.invocation.context)
- if (provider === undefined) {
- throw new TypertGatewayError(
- 'context-unavailable',
- endpoint,
- `Context provider ${JSON.stringify(marker.invocation.context)} is unavailable`,
- )
- }
- if (wires.has(provider.wire)) {
- throw new TypertGatewayError(
- 'signature-invalid',
- endpoint,
- `Context identity conflicts with wire field ${JSON.stringify(provider.wire)}`,
- { field: provider.wire },
- )
- }
- receiver = {
- kind: 'context',
- context: marker.invocation.context,
- wire: provider.wire,
- codec: { mode: 'src-json' },
- }
- }
- return {
- id: `src:${binding.serviceKey}#${endpoint}`,
- service: binding.serviceKey,
- namespace: binding.namespace,
- method,
- ...(marker.method === method ? {} : { implementation: marker.method }),
- invocation: receiver,
- parameters,
- ...(cancellation === undefined ? {} : { cancellation }),
- result: { mode: 'src-json' },
- }
- }
- private async resolveReceiverContext(
- descriptor: InvocationDescriptor,
- args: Readonly<Record<string, unknown>>,
- endpoint: string,
- ): Promise<Context> {
- if (descriptor.invocation.kind === 'direct') return this.ctx
- const invocation = descriptor.invocation
- const provider = this.ctx.typert.contexts.getHost(invocation.context)
- if (provider === undefined) {
- throw new TypertGatewayError(
- 'context-unavailable',
- endpoint,
- `Context provider ${JSON.stringify(invocation.context)} is unavailable`,
- )
- }
- if (provider.wire !== invocation.wire
- || (invocation.codec.mode === 'strict' && provider.wireTypeSymbol !== invocation.codec.typeSymbol)) {
- throw new TypertGatewayError(
- 'provider-mismatch',
- endpoint,
- `Context provider ${JSON.stringify(invocation.context)} does not match its strict definition`,
- { field: invocation.wire },
- )
- }
- const identity = decode(invocation.codec, args[invocation.wire], 'input-invalid', endpoint, invocation.wire)
- let context: Context | undefined
- try {
- context = await provider.resolve(identity)
- } catch (cause) {
- if (cause instanceof TypeRTLookupFailure) throw cause
- throw new TypertGatewayError(
- 'context-failed',
- endpoint,
- `Context provider ${JSON.stringify(invocation.context)} failed`,
- { cause, field: invocation.wire },
- )
- }
- if (context === undefined) {
- throw new TypertGatewayError(
- 'context-not-found',
- endpoint,
- `Context provider ${JSON.stringify(invocation.context)} did not resolve the requested identity`,
- { field: invocation.wire },
- )
- }
- return context
- }
- private async resolveParameter(
- parameter: InvocationParameterDescriptor,
- args: Readonly<Record<string, unknown>>,
- endpoint: string,
- ): Promise<unknown> {
- const value = decode(parameter.codec, args[parameter.wire], 'input-invalid', endpoint, parameter.wire)
- if (parameter.source === 'json') return value
- const key = parameter.lookup
- /* v8 ignore next -- registry validation rejects strict descriptors without a key, and SRC derivation always supplies one. */
- if (key === undefined) {
- throw new TypertGatewayError(
- 'lookup-unavailable',
- endpoint,
- `lookup parameter ${JSON.stringify(parameter.name)} has no provider key`,
- { field: parameter.wire },
- )
- }
- const provider = this.ctx.typert.lookups.get(key)
- if (provider === undefined) {
- throw new TypertGatewayError(
- 'lookup-unavailable',
- endpoint,
- `lookup provider ${JSON.stringify(key)} is unavailable`,
- { field: parameter.wire },
- )
- }
- if (provider.wire !== parameter.wire
- || (parameter.codec.mode === 'strict' && provider.wireTypeSymbol !== parameter.codec.typeSymbol)) {
- throw new TypertGatewayError(
- 'provider-mismatch',
- endpoint,
- `lookup provider ${JSON.stringify(key)} does not match its strict definition`,
- { field: parameter.wire },
- )
- }
- let resolved: unknown
- try {
- resolved = await provider.resolve(value)
- } catch (cause) {
- if (cause instanceof TypeRTLookupFailure) throw cause
- throw new TypertGatewayError(
- 'lookup-failed',
- endpoint,
- `lookup provider ${JSON.stringify(key)} failed`,
- { cause, field: parameter.wire },
- )
- }
- if (resolved === undefined) {
- throw new TypertGatewayError(
- 'lookup-not-found',
- endpoint,
- `lookup provider ${JSON.stringify(key)} did not resolve the requested identity`,
- { field: parameter.wire },
- )
- }
- return resolved
- }
- }
- function rpcFailure(error: unknown): ConnectionRpcResult {
- if (error instanceof TypeRTLookupFailure) {
- return { ok: false, error: error.failure as ConnectionRpcError }
- }
- return {
- ok: false,
- error: {
- code: 'internal',
- message: error instanceof Error ? error.message : String(error),
- details: {},
- },
- }
- }
- function endpointOf(namespace: string, method: string): string {
- return `${namespace}/${method}`
- }
- function validateBinding(
- receiver: object,
- serviceKey: string,
- namespace: string,
- endpoint: string,
- ): ResolvedBinding {
- const original = originalOf(receiver)
- const value = Reflect.get(original, 'typertGateway') as unknown
- if (value === undefined) {
- throw new TypertGatewayError(
- 'binding-invalid',
- endpoint,
- `Service ${JSON.stringify(serviceKey)} has no visible typertGateway binding`,
- )
- }
- return {
- binding: readBinding(value, original, serviceKey, endpoint, namespace),
- original,
- }
- }
- function readBinding(
- value: unknown,
- original: object,
- serviceKey: string,
- endpoint: string,
- namespace?: string,
- ): TypeRTGatewayBinding {
- if (!isObject(value)
- || Reflect.get(value, 'service') !== original
- || Reflect.get(value, 'serviceKey') !== serviceKey
- || typeof Reflect.get(value, 'namespace') !== 'string'
- || (namespace !== undefined && Reflect.get(value, 'namespace') !== namespace)) {
- throw new TypertGatewayError(
- 'binding-invalid',
- endpoint,
- `Service ${JSON.stringify(serviceKey)} has an inconsistent typertGateway binding`,
- )
- }
- return value as unknown as TypeRTGatewayBinding
- }
- function originalOf(receiver: object): object {
- const original = Reflect.get(receiver, symbols.original) as unknown
- return isObject(original) ? original : receiver
- }
- function methodParameterNames(service: object, method: string, endpoint: string): readonly string[] {
- let prototype: object | null = Object.getPrototypeOf(service) as object | null
- let implementation: ((this: object, ...args: never[]) => unknown) | undefined
- while (prototype !== null) {
- const descriptor = Object.getOwnPropertyDescriptor(prototype, method)
- if (descriptor !== undefined) {
- if ('value' in descriptor && typeof descriptor.value === 'function') {
- implementation = descriptor.value as (this: object, ...args: never[]) => unknown
- }
- break
- }
- prototype = Object.getPrototypeOf(prototype) as object | null
- }
- if (implementation === undefined) {
- throw new TypertGatewayError(
- 'method-unavailable',
- endpoint,
- `Remote marker has no prototype method ${JSON.stringify(method)}`,
- )
- }
- const source = Function.prototype.toString.call(implementation)
- const open = source.indexOf('(')
- const close = source.indexOf(')', open + 1)
- /* v8 ignore next -- standard public class-method syntax always contains a parenthesized parameter list. */
- if (open < 0 || close < 0) return invalidSignature(endpoint, method)
- const body = source.slice(open + 1, close).trim()
- if (body.length === 0) return []
- const parts = body.split(',').map(part => part.trim())
- const names = new Set<string>()
- for (const part of parts) {
- if (!/^[$A-Z_a-z][$\w]*$/u.test(part) || names.has(part)) return invalidSignature(endpoint, method)
- names.add(part)
- }
- return [...names]
- }
- function invalidSignature(endpoint: string, method: string): never {
- throw new TypertGatewayError(
- 'signature-invalid',
- endpoint,
- `SRC method ${JSON.stringify(method)} must use unique identifier parameters without destructuring, defaults, or rest`,
- )
- }
- function assertExactArguments(
- args: Readonly<Record<string, unknown>>,
- descriptor: InvocationDescriptor,
- endpoint: string,
- ): void {
- if (!isPlainObject(args)) {
- throw new TypertGatewayError('arguments-invalid', endpoint, 'args must be a plain object')
- }
- const expected = new Set(descriptor.parameters.map(parameter => parameter.wire))
- if (descriptor.invocation.kind === 'context') expected.add(descriptor.invocation.wire)
- const actual = Reflect.ownKeys(args)
- const extra = actual.filter(key => typeof key !== 'string' || !expected.has(key))
- const missing = [...expected].filter(key => !Object.hasOwn(args, key))
- if (extra.length === 0 && missing.length === 0) return
- const clauses: string[] = []
- if (missing.length > 0) clauses.push(`missing ${missing.map(key => JSON.stringify(key)).join(', ')}`)
- if (extra.length > 0) clauses.push(`unexpected ${extra.map(key => JSON.stringify(String(key))).join(', ')}`)
- throw new TypertGatewayError('arguments-invalid', endpoint, `args fields do not match the descriptor: ${clauses.join('; ')}`)
- }
- function decode(
- codec: TypeRTCodec,
- value: unknown,
- code: 'input-invalid' | 'result-invalid',
- endpoint: string,
- field: string,
- ): unknown {
- try {
- if (codec.mode === 'strict') value = codec.schema.parse(value)
- assertJsonValue(value, new Set())
- return value
- } catch (cause) {
- throw new TypertGatewayError(
- code,
- endpoint,
- code === 'input-invalid'
- ? `wire field ${JSON.stringify(field)} failed boundary validation`
- : 'business result failed boundary validation',
- { cause, field },
- )
- }
- }
- function assertJsonValue(value: unknown, ancestors: Set<object>): void {
- if (value === null || typeof value === 'string' || typeof value === 'boolean') return
- if (typeof value === 'number') {
- if (Number.isFinite(value)) return
- throw new TypeError('non-finite number is not JSON-safe')
- }
- if (!isObject(value)) throw new TypeError(`${typeof value} is not JSON-safe`)
- if (ancestors.has(value)) throw new TypeError('cyclic value is not JSON-safe')
- ancestors.add(value)
- try {
- if (Array.isArray(value)) {
- if (Object.getOwnPropertySymbols(value).length > 0 || Object.keys(value).length !== value.length) {
- throw new TypeError('sparse or decorated array is not JSON-safe')
- }
- for (let index = 0; index < value.length; index += 1) {
- if (!Object.hasOwn(value, index)) throw new TypeError('sparse array is not JSON-safe')
- assertJsonValue(value[index], ancestors)
- }
- return
- }
- if (!isPlainObject(value)) throw new TypeError('non-plain object is not JSON-safe')
- if (Object.getOwnPropertySymbols(value).length > 0) throw new TypeError('symbol property is not JSON-safe')
- for (const key of Reflect.ownKeys(value)) {
- const descriptor = Object.getOwnPropertyDescriptor(value, key)
- /* v8 ignore next -- ownKeys() just returned this key; only a hostile same-process Proxy can delete it between operations. */
- if (descriptor === undefined || !descriptor.enumerable || !('value' in descriptor)) {
- throw new TypeError('non-data property is not JSON-safe')
- }
- assertJsonValue(descriptor.value, ancestors)
- }
- } finally {
- ancestors.delete(value)
- }
- }
- function isPlainObject(value: object): value is Record<string, unknown> {
- if (Array.isArray(value)) return false
- const prototype = Object.getPrototypeOf(value) as object | null
- return prototype === null || prototype === Object.prototype
- }
- function isObject(value: unknown): value is object {
- return (typeof value === 'object' && value !== null) || typeof value === 'function'
- }
- export default TypertGatewayService
|