index.ts 44 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204
  1. /**
  2. * Live Typert Remote dispatch over Cordis Services and registered providers.
  3. * Unary transport and response envelopes belong to Connection; live Remote
  4. * streams use the Gateway-owned WebSocket mux.
  5. * @module @deepseek-ai/dsh-api-gateway
  6. */
  7. import { randomUUID } from 'node:crypto'
  8. import { Context, Service, symbols } from '@deepseek-ai/cordis'
  9. import type { ConnectionRpcHandler } from '@deepseek-ai/dsh-client-connection'
  10. import { Deque } from '@deepseek-ai/dsh-deque'
  11. import type { WebUpgradeRoute } from '@deepseek-ai/dsh-host-webserver'
  12. import { MAX_TIMER_DELAY_MS } from '@deepseek-ai/dsh-timeout'
  13. import z from '@deepseek-ai/schemastery'
  14. export type { TypertGatewayFaultDetails } from './remote-error-codes.ts'
  15. import {
  16. RemoteError,
  17. remoteErrorOf,
  18. remoteMethods,
  19. type InvocationDescriptor,
  20. type InvocationParameterDescriptor,
  21. type TypertCodec,
  22. type TypertGatewayBinding,
  23. } from '@deepseek-ai/dsh-typert-protocol'
  24. import type {
  25. InvokeRemoteRequest,
  26. TypertGateway,
  27. TypertGatewayErrorCode,
  28. TypertGatewayWireStream,
  29. TypertRemoteEventDispatch,
  30. TypertRemoteEventFrame,
  31. TypertRemoteEventInvocation,
  32. TypertRemoteEventOutcome,
  33. TypertRemoteEventSource,
  34. } from './types.ts'
  35. import {
  36. RemoteStreamMuxServer,
  37. rejectRemoteStreamUpgrade,
  38. } from './stream-server.ts'
  39. import {
  40. REMOTE_EVENT_STREAM_ENDPOINT,
  41. REMOTE_EVENT_STREAM_READY,
  42. REMOTE_EVENT_RESULT_ENDPOINT,
  43. REMOTE_STREAM_MUX_PATH,
  44. isRemoteEventAgentId,
  45. isRemoteJsonValue,
  46. parseRemoteEventResult,
  47. projectRemoteEventRequest,
  48. restoreRemoteEventRejection,
  49. type RemoteEventCancellationFrame,
  50. type RemoteEventClientId,
  51. type RemoteEventEmitFrame,
  52. type RemoteEventHostInfo,
  53. type RemoteEventId,
  54. type RemoteEventInvocationFrame,
  55. type RemoteEventReadyFrame,
  56. type RemoteStreamFailure,
  57. } from './stream-protocol.ts'
  58. export type {
  59. InvokeRemoteRequest,
  60. TypertGateway,
  61. TypertGatewayErrorCode,
  62. TypertGatewayWireStream,
  63. TypertRemoteEventContext,
  64. TypertRemoteEventDispatch,
  65. TypertRemoteEventFrame,
  66. TypertRemoteEventInvocation,
  67. TypertRemoteEventOutcome,
  68. TypertRemoteEventSource,
  69. } from './types.ts'
  70. export type { RemoteEventHostInfo } from './stream-protocol.ts'
  71. interface GatewayErrorOptions {
  72. readonly cause?: unknown
  73. readonly field?: string
  74. }
  75. interface ResolvedBinding {
  76. readonly binding: TypertGatewayBinding
  77. readonly original: object
  78. }
  79. interface PreparedInvocation {
  80. readonly endpoint: string
  81. readonly descriptor: InvocationDescriptor
  82. readonly receiver: object
  83. readonly args: readonly unknown[]
  84. readonly method: (...args: never[]) => unknown
  85. }
  86. interface RegisteredRemoteEventSource {
  87. readonly lifetime: AbortController
  88. readonly done: Promise<void>
  89. readonly host: RemoteEventHostInfo
  90. }
  91. interface RemoteEventClient {
  92. readonly id: RemoteEventClientId
  93. readonly queue: RemoteEventQueue
  94. readonly deliveries: Map<RemoteEventId, PendingRemoteEvent>
  95. }
  96. interface PendingRemoteEvent {
  97. readonly id: RemoteEventId
  98. readonly source: TypertRemoteEventInvocation
  99. readonly frame: RemoteEventInvocationFrame
  100. readonly deliveries: Set<RemoteEventClient>
  101. releaseContext: () => void
  102. releaseSignal: () => void
  103. }
  104. type ConnectionRpcResult = Awaited<ReturnType<ConnectionRpcHandler>>
  105. type ConnectionRpcError = Extract<ConnectionRpcResult, { readonly ok: false }>['error']
  106. const NEVER_ABORTED_SIGNAL = new AbortController().signal
  107. const DEFAULT_WEBSOCKET_HEARTBEAT_INTERVAL_MS = 2_000
  108. /** Gateway transport configuration. */
  109. export interface Config {
  110. /** WebSocket Ping interval from 1 through 2,147,483,647 milliseconds. @default 2000 */
  111. readonly websocketHeartbeatIntervalMs?: number
  112. }
  113. interface ResolvedConfig extends Config {
  114. readonly websocketHeartbeatIntervalMs: number
  115. }
  116. /**
  117. * Dispatch failure produced outside the invoked business method. Rides the
  118. * shared Remote failure vocabulary, so its code crosses the wire instead of
  119. * folding to `internal`.
  120. */
  121. export class TypertGatewayError extends RemoteError<TypertGatewayErrorCode> {
  122. /** Canonical `<namespace>/<method>` endpoint. */
  123. readonly endpoint: string
  124. /** Affected wire field when the failure is field-specific. */
  125. readonly field: string | undefined
  126. /**
  127. * Construct a Gateway failure without embedding boundary values in its message.
  128. * @param code - stable failure category.
  129. * @param endpoint - canonical Remote endpoint.
  130. * @param message - correction-oriented diagnostic without sensitive values.
  131. * @param options - optional field and contained cause.
  132. */
  133. constructor(
  134. code: TypertGatewayErrorCode,
  135. endpoint: string,
  136. message: string,
  137. options: GatewayErrorOptions = {},
  138. ) {
  139. super(
  140. code,
  141. `typert gateway: ${endpoint}: ${message}`,
  142. { endpoint, ...options.field === undefined ? {} : { field: options.field } },
  143. options.cause === undefined ? undefined : { cause: options.cause },
  144. )
  145. this.name = 'TypertGatewayError'
  146. this.endpoint = endpoint
  147. this.field = options.field
  148. }
  149. }
  150. /**
  151. * Resolve strict generated definitions or conservative SRC markers against
  152. * current Cordis Services and Typert providers.
  153. * @typert service typertGateway
  154. */
  155. export class TypertGatewayService extends Service implements TypertGateway {
  156. static inject = ['typert']
  157. static Config: z<Config> = z.object({
  158. websocketHeartbeatIntervalMs: z.number().step(1).min(1).max(MAX_TIMER_DELAY_MS)
  159. .default(DEFAULT_WEBSOCKET_HEARTBEAT_INTERVAL_MS),
  160. })
  161. /** Carrier adapter shared by the WebSocket mux and local Host transports. */
  162. readonly wireStream: TypertGatewayWireStream = {
  163. open: (endpoint, payload, signal) => this.openWireStream(endpoint, payload, signal),
  164. failure: error => rpcError(error),
  165. }
  166. private srcClaims: ReadonlySet<string> | undefined
  167. private remoteEvents: RegisteredRemoteEventSource | undefined
  168. private readonly remoteEventClients = new Map<RemoteEventClientId, RemoteEventClient>()
  169. private readonly pendingRemoteEvents = new Map<RemoteEventId, PendingRemoteEvent>()
  170. /**
  171. * Register the Gateway against the active Typert registry.
  172. * @param ctx - owning Host Context with Typert registry access.
  173. * @param config - validated Gateway transport configuration.
  174. */
  175. constructor(ctx: Context, config: Config) {
  176. super(ctx, 'typertGateway')
  177. const resolved = config as ResolvedConfig
  178. ctx.on('internal/service', () => {
  179. this.srcClaims = undefined
  180. })
  181. ctx.inject(['connection'], (connectionCtx) => {
  182. connectionCtx.connection.rpc.intercept(
  183. '/api',
  184. endpoint => this.claimsEndpoint(endpoint),
  185. (endpoint, payload, signal) => this.dispatchRpc(endpoint, payload, signal),
  186. )
  187. })
  188. ctx.inject(['connection', 'webServer'], (webCtx) => {
  189. const mux = new RemoteStreamMuxServer(
  190. (endpoint, payload, signal) => this.openWireStream(endpoint, payload, signal),
  191. this.wireStream.failure,
  192. resolved.websocketHeartbeatIntervalMs,
  193. )
  194. webCtx.effect(() => {
  195. const route: WebUpgradeRoute = {
  196. path: REMOTE_STREAM_MUX_PATH,
  197. handler: (req, socket, head) => {
  198. const rejection = webCtx.connection.requestRejection(req)
  199. if (rejection !== undefined) {
  200. rejectRemoteStreamUpgrade(socket, rejection)
  201. return
  202. }
  203. mux.handleUpgrade(req, socket, head)
  204. },
  205. }
  206. const unregister = webCtx.webServer.registerUpgrade(route)
  207. return async () => {
  208. unregister()
  209. await mux.close()
  210. }
  211. }, `api-gateway: ${REMOTE_STREAM_MUX_PATH} WebSocket`)
  212. })
  213. }
  214. /**
  215. * Register the sole application-selected forwarded-event source.
  216. * @param source - stream factory installed by the Remote assembly.
  217. * @param host - stable Host facts included in each Client generation's opening frame.
  218. * @returns disposer removing this source and cancelling its active streams.
  219. */
  220. registerRemoteEvents(
  221. source: TypertRemoteEventSource,
  222. host: RemoteEventHostInfo,
  223. ): () => Promise<void> {
  224. if (this.remoteEvents !== undefined) {
  225. throw new Error('typert gateway: forwarded Remote event source is already registered')
  226. }
  227. const lifetime = new AbortController()
  228. const stream = source(lifetime.signal)
  229. const done = this.consumeRemoteEvents(stream, lifetime.signal).catch((error: unknown) => {
  230. if (this.remoteEvents?.lifetime !== lifetime || lifetime.signal.aborted) return
  231. this.closeRemoteEvents(error)
  232. this.remoteEvents = undefined
  233. lifetime.abort(error)
  234. })
  235. const registration: RegisteredRemoteEventSource = { lifetime, done, host: { home: host.home } }
  236. this.remoteEvents = registration
  237. return async () => {
  238. if (this.remoteEvents === registration) {
  239. this.remoteEvents = undefined
  240. const error = new Error('typert gateway: forwarded Remote event source was removed')
  241. registration.lifetime.abort(error)
  242. this.closeRemoteEvents(error)
  243. }
  244. await registration.done
  245. }
  246. }
  247. private claimsEndpoint(endpoint: string): boolean {
  248. if (endpoint === REMOTE_EVENT_RESULT_ENDPOINT) return true
  249. const segments = endpoint.split('/')
  250. if (segments.length !== 2 || segments[0] === '' || segments[1] === '') return false
  251. if (this.ctx.typert.local.get(endpoint) !== undefined || this.ctx.typert.local.hasSeen(endpoint)) return true
  252. this.srcClaims ??= this.collectSrcClaims()
  253. return this.srcClaims.has(endpoint)
  254. }
  255. private collectSrcClaims(): ReadonlySet<string> {
  256. const claims = new Set<string>()
  257. for (const [serviceKey, definition] of Object.entries(this.ctx.reflect.props)) {
  258. if (definition.type !== 'service') continue
  259. const receiver = this.ctx.get(serviceKey) as unknown
  260. if (!isObject(receiver)) continue
  261. const original = originalOf(receiver)
  262. const binding = Reflect.get(original, 'typertRemote') as unknown
  263. if (!isObject(binding) || typeof Reflect.get(binding, 'namespace') !== 'string') continue
  264. const namespace = Reflect.get(binding, 'namespace') as string
  265. for (const candidate of remoteMethods(original)) {
  266. claims.add(endpointOf(namespace, candidate.exportName ?? candidate.method))
  267. }
  268. }
  269. return claims
  270. }
  271. /**
  272. * Invoke one live Remote method through strict generated reflection or SRC markers.
  273. * @param request - decoded endpoint and exact named wire arguments.
  274. * @returns the business result without output decoding.
  275. * @throws {@link TypertGatewayError} for dispatch, provider, or boundary failures; lookup-policy and business errors retain identity.
  276. */
  277. async invoke(request: InvokeRemoteRequest): Promise<unknown> {
  278. const prepared = await this.prepareInvocation(request)
  279. if (prepared.descriptor.mode === 'stream') {
  280. throw new TypertGatewayError(
  281. 'gateway/signature-invalid',
  282. prepared.endpoint,
  283. 'stream Remote methods must be opened through the stream carrier',
  284. )
  285. }
  286. try {
  287. return await Reflect.apply(prepared.method, prepared.receiver, prepared.args) as unknown
  288. } catch (error) {
  289. if (request.signal?.aborted === true) throw remoteCancelled(prepared.endpoint, error)
  290. throw error
  291. }
  292. }
  293. /**
  294. * Open one live stream Remote method without assuming a physical carrier.
  295. * @param request - decoded endpoint and named wire arguments.
  296. * @returns a cancellation-aware iterable over the business results.
  297. */
  298. async stream(request: InvokeRemoteRequest): Promise<AsyncIterable<unknown>> {
  299. const prepared = await this.prepareInvocation(request)
  300. if (prepared.descriptor.mode !== 'stream') {
  301. throw new TypertGatewayError(
  302. 'gateway/signature-invalid',
  303. prepared.endpoint,
  304. 'unary Remote methods cannot be opened through the stream carrier',
  305. )
  306. }
  307. let source: unknown
  308. try {
  309. source = Reflect.apply(prepared.method, prepared.receiver, prepared.args) as unknown
  310. } catch (error) {
  311. if (request.signal?.aborted === true) throw remoteCancelled(prepared.endpoint, error)
  312. throw error
  313. }
  314. if (!isIterable(source)) {
  315. throw new TypertGatewayError(
  316. 'gateway/result-invalid',
  317. prepared.endpoint,
  318. 'stream Remote method did not return Iterable or AsyncIterable',
  319. { field: 'result' },
  320. )
  321. }
  322. return cancellableStream(
  323. source,
  324. prepared.endpoint,
  325. request.signal ?? NEVER_ABORTED_SIGNAL,
  326. )
  327. }
  328. private async dispatchRpc(
  329. endpoint: string,
  330. payload: unknown,
  331. signal: AbortSignal,
  332. ): Promise<ConnectionRpcResult> {
  333. if (endpoint === REMOTE_EVENT_RESULT_ENDPOINT) {
  334. try {
  335. const result = parseRemoteEventResultPayload(payload)
  336. const client = this.remoteEventClients.get(result.clientId)
  337. if (client === undefined) {
  338. throw new Error('typert gateway: Remote event result identifies no active event stream')
  339. }
  340. this.receiveRemoteEventResult(client, result)
  341. return { ok: true, value: undefined }
  342. } catch (error) {
  343. return rpcFailure(error)
  344. }
  345. }
  346. return this.invokeRpc(endpoint, payload, signal)
  347. }
  348. private async openWireStream(
  349. endpoint: string,
  350. payload: unknown,
  351. signal: AbortSignal,
  352. ): Promise<AsyncIterable<unknown>> {
  353. if (endpoint === REMOTE_EVENT_STREAM_ENDPOINT) {
  354. return this.openRemoteEvents(payload, signal)
  355. }
  356. return this.stream(remoteRequest(endpoint, payload, signal))
  357. }
  358. private async *openRemoteEvents(
  359. payload: unknown,
  360. signal: AbortSignal,
  361. ): AsyncGenerator<
  362. RemoteEventEmitFrame | RemoteEventInvocationFrame | RemoteEventCancellationFrame
  363. | RemoteEventReadyFrame
  364. > {
  365. if (!isObject(payload)
  366. || !isPlainObject(payload)
  367. || Reflect.ownKeys(payload).length !== 1
  368. || !Object.hasOwn(payload, 'args')
  369. || !isObject(payload.args)
  370. || !isPlainObject(payload.args)
  371. || Reflect.ownKeys(payload.args).length !== 0) {
  372. throw new TypertGatewayError(
  373. 'gateway/arguments-invalid',
  374. REMOTE_EVENT_STREAM_ENDPOINT,
  375. 'forwarded Remote event stream requires an empty args object',
  376. )
  377. }
  378. const registration = this.remoteEvents
  379. if (registration === undefined) {
  380. throw new TypertGatewayError(
  381. 'gateway/service-unavailable',
  382. REMOTE_EVENT_STREAM_ENDPOINT,
  383. 'forwarded Remote event source is unavailable',
  384. )
  385. }
  386. const lifetime = AbortSignal.any([signal, registration.lifetime.signal])
  387. let clientId = randomUUID() as RemoteEventClientId
  388. while (this.remoteEventClients.has(clientId)) clientId = randomUUID() as RemoteEventClientId
  389. const client: RemoteEventClient = {
  390. id: clientId,
  391. queue: new RemoteEventQueue(),
  392. deliveries: new Map(),
  393. }
  394. this.remoteEventClients.set(clientId, client)
  395. for (const pending of this.pendingRemoteEvents.values()) this.deliverRemoteEvent(pending, client)
  396. try {
  397. yield { ...REMOTE_EVENT_STREAM_READY, clientId, host: registration.host }
  398. yield* client.queue.iterate(lifetime)
  399. } finally {
  400. this.removeRemoteEventClient(client)
  401. }
  402. }
  403. private async consumeRemoteEvents(
  404. source: AsyncIterable<TypertRemoteEventDispatch>,
  405. signal: AbortSignal,
  406. ): Promise<void> {
  407. for await (const dispatch of source) {
  408. if (signal.aborted) {
  409. if ('context' in dispatch) dispatch.reject(signal.reason)
  410. return
  411. }
  412. if ('context' in dispatch) this.startRemoteEvent(dispatch)
  413. else this.broadcastRemoteEvent(dispatch)
  414. }
  415. if (!signal.aborted) {
  416. throw new Error('typert gateway: forwarded Remote event source ended unexpectedly')
  417. }
  418. }
  419. private broadcastRemoteEvent(frame: TypertRemoteEventFrame): void {
  420. assertRemoteEventFrame(frame)
  421. const wire: RemoteEventEmitFrame = {
  422. type: 'emit',
  423. event: frame.event,
  424. args: frame.args,
  425. }
  426. for (const client of this.remoteEventClients.values()) client.queue.push(wire)
  427. }
  428. private startRemoteEvent(source: TypertRemoteEventInvocation): void {
  429. try {
  430. assertRemoteEventName(source)
  431. if (!isRemoteEventAgentId(source.context.agentId)) {
  432. throw new TypeError(
  433. 'typert gateway: scoped Remote events require a non-empty Agent identity',
  434. )
  435. }
  436. const projected = projectRemoteEventRequest(source.request, source.context.subject)
  437. let id = randomUUID() as RemoteEventId
  438. while (this.pendingRemoteEvents.has(id)) id = randomUUID() as RemoteEventId
  439. let releaseContext: () => void
  440. try {
  441. const dispose = source.context.value.effect(
  442. () => () => {
  443. this.cancelRemoteEvent(
  444. pending,
  445. new Error('typert gateway: Remote event Agent Context was released'),
  446. )
  447. },
  448. `api-gateway: Remote event ${JSON.stringify(source.event)}`,
  449. )
  450. releaseContext = () => { void dispose() }
  451. } catch {
  452. source.resolve({ kind: 'next' })
  453. return
  454. }
  455. const signals = new Set(projected.signal === undefined ? [] : [projected.signal])
  456. const abort = (): void => {
  457. const reason = [...signals].find(signal => signal.aborted)?.reason as unknown
  458. this.cancelRemoteEvent(pending, reason instanceof Error
  459. ? reason
  460. : new Error('typert gateway: Remote event was cancelled', { cause: reason }))
  461. }
  462. const pending: PendingRemoteEvent = {
  463. id,
  464. source,
  465. frame: {
  466. type: 'waterfall',
  467. event: source.event,
  468. eventId: id,
  469. agentId: source.context.agentId,
  470. request: projected.request,
  471. },
  472. deliveries: new Set(),
  473. releaseContext,
  474. releaseSignal: () => {
  475. for (const signal of signals) signal.removeEventListener('abort', abort)
  476. },
  477. }
  478. this.pendingRemoteEvents.set(id, pending)
  479. for (const signal of signals) signal.addEventListener('abort', abort, { once: true })
  480. if ([...signals].some(signal => signal.aborted)) abort()
  481. else for (const client of this.remoteEventClients.values()) this.deliverRemoteEvent(pending, client)
  482. } catch (error) {
  483. source.reject(error)
  484. }
  485. }
  486. private deliverRemoteEvent(pending: PendingRemoteEvent, client: RemoteEventClient): void {
  487. pending.deliveries.add(client)
  488. client.deliveries.set(pending.id, pending)
  489. client.queue.push(pending.frame)
  490. }
  491. private receiveRemoteEventResult(
  492. client: RemoteEventClient,
  493. result: ReturnType<typeof parseRemoteEventResult>,
  494. ): void {
  495. const pending = this.pendingRemoteEvents.get(result.eventId)
  496. // Settlement and Client replacement may race the result request. Results
  497. // from a completed event or a superseded delivery are idempotent no-ops.
  498. if (pending === undefined || !pending.deliveries.has(client)) return
  499. this.removeRemoteEventDelivery(pending, client)
  500. if (result.outcome.kind === 'result') {
  501. this.settleRemoteEvent(pending, {
  502. kind: 'result',
  503. value: result.outcome.value,
  504. })
  505. } else if (result.outcome.kind === 'rejected') {
  506. this.cancelRemoteEvent(pending, restoreRemoteEventRejection(result.outcome.error))
  507. } else if (pending.deliveries.size === 0) {
  508. this.settleRemoteEvent(pending, { kind: 'next' })
  509. }
  510. }
  511. private removeRemoteEventDelivery(pending: PendingRemoteEvent, client: RemoteEventClient): void {
  512. pending.deliveries.delete(client)
  513. client.deliveries.delete(pending.id)
  514. }
  515. private removeRemoteEventClient(client: RemoteEventClient): void {
  516. this.remoteEventClients.delete(client.id)
  517. for (const pending of [...client.deliveries.values()]) this.removeRemoteEventDelivery(pending, client)
  518. client.queue.end()
  519. }
  520. private settleRemoteEvent(pending: PendingRemoteEvent, outcome: TypertRemoteEventOutcome): void {
  521. this.finishRemoteEvent(pending)
  522. pending.source.resolve(outcome)
  523. }
  524. private cancelRemoteEvent(pending: PendingRemoteEvent, reason: unknown): void {
  525. if (this.pendingRemoteEvents.get(pending.id) !== pending) return
  526. this.finishRemoteEvent(pending)
  527. pending.source.reject(reason)
  528. }
  529. private finishRemoteEvent(pending: PendingRemoteEvent): void {
  530. this.pendingRemoteEvents.delete(pending.id)
  531. pending.releaseSignal()
  532. pending.releaseContext()
  533. const clients = new Set(pending.deliveries)
  534. for (const client of clients) this.removeRemoteEventDelivery(pending, client)
  535. const cancellation: RemoteEventCancellationFrame = {
  536. type: 'cancel',
  537. eventId: pending.id,
  538. }
  539. for (const client of clients) client.queue.push(cancellation)
  540. }
  541. private closeRemoteEvents(reason: unknown): void {
  542. for (const pending of [...this.pendingRemoteEvents.values()]) {
  543. this.cancelRemoteEvent(pending, reason)
  544. }
  545. for (const client of [...this.remoteEventClients.values()]) client.queue.end()
  546. }
  547. private async invokeRpc(endpoint: string, payload: unknown, signal: AbortSignal): Promise<ConnectionRpcResult> {
  548. try {
  549. const value = await this.invoke(remoteRequest(endpoint, payload, signal))
  550. // A void or explicitly absent business result carries no `value` field;
  551. // JSON has no `undefined`, and the envelope's optional slot is the one
  552. // representation of absence that both args and results already use.
  553. return { ok: true, value }
  554. } catch (error) {
  555. return rpcFailure(error)
  556. }
  557. }
  558. private async prepareInvocation(request: InvokeRemoteRequest): Promise<PreparedInvocation> {
  559. const endpoint = endpointOf(request.namespace, request.method)
  560. const descriptor = this.resolveDescriptor(request.namespace, request.method, endpoint)
  561. assertExactArguments(request.args, descriptor, endpoint)
  562. const receiverContext = await this.resolveReceiverContext(descriptor, request.args, endpoint)
  563. const receiver = receiverContext.get(descriptor.service) as unknown
  564. if (!isObject(receiver)) {
  565. throw new TypertGatewayError(
  566. 'gateway/service-unavailable',
  567. endpoint,
  568. `active Service ${JSON.stringify(descriptor.service)} is unavailable`,
  569. )
  570. }
  571. validateBinding(receiver, descriptor.service, descriptor.namespace, endpoint)
  572. const args = await Promise.all(descriptor.parameters.map(parameter =>
  573. this.resolveParameter(parameter, request.args, endpoint)))
  574. if (descriptor.cancellation !== undefined) args.push(request.signal ?? NEVER_ABORTED_SIGNAL)
  575. const implementation = descriptor.implementation ?? descriptor.method
  576. const method = Reflect.get(receiver, implementation) as unknown
  577. if (typeof method !== 'function') {
  578. throw new TypertGatewayError(
  579. 'gateway/method-unavailable',
  580. endpoint,
  581. `active Service ${JSON.stringify(descriptor.service)} has no callable method ${JSON.stringify(implementation)}`,
  582. )
  583. }
  584. return { endpoint, descriptor, receiver, args, method: method as (...args: never[]) => unknown }
  585. }
  586. private resolveDescriptor(namespace: string, method: string, endpoint: string): InvocationDescriptor {
  587. const strict = this.ctx.typert.local.get(endpoint)
  588. if (strict !== undefined) return strict
  589. if (this.ctx.typert.local.hasSeen(endpoint)) {
  590. throw new TypertGatewayError(
  591. 'gateway/definition-unavailable',
  592. endpoint,
  593. 'its strict definition was withdrawn and SRC fallback is forbidden',
  594. )
  595. }
  596. return this.resolveSrcDescriptor(namespace, method, endpoint)
  597. }
  598. private resolveSrcDescriptor(namespace: string, method: string, endpoint: string): InvocationDescriptor {
  599. const candidates: InvocationDescriptor[] = []
  600. for (const [serviceKey, definition] of Object.entries(this.ctx.reflect.props)) {
  601. if (definition.type !== 'service') continue
  602. const receiver = this.ctx.get(serviceKey) as unknown
  603. if (!isObject(receiver)) continue
  604. const original = originalOf(receiver)
  605. const value = Reflect.get(original, 'typertRemote') as unknown
  606. if (value === undefined) continue
  607. const binding = readBinding(value, original, serviceKey, endpoint)
  608. if (binding.namespace !== namespace) continue
  609. const marker = remoteMethods(original).find(candidate => (candidate.exportName ?? candidate.method) === method)
  610. if (marker === undefined) continue
  611. candidates.push(this.srcDescriptor(binding, marker, method, endpoint))
  612. }
  613. if (candidates.length === 0) {
  614. throw new TypertGatewayError('gateway/invocation-unavailable', endpoint, 'no active Remote method exports this endpoint')
  615. }
  616. if (candidates.length > 1) {
  617. throw new TypertGatewayError(
  618. 'gateway/ambiguous-endpoint',
  619. endpoint,
  620. `multiple active Services export this endpoint: ${candidates.map(candidate => candidate.service).sort().join(', ')}`,
  621. )
  622. }
  623. return candidates[0] as InvocationDescriptor
  624. }
  625. private srcDescriptor(
  626. binding: TypertGatewayBinding,
  627. marker: ReturnType<typeof remoteMethods>[number],
  628. method: string,
  629. endpoint: string,
  630. ): InvocationDescriptor {
  631. const names = methodParameterNames(binding.service, marker.method, endpoint)
  632. const signalIndex = names.indexOf('signal')
  633. if (signalIndex >= 0 && signalIndex !== names.length - 1) {
  634. throw new TypertGatewayError(
  635. 'gateway/signature-invalid',
  636. endpoint,
  637. 'SRC cancellation parameter signal must be the final parameter',
  638. { field: 'signal' },
  639. )
  640. }
  641. const cancellation = signalIndex >= 0
  642. ? { parameter: 'signal' as const }
  643. : undefined
  644. const businessNames = cancellation === undefined ? names : names.slice(0, -1)
  645. const parameters: InvocationParameterDescriptor[] = []
  646. const wires = new Set<string>()
  647. for (const name of businessNames) {
  648. const matches = this.ctx.typert.lookups.definitions()
  649. .filter(definition => definition.parameter === name)
  650. if (matches.length > 1) {
  651. throw new TypertGatewayError(
  652. 'gateway/signature-invalid',
  653. endpoint,
  654. `parameter ${JSON.stringify(name)} matches multiple lookup providers`,
  655. { field: name },
  656. )
  657. }
  658. const match = matches[0]
  659. const parameter: InvocationParameterDescriptor = match === undefined
  660. ? { name, wire: name, source: 'json', codec: { mode: 'src-json' } }
  661. : {
  662. name,
  663. wire: match.wire,
  664. source: 'lookup',
  665. lookup: match.key,
  666. codec: { mode: 'src-json' },
  667. }
  668. if (wires.has(parameter.wire)) {
  669. throw new TypertGatewayError(
  670. 'gateway/signature-invalid',
  671. endpoint,
  672. `multiple parameters use wire field ${JSON.stringify(parameter.wire)}`,
  673. { field: parameter.wire },
  674. )
  675. }
  676. wires.add(parameter.wire)
  677. parameters.push(parameter)
  678. }
  679. let receiver: InvocationDescriptor['invocation'] = { kind: 'direct' }
  680. if (marker.invocation.kind === 'context') {
  681. const provider = this.ctx.typert.contexts.getHost(marker.invocation.context)
  682. if (provider === undefined) {
  683. throw new TypertGatewayError(
  684. 'gateway/context-unavailable',
  685. endpoint,
  686. `Context provider ${JSON.stringify(marker.invocation.context)} is unavailable`,
  687. )
  688. }
  689. if (wires.has(provider.wire)) {
  690. throw new TypertGatewayError(
  691. 'gateway/signature-invalid',
  692. endpoint,
  693. `Context identity conflicts with wire field ${JSON.stringify(provider.wire)}`,
  694. { field: provider.wire },
  695. )
  696. }
  697. receiver = {
  698. kind: 'context',
  699. context: marker.invocation.context,
  700. wire: provider.wire,
  701. codec: { mode: 'src-json' },
  702. }
  703. }
  704. return {
  705. id: `src:${binding.serviceKey}#${endpoint}`,
  706. service: binding.serviceKey,
  707. namespace: binding.namespace,
  708. method,
  709. ...(marker.method === method ? {} : { implementation: marker.method }),
  710. ...(marker.mode === undefined ? {} : { mode: marker.mode }),
  711. invocation: receiver,
  712. parameters,
  713. ...(cancellation === undefined ? {} : { cancellation }),
  714. result: { mode: 'src-json' },
  715. }
  716. }
  717. private async resolveReceiverContext(
  718. descriptor: InvocationDescriptor,
  719. args: Readonly<Record<string, unknown>>,
  720. endpoint: string,
  721. ): Promise<Context> {
  722. if (descriptor.invocation.kind === 'direct') return this.ctx
  723. const invocation = descriptor.invocation
  724. const provider = this.ctx.typert.contexts.getHost(invocation.context)
  725. if (provider === undefined) {
  726. throw new TypertGatewayError(
  727. 'gateway/context-unavailable',
  728. endpoint,
  729. `Context provider ${JSON.stringify(invocation.context)} is unavailable`,
  730. )
  731. }
  732. if (provider.wire !== invocation.wire
  733. || (invocation.codec.mode === 'strict' && provider.wireTypeSymbol !== invocation.codec.typeSymbol)) {
  734. throw new TypertGatewayError(
  735. 'gateway/provider-mismatch',
  736. endpoint,
  737. `Context provider ${JSON.stringify(invocation.context)} does not match its strict definition`,
  738. { field: invocation.wire },
  739. )
  740. }
  741. const identity = decode(invocation.codec, args[invocation.wire], endpoint, invocation.wire)
  742. let context: Context | undefined
  743. try {
  744. context = await provider.resolve(identity)
  745. } catch (cause) {
  746. if (remoteErrorOf(cause) !== undefined) throw cause
  747. throw new TypertGatewayError(
  748. 'gateway/context-failed',
  749. endpoint,
  750. `Context provider ${JSON.stringify(invocation.context)} failed`,
  751. { cause, field: invocation.wire },
  752. )
  753. }
  754. if (context === undefined) {
  755. throw new TypertGatewayError(
  756. 'gateway/context-not-found',
  757. endpoint,
  758. `Context provider ${JSON.stringify(invocation.context)} did not resolve the requested identity`,
  759. { field: invocation.wire },
  760. )
  761. }
  762. return context
  763. }
  764. private async resolveParameter(
  765. parameter: InvocationParameterDescriptor,
  766. args: Readonly<Record<string, unknown>>,
  767. endpoint: string,
  768. ): Promise<unknown> {
  769. // An absent field reached assertExactArguments' allowance, so this parameter
  770. // takes undefined; a present-but-undefined field is not JSON-safe input and
  771. // still fails decode. Lookup ids are never omissible, so absence here only
  772. // ever belongs to a json parameter.
  773. if (!Object.hasOwn(args, parameter.wire)) return undefined
  774. const value = decode(parameter.codec, args[parameter.wire], endpoint, parameter.wire)
  775. if (parameter.source === 'json') return value
  776. const key = parameter.lookup
  777. /* v8 ignore next -- registry validation rejects strict descriptors without a key, and SRC derivation always supplies one. */
  778. if (key === undefined) {
  779. throw new TypertGatewayError(
  780. 'gateway/lookup-unavailable',
  781. endpoint,
  782. `lookup parameter ${JSON.stringify(parameter.name)} has no provider key`,
  783. { field: parameter.wire },
  784. )
  785. }
  786. const provider = this.ctx.typert.lookups.get(key)
  787. if (provider === undefined) {
  788. throw new TypertGatewayError(
  789. 'gateway/lookup-unavailable',
  790. endpoint,
  791. `lookup provider ${JSON.stringify(key)} is unavailable`,
  792. { field: parameter.wire },
  793. )
  794. }
  795. if (provider.wire !== parameter.wire
  796. || (parameter.codec.mode === 'strict' && provider.wireTypeSymbol !== parameter.codec.typeSymbol)) {
  797. throw new TypertGatewayError(
  798. 'gateway/provider-mismatch',
  799. endpoint,
  800. `lookup provider ${JSON.stringify(key)} does not match its strict definition`,
  801. { field: parameter.wire },
  802. )
  803. }
  804. let resolved: unknown
  805. try {
  806. resolved = await provider.resolve(value)
  807. } catch (cause) {
  808. if (remoteErrorOf(cause) !== undefined) throw cause
  809. throw new TypertGatewayError(
  810. 'gateway/lookup-failed',
  811. endpoint,
  812. `lookup provider ${JSON.stringify(key)} failed`,
  813. { cause, field: parameter.wire },
  814. )
  815. }
  816. if (resolved === undefined) {
  817. throw new TypertGatewayError(
  818. 'gateway/lookup-not-found',
  819. endpoint,
  820. `lookup provider ${JSON.stringify(key)} did not resolve the requested identity`,
  821. { field: parameter.wire },
  822. )
  823. }
  824. return resolved
  825. }
  826. }
  827. type RemoteEventWireFrame =
  828. | RemoteEventEmitFrame
  829. | RemoteEventInvocationFrame
  830. | RemoteEventCancellationFrame
  831. /** Pull-driven queue owned by one connected Client event generation. */
  832. class RemoteEventQueue {
  833. private readonly frames = new Deque<RemoteEventWireFrame>()
  834. private waiter: (() => void) | undefined
  835. private closed = false
  836. push(frame: RemoteEventWireFrame): void {
  837. if (this.closed) return
  838. this.frames.pushBack(frame)
  839. this.waiter?.()
  840. }
  841. end(): void {
  842. if (this.closed) return
  843. this.closed = true
  844. this.waiter?.()
  845. }
  846. async *iterate(signal: AbortSignal): AsyncGenerator<RemoteEventWireFrame> {
  847. const abort = (): void => { this.end() }
  848. signal.addEventListener('abort', abort, { once: true })
  849. try {
  850. while (true) {
  851. while (this.frames.size > 0) yield this.frames.popFront() as RemoteEventWireFrame
  852. if (this.closed || signal.aborted) return
  853. await new Promise<void>((resolve) => { this.waiter = resolve })
  854. this.waiter = undefined
  855. }
  856. } finally {
  857. signal.removeEventListener('abort', abort)
  858. }
  859. }
  860. }
  861. function assertRemoteEventFrame(frame: TypertRemoteEventFrame): void {
  862. assertRemoteEventName(frame)
  863. if (!Array.isArray(frame.args) || !isRemoteJsonValue(frame.args)) {
  864. throw new TypeError(`typert gateway: Remote event ${JSON.stringify(frame.event)} arguments are not lossless JSON data`)
  865. }
  866. }
  867. function assertRemoteEventName(frame: { readonly event: unknown }): void {
  868. if (typeof frame.event !== 'string' || frame.event.length === 0) {
  869. throw new TypeError('typert gateway: Remote event name must be a nonempty string')
  870. }
  871. }
  872. function parseRemoteEventResultPayload(payload: unknown): ReturnType<typeof parseRemoteEventResult> {
  873. if (!isObject(payload)
  874. || !isPlainObject(payload)
  875. || Reflect.ownKeys(payload).length !== 1
  876. || !Object.hasOwn(payload, 'args')) {
  877. throw new Error('typert gateway: Remote event result requires exactly one plain-object args field')
  878. }
  879. return parseRemoteEventResult(payload.args)
  880. }
  881. function remoteRequest(endpoint: string, payload: unknown, signal: AbortSignal): InvokeRemoteRequest {
  882. const segments = endpoint.split('/')
  883. if (segments.length !== 2 || segments[0] === '' || segments[1] === '') {
  884. throw new Error(`invalid Remote endpoint ${JSON.stringify(endpoint)}`)
  885. }
  886. const [namespace, method] = segments as [string, string]
  887. if (!isObject(payload)
  888. || !isPlainObject(payload)
  889. || Reflect.ownKeys(payload).length !== 1
  890. || !Object.hasOwn(payload, 'args')
  891. || !isObject(payload.args)
  892. || !isPlainObject(payload.args)) {
  893. throw new Error('Remote payload must contain exactly one plain-object args field')
  894. }
  895. return { namespace, method, args: payload.args, signal }
  896. }
  897. function isIterable(value: unknown): value is Iterable<unknown> | AsyncIterable<unknown> {
  898. return isObject(value)
  899. && (typeof Reflect.get(value, Symbol.iterator) === 'function'
  900. || typeof Reflect.get(value, Symbol.asyncIterator) === 'function')
  901. }
  902. async function *cancellableStream(
  903. source: Iterable<unknown> | AsyncIterable<unknown>,
  904. endpoint: string,
  905. signal: AbortSignal,
  906. ): AsyncGenerator {
  907. const asyncFactory = Reflect.get(source, Symbol.asyncIterator) as unknown
  908. const syncFactory = Reflect.get(source, Symbol.iterator) as unknown
  909. const iterator = typeof asyncFactory === 'function'
  910. ? Reflect.apply(asyncFactory, source, []) as AsyncIterator<unknown>
  911. : Reflect.apply(syncFactory as (...args: never[]) => Iterator<unknown>, source, [])
  912. let rejectAbort: ((error: unknown) => void) | undefined
  913. const aborted = new Promise<never>((_resolve, reject) => { rejectAbort = reject })
  914. const onAbort = (): void => {
  915. rejectAbort?.(remoteCancelled(endpoint, signal.reason))
  916. }
  917. signal.addEventListener('abort', onAbort, { once: true })
  918. try {
  919. if (signal.aborted) throw remoteCancelled(endpoint, signal.reason)
  920. while (true) {
  921. const next = await Promise.race([Promise.resolve(iterator.next()), aborted])
  922. if (next.done === true) return
  923. yield next.value
  924. }
  925. } finally {
  926. signal.removeEventListener('abort', onAbort)
  927. await iterator.return?.()
  928. }
  929. }
  930. /** Carrier-signal cancellation as the shared failure vocabulary expresses it. */
  931. function remoteCancelled(endpoint: string, cause: unknown): RemoteError<'gateway/cancelled'> {
  932. return new RemoteError('gateway/cancelled', `Remote invocation "${endpoint}" was aborted`, {}, { cause })
  933. }
  934. function rpcFailure(error: unknown): ConnectionRpcResult {
  935. const remote = remoteErrorOf(error)
  936. if (remote !== undefined) {
  937. return { ok: false, error: { code: remote.code, message: remote.message, details: remote.details } }
  938. }
  939. return {
  940. ok: false,
  941. error: {
  942. code: 'gateway/internal',
  943. message: error instanceof Error ? error.message : String(error),
  944. details: {},
  945. },
  946. }
  947. }
  948. function rpcError(error: unknown): ConnectionRpcError & RemoteStreamFailure {
  949. return (rpcFailure(error) as Extract<ConnectionRpcResult, { readonly ok: false }>).error
  950. }
  951. function endpointOf(namespace: string, method: string): string {
  952. return `${namespace}/${method}`
  953. }
  954. function validateBinding(
  955. receiver: object,
  956. serviceKey: string,
  957. namespace: string,
  958. endpoint: string,
  959. ): ResolvedBinding {
  960. const original = originalOf(receiver)
  961. const value = Reflect.get(original, 'typertRemote') as unknown
  962. if (value === undefined) {
  963. throw new TypertGatewayError(
  964. 'gateway/binding-invalid',
  965. endpoint,
  966. `Service ${JSON.stringify(serviceKey)} has no visible typertRemote binding`,
  967. )
  968. }
  969. return {
  970. binding: readBinding(value, original, serviceKey, endpoint, namespace),
  971. original,
  972. }
  973. }
  974. function readBinding(
  975. value: unknown,
  976. original: object,
  977. serviceKey: string,
  978. endpoint: string,
  979. namespace?: string,
  980. ): TypertGatewayBinding {
  981. if (!isObject(value)
  982. || Reflect.get(value, 'service') !== original
  983. || Reflect.get(value, 'serviceKey') !== serviceKey
  984. || typeof Reflect.get(value, 'namespace') !== 'string'
  985. || (namespace !== undefined && Reflect.get(value, 'namespace') !== namespace)) {
  986. throw new TypertGatewayError(
  987. 'gateway/binding-invalid',
  988. endpoint,
  989. `Service ${JSON.stringify(serviceKey)} has an inconsistent typertRemote binding`,
  990. )
  991. }
  992. return value as unknown as TypertGatewayBinding
  993. }
  994. function originalOf(receiver: object): object {
  995. const original = Reflect.get(receiver, symbols.original) as unknown
  996. return isObject(original) ? original : receiver
  997. }
  998. function methodParameterNames(service: object, method: string, endpoint: string): readonly string[] {
  999. let prototype: object | null = Object.getPrototypeOf(service) as object | null
  1000. let implementation: ((this: object, ...args: never[]) => unknown) | undefined
  1001. while (prototype !== null) {
  1002. const descriptor = Object.getOwnPropertyDescriptor(prototype, method)
  1003. if (descriptor !== undefined) {
  1004. if ('value' in descriptor && typeof descriptor.value === 'function') {
  1005. implementation = descriptor.value as (this: object, ...args: never[]) => unknown
  1006. }
  1007. break
  1008. }
  1009. prototype = Object.getPrototypeOf(prototype) as object | null
  1010. }
  1011. if (implementation === undefined) {
  1012. throw new TypertGatewayError(
  1013. 'gateway/method-unavailable',
  1014. endpoint,
  1015. `Remote marker has no prototype method ${JSON.stringify(method)}`,
  1016. )
  1017. }
  1018. const source = Function.prototype.toString.call(implementation)
  1019. const open = source.indexOf('(')
  1020. const close = source.indexOf(')', open + 1)
  1021. /* v8 ignore next -- standard public class-method syntax always contains a parenthesized parameter list. */
  1022. if (open < 0 || close < 0) return invalidSignature(endpoint, method)
  1023. const body = source.slice(open + 1, close).trim()
  1024. if (body.length === 0) return []
  1025. const parts = body.split(',').map(part => part.trim())
  1026. const names = new Set<string>()
  1027. for (const part of parts) {
  1028. if (!/^[$A-Z_a-z][$\w]*$/u.test(part) || names.has(part)) return invalidSignature(endpoint, method)
  1029. names.add(part)
  1030. }
  1031. return [...names]
  1032. }
  1033. function invalidSignature(endpoint: string, method: string): never {
  1034. throw new TypertGatewayError(
  1035. 'gateway/signature-invalid',
  1036. endpoint,
  1037. `SRC method ${JSON.stringify(method)} must use unique identifier parameters without destructuring, defaults, or rest`,
  1038. )
  1039. }
  1040. function assertExactArguments(
  1041. args: Readonly<Record<string, unknown>>,
  1042. descriptor: InvocationDescriptor,
  1043. endpoint: string,
  1044. ): void {
  1045. if (!isPlainObject(args)) {
  1046. throw new TypertGatewayError('gateway/arguments-invalid', endpoint, 'args must be a plain object')
  1047. }
  1048. const expected = new Set(descriptor.parameters.map(parameter => parameter.wire))
  1049. if (descriptor.invocation.kind === 'context') expected.add(descriptor.invocation.wire)
  1050. const actual = Reflect.ownKeys(args)
  1051. const extra = actual.filter(key => typeof key !== 'string' || !expected.has(key))
  1052. // A JSON field may be omitted when the strict descriptor declares absence,
  1053. // and always under SRC: a weak descriptor reads parameter names from the
  1054. // JavaScript signature and cannot see which are optional, so LIB is where an
  1055. // omitted required argument is caught. Lookup ids are never omissible.
  1056. const acceptsMissing = new Set(descriptor.parameters
  1057. .filter(parameter => parameter.source === 'json'
  1058. && (parameter.acceptsUndefined === true || parameter.codec.mode === 'src-json'))
  1059. .map(parameter => parameter.wire))
  1060. const missing = [...expected].filter(key => !Object.hasOwn(args, key) && !acceptsMissing.has(key))
  1061. if (extra.length === 0 && missing.length === 0) return
  1062. const clauses: string[] = []
  1063. if (missing.length > 0) clauses.push(`missing ${missing.map(key => JSON.stringify(key)).join(', ')}`)
  1064. if (extra.length > 0) clauses.push(`unexpected ${extra.map(key => JSON.stringify(String(key))).join(', ')}`)
  1065. throw new TypertGatewayError('gateway/arguments-invalid', endpoint, `args fields do not match the descriptor: ${clauses.join('; ')}`)
  1066. }
  1067. function decode(
  1068. codec: TypertCodec,
  1069. value: unknown,
  1070. endpoint: string,
  1071. field: string,
  1072. ): unknown {
  1073. try {
  1074. if (codec.mode === 'strict') {
  1075. value = codec.create().parse(value)
  1076. /* v8 ignore next -- generated optional-input codecs are the only strict codecs that return undefined. */
  1077. if (value === undefined) return value
  1078. }
  1079. assertJsonValue(value, new Set())
  1080. return value
  1081. } catch (cause) {
  1082. throw new TypertGatewayError(
  1083. 'gateway/input-invalid',
  1084. endpoint,
  1085. `wire field ${JSON.stringify(field)} failed boundary validation`,
  1086. { cause, field },
  1087. )
  1088. }
  1089. }
  1090. function assertJsonValue(value: unknown, ancestors: Set<object>): void {
  1091. if (value === null || typeof value === 'string' || typeof value === 'boolean') return
  1092. if (typeof value === 'number') {
  1093. if (Number.isFinite(value)) return
  1094. throw new TypeError('non-finite number is not JSON-safe')
  1095. }
  1096. if (!isObject(value)) throw new TypeError(`${typeof value} is not JSON-safe`)
  1097. if (ancestors.has(value)) throw new TypeError('cyclic value is not JSON-safe')
  1098. ancestors.add(value)
  1099. try {
  1100. if (Array.isArray(value)) {
  1101. if (Object.getOwnPropertySymbols(value).length > 0 || Object.keys(value).length !== value.length) {
  1102. throw new TypeError('sparse or decorated array is not JSON-safe')
  1103. }
  1104. for (let index = 0; index < value.length; index += 1) {
  1105. if (!Object.hasOwn(value, index)) throw new TypeError('sparse array is not JSON-safe')
  1106. assertJsonValue(value[index], ancestors)
  1107. }
  1108. return
  1109. }
  1110. if (!isPlainObject(value)) throw new TypeError('non-plain object is not JSON-safe')
  1111. if (Object.getOwnPropertySymbols(value).length > 0) throw new TypeError('symbol property is not JSON-safe')
  1112. for (const key of Reflect.ownKeys(value)) {
  1113. const descriptor = Object.getOwnPropertyDescriptor(value, key)
  1114. /* v8 ignore next -- ownKeys() just returned this key; only a hostile same-process Proxy can delete it between operations. */
  1115. if (descriptor === undefined || !descriptor.enumerable || !('value' in descriptor)) {
  1116. throw new TypeError('non-data property is not JSON-safe')
  1117. }
  1118. assertJsonValue(descriptor.value, ancestors)
  1119. }
  1120. } finally {
  1121. ancestors.delete(value)
  1122. }
  1123. }
  1124. function isPlainObject(value: object): value is Record<string, unknown> {
  1125. if (Array.isArray(value)) return false
  1126. const prototype = Object.getPrototypeOf(value) as object | null
  1127. return prototype === null || prototype === Object.prototype
  1128. }
  1129. function isObject(value: unknown): value is object {
  1130. return (typeof value === 'object' && value !== null) || typeof value === 'function'
  1131. }
  1132. export default TypertGatewayService