client-apply.client.spec.ts 8.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237
  1. import { Context } from '@deepseek-ai/cordis'
  2. import type { Fiber } from '@deepseek-ai/cordis'
  3. import type {
  4. ConnectionHandle,
  5. HostDescription,
  6. } from '@deepseek-ai/dsh-client-connection/client'
  7. import {
  8. RemoteStreamCarrierError,
  9. RemoteStream,
  10. type RemoteStreamOptions,
  11. } from '@deepseek-ai/dsh-api-gateway/client'
  12. import type { SessionId } from '@deepseek-ai/dsh-session/types'
  13. import TypertRegistry from '@deepseek-ai/dsh-typert-registry'
  14. import { afterEach, describe, expect, it, vi } from 'vitest'
  15. import * as SessionClient from '../src/client/index.ts'
  16. import { ClientSessions } from '../src/client/sessions/service.ts'
  17. import { FakeApiClient, fakeRemote } from './fake-api.client.ts'
  18. const DESCRIPTION: HostDescription = {
  19. version: 'fixture',
  20. cwd: '/fixture',
  21. attachedSessions: 0,
  22. home: '/home/fixture',
  23. canOpenPath: true,
  24. }
  25. const sid = (value: string): SessionId => value as SessionId
  26. type RemoteListener = (...args: never[]) => void
  27. interface Bench {
  28. readonly ctx: Context
  29. readonly api: FakeApiClient
  30. readonly fiber: Fiber
  31. readonly sessions: ClientSessions
  32. dispatch(event: string, ...args: unknown[]): void
  33. publishHost(description: HostDescription | undefined): void
  34. }
  35. const contexts = new Set<Context>()
  36. afterEach(async () => {
  37. vi.restoreAllMocks()
  38. await Promise.all([...contexts].map(async (ctx) => { await ctx.fiber.dispose() }))
  39. contexts.clear()
  40. })
  41. async function mount(initialHost?: HostDescription): Promise<Bench> {
  42. const ctx = new Context()
  43. contexts.add(ctx)
  44. await ctx.plugin(TypertRegistry)
  45. const api = new FakeApiClient()
  46. const remote = fakeRemote(api)
  47. const listeners = new Map<string, Set<RemoteListener>>()
  48. const hostListeners = new Set<() => void>()
  49. let host = initialHost
  50. const connection: ConnectionHandle = {
  51. api,
  52. isLoopback: true,
  53. hostDescription: {
  54. getSnapshot: () => host,
  55. subscribe: (listener) => {
  56. hostListeners.add(listener)
  57. return () => { hostListeners.delete(listener) }
  58. },
  59. },
  60. generation: {
  61. getSnapshot: () => host === undefined
  62. ? undefined
  63. : { id: 1, host: { home: host.home } },
  64. subscribe: (listener) => {
  65. hostListeners.add(listener)
  66. return () => { hostListeners.delete(listener) }
  67. },
  68. },
  69. generation: {
  70. getSnapshot: () => host === undefined
  71. ? undefined
  72. : { id: 1, host: { home: host.home } },
  73. subscribe: (listener) => {
  74. hostListeners.add(listener)
  75. return () => { hostListeners.delete(listener) }
  76. },
  77. },
  78. rpc: {
  79. call: () => Promise.reject(new Error('unexpected generic RPC call')),
  80. },
  81. registerGenerationSource: () => () => {},
  82. start: () => ({ stop: () => {} }),
  83. }
  84. ctx.reflect.provide('connection', connection)
  85. ctx.reflect.provide('remote', {
  86. ...remote,
  87. $stream: <Item>(options: RemoteStreamOptions<Item>) => (
  88. new RemoteStream(connection, options)
  89. ),
  90. $on: (event: string, listener: RemoteListener) => {
  91. const eventListeners = listeners.get(event) ?? new Set<RemoteListener>()
  92. eventListeners.add(listener)
  93. listeners.set(event, eventListeners)
  94. return () => { eventListeners.delete(listener) }
  95. },
  96. })
  97. ctx.reflect.provide('remote.commands', remote.commands)
  98. ctx.reflect.provide('remote.session', remote.session)
  99. ctx.reflect.provide('remote.subagents', remote.subagents)
  100. const fiber = ctx.plugin(SessionClient)
  101. await fiber
  102. const sessions = ctx.sessions as ClientSessions
  103. return {
  104. ctx,
  105. api,
  106. fiber,
  107. sessions,
  108. dispatch: (event, ...args) => {
  109. for (const listener of listeners.get(event) ?? []) listener(...args as never[])
  110. },
  111. publishHost: (description) => {
  112. host = description
  113. for (const listener of [...hostListeners]) listener()
  114. },
  115. }
  116. }
  117. async function flush(): Promise<void> {
  118. for (let index = 0; index < 12; index++) await Promise.resolve()
  119. }
  120. describe('Session Controller Client apply', () => {
  121. it('routes Session Remote Events and connection generations into the object layer', async () => {
  122. const connected = vi.spyOn(ClientSessions.prototype, 'handleConnected')
  123. const error = vi.spyOn(ClientSessions.prototype, 'handleSessionError')
  124. const bench = await mount()
  125. expect(connected).not.toHaveBeenCalled()
  126. bench.dispatch('api-session/added', {
  127. sessionId: sid('session-1'),
  128. updatedAt: 1,
  129. running: false,
  130. blank: true,
  131. })
  132. await flush()
  133. expect(bench.sessions.list.getSnapshot().byId[sid('session-1')]).toMatchObject({
  134. running: false,
  135. updatedAt: 1,
  136. })
  137. bench.dispatch('api-session/status', sid('session-1'), true)
  138. bench.dispatch('api-session/activity', sid('session-1'), 9)
  139. bench.dispatch('api-session/error', sid('session-1'), 'agent failed')
  140. await flush()
  141. expect(bench.sessions.list.getSnapshot().byId[sid('session-1')]).toMatchObject({
  142. running: true,
  143. updatedAt: 9,
  144. })
  145. expect(error).toHaveBeenCalledWith(sid('session-1'), 'agent failed')
  146. bench.dispatch('api-session/removed', sid('session-1'))
  147. await flush()
  148. expect(bench.sessions.list.getSnapshot().byId[sid('session-1')]).toBeUndefined()
  149. bench.ctx.emit('connection/reset')
  150. expect(connected).toHaveBeenCalledOnce()
  151. })
  152. it('accepts the control baseline, retries a carrier generation, and reports terminal protocol failure', async () => {
  153. const accept = vi.spyOn(ClientSessions.prototype, 'handleControlFrame')
  154. const logged = vi.spyOn(console, 'error').mockImplementation(() => {})
  155. const bench = await mount(DESCRIPTION)
  156. await flush()
  157. expect(accept).toHaveBeenCalledWith({
  158. type: 'baseline',
  159. value: { queues: {}, jobs: {}, projections: {} },
  160. })
  161. bench.api.failStreams(new RemoteStreamCarrierError('generation lost'))
  162. await flush()
  163. expect(accept.mock.calls.filter(([frame]) => frame.type === 'baseline')).toHaveLength(2)
  164. bench.api.pushControl({ type: 'baseline', value: bench.api.controlBaseline } as never)
  165. await vi.waitFor(() => {
  166. expect(logged).toHaveBeenCalledWith(
  167. '[session-controller] control stream failed:',
  168. expect.objectContaining({ message: 'session control stream emitted more than one opening snapshot' }),
  169. )
  170. })
  171. })
  172. it('materializes Host-addressed Agent scopes before the Session list arrives', async () => {
  173. const bench = await mount()
  174. const adapter = bench.ctx.typert.contexts.getClient('agent')
  175. const first = adapter?.resolve(sid('agent-early'))
  176. expect(first).toBeDefined()
  177. expect(bench.sessions.scopeOf(first as Context)).toBe(sid('agent-early'))
  178. expect(adapter?.resolve(sid('agent-early'))).toBe(first)
  179. })
  180. it('projects Agent Context identity in both directions and withdraws the adapter on disposal', async () => {
  181. const bench = await mount(DESCRIPTION)
  182. await flush()
  183. expect(bench.sessions.list.getSnapshot().phase).toBe('ready')
  184. bench.dispatch('api-session/added', {
  185. sessionId: sid('agent-1'),
  186. updatedAt: 1,
  187. running: false,
  188. blank: true,
  189. })
  190. await flush()
  191. const scoped = bench.sessions.scope(sid('agent-1'))
  192. const adapter = bench.ctx.typert.contexts.getClient('agent')
  193. expect(scoped).toBeDefined()
  194. expect(adapter?.identity(bench.ctx)).toBeUndefined()
  195. expect(adapter?.identity(scoped!)).toBe(sid('agent-1'))
  196. expect(adapter?.resolve(sid('agent-1'))).toBe(scoped)
  197. await bench.fiber.dispose()
  198. expect(bench.ctx.typert.contexts.getClient('agent')).toBeUndefined()
  199. })
  200. it('waits for a Host generation before retrying the control stream', async () => {
  201. const accept = vi.spyOn(ClientSessions.prototype, 'handleControlFrame')
  202. const bench = await mount()
  203. await flush()
  204. expect(accept.mock.calls.filter(([frame]) => frame.type === 'baseline')).toHaveLength(1)
  205. bench.api.failStreams(new RemoteStreamCarrierError('offline'))
  206. await flush()
  207. expect(accept.mock.calls.filter(([frame]) => frame.type === 'baseline')).toHaveLength(1)
  208. bench.publishHost(DESCRIPTION)
  209. await flush()
  210. expect(accept.mock.calls.filter(([frame]) => frame.type === 'baseline')).toHaveLength(2)
  211. })
  212. })