1
/**2
* Live Typert Remote dispatch over Cordis Services and registered providers.3
* Unary transport and response envelopes belong to Connection; live Remote4
* streams use the Gateway-owned WebSocket mux.5
* @module @deepseek-ai/dsh-api-gateway6
*/8
import { randomUUID } from 'node:crypto'9
import { Context, Service, symbols } from '@deepseek-ai/cordis'10
import {11
OperatorPeer,12
type ConnectionRpcAttachment,13
type ConnectionRpcHandler,14
} from '@deepseek-ai/dsh-client-connection'15
import { Deque } from '@deepseek-ai/dsh-deque'16
import type { WebUpgradeRoute } from '@deepseek-ai/dsh-host-webserver'17
import { MAX_TIMER_DELAY_MS } from '@deepseek-ai/dsh-timeout'18
import type {} from '@deepseek-ai/dsh-cmdline'19
import z from '@deepseek-ai/schemastery'20
export type { TypertGatewayFaultDetails } from './remote-error-codes.ts'21
import {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'33
import type {34
InvokeRemoteRequest,35
TypertGateway,36
TypertGatewayErrorCode,37
TypertGatewayWireStream,38
TypertRemoteEventDispatch,39
TypertRemoteEventFrame,40
TypertRemoteEventInvocation,41
TypertRemoteEventOutcome,42
TypertRemoteEventSource,43
} from './types.ts'44
import {45
RemoteStreamMuxServer,46
rejectRemoteStreamUpgrade,47
} from './stream-server.ts'48
import {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'67
export 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'79
export type { RemoteEventHostInfo } from './stream-protocol.ts'81
interface GatewayErrorOptions {82
readonly cause?: unknown83
readonly field?: string84
}86
interface ResolvedBinding {87
readonly binding: TypertGatewayBinding88
readonly original: object89
}91
interface PreparedInvocation {92
readonly endpoint: string93
readonly descriptor: InvocationDescriptor94
/** Service view bound to a Context carrying `invocation`, so the method reads it as `this.ctx.invocation`. */95
readonly receiver: object96
readonly args: readonly unknown[]97
readonly method: (...args: never[]) => unknown98
readonly invocation: GatewayInvocation99
}101
/** Carrier inputs `GatewayInvocation.uplink()` decodes on first use. */102
interface 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: TypertCodec107
readonly endpoint: string108
/** Fail the logical stream with a Remote failure as the reason. */109
readonly abort: (reason: unknown) => void110
}112
interface RegisteredRemoteEventSource {113
readonly lifetime: AbortController114
readonly done: Promise<void>115
readonly host: RemoteEventHostInfo116
}118
interface RemoteEventClient {119
readonly id: RemoteEventClientId120
readonly queue: RemoteEventQueue121
readonly signal: AbortSignal122
readonly deliveries: Map<RemoteEventId, PendingRemoteEvent>123
}125
interface PendingRemoteEvent {126
readonly id: RemoteEventId127
readonly source: TypertRemoteEventInvocation128
readonly frame: RemoteEventInvocationFrame129
readonly deliveries: Set<RemoteEventClient>130
releaseContext: () => void131
releaseSignal: () => void132
}134
type ConnectionRpcResult = Awaited<ReturnType<ConnectionRpcHandler>>135
type ConnectionRpcError = Extract<ConnectionRpcResult, { readonly ok: false }>['error']136
const NEVER_ABORTED_SIGNAL = new AbortController().signal137
const DEFAULT_WEBSOCKET_HEARTBEAT_INTERVAL_MS = 2_000138
const DEFAULT_STREAM_INBOX_BYTES = 262_144139
const EMPTY_ASYNC_ITERABLE: AsyncIterable<never> = {140
[Symbol.asyncIterator]: () => ({ next: () => Promise.resolve({ value: undefined, done: true }) }),141
}142
const UPLINK_DONE: IteratorReturnResult<undefined> = { value: undefined, done: true }143
const SRC_JSON_CODEC: TypertCodec = { mode: 'src-json' }145
/** Gateway transport configuration. */146
export interface Config {147
/** WebSocket Ping interval from 1 through 2,147,483,647 milliseconds. @default 2000 */148
readonly websocketHeartbeatIntervalMs?: number149
/** Buffered uplink frame bytes one logical stream may hold before it fails with `gateway/uplink-overflow`. @default 262144 */150
readonly streamInboxBytes?: number151
}153
interface ResolvedConfig extends Config {154
readonly websocketHeartbeatIntervalMs: number155
readonly streamInboxBytes: number156
}158
/**159
* Dispatch failure produced outside the invoked business method. Rides the160
* shared Remote failure vocabulary, so its code crosses the wire instead of161
* folding to `internal`.162
*/163
export class TypertGatewayError extends RemoteError<TypertGatewayErrorCode> {164
/** Canonical `<namespace>/<method>` endpoint. */165
readonly endpoint: string166
/** Affected wire field when the failure is field-specific. */167
readonly field: string | undefined169
/**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 = endpoint190
this.field = options.field191
}192
}194
/**195
* Resolve strict generated definitions or conservative SRC markers against196
* current Cordis Services and Typert providers.197
* @typert service typertGateway198
*/199
export 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
})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
}214
private srcClaims: ReadonlySet<string> | undefined215
private inProcessOperator: PeerScope | undefined216
private remoteEvents: RegisteredRemoteEventSource | undefined217
private readonly remoteEventClients = new Map<RemoteEventClientId, RemoteEventClient>()218
private readonly pendingRemoteEvents = new Map<RemoteEventId, PendingRemoteEvent>()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 ResolvedConfig230
ctx.on('internal/service', () => {231
this.srcClaims = undefined232
})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
return258
}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 stream266
// 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 = false271
const cancel = ready.onReady(() => { if (!closed) listen() })272
return () => {273
closed = true274
cancel()275
}276
}, 'api-gateway: application readiness')277
})278
}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 true287
}288
return false289
}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) return308
this.closeRemoteEvents(error)309
this.remoteEvents = undefined310
lifetime.abort(error)311
})312
const registration: RegisteredRemoteEventSource = { lifetime, done, host: { home: host.home } }313
this.remoteEvents = registration314
return async () => {315
if (this.remoteEvents === registration) {316
this.remoteEvents = undefined317
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.done322
}323
}325
private claimsEndpoint(endpoint: string): boolean {326
if (endpoint === REMOTE_EVENT_RESULT_ENDPOINT) return true327
const segments = endpoint.split('/')328
if (segments.length !== 2 || segments[0] === '' || segments[1] === '') return false329
if (this.ctx.typert.local.get(endpoint) !== undefined || this.ctx.typert.local.hasSeen(endpoint)) return true330
this.srcClaims ??= this.collectSrcClaims()331
return this.srcClaims.has(endpoint)332
}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') continue338
const receiver: unknown = this.ctx.get(serviceKey)339
if (!isObject(receiver)) continue340
const original = originalOf(receiver)341
const binding: unknown = Reflect.get(original, 'typertRemote')342
if (!isObject(binding) || typeof Reflect.get(binding, 'namespace') !== 'string') continue343
const namespace = Reflect.get(binding, 'namespace') as string344
for (const candidate of remoteMethods(original)) {345
claims.add(endpointOf(namespace, candidate.exportName ?? candidate.method))346
}347
}348
return claims349
}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
}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
}370
try {371
return await Reflect.apply(prepared.method, prepared.receiver, prepared.args) as unknown372
} catch (error) {373
if (prepared.invocation.signal.aborted) throw remoteCancelled(prepared.endpoint, error)374
throw error375
} finally {376
// A unary call's uplink is readable only while the method runs.377
await prepared.invocation.close()378
}379
}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
}390
/**391
* `control` belongs to the logical stream: a rejected uplink item aborts it392
* 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: unknown405
try {406
source = Reflect.apply(prepared.method, prepared.receiver, prepared.args) as unknown407
} catch (error) {408
await prepared.invocation.close()409
if (prepared.invocation.signal.aborted) throw remoteCancelled(prepared.endpoint, error)410
throw error411
}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
}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
}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
}462
/**463
* The Peer an in-process carrier speaks for when it names none: the464
* operator's Peer when Connection is mounted, otherwise an operator scope the465
* 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.operator471
this.inProcessOperator ??= new OperatorPeer(this.ctx)472
return this.inProcessOperator473
}475
private async *openRemoteEvents(476
payload: unknown,477
signal: AbortSignal,478
): AsyncGenerator<479
RemoteEventEmitFrame | RemoteEventInvocationFrame | RemoteEventCancellationFrame480
| RemoteEventReadyFrame481
> {482
if (!isObject(payload)483
|| !isPlainObject(payload)484
|| Reflect.ownKeys(payload).length !== 1485
|| !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.remoteEvents496
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 RemoteEventClientId505
while (this.remoteEventClients.has(clientId)) clientId = randomUUID() as RemoteEventClientId506
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
}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
return530
}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
}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
}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 RemoteEventId559
while (this.pendingRemoteEvents.has(id)) id = randomUUID() as RemoteEventId560
let releaseContext: () => void561
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
return575
}576
const signals = new Set(projected.signal === undefined ? [] : [projected.signal])577
const abort = (): void => {578
const reason: unknown = [...signals].find(signal => signal.aborted)?.reason579
this.cancelRemoteEvent(pending, reason instanceof Error580
? reason581
: 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
}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
}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. Results620
// from a completed event or a superseded delivery are idempotent no-ops.621
if (pending === undefined || !pending.deliveries.has(client)) return622
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
}635
private removeRemoteEventDelivery(pending: PendingRemoteEvent, client: RemoteEventClient): void {636
pending.deliveries.delete(client)637
client.deliveries.delete(pending.id)638
}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
}646
private settleRemoteEvent(pending: PendingRemoteEvent, outcome: TypertRemoteEventOutcome): void {647
this.finishRemoteEvent(pending)648
pending.source.resolve(outcome)649
}651
private cancelRemoteEvent(pending: PendingRemoteEvent, reason: unknown): void {652
if (this.pendingRemoteEvents.get(pending.id) !== pending) return653
this.finishRemoteEvent(pending)654
pending.source.reject(reason)655
}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
}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
}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 one691
// 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
}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 view735
// resolves the Service the plain read above already found.736
const callReceiver = receiverContext.extend({ invocation }).get(descriptor.service) as object737
const implementation = descriptor.implementation ?? descriptor.method738
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
}756
private resolveDescriptor(namespace: string, method: string, endpoint: string): InvocationDescriptor {757
const strict = this.ctx.typert.local.get(endpoint)758
if (strict !== undefined) return strict759
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
}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') continue773
const receiver: unknown = this.ctx.get(serviceKey)774
if (!isObject(receiver)) continue775
const original = originalOf(receiver)776
const value: unknown = Reflect.get(original, 'typertRemote')777
if (value === undefined) continue778
const binding = readBinding(value, original, serviceKey, endpoint)779
if (binding.namespace !== namespace) continue780
const marker = remoteMethods(original).find(candidate => (candidate.exportName ?? candidate.method) === method)781
if (marker === undefined) continue782
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 InvocationDescriptor795
}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 >= 0814
? { parameter: 'signal' as const }815
: undefined816
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 === undefined832
? { 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
}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
}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
}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.ctx898
const invocation = descriptor.invocation899
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.wire908
|| (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 | undefined918
try {919
context = await provider.resolve(identity)920
} catch (cause) {921
if (remoteErrorOf(cause) !== undefined) throw cause922
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 context938
}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 parameter946
// takes undefined; a present-but-undefined field is not JSON-safe input and947
// still fails decode. Lookup ids are never omissible, so absence here only948
// ever belongs to a json parameter.949
if (!Object.hasOwn(args, parameter.wire)) return undefined950
const value = decode(parameter.codec, args[parameter.wire], endpoint, parameter.wire)951
if (parameter.source === 'json') return value952
const key = parameter.lookup953
/* 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.wire972
|| (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: unknown981
try {982
resolved = await provider.resolve(value)983
} catch (cause) {984
if (remoteErrorOf(cause) !== undefined) throw cause985
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 resolved1001
}1002
}1004
function 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 null1009
}1010
const encoded = codec.mode === 'strict'1011
? codec.encode?.(value, writeBytes) ?? value1012
: encodeRuntimeResult(value, writeBytes)1013
return { ok: true, value: encoded, ...(attachments.length === 0 ? {} : { attachments }) }1014
}1016
function 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 = input1025
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 value1031
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: object1035
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 = items1040
} 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') continue1046
const extracted = child(item, key)1047
if (key === '__proto__') Object.defineProperty(fields, key, { value: extracted, enumerable: true })1048
else fields[key] = extracted1049
}1050
copy = fields1051
}1052
ancestors.delete(value)1053
return copy1054
}1055
const child = (value: unknown, key: string | number): unknown => {1056
if (typeof value !== 'object' || value === null) return value1057
path.push(key)1058
const extracted = extract(value, String(key))1059
path.pop()1060
return extracted1061
}1062
return extract(input, 'value')1063
}1065
type RemoteEventWireFrame =1066
| RemoteEventEmitFrame1067
| RemoteEventInvocationFrame1068
| RemoteEventCancellationFrame1070
/** Pull-driven queue owned by one connected Client event generation. */1071
class RemoteEventQueue {1072
private readonly frames = new Deque<RemoteEventWireFrame>()1073
private waiter: (() => void) | undefined1074
private closed = false1076
push(frame: RemoteEventWireFrame): void {1077
if (this.closed) return1078
this.frames.pushBack(frame)1079
this.waiter?.()1080
}1082
end(): void {1083
if (this.closed) return1084
this.closed = true1085
this.waiter?.()1086
}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 RemoteEventWireFrame1094
if (this.closed || signal.aborted) return1095
await new Promise<void>((resolve) => { this.waiter = resolve })1096
this.waiter = undefined1097
}1098
} finally {1099
signal.removeEventListener('abort', abort)1100
}1101
}1102
}1104
function 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
}1111
function 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
}1117
function parseRemoteEventResultPayload(payload: unknown): ReturnType<typeof parseRemoteEventResult> {1118
if (!isObject(payload)1119
|| !isPlainObject(payload)1120
|| Reflect.ownKeys(payload).length !== 11121
|| !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
}1127
function 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 !== 11141
|| !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
}1149
/**1150
* The signal a method observes. A carrier that supplies an uplink fails the1151
* stream through `control` when an item is rejected, so that invocation joins1152
* `control` with the carrier signal; every other invocation keeps the carrier1153
* signal's identity.1154
*/1155
function methodSignal(request: InvokeRemoteRequest, control: AbortController): AbortSignal {1156
const carrier = request.signal1157
if (request.uplink === undefined) return carrier ?? NEVER_ABORTED_SIGNAL1158
if (carrier === undefined || carrier === control.signal) return control.signal1159
return AbortSignal.any([carrier, control.signal])1160
}1162
function 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
}1168
async function *cancellableStream(1169
source: Iterable<unknown> | AsyncIterable<unknown>,1170
endpoint: string,1171
invocation: GatewayInvocation,1172
): AsyncGenerator {1173
const { signal } = invocation1174
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) | undefined1180
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 = undefined1193
if (next.done === true) return1194
yield next.value1195
}1196
} finally {1197
rejectAbort = undefined1198
signal.removeEventListener('abort', onAbort)1199
// The uplink closes first so a method blocked on `uplink.next()` unwinds1200
// before its iterator is asked to return.1201
await invocation.close()1202
await iterator.return?.()1203
}1204
}1206
/**1207
* The failure a stream surfaces for its abort: a Remote failure used as the1208
* 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
*/1211
function streamAbortFailure(endpoint: string, reason: unknown): unknown {1212
return remoteErrorOf(reason) === undefined ? remoteCancelled(endpoint, reason) : reason1213
}1215
/** Carrier-signal cancellation as the shared failure vocabulary expresses it. */1216
function remoteCancelled(endpoint: string, cause: unknown): RemoteError<'gateway/cancelled'> {1217
return new RemoteError('gateway/cancelled', `Remote invocation "${endpoint}" was aborted`, {}, { cause })1218
}1220
/**1221
* The iterable `invocation.uplink()` returns. Each uplink item passes the1222
* descriptor codec, or the JSON-safety check when the descriptor declares no1223
* uplink, before delivery; a rejected item fails the whole logical stream.1224
* Iteration ends when the Client half-closes or the downlink finishes, and1225
* fails when the stream is cancelled, so a method blocked on the uplink always1226
* wakes.1227
*/1228
class 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 = false1237
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
}1248
[Symbol.asyncIterator](): AsyncIterator<unknown> {1249
return this1250
}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_DONE1256
// 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_DONE1270
}1271
let value: unknown1272
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
? undefined1276
: 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 failure1281
}1282
return { value, done: false }1283
}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
}1296
private finish(): void {1297
this.closed = true1298
this.signal.removeEventListener('abort', this.onAbort)1299
for (const read of this.interrupted) read.resolve(UPLINK_DONE)1300
}1301
}1303
/**1304
* The context of one Remote call as the receiving method reads it through1305
* `this.ctx.invocation`. The uplink is decoded on first use and released when1306
* the call's downlink finishes.1307
*/1308
class GatewayInvocation implements RemoteInvocation {1309
private decoder: UplinkDecoder | undefined1310
private taken = false1312
/**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
) {}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 = true1332
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
}1338
/**1339
* The downlink finished: release the uplink. Unread items are dropped, and a1340
* 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
return1348
}1349
this.taken = true1350
releaseUplink(this.uplink_.source)1351
}1352
}1354
/**1355
* Return a carrier uplink nobody will read, so it drops later items instead of1356
* 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
*/1360
function releaseUplink(source: AsyncIterable<unknown>): void {1361
void Promise.resolve().then(() => source[Symbol.asyncIterator]().return?.()).catch(() => undefined)1362
}1364
function 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
}1379
function rpcError(error: unknown): ConnectionRpcError & RemoteStreamFailure {1380
return (rpcFailure(error) as Extract<ConnectionRpcResult, { readonly ok: false }>).error1381
}1383
function endpointOf(namespace: string, method: string): string {1384
return `${namespace}/${method}`1385
}1387
function 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
}1408
function 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') !== original1417
|| Reflect.get(value, 'serviceKey') !== serviceKey1418
|| 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 TypertGatewayBinding1427
}1429
function originalOf(receiver: object): object {1430
const original: unknown = Reflect.get(receiver, symbols.original)1431
return isObject(original) ? original : receiver1432
}1434
function methodParameterNames(service: object, method: string, endpoint: string): readonly string[] {1435
let prototype: object | null = Object.getPrototypeOf(service) as object | null1436
let implementation: ((this: object, ...args: never[]) => unknown) | undefined1437
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[]) => unknown1442
}1443
break1444
}1445
prototype = Object.getPrototypeOf(prototype) as object | null1446
}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
}1470
function 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
}1478
function 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 the1492
// JavaScript signature and cannot see which are optional, so LIB is where an1493
// omitted required argument is caught. Lookup ids are never omissible.1494
const acceptsMissing = new Set(descriptor.parameters1495
.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) return1500
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
}1506
function 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 value1517
}1518
assertJsonValue(value, new Set())1519
return value1520
} 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
}1530
function assertJsonValue(value: unknown, ancestors: Set<object>): void {1531
if (value === null || typeof value === 'string' || typeof value === 'boolean') return1532
if (typeof value === 'number') {1533
if (Number.isFinite(value)) return1534
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
return1549
}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
}1565
function isPlainObject(value: object): value is Record<string, unknown> {1566
if (Array.isArray(value)) return false1567
const prototype = Object.getPrototypeOf(value) as object | null1568
return prototype === null || prototype === Object.prototype1569
}1571
function isObject(value: unknown): value is object {1572
return (typeof value === 'object' && value !== null) || typeof value === 'function'1573
}1575
export default TypertGatewayService