返回源码地图

packages/api/gateway/src/index.ts

main snapshot · da00f7f5358f · 正文引用章节 15 / 15;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

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
8import { randomUUID } from 'node:crypto'
9import { Context, Service, symbols } from '@deepseek-ai/cordis'
10import {
11 OperatorPeer,
12 type ConnectionRpcAttachment,
13 type ConnectionRpcHandler,
14} from '@deepseek-ai/dsh-client-connection'
15import { Deque } from '@deepseek-ai/dsh-deque'
16import type { WebUpgradeRoute } from '@deepseek-ai/dsh-host-webserver'
17import { MAX_TIMER_DELAY_MS } from '@deepseek-ai/dsh-timeout'
18import type {} from '@deepseek-ai/dsh-cmdline'
19import z from '@deepseek-ai/schemastery'
20export type { TypertGatewayFaultDetails } from './remote-error-codes.ts'
21import {
22 RemoteError,
23 isRemoteJsonValue,
24 remoteErrorOf,
25 remoteMethods,
26 type InvocationDescriptor,
27 type InvocationParameterDescriptor,
28 type PeerScope,
29 type RemoteInvocation,
30 type TypertCodec,
31 type TypertGatewayBinding,
32} from '@deepseek-ai/dsh-typert-protocol'
33import type {
34 InvokeRemoteRequest,
35 TypertGateway,
36 TypertGatewayErrorCode,
37 TypertGatewayWireStream,
38 TypertRemoteEventDispatch,
39 TypertRemoteEventFrame,
40 TypertRemoteEventInvocation,
41 TypertRemoteEventOutcome,
42 TypertRemoteEventSource,
43} from './types.ts'
44import {
45 RemoteStreamMuxServer,
46 rejectRemoteStreamUpgrade,
47} from './stream-server.ts'
48import {
49 REMOTE_EVENT_STREAM_ENDPOINT,
50 REMOTE_EVENT_STREAM_READY,
51 REMOTE_EVENT_RESULT_ENDPOINT,
52 REMOTE_STREAM_MUX_PATH,
53 isRemoteEventAgentId,
54 parseRemoteEventResult,
55 projectRemoteEventRequest,
56 restoreRemoteEventRejection,
57 type RemoteEventCancellationFrame,
58 type RemoteEventClientId,
59 type RemoteEventEmitFrame,
60 type RemoteEventHostInfo,
61 type RemoteEventId,
62 type RemoteEventInvocationFrame,
63 type RemoteEventReadyFrame,
64 type RemoteStreamFailure,
65} from './stream-protocol.ts'
66
67export type {
68 InvokeRemoteRequest,
69 TypertGateway,
70 TypertGatewayErrorCode,
71 TypertGatewayWireStream,
72 TypertRemoteEventContext,
73 TypertRemoteEventDispatch,
74 TypertRemoteEventFrame,
75 TypertRemoteEventInvocation,
76 TypertRemoteEventOutcome,
77 TypertRemoteEventSource,
78} from './types.ts'
79export type { RemoteEventHostInfo } from './stream-protocol.ts'
80
81interface GatewayErrorOptions {
82 readonly cause?: unknown
83 readonly field?: string
84}
85
86interface ResolvedBinding {
87 readonly binding: TypertGatewayBinding
88 readonly original: object
89}
90
91interface PreparedInvocation {
92 readonly endpoint: string
93 readonly descriptor: InvocationDescriptor
94 /** Service view bound to a Context carrying `invocation`, so the method reads it as `this.ctx.invocation`. */
95 readonly receiver: object
96 readonly args: readonly unknown[]
97 readonly method: (...args: never[]) => unknown
98 readonly invocation: GatewayInvocation
99}
100
101/** Carrier inputs `GatewayInvocation.uplink()` decodes on first use. */
102interface UplinkSource {
103 /** Carrier items; an immediately ended iterable when the carrier has none. */
104 readonly source: AsyncIterable<unknown>
105 /** Descriptor codec, or the JSON-safety codec when the descriptor declares no uplink. */
106 readonly codec: TypertCodec
107 readonly endpoint: string
108 /** Fail the logical stream with a Remote failure as the reason. */
109 readonly abort: (reason: unknown) => void
110}
111
112interface RegisteredRemoteEventSource {
113 readonly lifetime: AbortController
114 readonly done: Promise<void>
115 readonly host: RemoteEventHostInfo
116}
117
118interface RemoteEventClient {
119 readonly id: RemoteEventClientId
120 readonly queue: RemoteEventQueue
121 readonly signal: AbortSignal
122 readonly deliveries: Map<RemoteEventId, PendingRemoteEvent>
123}
124
125interface PendingRemoteEvent {
126 readonly id: RemoteEventId
127 readonly source: TypertRemoteEventInvocation
128 readonly frame: RemoteEventInvocationFrame
129 readonly deliveries: Set<RemoteEventClient>
130 releaseContext: () => void
131 releaseSignal: () => void
132}
133
134type ConnectionRpcResult = Awaited<ReturnType<ConnectionRpcHandler>>
135type ConnectionRpcError = Extract<ConnectionRpcResult, { readonly ok: false }>['error']
136const NEVER_ABORTED_SIGNAL = new AbortController().signal
137const DEFAULT_WEBSOCKET_HEARTBEAT_INTERVAL_MS = 2_000
138const DEFAULT_STREAM_INBOX_BYTES = 262_144
139const EMPTY_ASYNC_ITERABLE: AsyncIterable<never> = {
140 [Symbol.asyncIterator]: () => ({ next: () => Promise.resolve({ value: undefined, done: true }) }),
141}
142const UPLINK_DONE: IteratorReturnResult<undefined> = { value: undefined, done: true }
143const SRC_JSON_CODEC: TypertCodec = { mode: 'src-json' }
144
145/** Gateway transport configuration. */
146export interface Config {
147 /** WebSocket Ping interval from 1 through 2,147,483,647 milliseconds. @default 2000 */
148 readonly websocketHeartbeatIntervalMs?: number
149 /** Buffered uplink frame bytes one logical stream may hold before it fails with `gateway/uplink-overflow`. @default 262144 */
150 readonly streamInboxBytes?: number
151}
152
153interface ResolvedConfig extends Config {
154 readonly websocketHeartbeatIntervalMs: number
155 readonly streamInboxBytes: number
156}
157
158/**
159 * Dispatch failure produced outside the invoked business method. Rides the
160 * shared Remote failure vocabulary, so its code crosses the wire instead of
161 * folding to `internal`.
162 */
163export class TypertGatewayError extends RemoteError<TypertGatewayErrorCode> {
164 /** Canonical `<namespace>/<method>` endpoint. */
165 readonly endpoint: string
166 /** Affected wire field when the failure is field-specific. */
167 readonly field: string | undefined
168
169 /**
170 * Construct a Gateway failure without embedding boundary values in its message.
171 * @param code - stable failure category.
172 * @param endpoint - canonical Remote endpoint.
173 * @param message - correction-oriented diagnostic without sensitive values.
174 * @param options - optional field and contained cause.
175 */
176 constructor(
177 code: TypertGatewayErrorCode,
178 endpoint: string,
179 message: string,
180 options: GatewayErrorOptions = {},
181 ) {
182 super(
183 code,
184 `typert gateway: ${endpoint}: ${message}`,
185 { endpoint, ...options.field === undefined ? {} : { field: options.field } },
186 options.cause === undefined ? undefined : { cause: options.cause },
187 )
188 this.name = 'TypertGatewayError'
189 this.endpoint = endpoint
190 this.field = options.field
191 }
192}
193
194/**
195 * Resolve strict generated definitions or conservative SRC markers against
196 * current Cordis Services and Typert providers.
197 * @typert service typertGateway
198 */
199export class TypertGatewayService extends Service implements TypertGateway {
200 static inject = ['typert']
201 static Config: z<Config> = z.object({
202 websocketHeartbeatIntervalMs: z.number().step(1).min(1).max(MAX_TIMER_DELAY_MS)
203 .default(DEFAULT_WEBSOCKET_HEARTBEAT_INTERVAL_MS),
204 streamInboxBytes: z.number().step(1).min(1).default(DEFAULT_STREAM_INBOX_BYTES),
205 })
206
207 /** Carrier adapter shared by the WebSocket mux and local Host transports. */
208 readonly wireStream: TypertGatewayWireStream = {
209 open: (endpoint, payload, uplink, peer, signal) =>
210 this.openWireStream(endpoint, payload, uplink, peer, signal, new AbortController()),
211 failure: error => rpcError(error),
212 }
213
214 private srcClaims: ReadonlySet<string> | undefined
215 private inProcessOperator: PeerScope | undefined
216 private remoteEvents: RegisteredRemoteEventSource | undefined
217 private readonly remoteEventClients = new Map<RemoteEventClientId, RemoteEventClient>()
218 private readonly pendingRemoteEvents = new Map<RemoteEventId, PendingRemoteEvent>()
219
220 /**
221 * Register the Gateway against the active Typert registry.
222 * WebSocket admission waits for launcher-owned application readiness when supplied;
223 * direct invocation and in-process streams remain available independently.
224 * @param ctx - owning Host Context with Typert registry access.
225 * @param config - validated Gateway transport configuration.
226 */
227 constructor(ctx: Context, config: Config) {
228 super(ctx, 'typertGateway')
229 const resolved = config as ResolvedConfig
230 ctx.on('internal/service', () => {
231 this.srcClaims = undefined
232 })
233 ctx.inject(['connection'], (connectionCtx) => {
234 connectionCtx.connection.rpc.intercept(
235 '/api',
236 endpoint => this.claimsEndpoint(endpoint),
237 (endpoint, payload, signal, peer) => this.dispatchRpc(endpoint, payload, signal, peer),
238 )
239 })
240 ctx.inject(['connection', 'webServer'], (webCtx) => {
241 const listen = (): void => {
242 const mux = new RemoteStreamMuxServer(
243 (endpoint, payload, uplink, peer, control) =>
244 this.openWireStream(endpoint, payload, uplink, peer, control.signal, control),
245 this.wireStream.failure,
246 resolved.websocketHeartbeatIntervalMs,
247 resolved.streamInboxBytes,
248 )
249 webCtx.effect(function* () {
250 yield () => mux.close()
251 const route: WebUpgradeRoute = {
252 path: REMOTE_STREAM_MUX_PATH,
253 handler: (req, socket, head) => {
254 const admission = webCtx.connection.admit(req)
255 if ('rejection' in admission) {
256 rejectRemoteStreamUpgrade(socket, admission.rejection)
257 return
258 }
259 mux.handleUpgrade(req, socket, head, admission.peer)
260 },
261 }
262 yield webCtx.webServer.registerUpgrade(route)
263 }, `api-gateway: ${REMOTE_STREAM_MUX_PATH} WebSocket`)
264 }
265 // Existing pages reconnect before the new Host prints its URL. No stream
266 // may enter until the launcher has activated and audited its controllers.
267 const ready = webCtx.get('appReady')
268 if (ready === undefined) listen()
269 else webCtx.effect(() => {
270 let closed = false
271 const cancel = ready.onReady(() => { if (!closed) listen() })
272 return () => {
273 closed = true
274 cancel()
275 }
276 }, 'api-gateway: application readiness')
277 })
278 }
279
280 /**
281 * Check for an active Client event stream.
282 * @returns whether a stream is open and has not been cancelled.
283 */
284 hasLiveClient(): boolean {
285 for (const client of this.remoteEventClients.values()) {
286 if (!client.signal.aborted) return true
287 }
288 return false
289 }
290
291 /**
292 * Register the sole application-selected forwarded-event source.
293 * @param source - stream factory installed by the Remote assembly.
294 * @param host - stable Host facts included in each Client generation's opening frame.
295 * @returns disposer removing this source and cancelling its active streams.
296 */
297 registerRemoteEvents(
298 source: TypertRemoteEventSource,
299 host: RemoteEventHostInfo,
300 ): () => Promise<void> {
301 if (this.remoteEvents !== undefined) {
302 throw new Error('typert gateway: forwarded Remote event source is already registered')
303 }
304 const lifetime = new AbortController()
305 const stream = source(lifetime.signal)
306 const done = this.consumeRemoteEvents(stream, lifetime.signal).catch((error: unknown) => {
307 if (this.remoteEvents?.lifetime !== lifetime || lifetime.signal.aborted) return
308 this.closeRemoteEvents(error)
309 this.remoteEvents = undefined
310 lifetime.abort(error)
311 })
312 const registration: RegisteredRemoteEventSource = { lifetime, done, host: { home: host.home } }
313 this.remoteEvents = registration
314 return async () => {
315 if (this.remoteEvents === registration) {
316 this.remoteEvents = undefined
317 const error = new Error('typert gateway: forwarded Remote event source was removed')
318 registration.lifetime.abort(error)
319 this.closeRemoteEvents(error)
320 }
321 await registration.done
322 }
323 }
324
325 private claimsEndpoint(endpoint: string): boolean {
326 if (endpoint === REMOTE_EVENT_RESULT_ENDPOINT) return true
327 const segments = endpoint.split('/')
328 if (segments.length !== 2 || segments[0] === '' || segments[1] === '') return false
329 if (this.ctx.typert.local.get(endpoint) !== undefined || this.ctx.typert.local.hasSeen(endpoint)) return true
330 this.srcClaims ??= this.collectSrcClaims()
331 return this.srcClaims.has(endpoint)
332 }
333
334 private collectSrcClaims(): ReadonlySet<string> {
335 const claims = new Set<string>()
336 for (const [serviceKey, definition] of Object.entries(this.ctx.reflect.props)) {
337 if (definition.type !== 'service') continue
338 const receiver: unknown = this.ctx.get(serviceKey)
339 if (!isObject(receiver)) continue
340 const original = originalOf(receiver)
341 const binding: unknown = Reflect.get(original, 'typertRemote')
342 if (!isObject(binding) || typeof Reflect.get(binding, 'namespace') !== 'string') continue
343 const namespace = Reflect.get(binding, 'namespace') as string
344 for (const candidate of remoteMethods(original)) {
345 claims.add(endpointOf(namespace, candidate.exportName ?? candidate.method))
346 }
347 }
348 return claims
349 }
350
351 /**
352 * Invoke one live Remote method through strict generated reflection or SRC markers.
353 * @param request - decoded endpoint and exact named wire arguments.
354 * @returns the business result without output decoding.
355 * @throws {@link TypertGatewayError} for dispatch, provider, or boundary failures; lookup-policy and business errors retain identity.
356 */
357 async invoke(request: InvokeRemoteRequest): Promise<unknown> {
358 return this.invokePrepared(await this.prepareInvocation(request, new AbortController()))
359 }
360
361 private async invokePrepared(prepared: PreparedInvocation): Promise<unknown> {
362 if (prepared.descriptor.mode !== undefined) {
363 throw new TypertGatewayError(
364 'gateway/signature-invalid',
365 prepared.endpoint,
366 'stream Remote methods must be opened through the stream carrier',
367 )
368 }
369
370 try {
371 return await Reflect.apply(prepared.method, prepared.receiver, prepared.args) as unknown
372 } catch (error) {
373 if (prepared.invocation.signal.aborted) throw remoteCancelled(prepared.endpoint, error)
374 throw error
375 } finally {
376 // A unary call's uplink is readable only while the method runs.
377 await prepared.invocation.close()
378 }
379 }
380
381 /**
382 * Open one live stream Remote method without assuming a physical carrier.
383 * @param request - decoded endpoint, named wire arguments, and the Client uplink when the carrier has one.
384 * @returns a cancellation-aware iterable over the business results.
385 */
386 async stream(request: InvokeRemoteRequest): Promise<AsyncIterable<unknown>> {
387 return this.openStream(request, new AbortController())
388 }
389
390 /**
391 * `control` belongs to the logical stream: a rejected uplink item aborts it
392 * with the Remote failure as the reason so the carrier delivers that failure.
393 */
394 private async openStream(request: InvokeRemoteRequest, control: AbortController): Promise<AsyncIterable<unknown>> {
395 const prepared = await this.prepareInvocation(request, control)
396 if (prepared.descriptor.mode === undefined) {
397 await prepared.invocation.close()
398 throw new TypertGatewayError(
399 'gateway/signature-invalid',
400 prepared.endpoint,
401 'unary Remote methods cannot be opened through the stream carrier',
402 )
403 }
404 let source: unknown
405 try {
406 source = Reflect.apply(prepared.method, prepared.receiver, prepared.args) as unknown
407 } catch (error) {
408 await prepared.invocation.close()
409 if (prepared.invocation.signal.aborted) throw remoteCancelled(prepared.endpoint, error)
410 throw error
411 }
412 if (!isIterable(source)) {
413 await prepared.invocation.close()
414 throw new TypertGatewayError(
415 'gateway/result-invalid',
416 prepared.endpoint,
417 'stream Remote method did not return Iterable or AsyncIterable',
418 { field: 'result' },
419 )
420 }
421 return cancellableStream(source, prepared.endpoint, prepared.invocation)
422 }
423
424 private async dispatchRpc(
425 endpoint: string,
426 payload: unknown,
427 signal: AbortSignal,
428 peer: PeerScope,
429 ): Promise<ConnectionRpcResult> {
430 if (endpoint === REMOTE_EVENT_RESULT_ENDPOINT) {
431 try {
432 const result = parseRemoteEventResultPayload(payload)
433 const client = this.remoteEventClients.get(result.clientId)
434 if (client === undefined) {
435 throw new Error('typert gateway: Remote event result identifies no active event stream')
436 }
437 this.receiveRemoteEventResult(client, result)
438 return { ok: true, value: undefined }
439 } catch (error) {
440 return rpcFailure(error)
441 }
442 }
443 return this.invokeRpc(endpoint, payload, signal, peer)
444 }
445
446 private async openWireStream(
447 endpoint: string,
448 payload: unknown,
449 uplink: AsyncIterable<unknown>,
450 peer: PeerScope | undefined,
451 signal: AbortSignal,
452 control: AbortController,
453 ): Promise<AsyncIterable<unknown>> {
454 if (endpoint === REMOTE_EVENT_STREAM_ENDPOINT) {
455 // A Gateway-owned stream reads no uplink: releasing it now keeps its items out of the bounded inbox.
456 releaseUplink(uplink)
457 return this.openRemoteEvents(payload, signal)
458 }
459 return this.openStream({ ...remoteRequest(endpoint, payload, signal, peer), uplink }, control)
460 }
461
462 /**
463 * The Peer an in-process carrier speaks for when it names none: the
464 * operator's Peer when Connection is mounted, otherwise an operator scope the
465 * Gateway owns for its own lifetime.
466 * @returns the operator Peer.
467 */
468 private operatorPeer(): PeerScope {
469 const connection = this.ctx.get('connection')
470 if (connection !== undefined) return connection.operator
471 this.inProcessOperator ??= new OperatorPeer(this.ctx)
472 return this.inProcessOperator
473 }
474
475 private async *openRemoteEvents(
476 payload: unknown,
477 signal: AbortSignal,
478 ): AsyncGenerator<
479 RemoteEventEmitFrame | RemoteEventInvocationFrame | RemoteEventCancellationFrame
480 | RemoteEventReadyFrame
481 > {
482 if (!isObject(payload)
483 || !isPlainObject(payload)
484 || Reflect.ownKeys(payload).length !== 1
485 || !Object.hasOwn(payload, 'args')
486 || !isObject(payload.args)
487 || !isPlainObject(payload.args)
488 || Reflect.ownKeys(payload.args).length !== 0) {
489 throw new TypertGatewayError(
490 'gateway/arguments-invalid',
491 REMOTE_EVENT_STREAM_ENDPOINT,
492 'forwarded Remote event stream requires an empty args object',
493 )
494 }
495 const registration = this.remoteEvents
496 if (registration === undefined) {
497 throw new TypertGatewayError(
498 'gateway/service-unavailable',
499 REMOTE_EVENT_STREAM_ENDPOINT,
500 'forwarded Remote event source is unavailable',
501 )
502 }
503 const lifetime = AbortSignal.any([signal, registration.lifetime.signal])
504 let clientId = randomUUID() as RemoteEventClientId
505 while (this.remoteEventClients.has(clientId)) clientId = randomUUID() as RemoteEventClientId
506 const client: RemoteEventClient = {
507 id: clientId,
508 queue: new RemoteEventQueue(),
509 signal: lifetime,
510 deliveries: new Map(),
511 }
512 this.remoteEventClients.set(clientId, client)
513 for (const pending of this.pendingRemoteEvents.values()) this.deliverRemoteEvent(pending, client)
514 try {
515 yield { ...REMOTE_EVENT_STREAM_READY, clientId, host: registration.host }
516 yield* client.queue.iterate(lifetime)
517 } finally {
518 this.removeRemoteEventClient(client)
519 }
520 }
521
522 private async consumeRemoteEvents(
523 source: AsyncIterable<TypertRemoteEventDispatch>,
524 signal: AbortSignal,
525 ): Promise<void> {
526 for await (const dispatch of source) {
527 if (signal.aborted) {
528 if ('context' in dispatch) dispatch.reject(signal.reason)
529 return
530 }
531 if ('context' in dispatch) this.startRemoteEvent(dispatch)
532 else this.broadcastRemoteEvent(dispatch)
533 }
534 if (!signal.aborted) {
535 throw new Error('typert gateway: forwarded Remote event source ended unexpectedly')
536 }
537 }
538
539 private broadcastRemoteEvent(frame: TypertRemoteEventFrame): void {
540 assertRemoteEventFrame(frame)
541 const wire: RemoteEventEmitFrame = {
542 type: 'emit',
543 event: frame.event,
544 args: frame.args,
545 }
546 for (const client of this.remoteEventClients.values()) client.queue.push(wire)
547 }
548
549 private startRemoteEvent(source: TypertRemoteEventInvocation): void {
550 try {
551 assertRemoteEventName(source)
552 if (!isRemoteEventAgentId(source.context.agentId)) {
553 throw new TypeError(
554 'typert gateway: scoped Remote events require a non-empty Agent identity',
555 )
556 }
557 const projected = projectRemoteEventRequest(source.request, source.context.subject)
558 let id = randomUUID() as RemoteEventId
559 while (this.pendingRemoteEvents.has(id)) id = randomUUID() as RemoteEventId
560 let releaseContext: () => void
561 try {
562 const dispose = source.context.value.effect(
563 () => () => {
564 this.cancelRemoteEvent(
565 pending,
566 new Error('typert gateway: Remote event Agent Context was released'),
567 )
568 },
569 `api-gateway: Remote event ${JSON.stringify(source.event)}`,
570 )
571 releaseContext = () => { void dispose() }
572 } catch {
573 source.resolve({ kind: 'next' })
574 return
575 }
576 const signals = new Set(projected.signal === undefined ? [] : [projected.signal])
577 const abort = (): void => {
578 const reason: unknown = [...signals].find(signal => signal.aborted)?.reason
579 this.cancelRemoteEvent(pending, reason instanceof Error
580 ? reason
581 : new Error('typert gateway: Remote event was cancelled', { cause: reason }))
582 }
583 const pending: PendingRemoteEvent = {
584 id,
585 source,
586 frame: {
587 type: 'waterfall',
588 event: source.event,
589 eventId: id,
590 agentId: source.context.agentId,
591 request: projected.request,
592 },
593 deliveries: new Set(),
594 releaseContext,
595 releaseSignal: () => {
596 for (const signal of signals) signal.removeEventListener('abort', abort)
597 },
598 }
599 this.pendingRemoteEvents.set(id, pending)
600 for (const signal of signals) signal.addEventListener('abort', abort, { once: true })
601 if ([...signals].some(signal => signal.aborted)) abort()
602 else for (const client of this.remoteEventClients.values()) this.deliverRemoteEvent(pending, client)
603 } catch (error) {
604 source.reject(error)
605 }
606 }
607
608 private deliverRemoteEvent(pending: PendingRemoteEvent, client: RemoteEventClient): void {
609 pending.deliveries.add(client)
610 client.deliveries.set(pending.id, pending)
611 client.queue.push(pending.frame)
612 }
613
614 private receiveRemoteEventResult(
615 client: RemoteEventClient,
616 result: ReturnType<typeof parseRemoteEventResult>,
617 ): void {
618 const pending = this.pendingRemoteEvents.get(result.eventId)
619 // Settlement and Client replacement may race the result request. Results
620 // from a completed event or a superseded delivery are idempotent no-ops.
621 if (pending === undefined || !pending.deliveries.has(client)) return
622 this.removeRemoteEventDelivery(pending, client)
623 if (result.outcome.kind === 'result') {
624 this.settleRemoteEvent(pending, {
625 kind: 'result',
626 value: result.outcome.value,
627 })
628 } else if (result.outcome.kind === 'rejected') {
629 this.cancelRemoteEvent(pending, restoreRemoteEventRejection(result.outcome.error))
630 } else if (pending.deliveries.size === 0) {
631 this.settleRemoteEvent(pending, { kind: 'next' })
632 }
633 }
634
635 private removeRemoteEventDelivery(pending: PendingRemoteEvent, client: RemoteEventClient): void {
636 pending.deliveries.delete(client)
637 client.deliveries.delete(pending.id)
638 }
639
640 private removeRemoteEventClient(client: RemoteEventClient): void {
641 this.remoteEventClients.delete(client.id)
642 for (const pending of [...client.deliveries.values()]) this.removeRemoteEventDelivery(pending, client)
643 client.queue.end()
644 }
645
646 private settleRemoteEvent(pending: PendingRemoteEvent, outcome: TypertRemoteEventOutcome): void {
647 this.finishRemoteEvent(pending)
648 pending.source.resolve(outcome)
649 }
650
651 private cancelRemoteEvent(pending: PendingRemoteEvent, reason: unknown): void {
652 if (this.pendingRemoteEvents.get(pending.id) !== pending) return
653 this.finishRemoteEvent(pending)
654 pending.source.reject(reason)
655 }
656
657 private finishRemoteEvent(pending: PendingRemoteEvent): void {
658 this.pendingRemoteEvents.delete(pending.id)
659 pending.releaseSignal()
660 pending.releaseContext()
661 const clients = new Set(pending.deliveries)
662 for (const client of clients) this.removeRemoteEventDelivery(pending, client)
663 const cancellation: RemoteEventCancellationFrame = {
664 type: 'cancel',
665 eventId: pending.id,
666 }
667 for (const client of clients) client.queue.push(cancellation)
668 }
669
670 private closeRemoteEvents(reason: unknown): void {
671 for (const pending of [...this.pendingRemoteEvents.values()]) {
672 this.cancelRemoteEvent(pending, reason)
673 }
674 for (const client of [...this.remoteEventClients.values()]) client.queue.end()
675 }
676
677 private async invokeRpc(
678 endpoint: string,
679 payload: unknown,
680 signal: AbortSignal,
681 peer: PeerScope,
682 ): Promise<ConnectionRpcResult> {
683 try {
684 const prepared = await this.prepareInvocation(
685 remoteRequest(endpoint, payload, signal, peer),
686 new AbortController(),
687 )
688 const value = await this.invokePrepared(prepared)
689 // A void or explicitly absent business result carries no `value` field;
690 // JSON has no `undefined`, and the envelope's optional slot is the one
691 // representation of absence that both args and results already use.
692 return encodeRpcResult(value, prepared.descriptor.result)
693 } catch (error) {
694 return rpcFailure(error)
695 }
696 }
697
698 /** `control` fails the logical stream when an uplink item is rejected; unary calls hand over an inert one. */
699 private async prepareInvocation(
700 request: InvokeRemoteRequest,
701 control: AbortController,
702 ): Promise<PreparedInvocation> {
703 const endpoint = endpointOf(request.namespace, request.method)
704 const descriptor = this.resolveDescriptor(request.namespace, request.method, endpoint)
705 assertExactArguments(request.args, descriptor, endpoint)
706 const receiverContext = await this.resolveReceiverContext(descriptor, request.args, endpoint)
707 const receiver: unknown = receiverContext.get(descriptor.service)
708 if (!isObject(receiver)) {
709 throw new TypertGatewayError(
710 'gateway/service-unavailable',
711 endpoint,
712 `active Service ${JSON.stringify(descriptor.service)} is unavailable`,
713 )
714 }
715 validateBinding(receiver, descriptor.service, descriptor.namespace, endpoint)
716 const args = await Promise.all(descriptor.parameters.map(parameter =>
717 this.resolveParameter(parameter, request.args, endpoint)))
718 const signal = methodSignal(request, control)
719 const invocation = new GatewayInvocation(
720 { namespace: request.namespace, method: request.method, args: request.args },
721 descriptor.service,
722 request.peer ?? this.operatorPeer(),
723 signal,
724 {
725 source: request.uplink ?? EMPTY_ASYNC_ITERABLE,
726 codec: descriptor.uplink?.codec ?? SRC_JSON_CODEC,
727 endpoint,
728 abort: (reason) => { control.abort(reason) },
729 },
730 )
731 if (descriptor.cancellation !== undefined) args.push(signal)
732 // The method runs on a Service view bound to a Context carrying this call:
733 // Cordis rebinds `this.ctx` to the accessing Context, so `this.ctx.invocation`
734 // is this call and nothing travels through the parameter list. The view
735 // resolves the Service the plain read above already found.
736 const callReceiver = receiverContext.extend({ invocation }).get(descriptor.service) as object
737 const implementation = descriptor.implementation ?? descriptor.method
738 const method: unknown = Reflect.get(callReceiver, implementation)
739 if (typeof method !== 'function') {
740 throw new TypertGatewayError(
741 'gateway/method-unavailable',
742 endpoint,
743 `active Service ${JSON.stringify(descriptor.service)} has no callable method ${JSON.stringify(implementation)}`,
744 )
745 }
746 return {
747 endpoint,
748 descriptor,
749 receiver: callReceiver,
750 args,
751 method: method as (...args: never[]) => unknown,
752 invocation,
753 }
754 }
755
756 private resolveDescriptor(namespace: string, method: string, endpoint: string): InvocationDescriptor {
757 const strict = this.ctx.typert.local.get(endpoint)
758 if (strict !== undefined) return strict
759 if (this.ctx.typert.local.hasSeen(endpoint)) {
760 throw new TypertGatewayError(
761 'gateway/definition-unavailable',
762 endpoint,
763 'its strict definition was withdrawn and SRC fallback is forbidden',
764 )
765 }
766 return this.resolveSrcDescriptor(namespace, method, endpoint)
767 }
768
769 private resolveSrcDescriptor(namespace: string, method: string, endpoint: string): InvocationDescriptor {
770 const candidates: InvocationDescriptor[] = []
771 for (const [serviceKey, definition] of Object.entries(this.ctx.reflect.props)) {
772 if (definition.type !== 'service') continue
773 const receiver: unknown = this.ctx.get(serviceKey)
774 if (!isObject(receiver)) continue
775 const original = originalOf(receiver)
776 const value: unknown = Reflect.get(original, 'typertRemote')
777 if (value === undefined) continue
778 const binding = readBinding(value, original, serviceKey, endpoint)
779 if (binding.namespace !== namespace) continue
780 const marker = remoteMethods(original).find(candidate => (candidate.exportName ?? candidate.method) === method)
781 if (marker === undefined) continue
782 candidates.push(this.srcDescriptor(binding, marker, method, endpoint))
783 }
784 if (candidates.length === 0) {
785 throw new TypertGatewayError('gateway/invocation-unavailable', endpoint, 'no active Remote method exports this endpoint')
786 }
787 if (candidates.length > 1) {
788 throw new TypertGatewayError(
789 'gateway/ambiguous-endpoint',
790 endpoint,
791 `multiple active Services export this endpoint: ${candidates.map(candidate => candidate.service).sort().join(', ')}`,
792 )
793 }
794 return candidates[0] as InvocationDescriptor
795 }
796
797 private srcDescriptor(
798 binding: TypertGatewayBinding,
799 marker: ReturnType<typeof remoteMethods>[number],
800 method: string,
801 endpoint: string,
802 ): InvocationDescriptor {
803 const names = methodParameterNames(binding.service, marker.method, endpoint)
804 const signalIndex = names.indexOf('signal')
805 if (signalIndex >= 0 && signalIndex !== names.length - 1) {
806 throw new TypertGatewayError(
807 'gateway/signature-invalid',
808 endpoint,
809 'SRC cancellation parameter signal must be the final parameter',
810 { field: 'signal' },
811 )
812 }
813 const cancellation = signalIndex >= 0
814 ? { parameter: 'signal' as const }
815 : undefined
816 const businessNames = cancellation === undefined ? names : names.slice(0, -1)
817 const parameters: InvocationParameterDescriptor[] = []
818 const wires = new Set<string>()
819 for (const name of businessNames) {
820 const matches = this.ctx.typert.lookups.definitions()
821 .filter(definition => definition.parameter === name)
822 if (matches.length > 1) {
823 throw new TypertGatewayError(
824 'gateway/signature-invalid',
825 endpoint,
826 `parameter ${JSON.stringify(name)} matches multiple lookup providers`,
827 { field: name },
828 )
829 }
830 const match = matches[0]
831 const parameter: InvocationParameterDescriptor = match === undefined
832 ? { name, wire: name, source: 'json', codec: { mode: 'src-json' } }
833 : {
834 name,
835 wire: match.wire,
836 source: 'lookup',
837 lookup: match.key,
838 codec: { mode: 'src-json' },
839 }
840 if (wires.has(parameter.wire)) {
841 throw new TypertGatewayError(
842 'gateway/signature-invalid',
843 endpoint,
844 `multiple parameters use wire field ${JSON.stringify(parameter.wire)}`,
845 { field: parameter.wire },
846 )
847 }
848 wires.add(parameter.wire)
849 parameters.push(parameter)
850 }
851
852 let receiver: InvocationDescriptor['invocation'] = { kind: 'direct' }
853 if (marker.invocation.kind === 'context') {
854 const provider = this.ctx.typert.contexts.getHost(marker.invocation.context)
855 if (provider === undefined) {
856 throw new TypertGatewayError(
857 'gateway/context-unavailable',
858 endpoint,
859 `Context provider ${JSON.stringify(marker.invocation.context)} is unavailable`,
860 )
861 }
862 if (wires.has(provider.wire)) {
863 throw new TypertGatewayError(
864 'gateway/signature-invalid',
865 endpoint,
866 `Context identity conflicts with wire field ${JSON.stringify(provider.wire)}`,
867 { field: provider.wire },
868 )
869 }
870 receiver = {
871 kind: 'context',
872 context: marker.invocation.context,
873 wire: provider.wire,
874 codec: { mode: 'src-json' },
875 }
876 }
877
878 return {
879 id: `src:${binding.serviceKey}#${endpoint}`,
880 service: binding.serviceKey,
881 namespace: binding.namespace,
882 method,
883 ...(marker.method === method ? {} : { implementation: marker.method }),
884 ...(marker.mode === undefined ? {} : { mode: marker.mode }),
885 invocation: receiver,
886 parameters,
887 ...(cancellation === undefined ? {} : { cancellation }),
888 result: { mode: 'src-json' },
889 }
890 }
891
892 private async resolveReceiverContext(
893 descriptor: InvocationDescriptor,
894 args: Readonly<Record<string, unknown>>,
895 endpoint: string,
896 ): Promise<Context> {
897 if (descriptor.invocation.kind === 'direct') return this.ctx
898 const invocation = descriptor.invocation
899 const provider = this.ctx.typert.contexts.getHost(invocation.context)
900 if (provider === undefined) {
901 throw new TypertGatewayError(
902 'gateway/context-unavailable',
903 endpoint,
904 `Context provider ${JSON.stringify(invocation.context)} is unavailable`,
905 )
906 }
907 if (provider.wire !== invocation.wire
908 || (invocation.codec.mode === 'strict' && provider.wireTypeSymbol !== invocation.codec.typeSymbol)) {
909 throw new TypertGatewayError(
910 'gateway/provider-mismatch',
911 endpoint,
912 `Context provider ${JSON.stringify(invocation.context)} does not match its strict definition`,
913 { field: invocation.wire },
914 )
915 }
916 const identity = decode(invocation.codec, args[invocation.wire], endpoint, invocation.wire)
917 let context: Context | undefined
918 try {
919 context = await provider.resolve(identity)
920 } catch (cause) {
921 if (remoteErrorOf(cause) !== undefined) throw cause
922 throw new TypertGatewayError(
923 'gateway/context-failed',
924 endpoint,
925 `Context provider ${JSON.stringify(invocation.context)} failed`,
926 { cause, field: invocation.wire },
927 )
928 }
929 if (context === undefined) {
930 throw new TypertGatewayError(
931 'gateway/context-not-found',
932 endpoint,
933 `Context provider ${JSON.stringify(invocation.context)} did not resolve the requested identity`,
934 { field: invocation.wire },
935 )
936 }
937 return context
938 }
939
940 private async resolveParameter(
941 parameter: InvocationParameterDescriptor,
942 args: Readonly<Record<string, unknown>>,
943 endpoint: string,
944 ): Promise<unknown> {
945 // An absent field reached assertExactArguments' allowance, so this parameter
946 // takes undefined; a present-but-undefined field is not JSON-safe input and
947 // still fails decode. Lookup ids are never omissible, so absence here only
948 // ever belongs to a json parameter.
949 if (!Object.hasOwn(args, parameter.wire)) return undefined
950 const value = decode(parameter.codec, args[parameter.wire], endpoint, parameter.wire)
951 if (parameter.source === 'json') return value
952 const key = parameter.lookup
953 /* v8 ignore next -- registry validation rejects strict descriptors without a key, and SRC derivation always supplies one. */
954 if (key === undefined) {
955 throw new TypertGatewayError(
956 'gateway/lookup-unavailable',
957 endpoint,
958 `lookup parameter ${JSON.stringify(parameter.name)} has no provider key`,
959 { field: parameter.wire },
960 )
961 }
962 const provider = this.ctx.typert.lookups.get(key)
963 if (provider === undefined) {
964 throw new TypertGatewayError(
965 'gateway/lookup-unavailable',
966 endpoint,
967 `lookup provider ${JSON.stringify(key)} is unavailable`,
968 { field: parameter.wire },
969 )
970 }
971 if (provider.wire !== parameter.wire
972 || (parameter.codec.mode === 'strict' && provider.wireTypeSymbol !== parameter.codec.typeSymbol)) {
973 throw new TypertGatewayError(
974 'gateway/provider-mismatch',
975 endpoint,
976 `lookup provider ${JSON.stringify(key)} does not match its strict definition`,
977 { field: parameter.wire },
978 )
979 }
980 let resolved: unknown
981 try {
982 resolved = await provider.resolve(value)
983 } catch (cause) {
984 if (remoteErrorOf(cause) !== undefined) throw cause
985 throw new TypertGatewayError(
986 'gateway/lookup-failed',
987 endpoint,
988 `lookup provider ${JSON.stringify(key)} failed`,
989 { cause, field: parameter.wire },
990 )
991 }
992 if (resolved === undefined) {
993 throw new TypertGatewayError(
994 'gateway/lookup-not-found',
995 endpoint,
996 `lookup provider ${JSON.stringify(key)} did not resolve the requested identity`,
997 { field: parameter.wire },
998 )
999 }
1000 return resolved
1001 }
1002}
1003
1004function encodeRpcResult(value: unknown, codec: TypertCodec): ConnectionRpcResult {
1005 const attachments: ConnectionRpcAttachment[] = []
1006 const writeBytes = (bytes: Uint8Array, path: readonly (string | number)[]): null => {
1007 attachments.push({ path: [...path], bytes })
1008 return null
1009 }
1010 const encoded = codec.mode === 'strict'
1011 ? codec.encode?.(value, writeBytes) ?? value
1012 : encodeRuntimeResult(value, writeBytes)
1013 return { ok: true, value: encoded, ...(attachments.length === 0 ? {} : { attachments }) }
1014}
1015
1016function encodeRuntimeResult(
1017 input: unknown,
1018 writeBytes: (bytes: Uint8Array, path: readonly (string | number)[]) => null,
1019): unknown {
1020 const path: (string | number)[] = []
1021 const ancestors = new Set<object>()
1022 const extract = (input: unknown, key: string): unknown => {
1023 // Capture accessors and toJSON once, before materializing the JSON metadata.
1024 let value = input
1025 if (input !== null && typeof input === 'object' && !(input instanceof Uint8Array)) {
1026 const toJSON: unknown = Reflect.get(input, 'toJSON')
1027 if (typeof toJSON === 'function') value = Reflect.apply(toJSON, input, [key])
1028 }
1029 if (value instanceof Uint8Array) return writeBytes(value, path)
1030 if (typeof value !== 'object' || value === null) return value
1031 if (value instanceof Number || value instanceof String || value instanceof Boolean) return value.valueOf()
1032 if (ancestors.has(value)) throw new TypeError('gateway: circular RPC result')
1033 ancestors.add(value)
1034 let copy: object
1035 if (Array.isArray(value)) {
1036 const items: unknown[] = []
1037 // JSON arrays include every index, even holes and non-enumerable elements.
1038 for (let index = 0, length = value.length; index < length; index++) items.push(child(value[index], index))
1039 copy = items
1040 } else {
1041 const fields: Record<string, unknown> = {}
1042 for (const key of Object.keys(value)) {
1043 const item: unknown = Reflect.get(value, key)
1044 // A toJSON method on the projected object must not run a second time.
1045 if (key === 'toJSON' && typeof item === 'function') continue
1046 const extracted = child(item, key)
1047 if (key === '__proto__') Object.defineProperty(fields, key, { value: extracted, enumerable: true })
1048 else fields[key] = extracted
1049 }
1050 copy = fields
1051 }
1052 ancestors.delete(value)
1053 return copy
1054 }
1055 const child = (value: unknown, key: string | number): unknown => {
1056 if (typeof value !== 'object' || value === null) return value
1057 path.push(key)
1058 const extracted = extract(value, String(key))
1059 path.pop()
1060 return extracted
1061 }
1062 return extract(input, 'value')
1063}
1064
1065type RemoteEventWireFrame =
1066 | RemoteEventEmitFrame
1067 | RemoteEventInvocationFrame
1068 | RemoteEventCancellationFrame
1069
1070/** Pull-driven queue owned by one connected Client event generation. */
1071class RemoteEventQueue {
1072 private readonly frames = new Deque<RemoteEventWireFrame>()
1073 private waiter: (() => void) | undefined
1074 private closed = false
1075
1076 push(frame: RemoteEventWireFrame): void {
1077 if (this.closed) return
1078 this.frames.pushBack(frame)
1079 this.waiter?.()
1080 }
1081
1082 end(): void {
1083 if (this.closed) return
1084 this.closed = true
1085 this.waiter?.()
1086 }
1087
1088 async *iterate(signal: AbortSignal): AsyncGenerator<RemoteEventWireFrame> {
1089 const abort = (): void => { this.end() }
1090 signal.addEventListener('abort', abort, { once: true })
1091 try {
1092 while (true) {
1093 while (this.frames.size > 0) yield this.frames.popFront() as RemoteEventWireFrame
1094 if (this.closed || signal.aborted) return
1095 await new Promise<void>((resolve) => { this.waiter = resolve })
1096 this.waiter = undefined
1097 }
1098 } finally {
1099 signal.removeEventListener('abort', abort)
1100 }
1101 }
1102}
1103
1104function assertRemoteEventFrame(frame: TypertRemoteEventFrame): void {
1105 assertRemoteEventName(frame)
1106 if (!Array.isArray(frame.args) || !isRemoteJsonValue(frame.args)) {
1107 throw new TypeError(`typert gateway: Remote event ${JSON.stringify(frame.event)} arguments are not lossless JSON data`)
1108 }
1109}
1110
1111function assertRemoteEventName(frame: { readonly event: unknown }): void {
1112 if (typeof frame.event !== 'string' || frame.event.length === 0) {
1113 throw new TypeError('typert gateway: Remote event name must be a nonempty string')
1114 }
1115}
1116
1117function parseRemoteEventResultPayload(payload: unknown): ReturnType<typeof parseRemoteEventResult> {
1118 if (!isObject(payload)
1119 || !isPlainObject(payload)
1120 || Reflect.ownKeys(payload).length !== 1
1121 || !Object.hasOwn(payload, 'args')) {
1122 throw new Error('typert gateway: Remote event result requires exactly one plain-object args field')
1123 }
1124 return parseRemoteEventResult(payload.args)
1125}
1126
1127function remoteRequest(
1128 endpoint: string,
1129 payload: unknown,
1130 signal: AbortSignal,
1131 peer?: PeerScope,
1132): InvokeRemoteRequest {
1133 const segments = endpoint.split('/')
1134 if (segments.length !== 2 || segments[0] === '' || segments[1] === '') {
1135 throw new Error(`invalid Remote endpoint ${JSON.stringify(endpoint)}`)
1136 }
1137 const [namespace, method] = segments as [string, string]
1138 if (!isObject(payload)
1139 || !isPlainObject(payload)
1140 || Reflect.ownKeys(payload).length !== 1
1141 || !Object.hasOwn(payload, 'args')
1142 || !isObject(payload.args)
1143 || !isPlainObject(payload.args)) {
1144 throw new Error('Remote payload must contain exactly one plain-object args field')
1145 }
1146 return { namespace, method, args: payload.args, signal, ...(peer === undefined ? {} : { peer }) }
1147}
1148
1149/**
1150 * The signal a method observes. A carrier that supplies an uplink fails the
1151 * stream through `control` when an item is rejected, so that invocation joins
1152 * `control` with the carrier signal; every other invocation keeps the carrier
1153 * signal's identity.
1154 */
1155function methodSignal(request: InvokeRemoteRequest, control: AbortController): AbortSignal {
1156 const carrier = request.signal
1157 if (request.uplink === undefined) return carrier ?? NEVER_ABORTED_SIGNAL
1158 if (carrier === undefined || carrier === control.signal) return control.signal
1159 return AbortSignal.any([carrier, control.signal])
1160}
1161
1162function isIterable(value: unknown): value is Iterable<unknown> | AsyncIterable<unknown> {
1163 return isObject(value)
1164 && (typeof Reflect.get(value, Symbol.iterator) === 'function'
1165 || typeof Reflect.get(value, Symbol.asyncIterator) === 'function')
1166}
1167
1168async function *cancellableStream(
1169 source: Iterable<unknown> | AsyncIterable<unknown>,
1170 endpoint: string,
1171 invocation: GatewayInvocation,
1172): AsyncGenerator {
1173 const { signal } = invocation
1174 const asyncFactory: unknown = Reflect.get(source, Symbol.asyncIterator)
1175 const syncFactory: unknown = Reflect.get(source, Symbol.iterator)
1176 const iterator = typeof asyncFactory === 'function'
1177 ? Reflect.apply(asyncFactory, source, []) as AsyncIterator<unknown>
1178 : Reflect.apply(syncFactory as (...args: never[]) => Iterator<unknown>, source, [])
1179 let rejectAbort: ((error: unknown) => void) | undefined
1180 const onAbort = (): void => {
1181 rejectAbort?.(streamAbortFailure(endpoint, signal.reason))
1182 }
1183 signal.addEventListener('abort', onAbort, { once: true })
1184 try {
1185 while (true) {
1186 if (signal.aborted) throw streamAbortFailure(endpoint, signal.reason)
1187 // Reusing a pending cancellation promise retains every completed race.
1188 const aborted = new Promise<never>((_resolve, reject) => { rejectAbort = reject })
1189 // next() can abort and throw synchronously before the race subscribes.
1190 void aborted.catch(() => undefined)
1191 const next = await Promise.race([Promise.resolve(iterator.next()), aborted])
1192 rejectAbort = undefined
1193 if (next.done === true) return
1194 yield next.value
1195 }
1196 } finally {
1197 rejectAbort = undefined
1198 signal.removeEventListener('abort', onAbort)
1199 // The uplink closes first so a method blocked on `uplink.next()` unwinds
1200 // before its iterator is asked to return.
1201 await invocation.close()
1202 await iterator.return?.()
1203 }
1204}
1205
1206/**
1207 * The failure a stream surfaces for its abort: a Remote failure used as the
1208 * reason is the Gateway or carrier failing the stream itself (a rejected,
1209 * overflowing, or misplaced uplink item); any other abort is a cancellation.
1210 */
1211function streamAbortFailure(endpoint: string, reason: unknown): unknown {
1212 return remoteErrorOf(reason) === undefined ? remoteCancelled(endpoint, reason) : reason
1213}
1214
1215/** Carrier-signal cancellation as the shared failure vocabulary expresses it. */
1216function remoteCancelled(endpoint: string, cause: unknown): RemoteError<'gateway/cancelled'> {
1217 return new RemoteError('gateway/cancelled', `Remote invocation "${endpoint}" was aborted`, {}, { cause })
1218}
1219
1220/**
1221 * The iterable `invocation.uplink()` returns. Each uplink item passes the
1222 * descriptor codec, or the JSON-safety check when the descriptor declares no
1223 * uplink, before delivery; a rejected item fails the whole logical stream.
1224 * Iteration ends when the Client half-closes or the downlink finishes, and
1225 * fails when the stream is cancelled, so a method blocked on the uplink always
1226 * wakes.
1227 */
1228class UplinkDecoder implements AsyncIterable<unknown>, AsyncIterator<unknown> {
1229 private readonly source: AsyncIterator<unknown>
1230 private readonly interrupted = new Set<PromiseWithResolvers<IteratorResult<unknown>>>()
1231 private readonly onAbort = (): void => {
1232 const failure = streamAbortFailure(this.endpoint, this.signal.reason)
1233 for (const read of this.interrupted) read.reject(failure)
1234 }
1235 private closed = false
1236
1237 constructor(
1238 uplink: AsyncIterable<unknown>,
1239 private readonly codec: TypertCodec,
1240 private readonly endpoint: string,
1241 private readonly signal: AbortSignal,
1242 private readonly abort: (reason: unknown) => void,
1243 ) {
1244 this.source = uplink[Symbol.asyncIterator]()
1245 signal.addEventListener('abort', this.onAbort, { once: true })
1246 }
1247
1248 [Symbol.asyncIterator](): AsyncIterator<unknown> {
1249 return this
1250 }
1251
1252 async next(): Promise<IteratorResult<unknown>> {
1253 // A cancelled stream reports the cancellation on every read, closed or not.
1254 if (this.signal.aborted) throw streamAbortFailure(this.endpoint, this.signal.reason)
1255 if (this.closed) return UPLINK_DONE
1256 // Only pending reads belong to the decoder; finish() must wake all of them.
1257 const interrupted = Promise.withResolvers<IteratorResult<unknown>>()
1258 this.interrupted.add(interrupted)
1259 // The source can abort and throw synchronously before the race subscribes.
1260 void interrupted.promise.catch(() => undefined)
1261 let next: IteratorResult<unknown>
1262 try {
1263 next = await Promise.race([this.source.next(), interrupted.promise])
1264 } finally {
1265 this.interrupted.delete(interrupted)
1266 }
1267 if (next.done === true) {
1268 this.finish()
1269 return UPLINK_DONE
1270 }
1271 let value: unknown
1272 try {
1273 // A top-level `undefined` is the absent `value` of its frame: without a codec it is the item itself.
1274 value = next.value === undefined && this.codec.mode === 'src-json'
1275 ? undefined
1276 : decode(this.codec, next.value, this.endpoint, 'uplink')
1277 } catch (failure) {
1278 // The codec failure (`gateway/input-invalid`, field `uplink`) fails the whole logical stream.
1279 this.abort(failure)
1280 throw failure
1281 }
1282 return { value, done: false }
1283 }
1284
1285 /** Downlink finished or the method stopped reading: unread uplink items are dropped. */
1286 return(): Promise<IteratorResult<unknown>> {
1287 if (!this.closed) {
1288 this.finish()
1289 // The carrier owns the source iterator; a generator blocked in next()
1290 // completes this return once it yields, so it is not awaited here.
1291 void Promise.resolve().then(() => this.source.return?.()).catch(() => undefined)
1292 }
1293 return Promise.resolve(UPLINK_DONE)
1294 }
1295
1296 private finish(): void {
1297 this.closed = true
1298 this.signal.removeEventListener('abort', this.onAbort)
1299 for (const read of this.interrupted) read.resolve(UPLINK_DONE)
1300 }
1301}
1302
1303/**
1304 * The context of one Remote call as the receiving method reads it through
1305 * `this.ctx.invocation`. The uplink is decoded on first use and released when
1306 * the call's downlink finishes.
1307 */
1308class GatewayInvocation implements RemoteInvocation {
1309 private decoder: UplinkDecoder | undefined
1310 private taken = false
1311
1312 /**
1313 * @param request - decoded endpoint and wire arguments.
1314 * @param service - Cordis service key of the receiver.
1315 * @param peer - Peer the call speaks for.
1316 * @param signal - the signal the method observes.
1317 * @param uplink - carrier items and the codec that decodes them.
1318 */
1319 constructor(
1320 readonly request: RemoteInvocation['request'],
1321 readonly service: string,
1322 readonly peer: PeerScope,
1323 readonly signal: AbortSignal,
1324 private readonly uplink_: UplinkSource,
1325 ) {}
1326
1327 uplink<In = unknown>(): AsyncIterable<In> {
1328 if (this.taken) {
1329 throw new Error(`typert gateway: ${this.uplink_.endpoint}: invocation.uplink() is available once per call`)
1330 }
1331 this.taken = true
1332 const { source, codec, endpoint, abort } = this.uplink_
1333 this.decoder = new UplinkDecoder(source, codec, endpoint, this.signal, abort)
1334 // The descriptor codec decides what arrives; `In` is the caller's assertion.
1335 return this.decoder as AsyncIterable<In>
1336 }
1337
1338 /**
1339 * The downlink finished: release the uplink. Unread items are dropped, and a
1340 * carrier iterable the method never took is returned so it stops producing;
1341 * a later `uplink()` throws like a second one would.
1342 * @returns settles once a taken uplink has closed.
1343 */
1344 async close(): Promise<void> {
1345 if (this.decoder !== undefined) {
1346 await this.decoder.return()
1347 return
1348 }
1349 this.taken = true
1350 releaseUplink(this.uplink_.source)
1351 }
1352}
1353
1354/**
1355 * Return a carrier uplink nobody will read, so it drops later items instead of
1356 * buffering them. The carrier owns the iterator and `return()` is not awaited:
1357 * a generator blocked in `next()` completes it only once it yields.
1358 * @param source - the carrier's uplink iterable.
1359 */
1360function releaseUplink(source: AsyncIterable<unknown>): void {
1361 void Promise.resolve().then(() => source[Symbol.asyncIterator]().return?.()).catch(() => undefined)
1362}
1363
1364function rpcFailure(error: unknown): ConnectionRpcResult {
1365 const remote = remoteErrorOf(error)
1366 if (remote !== undefined) {
1367 return { ok: false, error: { code: remote.code, message: remote.message, details: remote.details } }
1368 }
1369 return {
1370 ok: false,
1371 error: {
1372 code: 'gateway/internal',
1373 message: error instanceof Error ? error.message : String(error),
1374 details: {},
1375 },
1376 }
1377}
1378
1379function rpcError(error: unknown): ConnectionRpcError & RemoteStreamFailure {
1380 return (rpcFailure(error) as Extract<ConnectionRpcResult, { readonly ok: false }>).error
1381}
1382
1383function endpointOf(namespace: string, method: string): string {
1384 return `${namespace}/${method}`
1385}
1386
1387function validateBinding(
1388 receiver: object,
1389 serviceKey: string,
1390 namespace: string,
1391 endpoint: string,
1392): ResolvedBinding {
1393 const original = originalOf(receiver)
1394 const value: unknown = Reflect.get(original, 'typertRemote')
1395 if (value === undefined) {
1396 throw new TypertGatewayError(
1397 'gateway/binding-invalid',
1398 endpoint,
1399 `Service ${JSON.stringify(serviceKey)} has no visible typertRemote binding`,
1400 )
1401 }
1402 return {
1403 binding: readBinding(value, original, serviceKey, endpoint, namespace),
1404 original,
1405 }
1406}
1407
1408function readBinding(
1409 value: unknown,
1410 original: object,
1411 serviceKey: string,
1412 endpoint: string,
1413 namespace?: string,
1414): TypertGatewayBinding {
1415 if (!isObject(value)
1416 || Reflect.get(value, 'service') !== original
1417 || Reflect.get(value, 'serviceKey') !== serviceKey
1418 || typeof Reflect.get(value, 'namespace') !== 'string'
1419 || (namespace !== undefined && Reflect.get(value, 'namespace') !== namespace)) {
1420 throw new TypertGatewayError(
1421 'gateway/binding-invalid',
1422 endpoint,
1423 `Service ${JSON.stringify(serviceKey)} has an inconsistent typertRemote binding`,
1424 )
1425 }
1426 return value as TypertGatewayBinding
1427}
1428
1429function originalOf(receiver: object): object {
1430 const original: unknown = Reflect.get(receiver, symbols.original)
1431 return isObject(original) ? original : receiver
1432}
1433
1434function methodParameterNames(service: object, method: string, endpoint: string): readonly string[] {
1435 let prototype: object | null = Object.getPrototypeOf(service) as object | null
1436 let implementation: ((this: object, ...args: never[]) => unknown) | undefined
1437 while (prototype !== null) {
1438 const descriptor = Object.getOwnPropertyDescriptor(prototype, method)
1439 if (descriptor !== undefined) {
1440 if ('value' in descriptor && typeof descriptor.value === 'function') {
1441 implementation = descriptor.value as (this: object, ...args: never[]) => unknown
1442 }
1443 break
1444 }
1445 prototype = Object.getPrototypeOf(prototype) as object | null
1446 }
1447 if (implementation === undefined) {
1448 throw new TypertGatewayError(
1449 'gateway/method-unavailable',
1450 endpoint,
1451 `Remote marker has no prototype method ${JSON.stringify(method)}`,
1452 )
1453 }
1454 const source = Function.prototype.toString.call(implementation)
1455 const open = source.indexOf('(')
1456 const close = source.indexOf(')', open + 1)
1457 /* v8 ignore next -- standard public class-method syntax always contains a parenthesized parameter list. */
1458 if (open < 0 || close < 0) return invalidSignature(endpoint, method)
1459 const body = source.slice(open + 1, close).trim()
1460 if (body.length === 0) return []
1461 const parts = body.split(',').map(part => part.trim())
1462 const names = new Set<string>()
1463 for (const part of parts) {
1464 if (!/^[$A-Z_a-z][$\w]*$/u.test(part) || names.has(part)) return invalidSignature(endpoint, method)
1465 names.add(part)
1466 }
1467 return [...names]
1468}
1469
1470function invalidSignature(endpoint: string, method: string): never {
1471 throw new TypertGatewayError(
1472 'gateway/signature-invalid',
1473 endpoint,
1474 `SRC method ${JSON.stringify(method)} must use unique identifier parameters without destructuring, defaults, or rest`,
1475 )
1476}
1477
1478function assertExactArguments(
1479 args: Readonly<Record<string, unknown>>,
1480 descriptor: InvocationDescriptor,
1481 endpoint: string,
1482): void {
1483 if (!isPlainObject(args)) {
1484 throw new TypertGatewayError('gateway/arguments-invalid', endpoint, 'args must be a plain object')
1485 }
1486 const expected = new Set(descriptor.parameters.map(parameter => parameter.wire))
1487 if (descriptor.invocation.kind === 'context') expected.add(descriptor.invocation.wire)
1488 const actual = Reflect.ownKeys(args)
1489 const extra = actual.filter(key => typeof key !== 'string' || !expected.has(key))
1490 // A JSON field may be omitted when the strict descriptor declares absence,
1491 // and always under SRC: a weak descriptor reads parameter names from the
1492 // JavaScript signature and cannot see which are optional, so LIB is where an
1493 // omitted required argument is caught. Lookup ids are never omissible.
1494 const acceptsMissing = new Set(descriptor.parameters
1495 .filter(parameter => parameter.source === 'json'
1496 && (parameter.acceptsUndefined === true || parameter.codec.mode === 'src-json'))
1497 .map(parameter => parameter.wire))
1498 const missing = [...expected].filter(key => !Object.hasOwn(args, key) && !acceptsMissing.has(key))
1499 if (extra.length === 0 && missing.length === 0) return
1500 const clauses: string[] = []
1501 if (missing.length > 0) clauses.push(`missing ${missing.map(key => JSON.stringify(key)).join(', ')}`)
1502 if (extra.length > 0) clauses.push(`unexpected ${extra.map(key => JSON.stringify(String(key))).join(', ')}`)
1503 throw new TypertGatewayError('gateway/arguments-invalid', endpoint, `args fields do not match the descriptor: ${clauses.join('; ')}`)
1504}
1505
1506function decode(
1507 codec: TypertCodec,
1508 value: unknown,
1509 endpoint: string,
1510 field: string,
1511): unknown {
1512 try {
1513 if (codec.mode === 'strict') {
1514 value = codec.create().parse(value)
1515 /* v8 ignore next -- generated optional-input codecs are the only strict codecs that return undefined. */
1516 if (value === undefined) return value
1517 }
1518 assertJsonValue(value, new Set())
1519 return value
1520 } catch (cause) {
1521 throw new TypertGatewayError(
1522 'gateway/input-invalid',
1523 endpoint,
1524 `wire field ${JSON.stringify(field)} failed boundary validation`,
1525 { cause, field },
1526 )
1527 }
1528}
1529
1530function assertJsonValue(value: unknown, ancestors: Set<object>): void {
1531 if (value === null || typeof value === 'string' || typeof value === 'boolean') return
1532 if (typeof value === 'number') {
1533 if (Number.isFinite(value)) return
1534 throw new TypeError('non-finite number is not JSON-safe')
1535 }
1536 if (!isObject(value)) throw new TypeError(`${typeof value} is not JSON-safe`)
1537 if (ancestors.has(value)) throw new TypeError('cyclic value is not JSON-safe')
1538 ancestors.add(value)
1539 try {
1540 if (Array.isArray(value)) {
1541 if (Object.getOwnPropertySymbols(value).length > 0 || Object.keys(value).length !== value.length) {
1542 throw new TypeError('sparse or decorated array is not JSON-safe')
1543 }
1544 for (let index = 0; index < value.length; index += 1) {
1545 if (!Object.hasOwn(value, index)) throw new TypeError('sparse array is not JSON-safe')
1546 assertJsonValue(value[index], ancestors)
1547 }
1548 return
1549 }
1550 if (!isPlainObject(value)) throw new TypeError('non-plain object is not JSON-safe')
1551 if (Object.getOwnPropertySymbols(value).length > 0) throw new TypeError('symbol property is not JSON-safe')
1552 for (const key of Reflect.ownKeys(value)) {
1553 const descriptor = Object.getOwnPropertyDescriptor(value, key)
1554 /* v8 ignore next -- ownKeys() just returned this key; only a hostile same-process Proxy can delete it between operations. */
1555 if (descriptor === undefined || !descriptor.enumerable || !('value' in descriptor)) {
1556 throw new TypeError('non-data property is not JSON-safe')
1557 }
1558 assertJsonValue(descriptor.value, ancestors)
1559 }
1560 } finally {
1561 ancestors.delete(value)
1562 }
1563}
1564
1565function isPlainObject(value: object): value is Record<string, unknown> {
1566 if (Array.isArray(value)) return false
1567 const prototype = Object.getPrototypeOf(value) as object | null
1568 return prototype === null || prototype === Object.prototype
1569}
1570
1571function isObject(value: unknown): value is object {
1572 return (typeof value === 'object' && value !== null) || typeof value === 'function'
1573}
1574
1575export default TypertGatewayService