remote-events.host.spec.ts 7.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216
  1. import { Context } from '@deepseek-ai/cordis'
  2. import type { Fiber } from '@deepseek-ai/cordis'
  3. import type {
  4. TypertRemoteEventInvocation,
  5. TypertRemoteEventSource,
  6. } from '@deepseek-ai/dsh-api-gateway'
  7. import { scopeTarget } from '@deepseek-ai/dsh-scope'
  8. import { describe, expect, it } from 'vitest'
  9. import { apply, inject } from '../src/index.ts'
  10. interface GatewayProbe {
  11. source: TypertRemoteEventSource | undefined
  12. removals: number
  13. registerRemoteEvents(source: TypertRemoteEventSource): () => Promise<void>
  14. }
  15. async function setup(): Promise<{
  16. readonly ctx: Context
  17. readonly gateway: GatewayProbe
  18. readonly fiber: Fiber
  19. }> {
  20. const ctx = new Context()
  21. const gateway: GatewayProbe = {
  22. source: undefined,
  23. removals: 0,
  24. registerRemoteEvents(source) {
  25. gateway.source = source
  26. return async () => {
  27. if (gateway.source !== source) return
  28. gateway.source = undefined
  29. gateway.removals += 1
  30. }
  31. },
  32. }
  33. ctx.reflect.provide('typertGateway', gateway)
  34. const fiber = ctx.plugin({ inject: [...inject], apply })
  35. await fiber
  36. return { ctx, gateway, fiber }
  37. }
  38. function sourceOf(gateway: GatewayProbe): TypertRemoteEventSource {
  39. if (gateway.source === undefined) throw new Error('fixture Gateway has no Remote event source')
  40. return gateway.source
  41. }
  42. function emitRaw(ctx: Context, event: string, args: readonly unknown[]): void {
  43. const emit = ctx.emit.bind(ctx) as unknown as (name: string, ...values: readonly unknown[]) => void
  44. emit(event, ...args)
  45. }
  46. function waterfallRaw(
  47. ctx: Context,
  48. target: object,
  49. event: string,
  50. args: readonly unknown[],
  51. next: () => Promise<unknown>,
  52. ): Promise<unknown> {
  53. const waterfall = ctx.waterfall.bind(ctx) as unknown as (
  54. receiver: object,
  55. name: string,
  56. ...values: readonly unknown[]
  57. ) => Promise<unknown>
  58. return waterfall(target, event, ...args, next)
  59. }
  60. function invocationOf(value: unknown): TypertRemoteEventInvocation {
  61. if (typeof value !== 'object' || value === null || !Object.hasOwn(value, 'context')) {
  62. throw new Error('fixture did not receive a scoped Remote Event invocation')
  63. }
  64. return value as TypertRemoteEventInvocation
  65. }
  66. describe('Remote event Host source', () => {
  67. it('gives each Client stream an independent allowlisted event queue', async () => {
  68. const { ctx, gateway, fiber } = await setup()
  69. const firstAbort = new AbortController()
  70. const secondAbort = new AbortController()
  71. const first = sourceOf(gateway)(firstAbort.signal)[Symbol.asyncIterator]()
  72. const second = sourceOf(gateway)(secondAbort.signal)[Symbol.asyncIterator]()
  73. emitRaw(ctx, 'settings/document-updated', ['ui-theme', 1])
  74. await expect(first.next()).resolves.toEqual({
  75. done: false,
  76. value: { event: 'settings/document-updated', args: ['ui-theme', 1] },
  77. })
  78. await expect(second.next()).resolves.toEqual({
  79. done: false,
  80. value: { event: 'settings/document-updated', args: ['ui-theme', 1] },
  81. })
  82. const firstDone = first.next()
  83. firstAbort.abort(new Error('first Client disconnected'))
  84. emitRaw(ctx, 'commands/change', [])
  85. await expect(firstDone).resolves.toEqual({ done: true, value: undefined })
  86. await expect(second.next()).resolves.toEqual({
  87. done: false,
  88. value: { event: 'commands/change', args: [] },
  89. })
  90. const secondDone = second.next()
  91. secondAbort.abort(new Error('second Client disconnected'))
  92. await expect(secondDone).resolves.toEqual({ done: true, value: undefined })
  93. await fiber.dispose()
  94. expect(gateway.source).toBeUndefined()
  95. expect(gateway.removals).toBe(1)
  96. await ctx.fiber.dispose()
  97. })
  98. it('rejects a non-JSON argument without poisoning the stream', async () => {
  99. const { ctx, gateway } = await setup()
  100. const abort = new AbortController()
  101. const iterator = sourceOf(gateway)(abort.signal)[Symbol.asyncIterator]()
  102. const pending = iterator.next()
  103. expect(() => {
  104. emitRaw(ctx, 'settings/document-updated', ['ui-theme', 1n])
  105. }).toThrow('argument 1 is not lossless JSON data')
  106. emitRaw(ctx, 'settings/document-updated', ['ui-theme', 2])
  107. await expect(pending).resolves.toEqual({
  108. done: false,
  109. value: { event: 'settings/document-updated', args: ['ui-theme', 2] },
  110. })
  111. const done = iterator.next()
  112. abort.abort()
  113. await expect(done).resolves.toEqual({ done: true, value: undefined })
  114. const alreadyAborted = new AbortController()
  115. alreadyAborted.abort()
  116. await expect(sourceOf(gateway)(alreadyAborted.signal)[Symbol.asyncIterator]().next())
  117. .resolves.toEqual({ done: true, value: undefined })
  118. await ctx.fiber.dispose()
  119. })
  120. it('bridges scoped waterfall result, next delegation, and rejection', async () => {
  121. const { ctx, gateway } = await setup()
  122. const abort = new AbortController()
  123. const iterator = sourceOf(gateway)(abort.signal)[Symbol.asyncIterator]()
  124. const agentCtx = ctx.extend()
  125. const agent = { ctx: agentCtx }
  126. const target = scopeTarget(ctx, agent)
  127. const request = { questions: [], agent }
  128. const claimed = waterfallRaw(
  129. ctx,
  130. target,
  131. 'user-questions/request',
  132. [request],
  133. () => Promise.resolve('host fallback'),
  134. )
  135. const claimedDispatch = invocationOf((await iterator.next()).value)
  136. expect(claimedDispatch).toMatchObject({
  137. event: 'user-questions/request',
  138. request,
  139. context: { value: agentCtx, subject: agent },
  140. })
  141. claimedDispatch.resolve({ kind: 'result', value: 'client answer' })
  142. await expect(claimed).resolves.toBe('client answer')
  143. const delegated = waterfallRaw(
  144. ctx,
  145. target,
  146. 'user-questions/request',
  147. [request],
  148. () => Promise.resolve('host fallback'),
  149. )
  150. const delegatedDispatch = invocationOf((await iterator.next()).value)
  151. delegatedDispatch.resolve({ kind: 'next' })
  152. await expect(delegated).resolves.toBe('host fallback')
  153. const rejection = Object.assign(new Error('the user cancelled ask_user_question'), {
  154. code: 'ASK_CANCELLED',
  155. })
  156. const rejected = waterfallRaw(
  157. ctx,
  158. target,
  159. 'user-questions/request',
  160. [request],
  161. () => Promise.resolve('host fallback'),
  162. )
  163. const rejectedAssertion = expect(rejected).rejects.toBe(rejection)
  164. const rejectedDispatch = invocationOf((await iterator.next()).value)
  165. rejectedDispatch.reject(rejection)
  166. await rejectedAssertion
  167. const done = iterator.next()
  168. abort.abort()
  169. await expect(done).resolves.toEqual({ done: true, value: undefined })
  170. await ctx.fiber.dispose()
  171. })
  172. it('rejects a queued scoped waterfall when its source is withdrawn', async () => {
  173. const { ctx, gateway, fiber } = await setup()
  174. const abort = new AbortController()
  175. const iterator = sourceOf(gateway)(abort.signal)[Symbol.asyncIterator]()
  176. const delivery = iterator.next()
  177. const agent = { ctx: ctx.extend() }
  178. const reason = new Error('forwarded event source removed')
  179. const pending = waterfallRaw(
  180. ctx,
  181. scopeTarget(ctx, agent),
  182. 'user-questions/request',
  183. [{ questions: [], agent }],
  184. () => Promise.resolve('host fallback'),
  185. )
  186. const rejected = expect(pending).rejects.toBe(reason)
  187. abort.abort(reason)
  188. await rejected
  189. await expect(delivery).resolves.toEqual({ done: true, value: undefined })
  190. await fiber.dispose()
  191. await ctx.fiber.dispose()
  192. })
  193. })