index.ts 6.1 KB

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