manager.spec.ts 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283
  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', payload: { 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.key)).toEqual(Array.from({ length: 32 }, (_, i) => `q: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. it('retains monotonic title snapshots before list arrival, merges recency, and clears them on removal', async () => {
  84. const api = new FakeApiClient()
  85. const manager = new SessionManager(api)
  86. manager.handleMuxEnvelope({
  87. rpcId: 'title-new' as never,
  88. payload: { type: 'session/title', sessionId: S1, title: 'Newest', eventSeq: 4, updatedAt: 300 },
  89. })
  90. manager.handleMuxEnvelope({
  91. rpcId: 'title-stale' as never,
  92. payload: { type: 'session/title', sessionId: S1, title: 'Stale', eventSeq: 3, updatedAt: 900 },
  93. })
  94. manager.handleMuxEnvelope({
  95. rpcId: 'title-equal' as never,
  96. payload: { type: 'session/title', sessionId: S1, title: 'Equal', eventSeq: 4, updatedAt: 901 },
  97. })
  98. api.onList = () => Promise.resolve(ok({
  99. items: [summary(S1), summary(S2, { updatedAt: 200 })] as never[],
  100. }))
  101. await manager.refreshList()
  102. const titled = manager.getListSnapshot()
  103. expect(titled.items.map(item => item.sessionId)).toEqual([S1, S2])
  104. expect(titled.items[0]).toMatchObject({ title: 'Newest', updatedAt: 300 })
  105. expect(titled.items[1]?.title).toBeUndefined()
  106. manager.handleHostEnvelope({ rpcId: 'removed' as never, payload: { type: 'host/session-removed', sessionId: S1 } })
  107. manager.handleHostEnvelope({ rpcId: 'readded' as never, payload: { type: 'host/session-added', sessionId: S1 } })
  108. expect(manager.getListSnapshot().items.find(item => item.sessionId === S1)?.title).toBeUndefined()
  109. })
  110. it('drops a retained title beyond the subscription baseline before accepting its durable replay', async () => {
  111. const api = new FakeApiClient()
  112. api.onList = () => Promise.resolve(ok({ items: [summary(S1)] as never[] }))
  113. const manager = new SessionManager(api)
  114. await manager.refreshList()
  115. manager.handleMuxEnvelope({
  116. rpcId: 'title-unflushed' as never,
  117. payload: { type: 'session/title', sessionId: S1, title: 'Unflushed', eventSeq: 4, updatedAt: 400 },
  118. })
  119. manager.handleMuxEnvelope({
  120. rpcId: 'subscribed-recovered' as never,
  121. payload: { type: 'session/subscribed', sessionId: S1, lastSeq: 2 },
  122. })
  123. expect(manager.getListSnapshot().items[0]?.title).toBeUndefined()
  124. expect(manager.getListSnapshot().items[0]?.updatedAt).toBe(100)
  125. manager.handleMuxEnvelope({
  126. rpcId: 'title-durable' as never,
  127. payload: { type: 'session/title', sessionId: S1, title: 'Durable', eventSeq: 2, updatedAt: 200 },
  128. })
  129. expect(manager.getListSnapshot().items[0]).toMatchObject({ title: 'Durable', updatedAt: 200 })
  130. manager.handleMuxEnvelope({
  131. rpcId: 'subscribed-current' as never,
  132. payload: { type: 'session/subscribed', sessionId: S1, lastSeq: 2 },
  133. })
  134. expect(manager.getListSnapshot().items[0]).toMatchObject({ title: 'Durable', updatedAt: 200 })
  135. })
  136. })
  137. describe('host frame routing', () => {
  138. it('adds/removes/flips sessions from host frames and keeps removed instances resident', async () => {
  139. const api = new FakeApiClient()
  140. const manager = new SessionManager(api)
  141. manager.handleHostEnvelope({ rpcId: 'h1' as never, payload: { type: 'host/session-added', sessionId: S1 } })
  142. manager.handleHostEnvelope({ rpcId: 'h2' as never, payload: { type: 'host/session-added', sessionId: S1 } }) // dup: ignored
  143. expect(manager.getListSnapshot().items).toHaveLength(1)
  144. const session = manager.get(S1)
  145. manager.handleHostEnvelope({ rpcId: 'h3' as never, payload: { type: 'host/session-status', sessionId: S1, running: true } })
  146. expect(session.getSnapshot().running).toBe(true)
  147. expect(manager.getListSnapshot().items[0]?.running).toBe(true)
  148. manager.handleHostEnvelope({ rpcId: 'h4' as never, payload: { type: 'host/agent-error', sessionId: S1, message: '炸了' } })
  149. expect(session.getSnapshot().lastAgentError).toBe('炸了')
  150. manager.handleHostEnvelope({ rpcId: 'h5' as never, payload: { type: 'host/session-removed', sessionId: S1 } })
  151. expect(manager.getListSnapshot().items).toHaveLength(0)
  152. expect(session.getSnapshot().removed).toBe(true)
  153. expect(manager.get(S1)).toBe(session) // resident-instance rule survives removal
  154. })
  155. })
  156. describe('remaining branches', () => {
  157. it('refreshList folds a transport throw into the error state', async () => {
  158. const api = new FakeApiClient()
  159. api.onList = () => Promise.reject(new Error('list wire down'))
  160. const manager = new SessionManager(api)
  161. await manager.refreshList()
  162. expect(manager.getListSnapshot()).toMatchObject({ state: 'error', error: { code: 'internal', message: 'list wire down' } })
  163. })
  164. it('refreshList pushes running bits down to already-instantiated sessions', async () => {
  165. const api = new FakeApiClient()
  166. const manager = new SessionManager(api)
  167. const session = manager.get(S1)
  168. api.onList = () => Promise.resolve(ok({ items: [summary(S1, { running: true })] as never[] }))
  169. await manager.refreshList()
  170. expect(session.getSnapshot().running).toBe(true)
  171. })
  172. it('create passes cwd through, folds transport throws, and skips the merge when already listed', async () => {
  173. const api = new FakeApiClient()
  174. api.onCreate = () => Promise.resolve(ok({ sessionId: S1 }))
  175. const manager = new SessionManager(api)
  176. await manager.create('/tmp/w')
  177. expect(api.callsOf('session.create')).toEqual([{ cwd: '/tmp/w' }])
  178. expect(manager.getListSnapshot().items[0]).toMatchObject({ sessionId: S1, cwd: '/tmp/w' })
  179. await manager.create('/tmp/w') // same id returned: no duplicate row
  180. expect(manager.getListSnapshot().items).toHaveLength(1)
  181. api.onCreate = () => Promise.reject(new Error('create wire down'))
  182. expect(await manager.create()).toMatchObject({ ok: false, error: { code: 'internal' } })
  183. // Business error passes through untouched.
  184. api.onCreate = () => Promise.resolve(err({ code: 'internal', message: 'no', details: {} }))
  185. expect(await manager.create()).toMatchObject({ ok: false })
  186. })
  187. it('subscribe notifies on list changes and stops after unsubscribe', async () => {
  188. const api = new FakeApiClient()
  189. const manager = new SessionManager(api)
  190. let notified = 0
  191. const unsubscribe = manager.subscribe(() => { notified++ })
  192. await manager.refreshList()
  193. await new Promise(resolve => setTimeout(resolve, 0))
  194. expect(notified).toBeGreaterThan(0)
  195. const seen = notified
  196. unsubscribe()
  197. manager.handleHostEnvelope({ rpcId: 'h' as never, payload: { type: 'host/session-added', sessionId: S1 } })
  198. await new Promise(resolve => setTimeout(resolve, 0))
  199. expect(notified).toBe(seen)
  200. })
  201. it('routes stream/error and unknown frames to the documented drops, and dispatches to instantiated sessions', () => {
  202. const api = new FakeApiClient()
  203. const manager = new SessionManager(api)
  204. manager.handleMuxEnvelope({ rpcId: 'e' as never, payload: { type: 'stream/error', error: { code: 'internal', message: 'x', details: {} } } })
  205. manager.handleHostEnvelope({ rpcId: 'e2' as never, payload: { type: 'stream/error', error: { code: 'internal', message: 'x', details: {} } } })
  206. manager.handleHostEnvelope({ rpcId: 'e3' as never, payload: { type: 'future/host-frame' } as never })
  207. const session = manager.get(S1)
  208. manager.handleMuxEnvelope({ rpcId: 'q1' as never, payload: { type: 'question/requested', sessionId: S1, questions: [] } })
  209. expect(session.getSnapshot().pending).toMatchObject([{ kind: 'question' }])
  210. // status flip for an unknown session only touches summaries (no crash).
  211. manager.handleHostEnvelope({ rpcId: 'h9' as never, payload: { type: 'host/session-status', sessionId: S2, running: true } })
  212. manager.handleHostEnvelope({ rpcId: 'ha' as never, payload: { type: 'host/agent-error', sessionId: S2, message: '无实例' } })
  213. })
  214. it('keeps list-entry identity for unchanged rows across an unrelated list change', async () => {
  215. const api = new FakeApiClient()
  216. api.onList = () => Promise.resolve(ok({ items: [summary(S1), summary(S2, { updatedAt: 200 })] as never[] }))
  217. const manager = new SessionManager(api)
  218. await manager.refreshList()
  219. const before = manager.getListSnapshot()
  220. manager.handleHostEnvelope({ rpcId: 'h' as never, payload: { type: 'host/session-status', sessionId: S2, running: true } })
  221. const after = manager.getListSnapshot()
  222. expect(after.items).not.toBe(before.items)
  223. const beforeS1 = before.items.find(e => e.sessionId === S1)
  224. const afterS1 = after.items.find(e => e.sessionId === S1)
  225. expect(afterS1).toBe(beforeS1) // untouched entry keeps identity (entryCache)
  226. // Same-order same-entries snapshot reuses the items array.
  227. manager.handleHostEnvelope({ rpcId: 'h2' as never, payload: { type: 'host/agent-error', sessionId: S1, message: 'x' } })
  228. expect(manager.getListSnapshot().items).toBe(after.items)
  229. })
  230. it('carries parentSessionId from host/session-added into the lineage row', () => {
  231. const api = new FakeApiClient()
  232. const manager = new SessionManager(api)
  233. manager.handleHostEnvelope({ rpcId: 'h1' as never, payload: { type: 'host/session-added', sessionId: S1 } })
  234. manager.handleHostEnvelope({ rpcId: 'h2' as never, payload: { type: 'host/session-added', sessionId: S2, parentSessionId: S1 } })
  235. const items = manager.getListSnapshot().items
  236. expect(items.find(e => e.sessionId === S2)).toMatchObject({ parentSessionId: S1, depth: 1 })
  237. })
  238. })
  239. describe('connected generation', () => {
  240. it('refreshes the list and resyncs only opened instances', async () => {
  241. const api = new FakeApiClient()
  242. api.onHistory = () => Promise.resolve(ok({ events: entries(plainTurn(0, 0, 'a', 'b')) as never[], hasMore: false }))
  243. const manager = new SessionManager(api)
  244. const openedSession = manager.get(S1)
  245. await openedSession.open()
  246. manager.get(S2) // instantiated but never opened
  247. const historyCallsBefore = api.callsOf('session.history').length
  248. manager.handleConnected()
  249. await vi.waitFor(() => {
  250. expect(api.callsOf('session.list').length).toBe(1)
  251. // Only the opened instance repulls history; the cold one stays silent.
  252. expect(api.callsOf('session.history').length).toBe(historyCallsBefore + 1)
  253. })
  254. })
  255. })