| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164 |
- /** Host BFF entry and Loader shell for the Remote contribution assembly. */
- import type { Context } from '@deepseek-ai/cordis'
- import type {
- TypertRemoteEventDispatch,
- TypertRemoteEventInvocation,
- TypertRemoteEventOutcome,
- TypertRemoteEventSource,
- } from '@deepseek-ai/dsh-api-gateway'
- import { carrierKeyOf } from '@deepseek-ai/dsh-scope'
- import { isJsonValue } from '@deepseek-ai/dsh-session'
- import type { JsonValue } from '@deepseek-ai/dsh-session'
- import { API_REMOTE_FORWARDED_EVENTS } from './remote-events.ts'
- // The owner packages' client-safe `./types` exports carry the cordis `Events`
- // declarations for every allowlisted event. Pulling them into this face is what
- // makes the shape assertion below judge real signatures rather than an empty
- // event vocabulary.
- import type {} from '@deepseek-ai/dsh-commands/types'
- import type {} from '@deepseek-ai/dsh-cordis-host-runner/types'
- import type {} from '@deepseek-ai/dsh-credentials/types'
- import type {} from '@deepseek-ai/dsh-llm/types'
- import type {} from '@deepseek-ai/dsh-agent-presets/types'
- import type {} from '@deepseek-ai/dsh-settings/types'
- import type {} from '@deepseek-ai/dsh-user-approval'
- import type {} from '@deepseek-ai/dsh-user-questions'
- export type {} from '@deepseek-ai/dsh-api-session-controller/types'
- export { API_REMOTE_FORWARDED_EVENTS } from './remote-events.ts'
- export type { ApiRemoteForwardedEvent } from './types.ts'
- /** Required Host service: the Gateway owns the physical Remote stream mux. */
- export const inject = ['typertGateway']
- /** Host plugin body registering this application's selected Cordis event source. */
- export function apply(ctx: Context): void {
- ctx.effect(
- () => ctx.typertGateway.registerRemoteEvents(remoteEventSource(ctx)),
- 'api-remotes: forwarded Cordis event source',
- )
- }
- /** Create the sole queue and listener set consumed by the registered Gateway. */
- function remoteEventSource(ctx: Context): TypertRemoteEventSource {
- return (signal) => {
- const queue = new RemoteEventQueue()
- const disposers = API_REMOTE_FORWARDED_EVENTS.map(({ event, mode }) => {
- if (mode === 'emit') {
- return ctx.on(event as never, ((...args: unknown[]) => {
- queue.push({ event, args: assertJsonArgs(event, args) })
- }) as never)
- }
- return ctx.on(event as never, (function (
- this: unknown,
- request: object,
- next: () => unknown,
- ) {
- const subject = carrierKeyOf(this)
- if (subject === undefined) return next()
- const value = Reflect.get(subject, 'ctx') as unknown
- if (typeof value !== 'object' || value === null) {
- throw new TypeError(`forwarded scoped event ${JSON.stringify(event)} has no live Context`)
- }
- return forwardWaterfall(
- queue,
- event,
- request,
- { value: value as Context, subject },
- next,
- )
- }) as never)
- })
- return queue.iterate(signal, () => {
- for (const dispose of disposers) dispose()
- })
- }
- }
- /** One pull-driven queue bridging synchronous Cordis listeners to an AsyncIterable. */
- class RemoteEventQueue {
- private readonly buffer: TypertRemoteEventDispatch[] = []
- private waiter: (() => void) | undefined
- private done = false
- push(frame: TypertRemoteEventDispatch): boolean {
- if (this.done) return false
- this.buffer.push(frame)
- this.waiter?.()
- return true
- }
- private end(reason: unknown): void {
- if (this.done) return
- this.done = true
- const buffered = this.buffer.splice(0)
- for (const dispatch of buffered) {
- if ('context' in dispatch) dispatch.reject(reason)
- }
- this.waiter?.()
- }
- async *iterate(signal: AbortSignal, cleanup: () => void): AsyncGenerator<TypertRemoteEventDispatch> {
- const abort = (): void => { this.end(remoteEventSourceEndReason(signal)) }
- signal.addEventListener('abort', abort, { once: true })
- try {
- while (true) {
- if (this.done || signal.aborted) return
- while (this.buffer.length > 0) yield this.buffer.shift() as TypertRemoteEventDispatch
- await new Promise<void>((resolve) => { this.waiter = resolve })
- this.waiter = undefined
- }
- } finally {
- signal.removeEventListener('abort', abort)
- this.end(remoteEventSourceEndReason(signal))
- cleanup()
- }
- }
- }
- /**
- * Normalize an event-source shutdown for pending Host waterfalls.
- * @param signal - source lifetime whose reason wins after cancellation.
- * @returns the cancellation reason or an unexpected-end failure.
- */
- function remoteEventSourceEndReason(signal: AbortSignal): unknown {
- if (signal.aborted) return signal.reason
- return new Error('api-remotes: forwarded Remote event source ended')
- }
- /** Bridge one Cordis waterfall listener through the Gateway-owned pending event. */
- function forwardWaterfall(
- queue: RemoteEventQueue,
- event: string,
- request: object,
- context: TypertRemoteEventInvocation['context'],
- next: () => unknown,
- ): Promise<unknown> {
- const settled = Promise.withResolvers<unknown>()
- const dispatch: TypertRemoteEventInvocation = {
- event,
- request,
- context,
- resolve: (outcome: TypertRemoteEventOutcome) => {
- if (outcome.kind === 'result') {
- settled.resolve(outcome.value)
- return
- }
- void Promise.resolve().then(next).then(settled.resolve, settled.reject)
- },
- reject: settled.reject,
- }
- if (!queue.push(dispatch)) void Promise.resolve().then(next).then(settled.resolve, settled.reject)
- return settled.promise
- }
- /** Reject an allowlisted event whose runtime arguments are not lossless JSON data. */
- function assertJsonArgs(event: string, args: readonly unknown[]): JsonValue[] {
- for (const [index, arg] of args.entries()) {
- if (!isJsonValue(arg)) {
- throw new Error(`forwarded host event "${event}" argument ${String(index)} is not lossless JSON data`)
- }
- }
- return args as JsonValue[]
- }
|