index.ts 5.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164
  1. /** Host BFF entry and Loader shell for the Remote contribution assembly. */
  2. import type { Context } from '@deepseek-ai/cordis'
  3. import type {
  4. TypertRemoteEventDispatch,
  5. TypertRemoteEventInvocation,
  6. TypertRemoteEventOutcome,
  7. TypertRemoteEventSource,
  8. } from '@deepseek-ai/dsh-api-gateway'
  9. import { carrierKeyOf } from '@deepseek-ai/dsh-scope'
  10. import { isJsonValue } from '@deepseek-ai/dsh-session'
  11. import type { JsonValue } from '@deepseek-ai/dsh-session'
  12. import { API_REMOTE_FORWARDED_EVENTS } from './remote-events.ts'
  13. // The owner packages' client-safe `./types` exports carry the cordis `Events`
  14. // declarations for every allowlisted event. Pulling them into this face is what
  15. // makes the shape assertion below judge real signatures rather than an empty
  16. // event vocabulary.
  17. import type {} from '@deepseek-ai/dsh-commands/types'
  18. import type {} from '@deepseek-ai/dsh-cordis-host-runner/types'
  19. import type {} from '@deepseek-ai/dsh-credentials/types'
  20. import type {} from '@deepseek-ai/dsh-llm/types'
  21. import type {} from '@deepseek-ai/dsh-agent-presets/types'
  22. import type {} from '@deepseek-ai/dsh-settings/types'
  23. import type {} from '@deepseek-ai/dsh-user-approval'
  24. import type {} from '@deepseek-ai/dsh-user-questions'
  25. export type {} from '@deepseek-ai/dsh-api-session-controller/types'
  26. export { API_REMOTE_FORWARDED_EVENTS } from './remote-events.ts'
  27. export type { ApiRemoteForwardedEvent } from './types.ts'
  28. /** Required Host service: the Gateway owns the physical Remote stream mux. */
  29. export const inject = ['typertGateway']
  30. /** Host plugin body registering this application's selected Cordis event source. */
  31. export function apply(ctx: Context): void {
  32. ctx.effect(
  33. () => ctx.typertGateway.registerRemoteEvents(remoteEventSource(ctx)),
  34. 'api-remotes: forwarded Cordis event source',
  35. )
  36. }
  37. /** Create the sole queue and listener set consumed by the registered Gateway. */
  38. function remoteEventSource(ctx: Context): TypertRemoteEventSource {
  39. return (signal) => {
  40. const queue = new RemoteEventQueue()
  41. const disposers = API_REMOTE_FORWARDED_EVENTS.map(({ event, mode }) => {
  42. if (mode === 'emit') {
  43. return ctx.on(event as never, ((...args: unknown[]) => {
  44. queue.push({ event, args: assertJsonArgs(event, args) })
  45. }) as never)
  46. }
  47. return ctx.on(event as never, (function (
  48. this: unknown,
  49. request: object,
  50. next: () => unknown,
  51. ) {
  52. const subject = carrierKeyOf(this)
  53. if (subject === undefined) return next()
  54. const value = Reflect.get(subject, 'ctx') as unknown
  55. if (typeof value !== 'object' || value === null) {
  56. throw new TypeError(`forwarded scoped event ${JSON.stringify(event)} has no live Context`)
  57. }
  58. return forwardWaterfall(
  59. queue,
  60. event,
  61. request,
  62. { value: value as Context, subject },
  63. next,
  64. )
  65. }) as never)
  66. })
  67. return queue.iterate(signal, () => {
  68. for (const dispose of disposers) dispose()
  69. })
  70. }
  71. }
  72. /** One pull-driven queue bridging synchronous Cordis listeners to an AsyncIterable. */
  73. class RemoteEventQueue {
  74. private readonly buffer: TypertRemoteEventDispatch[] = []
  75. private waiter: (() => void) | undefined
  76. private done = false
  77. push(frame: TypertRemoteEventDispatch): boolean {
  78. if (this.done) return false
  79. this.buffer.push(frame)
  80. this.waiter?.()
  81. return true
  82. }
  83. private end(reason: unknown): void {
  84. if (this.done) return
  85. this.done = true
  86. const buffered = this.buffer.splice(0)
  87. for (const dispatch of buffered) {
  88. if ('context' in dispatch) dispatch.reject(reason)
  89. }
  90. this.waiter?.()
  91. }
  92. async *iterate(signal: AbortSignal, cleanup: () => void): AsyncGenerator<TypertRemoteEventDispatch> {
  93. const abort = (): void => { this.end(remoteEventSourceEndReason(signal)) }
  94. signal.addEventListener('abort', abort, { once: true })
  95. try {
  96. while (true) {
  97. if (this.done || signal.aborted) return
  98. while (this.buffer.length > 0) yield this.buffer.shift() as TypertRemoteEventDispatch
  99. await new Promise<void>((resolve) => { this.waiter = resolve })
  100. this.waiter = undefined
  101. }
  102. } finally {
  103. signal.removeEventListener('abort', abort)
  104. this.end(remoteEventSourceEndReason(signal))
  105. cleanup()
  106. }
  107. }
  108. }
  109. /**
  110. * Normalize an event-source shutdown for pending Host waterfalls.
  111. * @param signal - source lifetime whose reason wins after cancellation.
  112. * @returns the cancellation reason or an unexpected-end failure.
  113. */
  114. function remoteEventSourceEndReason(signal: AbortSignal): unknown {
  115. if (signal.aborted) return signal.reason
  116. return new Error('api-remotes: forwarded Remote event source ended')
  117. }
  118. /** Bridge one Cordis waterfall listener through the Gateway-owned pending event. */
  119. function forwardWaterfall(
  120. queue: RemoteEventQueue,
  121. event: string,
  122. request: object,
  123. context: TypertRemoteEventInvocation['context'],
  124. next: () => unknown,
  125. ): Promise<unknown> {
  126. const settled = Promise.withResolvers<unknown>()
  127. const dispatch: TypertRemoteEventInvocation = {
  128. event,
  129. request,
  130. context,
  131. resolve: (outcome: TypertRemoteEventOutcome) => {
  132. if (outcome.kind === 'result') {
  133. settled.resolve(outcome.value)
  134. return
  135. }
  136. void Promise.resolve().then(next).then(settled.resolve, settled.reject)
  137. },
  138. reject: settled.reject,
  139. }
  140. if (!queue.push(dispatch)) void Promise.resolve().then(next).then(settled.resolve, settled.reject)
  141. return settled.promise
  142. }
  143. /** Reject an allowlisted event whose runtime arguments are not lossless JSON data. */
  144. function assertJsonArgs(event: string, args: readonly unknown[]): JsonValue[] {
  145. for (const [index, arg] of args.entries()) {
  146. if (!isJsonValue(arg)) {
  147. throw new Error(`forwarded host event "${event}" argument ${String(index)} is not lossless JSON data`)
  148. }
  149. }
  150. return args as JsonValue[]
  151. }