manager.spec.ts 20 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378
  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. type SummaryOver = Partial<{ updatedAt: number; running: boolean; blank: boolean; parentSessionId: SessionId }>
  13. function summary(sessionId: SessionId, over: SummaryOver = {}) {
  14. return { sessionId, updatedAt: 100, running: false, blank: false, ...over }
  15. }
  16. describe('instances', () => {
  17. it('lazily builds one resident instance per id and syncs the running bit from the list', async () => {
  18. const api = new FakeApiClient()
  19. api.onList = () => Promise.resolve(ok({ items: [summary(S1, { running: true })] as never[] }))
  20. const manager = new SessionManager(api)
  21. await manager.refreshList()
  22. const session = manager.get(S1)
  23. expect(manager.get(S1)).toBe(session) // resident: same instance forever
  24. expect(session.getSnapshot().running).toBe(true) // list preceded instantiation
  25. })
  26. it('replays buffered approval frames on instantiation and drops ordinary frames for uninstantiated sessions', () => {
  27. const api = new FakeApiClient()
  28. const manager = new SessionManager(api)
  29. // Uninstantiated: approval buffers, plain session/event drops.
  30. manager.handleMuxEnvelope({ rpcId: 'ra' as never, payload: { type: 'approval/requested', sessionId: S1, approvalId: 'ap1' as never, toolName: 'rm' } })
  31. manager.handleMuxEnvelope({ rpcId: 're' as never, payload: { type: 'session/event', sessionId: S1, event: plainTurn(0, 0, 'x', 'y')[0] as never } })
  32. const session = manager.get(S1)
  33. expect(session.getSnapshot().pending).toMatchObject([{ kind: 'approval', payload: { approvalId: 'ap1' } }])
  34. // Buffer cleared: a second instantiation of another id gets nothing.
  35. expect(manager.get(S2).getSnapshot().pending).toEqual([])
  36. })
  37. it('caps the pending buffer at 32 keeping the newest, and drops it on session-removed', () => {
  38. const api = new FakeApiClient()
  39. const manager = new SessionManager(api)
  40. // 40 distinct question frames for an uninstantiated session: only the newest 32 survive.
  41. for (let i = 0; i < 40; i++) {
  42. manager.handleMuxEnvelope({ rpcId: `q${i}` as never, payload: { type: 'question/requested', sessionId: S1, questions: [] } })
  43. }
  44. const pending = manager.get(S1).getSnapshot().pending
  45. expect(pending).toHaveLength(32)
  46. expect(pending.map(p => p.key)).toEqual(Array.from({ length: 32 }, (_, i) => `q:q${i + 8}`)) // oldest 8 dropped
  47. // Removed session: buffered frames must not replay on a future instantiation.
  48. manager.handleMuxEnvelope({ rpcId: 'qz' as never, payload: { type: 'question/requested', sessionId: S2, questions: [] } })
  49. manager.handleHostEnvelope({ rpcId: 'hz' as never, payload: { type: 'host/session-removed', sessionId: S2 } })
  50. expect(manager.get(S2).getSnapshot().pending).toEqual([])
  51. })
  52. })
  53. describe('list lifecycle', () => {
  54. it('single-flights refreshList and preserves the Host baseline order', async () => {
  55. const api = new FakeApiClient()
  56. const gate = deferred<Awaited<ReturnType<FakeApiClient['onList']>>>()
  57. api.onList = () => gate.promise
  58. const manager = new SessionManager(api)
  59. const first = manager.refreshList()
  60. const second = manager.refreshList()
  61. expect(manager.getListSnapshot().state).toBe('loading')
  62. gate.resolve(ok({ items: [summary(S2, { updatedAt: 200 }), summary(S1)] as never[] }))
  63. await Promise.all([first, second])
  64. expect(api.callsOf('session.list')).toHaveLength(1)
  65. const snapshot = manager.getListSnapshot()
  66. expect(snapshot.state).toBe('idle')
  67. expect(snapshot.items.map(i => i.sessionId)).toEqual([S2, S1])
  68. })
  69. it('replays incremental frames over hydration and never batch-reorders established ids', async () => {
  70. const api = new FakeApiClient()
  71. const first = deferred<Awaited<ReturnType<FakeApiClient['onList']>>>()
  72. api.onList = () => first.promise
  73. const manager = new SessionManager(api)
  74. const hydration = manager.refreshList()
  75. manager.handleHostEnvelope({
  76. rpcId: 'during-first' as never,
  77. payload: { type: 'host/session-added', blank: true, sessionId: S2 },
  78. })
  79. first.resolve(ok({ items: [summary(S1)] as never[] }))
  80. await hydration
  81. expect(manager.getListSnapshot().items.map(item => item.sessionId)).toEqual([S2, S1])
  82. api.onList = () => Promise.resolve(ok({
  83. items: [summary(S1, { updatedAt: 900 }), summary(S2, { updatedAt: 800 })] as never[],
  84. }))
  85. await manager.refreshList()
  86. expect(manager.getListSnapshot().items.map(item => item.sessionId)).toEqual([S2, S1])
  87. })
  88. it('keeps the error in the list snapshot on failure', async () => {
  89. const api = new FakeApiClient()
  90. api.onList = () => Promise.resolve(err({ code: 'internal', message: 'boom', details: {} }))
  91. const manager = new SessionManager(api)
  92. await manager.refreshList()
  93. expect(manager.getListSnapshot()).toMatchObject({ state: 'error', error: { code: 'internal' } })
  94. // A failed pull does not step the arrival phase: still pending.
  95. expect(manager.getListSnapshot().phase).toBe('pending')
  96. })
  97. it('phase steps pending → ready on the first successful pull and never returns', async () => {
  98. const api = new FakeApiClient()
  99. const manager = new SessionManager(api)
  100. expect(manager.getListSnapshot().phase).toBe('pending')
  101. await manager.refreshList()
  102. expect(manager.getListSnapshot().phase).toBe('ready')
  103. // Sticky across later failures: the pull-activity axis reports the error,
  104. // the arrival phase holds.
  105. api.onList = () => Promise.resolve(err({ code: 'internal', message: 'down', details: {} }))
  106. await manager.refreshList()
  107. expect(manager.getListSnapshot()).toMatchObject({ state: 'error', phase: 'ready' })
  108. // And across an empty re-pull (empty-with-ready = truly no sessions).
  109. api.onList = () => Promise.resolve(ok({ items: [] as never[] }))
  110. await manager.refreshList()
  111. expect(manager.getListSnapshot()).toMatchObject({ state: 'idle', phase: 'ready' })
  112. expect(manager.getListSnapshot().items).toEqual([])
  113. })
  114. it('merges create into the list immediately without waiting for a refresh', async () => {
  115. const api = new FakeApiClient()
  116. api.onCreate = () => Promise.resolve(ok({ sessionId: S2 }))
  117. const manager = new SessionManager(api)
  118. const result = await manager.create()
  119. expect(result).toMatchObject({ ok: true, value: { sessionId: S2 } })
  120. expect(manager.getListSnapshot().items.map(i => i.sessionId)).toEqual([S2])
  121. })
  122. it('retains title projections before list arrival, keeps last-wins by seq, and clears them on removal', async () => {
  123. const api = new FakeApiClient()
  124. const manager = new SessionManager(api)
  125. const titleFrame = (rpcId: string, title: string, seq: number) => {
  126. manager.handleMuxEnvelope({
  127. rpcId: rpcId as never,
  128. payload: { type: 'session/projection', sessionId: S1, key: 'title', value: title, seq } as never,
  129. })
  130. }
  131. titleFrame('title-new', 'Newest', 4)
  132. titleFrame('title-stale', 'Stale', 3)
  133. titleFrame('title-equal', 'Equal', 4)
  134. api.onList = () => Promise.resolve(ok({
  135. items: [summary(S1), summary(S2, { updatedAt: 200 })] as never[],
  136. }))
  137. await manager.refreshList()
  138. const titled = manager.getListSnapshot()
  139. expect(titled.items.map(item => item.sessionId)).toEqual([S1, S2])
  140. expect(titled.items[0]?.title).toBe('Newest')
  141. expect(titled.items[1]?.title).toBeUndefined()
  142. manager.handleHostEnvelope({ rpcId: 'removed' as never, payload: { type: 'host/session-removed', sessionId: S1 } })
  143. manager.handleHostEnvelope({ rpcId: 'readded' as never, payload: { type: 'host/session-added', blank: true, sessionId: S1 } })
  144. expect(manager.getListSnapshot().items.find(item => item.sessionId === S1)?.title).toBeUndefined()
  145. })
  146. it('seeds cold titles from the list rows\' projections block under higher-seq-wins', async () => {
  147. const api = new FakeApiClient()
  148. const manager = new SessionManager(api)
  149. // A push frame landed before the list (S2's title is newer than the block's cut).
  150. manager.handleMuxEnvelope({
  151. rpcId: 'push-newer' as never,
  152. payload: { type: 'session/projection', sessionId: S2, key: 'title', value: 'Pushed', seq: 9 } as never,
  153. })
  154. api.onList = () => Promise.resolve(ok({
  155. items: [
  156. { ...summary(S1), projections: { asOfSeq: 4, values: { title: 'Cold cached' } } },
  157. { ...summary(S2, { updatedAt: 200 }), projections: { asOfSeq: 5, values: { title: 'List stale' } } },
  158. ] as never[],
  159. }))
  160. await manager.refreshList()
  161. const items = manager.getListSnapshot().items
  162. // Cold row: title surfaces straight from the list block — no open, no history.
  163. expect(items.find(item => item.sessionId === S1)?.title).toBe('Cold cached')
  164. // The stale list block (seq 5) cannot overwrite the newer push frame (seq 9).
  165. expect(items.find(item => item.sessionId === S2)?.title).toBe('Pushed')
  166. })
  167. it('drops a projection row beyond the subscription baseline before accepting its durable replay', async () => {
  168. const api = new FakeApiClient()
  169. api.onList = () => Promise.resolve(ok({ items: [summary(S1)] as never[] }))
  170. const manager = new SessionManager(api)
  171. await manager.refreshList()
  172. const frame = (rpcId: string, payload: object) => {
  173. manager.handleMuxEnvelope({ rpcId: rpcId as never, payload: payload as never })
  174. }
  175. frame('title-unflushed', { type: 'session/projection', sessionId: S1, key: 'title', value: 'Unflushed', seq: 4 })
  176. // The durable baseline says the host only knows up to seq 2: the phantom
  177. // row rode lost state and must drop, or last-wins pins it forever.
  178. frame('subscribed-recovered', { type: 'session/subscribed', sessionId: S1, lastSeq: 2 })
  179. expect(manager.getListSnapshot().items[0]?.title).toBeUndefined()
  180. frame('title-durable', { type: 'session/projection', sessionId: S1, key: 'title', value: 'Durable', seq: 2 })
  181. expect(manager.getListSnapshot().items[0]?.title).toBe('Durable')
  182. // A baseline at or past the row's seq keeps it (nothing phantom to drop).
  183. frame('subscribed-current', { type: 'session/subscribed', sessionId: S1, lastSeq: 2 })
  184. expect(manager.getListSnapshot().items[0]?.title).toBe('Durable')
  185. })
  186. })
  187. describe('host frame routing', () => {
  188. it('adds/removes/flips sessions from host frames and keeps removed instances resident', async () => {
  189. const api = new FakeApiClient()
  190. const manager = new SessionManager(api)
  191. manager.handleHostEnvelope({ rpcId: 'h1' as never, payload: { type: 'host/session-added', blank: true, sessionId: S1 } })
  192. manager.handleHostEnvelope({ rpcId: 'h2' as never, payload: { type: 'host/session-added', blank: true, sessionId: S1 } }) // dup: ignored
  193. expect(manager.getListSnapshot().items).toHaveLength(1)
  194. const session = manager.get(S1)
  195. manager.handleHostEnvelope({ rpcId: 'h3' as never, payload: { type: 'host/session-status', sessionId: S1, running: true } })
  196. expect(session.getSnapshot().running).toBe(true)
  197. expect(manager.getListSnapshot().items[0]?.running).toBe(true)
  198. manager.handleHostEnvelope({ rpcId: 'h4' as never, payload: { type: 'host/agent-error', sessionId: S1, message: '炸了' } })
  199. expect(session.getSnapshot().lastAgentError).toBe('炸了')
  200. manager.handleHostEnvelope({ rpcId: 'h5' as never, payload: { type: 'host/session-removed', sessionId: S1 } })
  201. expect(manager.getListSnapshot().items).toHaveLength(0)
  202. expect(session.getSnapshot().removed).toBe(true)
  203. expect(manager.get(S1)).toBe(session) // resident-instance rule survives removal
  204. })
  205. })
  206. describe('remaining branches', () => {
  207. it('refreshList folds a transport throw into the error state', async () => {
  208. const api = new FakeApiClient()
  209. api.onList = () => Promise.reject(new Error('list wire down'))
  210. const manager = new SessionManager(api)
  211. await manager.refreshList()
  212. expect(manager.getListSnapshot()).toMatchObject({ state: 'error', error: { code: 'internal', message: 'list wire down' } })
  213. })
  214. it('refreshList pushes running bits down to already-instantiated sessions', async () => {
  215. const api = new FakeApiClient()
  216. const manager = new SessionManager(api)
  217. const session = manager.get(S1)
  218. api.onList = () => Promise.resolve(ok({ items: [summary(S1, { running: true })] as never[] }))
  219. await manager.refreshList()
  220. expect(session.getSnapshot().running).toBe(true)
  221. })
  222. it('create passes cwd and a preallocated id, folds transport throws, and deduplicates the echo', async () => {
  223. const api = new FakeApiClient()
  224. api.onCreate = () => Promise.resolve(ok({ sessionId: S1 }))
  225. const manager = new SessionManager(api)
  226. await manager.create({ cwd: '/tmp/w', sessionId: S1 })
  227. expect(api.callsOf('session.create')).toEqual([{ cwd: '/tmp/w', sessionId: S1 }])
  228. expect(manager.getListSnapshot().items[0]).toMatchObject({ sessionId: S1, cwd: '/tmp/w' })
  229. await manager.create({ cwd: '/tmp/w' }) // same id returned: no duplicate row
  230. expect(manager.getListSnapshot().items).toHaveLength(1)
  231. api.onCreate = () => Promise.reject(new Error('create wire down'))
  232. expect(await manager.create()).toMatchObject({ ok: false, error: { code: 'internal' } })
  233. // Business error passes through untouched.
  234. api.onCreate = () => Promise.resolve(err({ code: 'internal', message: 'no', details: {} }))
  235. expect(await manager.create()).toMatchObject({ ok: false })
  236. })
  237. it('publishes a real Ungrouped summary from workspace-attach-failed', async () => {
  238. const api = new FakeApiClient()
  239. api.onCreate = () => Promise.resolve(err({
  240. code: 'workspace-attach-failed',
  241. message: 'published but unattached',
  242. details: { sessionId: S1, workspaceId: 'w1' },
  243. } as never))
  244. const manager = new SessionManager(api)
  245. const result = await manager.create({ workspaceId: 'w1' as never, sessionId: S1 })
  246. expect(result).toMatchObject({ ok: false, error: { code: 'workspace-attach-failed' } })
  247. expect(manager.getListSnapshot().items).toEqual([expect.objectContaining({ sessionId: S1 })])
  248. expect(manager.getListSnapshot().items[0]).not.toHaveProperty('cwd')
  249. })
  250. it('reconciles a preallocated id after an ordinary transport failure', async () => {
  251. const api = new FakeApiClient()
  252. api.onCreate = () => Promise.reject(new Error('response lost'))
  253. const manager = new SessionManager(api)
  254. const failed = await manager.create({ workspaceId: 'w1' as never, sessionId: S1 })
  255. expect(failed).toMatchObject({ ok: false, error: { message: 'response lost' } })
  256. expect(manager.getListSnapshot().items).toEqual([])
  257. manager.handleHostEnvelope({
  258. rpcId: 'published-later' as never,
  259. payload: { type: 'host/session-added', blank: true, sessionId: S1, cwd: '/w/one' },
  260. })
  261. expect(manager.getListSnapshot().items).toEqual([
  262. expect.objectContaining({ sessionId: S1, cwd: '/w/one' }),
  263. ])
  264. manager.handleHostEnvelope({
  265. rpcId: 'duplicate-frame' as never,
  266. payload: { type: 'host/session-added', blank: true, sessionId: S1, cwd: '/w/one' },
  267. })
  268. expect(manager.getListSnapshot().items).toHaveLength(1)
  269. })
  270. it('subscribe notifies on list changes and stops after unsubscribe', async () => {
  271. const api = new FakeApiClient()
  272. const manager = new SessionManager(api)
  273. let notified = 0
  274. const unsubscribe = manager.subscribe(() => { notified++ })
  275. await manager.refreshList()
  276. await new Promise(resolve => setTimeout(resolve, 0))
  277. expect(notified).toBeGreaterThan(0)
  278. const seen = notified
  279. unsubscribe()
  280. manager.handleHostEnvelope({ rpcId: 'h' as never, payload: { type: 'host/session-added', blank: true, sessionId: S1 } })
  281. await new Promise(resolve => setTimeout(resolve, 0))
  282. expect(notified).toBe(seen)
  283. })
  284. it('routes stream/error and unknown frames to the documented drops, and dispatches to instantiated sessions', () => {
  285. const api = new FakeApiClient()
  286. const manager = new SessionManager(api)
  287. manager.handleMuxEnvelope({ rpcId: 'e' as never, payload: { type: 'stream/error', error: { code: 'internal', message: 'x', details: {} } } })
  288. manager.handleHostEnvelope({ rpcId: 'e2' as never, payload: { type: 'stream/error', error: { code: 'internal', message: 'x', details: {} } } })
  289. manager.handleHostEnvelope({ rpcId: 'e3' as never, payload: { type: 'future/host-frame' } as never })
  290. const session = manager.get(S1)
  291. manager.handleMuxEnvelope({ rpcId: 'q1' as never, payload: { type: 'question/requested', sessionId: S1, questions: [] } })
  292. expect(session.getSnapshot().pending).toMatchObject([{ kind: 'question' }])
  293. // status flip for an unknown session only touches summaries (no crash).
  294. manager.handleHostEnvelope({ rpcId: 'h9' as never, payload: { type: 'host/session-status', sessionId: S2, running: true } })
  295. manager.handleHostEnvelope({ rpcId: 'ha' as never, payload: { type: 'host/agent-error', sessionId: S2, message: '无实例' } })
  296. })
  297. it('keeps list-entry identity for unchanged rows across an unrelated list change', async () => {
  298. const api = new FakeApiClient()
  299. api.onList = () => Promise.resolve(ok({ items: [summary(S1), summary(S2, { updatedAt: 200 })] as never[] }))
  300. const manager = new SessionManager(api)
  301. await manager.refreshList()
  302. const before = manager.getListSnapshot()
  303. manager.handleHostEnvelope({ rpcId: 'h' as never, payload: { type: 'host/session-status', sessionId: S2, running: true } })
  304. const after = manager.getListSnapshot()
  305. expect(after.items).not.toBe(before.items)
  306. const beforeS1 = before.items.find(e => e.sessionId === S1)
  307. const afterS1 = after.items.find(e => e.sessionId === S1)
  308. expect(afterS1).toBe(beforeS1) // untouched entry keeps identity (entryCache)
  309. // Same-order same-entries snapshot reuses the items array.
  310. manager.handleHostEnvelope({ rpcId: 'h2' as never, payload: { type: 'host/agent-error', sessionId: S1, message: 'x' } })
  311. expect(manager.getListSnapshot().items).toBe(after.items)
  312. })
  313. it('carries parentSessionId from host/session-added into the lineage row', () => {
  314. const api = new FakeApiClient()
  315. const manager = new SessionManager(api)
  316. manager.handleHostEnvelope({ rpcId: 'h1' as never, payload: { type: 'host/session-added', blank: true, sessionId: S1 } })
  317. manager.handleHostEnvelope({ rpcId: 'h2' as never, payload: { type: 'host/session-added', blank: true, sessionId: S2, parentSessionId: S1 } })
  318. const items = manager.getListSnapshot().items
  319. expect(items.find(e => e.sessionId === S2)).toMatchObject({ parentSessionId: S1, depth: 1 })
  320. })
  321. })
  322. describe('connected generation', () => {
  323. it('refreshes the list and resyncs only opened instances', async () => {
  324. const api = new FakeApiClient()
  325. api.onHistory = () => Promise.resolve(ok({
  326. events: entries(plainTurn(0, 0, 'a', 'b')) as never[],
  327. hasMore: false,
  328. modelTarget: { provider: 'deepseek', model: 'deepseek-chat' },
  329. }))
  330. const manager = new SessionManager(api)
  331. const openedSession = manager.get(S1)
  332. await openedSession.open()
  333. manager.get(S2) // instantiated but never opened
  334. const historyCallsBefore = api.callsOf('session.history').length
  335. manager.handleConnected()
  336. await vi.waitFor(() => {
  337. expect(api.callsOf('session.list').length).toBe(1)
  338. // Only the opened instance repulls history; the cold one stays silent.
  339. expect(api.callsOf('session.history').length).toBe(historyCallsBefore + 1)
  340. })
  341. })
  342. })