manager.spec.ts 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223
  1. /**
  2. * SessionManager orchestration: lazy resident instances, list lifecycle, host
  3. * frame routing, and the pending-frame buffer for uninstantiated sessions.
  4. */
  5. import { describe, expect, it, vi } from 'vitest'
  6. import type { SessionId } from '@deepseek-ai/dsh-client-connection/client'
  7. import { SessionManager } from '../src/client/sessions/manager.ts'
  8. import { FakeApiClient, deferred, err, ok } from './fake-api.ts'
  9. import { entries, plainTurn } from './event-script.ts'
  10. const S1 = 'fk-m1' as SessionId
  11. const S2 = 'fk-m2' as SessionId
  12. function summary(sessionId: SessionId, over: Partial<{ updatedAt: number; running: boolean; parentSessionId: SessionId }> = {}) {
  13. return { sessionId, updatedAt: 100, running: false, ...over }
  14. }
  15. describe('instances', () => {
  16. it('lazily builds one resident instance per id and syncs the running bit from the list', async () => {
  17. const api = new FakeApiClient()
  18. api.onList = () => Promise.resolve(ok({ items: [summary(S1, { running: true })] as never[] }))
  19. const manager = new SessionManager(api)
  20. await manager.refreshList()
  21. const session = manager.get(S1)
  22. expect(manager.get(S1)).toBe(session) // resident: same instance forever
  23. expect(session.getSnapshot().running).toBe(true) // list preceded instantiation
  24. })
  25. it('replays buffered approval frames on instantiation and drops ordinary frames for uninstantiated sessions', () => {
  26. const api = new FakeApiClient()
  27. const manager = new SessionManager(api)
  28. // Uninstantiated: approval buffers, plain session/event drops.
  29. manager.handleMuxEnvelope({ rpcId: 'ra' as never, payload: { type: 'approval/requested', sessionId: S1, approvalId: 'ap1' as never, toolName: 'rm' } })
  30. manager.handleMuxEnvelope({ rpcId: 're' as never, payload: { type: 'session/event', sessionId: S1, event: plainTurn(0, 0, 'x', 'y')[0] as never } })
  31. const session = manager.get(S1)
  32. expect(session.getSnapshot().pending).toMatchObject([{ kind: 'approval', approvalId: 'ap1' }])
  33. // Buffer cleared: a second instantiation of another id gets nothing.
  34. expect(manager.get(S2).getSnapshot().pending).toEqual([])
  35. })
  36. it('caps the pending buffer at 32 keeping the newest, and drops it on session-removed', () => {
  37. const api = new FakeApiClient()
  38. const manager = new SessionManager(api)
  39. // 40 distinct question frames for an uninstantiated session: only the newest 32 survive.
  40. for (let i = 0; i < 40; i++) {
  41. manager.handleMuxEnvelope({ rpcId: `q${i}` as never, payload: { type: 'question/requested', sessionId: S1, questions: [] } })
  42. }
  43. const pending = manager.get(S1).getSnapshot().pending
  44. expect(pending).toHaveLength(32)
  45. expect(pending.map(p => p.rpcId)).toEqual(Array.from({ length: 32 }, (_, i) => `q${i + 8}`)) // oldest 8 dropped
  46. // Removed session: buffered frames must not replay on a future instantiation.
  47. manager.handleMuxEnvelope({ rpcId: 'qz' as never, payload: { type: 'question/requested', sessionId: S2, questions: [] } })
  48. manager.handleHostEnvelope({ rpcId: 'hz' as never, payload: { type: 'host/session-removed', sessionId: S2 } })
  49. expect(manager.get(S2).getSnapshot().pending).toEqual([])
  50. })
  51. })
  52. describe('list lifecycle', () => {
  53. it('single-flights refreshList and lands items sorted through lineage flattening', async () => {
  54. const api = new FakeApiClient()
  55. const gate = deferred<Awaited<ReturnType<FakeApiClient['onList']>>>()
  56. api.onList = () => gate.promise
  57. const manager = new SessionManager(api)
  58. const first = manager.refreshList()
  59. const second = manager.refreshList()
  60. expect(manager.getListSnapshot().state).toBe('loading')
  61. gate.resolve(ok({ items: [summary(S1), summary(S2, { updatedAt: 200 })] as never[] }))
  62. await Promise.all([first, second])
  63. expect(api.callsOf('session.list')).toHaveLength(1)
  64. const snapshot = manager.getListSnapshot()
  65. expect(snapshot.state).toBe('idle')
  66. expect(snapshot.items.map(i => i.sessionId)).toEqual([S2, S1]) // updatedAt desc
  67. })
  68. it('keeps the error in the list snapshot on failure', async () => {
  69. const api = new FakeApiClient()
  70. api.onList = () => Promise.resolve(err({ code: 'internal', message: 'boom', details: {} }))
  71. const manager = new SessionManager(api)
  72. await manager.refreshList()
  73. expect(manager.getListSnapshot()).toMatchObject({ state: 'error', error: { code: 'internal' } })
  74. })
  75. it('merges create into the list immediately without waiting for a refresh', async () => {
  76. const api = new FakeApiClient()
  77. api.onCreate = () => Promise.resolve(ok({ sessionId: S2 }))
  78. const manager = new SessionManager(api)
  79. const result = await manager.create()
  80. expect(result).toMatchObject({ ok: true, value: { sessionId: S2 } })
  81. expect(manager.getListSnapshot().items.map(i => i.sessionId)).toEqual([S2])
  82. })
  83. })
  84. describe('host frame routing', () => {
  85. it('adds/removes/flips sessions from host frames and keeps removed instances resident', async () => {
  86. const api = new FakeApiClient()
  87. const manager = new SessionManager(api)
  88. manager.handleHostEnvelope({ rpcId: 'h1' as never, payload: { type: 'host/session-added', sessionId: S1 } })
  89. manager.handleHostEnvelope({ rpcId: 'h2' as never, payload: { type: 'host/session-added', sessionId: S1 } }) // dup: ignored
  90. expect(manager.getListSnapshot().items).toHaveLength(1)
  91. const session = manager.get(S1)
  92. manager.handleHostEnvelope({ rpcId: 'h3' as never, payload: { type: 'host/session-status', sessionId: S1, running: true } })
  93. expect(session.getSnapshot().running).toBe(true)
  94. expect(manager.getListSnapshot().items[0]?.running).toBe(true)
  95. manager.handleHostEnvelope({ rpcId: 'h4' as never, payload: { type: 'host/agent-error', sessionId: S1, message: '炸了' } })
  96. expect(session.getSnapshot().lastAgentError).toBe('炸了')
  97. manager.handleHostEnvelope({ rpcId: 'h5' as never, payload: { type: 'host/session-removed', sessionId: S1 } })
  98. expect(manager.getListSnapshot().items).toHaveLength(0)
  99. expect(session.getSnapshot().removed).toBe(true)
  100. expect(manager.get(S1)).toBe(session) // resident-instance rule survives removal
  101. })
  102. })
  103. describe('remaining branches', () => {
  104. it('refreshList folds a transport throw into the error state', async () => {
  105. const api = new FakeApiClient()
  106. api.onList = () => Promise.reject(new Error('list wire down'))
  107. const manager = new SessionManager(api)
  108. await manager.refreshList()
  109. expect(manager.getListSnapshot()).toMatchObject({ state: 'error', error: { code: 'internal', message: 'list wire down' } })
  110. })
  111. it('refreshList pushes running bits down to already-instantiated sessions', async () => {
  112. const api = new FakeApiClient()
  113. const manager = new SessionManager(api)
  114. const session = manager.get(S1)
  115. api.onList = () => Promise.resolve(ok({ items: [summary(S1, { running: true })] as never[] }))
  116. await manager.refreshList()
  117. expect(session.getSnapshot().running).toBe(true)
  118. })
  119. it('create passes cwd through, folds transport throws, and skips the merge when already listed', async () => {
  120. const api = new FakeApiClient()
  121. api.onCreate = () => Promise.resolve(ok({ sessionId: S1 }))
  122. const manager = new SessionManager(api)
  123. await manager.create('/tmp/w')
  124. expect(api.callsOf('session.create')).toEqual([{ cwd: '/tmp/w' }])
  125. expect(manager.getListSnapshot().items[0]).toMatchObject({ sessionId: S1, cwd: '/tmp/w' })
  126. await manager.create('/tmp/w') // same id returned: no duplicate row
  127. expect(manager.getListSnapshot().items).toHaveLength(1)
  128. api.onCreate = () => Promise.reject(new Error('create wire down'))
  129. expect(await manager.create()).toMatchObject({ ok: false, error: { code: 'internal' } })
  130. // Business error passes through untouched.
  131. api.onCreate = () => Promise.resolve(err({ code: 'internal', message: 'no', details: {} }))
  132. expect(await manager.create()).toMatchObject({ ok: false })
  133. })
  134. it('subscribe notifies on list changes and stops after unsubscribe', async () => {
  135. const api = new FakeApiClient()
  136. const manager = new SessionManager(api)
  137. let notified = 0
  138. const unsubscribe = manager.subscribe(() => { notified++ })
  139. await manager.refreshList()
  140. await new Promise(resolve => setTimeout(resolve, 0))
  141. expect(notified).toBeGreaterThan(0)
  142. const seen = notified
  143. unsubscribe()
  144. manager.handleHostEnvelope({ rpcId: 'h' as never, payload: { type: 'host/session-added', sessionId: S1 } })
  145. await new Promise(resolve => setTimeout(resolve, 0))
  146. expect(notified).toBe(seen)
  147. })
  148. it('routes stream/error and unknown frames to the documented drops, and dispatches to instantiated sessions', () => {
  149. const api = new FakeApiClient()
  150. const manager = new SessionManager(api)
  151. manager.handleMuxEnvelope({ rpcId: 'e' as never, payload: { type: 'stream/error', error: { code: 'internal', message: 'x', details: {} } } })
  152. manager.handleHostEnvelope({ rpcId: 'e2' as never, payload: { type: 'stream/error', error: { code: 'internal', message: 'x', details: {} } } })
  153. manager.handleHostEnvelope({ rpcId: 'e3' as never, payload: { type: 'future/host-frame' } as never })
  154. const session = manager.get(S1)
  155. manager.handleMuxEnvelope({ rpcId: 'q1' as never, payload: { type: 'question/requested', sessionId: S1, questions: [] } })
  156. expect(session.getSnapshot().pending).toMatchObject([{ kind: 'question' }])
  157. // status flip for an unknown session only touches summaries (no crash).
  158. manager.handleHostEnvelope({ rpcId: 'h9' as never, payload: { type: 'host/session-status', sessionId: S2, running: true } })
  159. manager.handleHostEnvelope({ rpcId: 'ha' as never, payload: { type: 'host/agent-error', sessionId: S2, message: '无实例' } })
  160. })
  161. it('keeps list-entry identity for unchanged rows across an unrelated list change', async () => {
  162. const api = new FakeApiClient()
  163. api.onList = () => Promise.resolve(ok({ items: [summary(S1), summary(S2, { updatedAt: 200 })] as never[] }))
  164. const manager = new SessionManager(api)
  165. await manager.refreshList()
  166. const before = manager.getListSnapshot()
  167. manager.handleHostEnvelope({ rpcId: 'h' as never, payload: { type: 'host/session-status', sessionId: S2, running: true } })
  168. const after = manager.getListSnapshot()
  169. expect(after.items).not.toBe(before.items)
  170. const beforeS1 = before.items.find(e => e.sessionId === S1)
  171. const afterS1 = after.items.find(e => e.sessionId === S1)
  172. expect(afterS1).toBe(beforeS1) // untouched entry keeps identity (entryCache)
  173. // Same-order same-entries snapshot reuses the items array.
  174. manager.handleHostEnvelope({ rpcId: 'h2' as never, payload: { type: 'host/agent-error', sessionId: S1, message: 'x' } })
  175. expect(manager.getListSnapshot().items).toBe(after.items)
  176. })
  177. it('carries parentSessionId from host/session-added into the lineage row', () => {
  178. const api = new FakeApiClient()
  179. const manager = new SessionManager(api)
  180. manager.handleHostEnvelope({ rpcId: 'h1' as never, payload: { type: 'host/session-added', sessionId: S1 } })
  181. manager.handleHostEnvelope({ rpcId: 'h2' as never, payload: { type: 'host/session-added', sessionId: S2, parentSessionId: S1 } })
  182. const items = manager.getListSnapshot().items
  183. expect(items.find(e => e.sessionId === S2)).toMatchObject({ parentSessionId: S1, depth: 1 })
  184. })
  185. })
  186. describe('connected generation', () => {
  187. it('refreshes the list and resyncs only opened instances', async () => {
  188. const api = new FakeApiClient()
  189. api.onHistory = () => Promise.resolve(ok({ events: entries(plainTurn(0, 0, 'a', 'b')) as never[], hasMore: false }))
  190. const manager = new SessionManager(api)
  191. const openedSession = manager.get(S1)
  192. await openedSession.open()
  193. manager.get(S2) // instantiated but never opened
  194. const historyCallsBefore = api.callsOf('session.history').length
  195. manager.handleConnected()
  196. await vi.waitFor(() => {
  197. expect(api.callsOf('session.list').length).toBe(1)
  198. // Only the opened instance repulls history; the cold one stays silent.
  199. expect(api.callsOf('session.history').length).toBe(historyCallsBefore + 1)
  200. })
  201. })
  202. })