manager.client.spec.ts 57 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210
  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-api-remotes/client'
  7. import { SessionManager } from '../src/client/sessions/manager.ts'
  8. import { FakeApiClient, deferred, err, fakeRemote, ok } from './fake-api.client.ts'
  9. import { entries, ev, plainTurn } from './event-script.client.ts'
  10. const S1 = 'fk-m1' as SessionId
  11. const S2 = 'fk-m2' as SessionId
  12. type SummaryOver = Partial<{
  13. updatedAt: number
  14. running: boolean
  15. blank: boolean
  16. parentSessionId: SessionId
  17. origin: 'subagent'
  18. }>
  19. function summary(sessionId: SessionId, over: SummaryOver = {}) {
  20. return { sessionId, updatedAt: 100, running: false, blank: false, ...over }
  21. }
  22. describe('instances', () => {
  23. it('lazily builds one resident instance per id and syncs the running bit from the list', async () => {
  24. const api = new FakeApiClient()
  25. api.onList = () => Promise.resolve(ok({ items: [summary(S1, { running: true })] as never[] }))
  26. const manager = new SessionManager(api, fakeRemote())
  27. await manager.refreshList()
  28. const session = manager.get(S1)
  29. expect(manager.get(S1)).toBe(session) // resident: same instance forever
  30. expect(session.getSnapshot().running).toBe(true) // list preceded instantiation
  31. })
  32. it('replays buffered approval frames on instantiation and drops ordinary frames for uninstantiated sessions', () => {
  33. const api = new FakeApiClient()
  34. const manager = new SessionManager(api, fakeRemote())
  35. // Uninstantiated: approval buffers, plain session/event drops.
  36. manager.handleMuxEnvelope({ rpcId: 'ra' as never, payload: { type: 'approval/requested', sessionId: S1, approvalId: 'ap1' as never, toolName: 'rm' } })
  37. manager.handleMuxEnvelope({ rpcId: 'ra' as never, payload: { type: 'approval/requested', sessionId: S1, approvalId: 'ap1' as never, toolName: 'rm' } })
  38. manager.handleMuxEnvelope({ rpcId: 're' as never, payload: { type: 'session/event', sessionId: S1, event: plainTurn(0, 0, 'x', 'y')[0] as never } })
  39. const session = manager.get(S1)
  40. expect(session.getSnapshot().pending).toMatchObject([{ kind: 'approval', payload: { approvalId: 'ap1' } }])
  41. // Buffer cleared: a second instantiation of another id gets nothing.
  42. expect(manager.get(S2).getSnapshot().pending).toEqual([])
  43. })
  44. it('retains every live answerable request and compacts resolutions before instantiation', () => {
  45. const api = new FakeApiClient()
  46. const manager = new SessionManager(api, fakeRemote())
  47. manager.handleHostEnvelope({ rpcId: 'h1' as never, payload: { type: 'host/session-added', sessionId: S1, blank: false } })
  48. for (let i = 0; i < 40; i++) {
  49. manager.handleMuxEnvelope({ rpcId: `q${i}` as never, payload: { type: 'question/requested', sessionId: S1, questions: [] } })
  50. }
  51. expect(manager.getListSnapshot().items[0]?.pendingInteraction).toBe('question')
  52. for (let i = 0; i < 40; i++) {
  53. manager.handleMuxEnvelope({
  54. rpcId: `r${i}` as never,
  55. payload: { type: 'question/resolved', sessionId: S1, questionRpcId: `q${i}` as never, outcome: 'answered' },
  56. })
  57. }
  58. expect(manager.getListSnapshot().items[0]?.pendingInteraction).toBeUndefined()
  59. expect(manager.get(S1).getSnapshot().pending).toEqual([])
  60. })
  61. it('drops buffered answerable requests on session removal', () => {
  62. const manager = new SessionManager(new FakeApiClient(), fakeRemote())
  63. // Removed session: buffered frames must not replay on a future instantiation.
  64. manager.handleMuxEnvelope({ rpcId: 'qz' as never, payload: { type: 'question/requested', sessionId: S2, questions: [] } })
  65. manager.handleHostEnvelope({ rpcId: 'hz' as never, payload: { type: 'host/session-removed', sessionId: S2 } })
  66. expect(manager.get(S2).getSnapshot().pending).toEqual([])
  67. })
  68. })
  69. describe('list lifecycle', () => {
  70. it('single-flights refreshList and preserves the Host baseline order', async () => {
  71. const api = new FakeApiClient()
  72. const gate = deferred<Awaited<ReturnType<FakeApiClient['onList']>>>()
  73. api.onList = () => gate.promise
  74. const manager = new SessionManager(api, fakeRemote())
  75. const first = manager.refreshList()
  76. const second = manager.refreshList()
  77. expect(manager.getListSnapshot().state).toBe('loading')
  78. gate.resolve(ok({ items: [summary(S2, { updatedAt: 200 }), summary(S1)] as never[] }))
  79. await Promise.all([first, second])
  80. expect(api.callsOf('session.list')).toHaveLength(1)
  81. const snapshot = manager.getListSnapshot()
  82. expect(snapshot.state).toBe('idle')
  83. expect(snapshot.items.map(i => i.sessionId)).toEqual([S2, S1])
  84. })
  85. it('replays incremental frames over hydration and never batch-reorders established ids', async () => {
  86. const api = new FakeApiClient()
  87. const first = deferred<Awaited<ReturnType<FakeApiClient['onList']>>>()
  88. api.onList = () => first.promise
  89. const manager = new SessionManager(api, fakeRemote())
  90. const hydration = manager.refreshList()
  91. manager.handleHostEnvelope({
  92. rpcId: 'during-first' as never,
  93. payload: { type: 'host/session-added', blank: true, sessionId: S2 },
  94. })
  95. first.resolve(ok({ items: [summary(S1)] as never[] }))
  96. await hydration
  97. expect(manager.getListSnapshot().items.map(item => item.sessionId)).toEqual([S2, S1])
  98. api.onList = () => Promise.resolve(ok({
  99. items: [summary(S1, { updatedAt: 900 }), summary(S2, { updatedAt: 800 })] as never[],
  100. }))
  101. await manager.refreshList()
  102. expect(manager.getListSnapshot().items.map(item => item.sessionId)).toEqual([S2, S1])
  103. })
  104. it('advances list activity only for direct user messages', async () => {
  105. const api = new FakeApiClient()
  106. api.onList = () => Promise.resolve(ok({ items: [summary(S1)] as never[] }))
  107. const manager = new SessionManager(api, fakeRemote())
  108. await manager.refreshList()
  109. // Both a new prompt and an admitted steer land as a user-sourced message.
  110. const activity = { ...ev.user(10, 'new'), time: 500 }
  111. manager.handleMuxEnvelope({
  112. rpcId: 'activity' as never,
  113. payload: { type: 'session/event', sessionId: S1, event: activity },
  114. })
  115. expect(manager.getListSnapshot().items[0]?.updatedAt).toBe(500)
  116. manager.handleMuxEnvelope({
  117. rpcId: 'older' as never,
  118. payload: { type: 'session/event', sessionId: S1, event: { ...activity, time: 400 } },
  119. })
  120. manager.handleMuxEnvelope({
  121. rpcId: 'assistant' as never,
  122. payload: { type: 'session/event', sessionId: S1, event: { ...ev.assistant(11, 0, 'reply'), time: 600 } },
  123. })
  124. const injected = ev.user(12, 'context')
  125. if (injected.type !== 'user/message') throw new Error('user builder returned another event type')
  126. manager.handleMuxEnvelope({
  127. rpcId: 'injected' as never,
  128. payload: {
  129. type: 'session/event',
  130. sessionId: S1,
  131. event: {
  132. ...injected,
  133. time: 700,
  134. data: { ...injected.data, source: { kind: 'plugin', plugin: 'test' } },
  135. },
  136. },
  137. })
  138. expect(manager.getListSnapshot().items[0]?.updatedAt).toBe(500)
  139. })
  140. it('keeps the error in the list snapshot on failure', async () => {
  141. const api = new FakeApiClient()
  142. api.onList = () => Promise.resolve(err({ code: 'internal', message: 'boom', details: {} }))
  143. const manager = new SessionManager(api, fakeRemote())
  144. await manager.refreshList()
  145. expect(manager.getListSnapshot()).toMatchObject({ state: 'error', error: { code: 'internal' } })
  146. // A failed pull does not step the arrival phase: still pending.
  147. expect(manager.getListSnapshot().phase).toBe('pending')
  148. })
  149. it('phase steps pending → ready on the first successful pull and never returns', async () => {
  150. const api = new FakeApiClient()
  151. const manager = new SessionManager(api, fakeRemote())
  152. expect(manager.getListSnapshot().phase).toBe('pending')
  153. await manager.refreshList()
  154. expect(manager.getListSnapshot().phase).toBe('ready')
  155. // Sticky across later failures: the pull-activity axis reports the error,
  156. // the arrival phase holds.
  157. api.onList = () => Promise.resolve(err({ code: 'internal', message: 'down', details: {} }))
  158. await manager.refreshList()
  159. expect(manager.getListSnapshot()).toMatchObject({ state: 'error', phase: 'ready' })
  160. // And across an empty re-pull (empty-with-ready = truly no sessions).
  161. api.onList = () => Promise.resolve(ok({ items: [] as never[] }))
  162. await manager.refreshList()
  163. expect(manager.getListSnapshot()).toMatchObject({ state: 'idle', phase: 'ready' })
  164. expect(manager.getListSnapshot().items).toEqual([])
  165. })
  166. it('merges create into the list immediately without waiting for a refresh', async () => {
  167. const api = new FakeApiClient()
  168. api.onCreate = () => Promise.resolve(ok({ sessionId: S2 }))
  169. const manager = new SessionManager(api, fakeRemote())
  170. const result = await manager.create()
  171. expect(result).toMatchObject({ ok: true, value: { sessionId: S2 } })
  172. expect(manager.getListSnapshot().items.map(i => i.sessionId)).toEqual([S2])
  173. })
  174. it('retains title projections before list arrival, keeps last-wins by seq, and clears them on removal', async () => {
  175. const api = new FakeApiClient()
  176. const manager = new SessionManager(api, fakeRemote())
  177. const titleFrame = (rpcId: string, title: string, seq: number) => {
  178. manager.handleMuxEnvelope({
  179. rpcId: rpcId as never,
  180. payload: { type: 'session/projection', sessionId: S1, key: 'title', value: title, seq } as never,
  181. })
  182. }
  183. titleFrame('title-new', 'Newest', 4)
  184. titleFrame('title-stale', 'Stale', 3)
  185. titleFrame('title-equal', 'Equal', 4)
  186. api.onList = () => Promise.resolve(ok({
  187. items: [summary(S1), summary(S2, { updatedAt: 200 })] as never[],
  188. }))
  189. await manager.refreshList()
  190. const titled = manager.getListSnapshot()
  191. expect(titled.items.map(item => item.sessionId)).toEqual([S1, S2])
  192. expect(titled.items[0]?.title).toBe('Newest')
  193. expect(titled.items[1]?.title).toBeUndefined()
  194. manager.handleHostEnvelope({ rpcId: 'removed' as never, payload: { type: 'host/session-removed', sessionId: S1 } })
  195. manager.handleHostEnvelope({ rpcId: 'readded' as never, payload: { type: 'host/session-added', blank: true, sessionId: S1 } })
  196. expect(manager.getListSnapshot().items.find(item => item.sessionId === S1)?.title).toBeUndefined()
  197. })
  198. it('seeds cold titles from the list rows\' projections block under higher-seq-wins', async () => {
  199. const api = new FakeApiClient()
  200. const manager = new SessionManager(api, fakeRemote())
  201. // A push frame landed before the list (S2's title is newer than the block's cut).
  202. manager.handleMuxEnvelope({
  203. rpcId: 'push-newer' as never,
  204. payload: { type: 'session/projection', sessionId: S2, key: 'title', value: 'Pushed', seq: 9 } as never,
  205. })
  206. api.onList = () => Promise.resolve(ok({
  207. items: [
  208. { ...summary(S1), projections: { asOfSeq: 4, values: { title: 'Cold cached' } } },
  209. { ...summary(S2, { updatedAt: 200 }), projections: { asOfSeq: 5, values: { title: 'List stale' } } },
  210. ] as never[],
  211. }))
  212. await manager.refreshList()
  213. const items = manager.getListSnapshot().items
  214. // Cold row: title surfaces straight from the list block — no open, no history.
  215. expect(items.find(item => item.sessionId === S1)?.title).toBe('Cold cached')
  216. // The stale list block (seq 5) cannot overwrite the newer push frame (seq 9).
  217. expect(items.find(item => item.sessionId === S2)?.title).toBe('Pushed')
  218. })
  219. it('drops a projection row beyond the subscription baseline before accepting its durable replay', async () => {
  220. const api = new FakeApiClient()
  221. api.onList = () => Promise.resolve(ok({ items: [summary(S1)] as never[] }))
  222. const manager = new SessionManager(api, fakeRemote())
  223. await manager.refreshList()
  224. const frame = (rpcId: string, payload: object) => {
  225. manager.handleMuxEnvelope({ rpcId: rpcId as never, payload: payload as never })
  226. }
  227. frame('title-unflushed', { type: 'session/projection', sessionId: S1, key: 'title', value: 'Unflushed', seq: 4 })
  228. // The durable baseline says the host only knows up to seq 2: the phantom
  229. // row rode lost state and must drop, or last-wins pins it forever.
  230. frame('subscribed-recovered', { type: 'session/subscribed', sessionId: S1, lastSeq: 2 })
  231. expect(manager.getListSnapshot().items[0]?.title).toBeUndefined()
  232. frame('title-durable', { type: 'session/projection', sessionId: S1, key: 'title', value: 'Durable', seq: 2 })
  233. expect(manager.getListSnapshot().items[0]?.title).toBe('Durable')
  234. // A baseline at or past the row's seq keeps it (nothing phantom to drop).
  235. frame('subscribed-current', { type: 'session/subscribed', sessionId: S1, lastSeq: 2 })
  236. expect(manager.getListSnapshot().items[0]?.title).toBe('Durable')
  237. })
  238. })
  239. describe('search', () => {
  240. it('returns bounded Host results and forwards the caller signal', async () => {
  241. const api = new FakeApiClient()
  242. api.onSearch = () => Promise.resolve(ok({
  243. items: [{ sessionId: S1, snippet: 'matching excerpt' }],
  244. hasMore: true,
  245. }))
  246. const manager = new SessionManager(api, fakeRemote())
  247. const signal = new AbortController().signal
  248. await expect(manager.search('exact phrase', signal)).resolves.toEqual({
  249. ok: true,
  250. value: {
  251. items: [{ sessionId: S1, snippet: 'matching excerpt' }],
  252. hasMore: true,
  253. },
  254. })
  255. expect(api.callsOf('session.search')).toEqual([{ query: 'exact phrase' }])
  256. expect(api.lastSearchSignal).toBe(signal)
  257. })
  258. it('preserves business errors and folds transport failures', async () => {
  259. const api = new FakeApiClient()
  260. const manager = new SessionManager(api, fakeRemote())
  261. api.onSearch = () => Promise.resolve(err({
  262. code: 'internal',
  263. message: 'index unavailable',
  264. details: {},
  265. }))
  266. const signal = new AbortController().signal
  267. await expect(manager.search('first', signal)).resolves.toMatchObject({
  268. ok: false,
  269. error: { code: 'internal', message: 'index unavailable' },
  270. })
  271. api.onSearch = () => Promise.reject(new Error('wire down'))
  272. await expect(manager.search('second', signal)).resolves.toMatchObject({
  273. ok: false,
  274. error: { code: 'internal', message: 'wire down' },
  275. })
  276. })
  277. })
  278. describe('host frame routing', () => {
  279. it('adds/removes/flips sessions from host frames and keeps removed instances resident', async () => {
  280. const api = new FakeApiClient()
  281. const manager = new SessionManager(api, fakeRemote())
  282. manager.handleHostEnvelope({ rpcId: 'h1' as never, payload: { type: 'host/session-added', blank: true, sessionId: S1 } })
  283. manager.handleHostEnvelope({ rpcId: 'h2' as never, payload: { type: 'host/session-added', blank: true, sessionId: S1 } }) // dup: ignored
  284. expect(manager.getListSnapshot().items).toHaveLength(1)
  285. const session = manager.get(S1)
  286. manager.handleHostEnvelope({ rpcId: 'h3' as never, payload: { type: 'host/session-status', sessionId: S1, running: true } })
  287. expect(session.getSnapshot().running).toBe(true)
  288. expect(manager.getListSnapshot().items[0]?.running).toBe(true)
  289. manager.handleHostEnvelope({ rpcId: 'h4' as never, payload: { type: 'host/agent-error', sessionId: S1, message: '炸了' } })
  290. expect(session.getSnapshot().lastAgentError).toBe('炸了')
  291. manager.handleHostEnvelope({ rpcId: 'h5' as never, payload: { type: 'host/session-removed', sessionId: S1 } })
  292. expect(manager.getListSnapshot().items).toHaveLength(0)
  293. expect(session.getSnapshot().removed).toBe(true)
  294. expect(manager.get(S1)).toBe(session) // resident-instance rule survives removal
  295. })
  296. })
  297. describe('subagent catalogs', () => {
  298. it('keeps a catalog-discovered child address across ordinary selection and status frames', async () => {
  299. const api = new FakeApiClient()
  300. api.onList = () => Promise.resolve(ok({ items: [
  301. summary(S1),
  302. summary(S2, { parentSessionId: S1, origin: 'subagent' }),
  303. ] as never[] }))
  304. api.onSubagentList = () => Promise.resolve(ok({
  305. entries: [{
  306. kind: 'child', id: S2, mode: 'continuable', label: 'worker',
  307. activity: 'running', hasChildren: false,
  308. }] as never[],
  309. parentAvailable: true,
  310. }))
  311. const manager = new SessionManager(api, fakeRemote())
  312. await manager.refreshList()
  313. await manager.refreshSubagents(S1)
  314. manager.selectSubagent({ parentSessionId: S1, childSessionId: S2, mode: 'continuable' })
  315. expect(manager.getListSnapshot().currentAddress).toEqual({
  316. parentSessionId: S1, childSessionId: S2, mode: 'continuable',
  317. })
  318. expect(manager.get(S2).getSnapshot().subagent).toEqual({
  319. address: { parentSessionId: S1, childSessionId: S2, mode: 'continuable' },
  320. parentAvailable: true,
  321. })
  322. // Clicking the same child through an ordinary list-selection path must not
  323. // erase the catalog-derived address and fall back to session.* transport.
  324. manager.select(S2)
  325. expect(manager.getListSnapshot().currentAddress).toEqual({
  326. parentSessionId: S1, childSessionId: S2, mode: 'continuable',
  327. })
  328. expect(manager.get(S2).getSnapshot().subagent).toEqual({
  329. address: { parentSessionId: S1, childSessionId: S2, mode: 'continuable' },
  330. parentAvailable: true,
  331. })
  332. await manager.get(S2).open()
  333. await manager.get(S2).prompt([{ type: 'text', text: 'continue' }], 'queue')
  334. expect(api.callsOf('subagent.history')).toEqual([
  335. { parentSessionId: S1, childSessionId: S2, mode: 'continuable', maxMessages: 50 },
  336. ])
  337. expect(api.callsOf('subagent.prompt')).toEqual([
  338. {
  339. parentSessionId: S1, childSessionId: S2, mode: 'continuable',
  340. content: [{ type: 'text', text: 'continue' }],
  341. clientTimeZone: new Intl.DateTimeFormat().resolvedOptions().timeZone,
  342. },
  343. ])
  344. expect(api.callsOf('session.history')).toEqual([])
  345. expect(api.callsOf('session.prompt')).toEqual([])
  346. const listCalls = api.callsOf('subagent.list').length
  347. manager.handleHostEnvelope({
  348. rpcId: 'child-complete' as never,
  349. payload: { type: 'host/session-status', sessionId: S2, running: false },
  350. })
  351. expect(manager.getListSnapshot().subagentsByParent[S1]?.entries[0]).toMatchObject({
  352. kind: 'child', id: S2, activity: 'inactive',
  353. })
  354. expect(api.callsOf('subagent.list')).toHaveLength(listCalls)
  355. manager.handleHostEnvelope({
  356. rpcId: 'child-detached' as never,
  357. payload: { type: 'host/session-removed', sessionId: S2 },
  358. })
  359. expect(manager.getListSnapshot().items.find(item => item.sessionId === S2)).toMatchObject({
  360. origin: 'subagent', parentSessionId: S1, running: false,
  361. })
  362. expect(manager.get(S2).getSnapshot()).toMatchObject({
  363. removed: false,
  364. subagent: {
  365. address: { parentSessionId: S1, childSessionId: S2, mode: 'continuable' },
  366. },
  367. })
  368. })
  369. it('refetches debounced membership only while the parent catalog is open', async () => {
  370. vi.useFakeTimers()
  371. try {
  372. const api = new FakeApiClient()
  373. const manager = new SessionManager(api, fakeRemote())
  374. await manager.refreshSubagents(S1)
  375. manager.setSubagentCatalogOpen(S1, true)
  376. await Promise.resolve()
  377. const baseline = api.callsOf('subagent.list').length
  378. manager.handleHostEnvelope({
  379. rpcId: 'child-added' as never,
  380. payload: {
  381. type: 'host/session-added', sessionId: S2, parentSessionId: S1, blank: false,
  382. },
  383. })
  384. manager.handleHostEnvelope({
  385. rpcId: 'child-added-again' as never,
  386. payload: {
  387. type: 'host/session-added', sessionId: 'fk-m3' as SessionId, parentSessionId: S1, blank: false,
  388. },
  389. })
  390. await vi.advanceTimersByTimeAsync(50)
  391. expect(api.callsOf('subagent.list')).toHaveLength(baseline + 1)
  392. manager.setSubagentCatalogOpen(S1, false)
  393. manager.handleHostEnvelope({
  394. rpcId: 'child-added-closed' as never,
  395. payload: {
  396. type: 'host/session-added', sessionId: 'fk-m4' as SessionId, parentSessionId: S1, blank: false,
  397. },
  398. })
  399. await vi.advanceTimersByTimeAsync(50)
  400. expect(api.callsOf('subagent.list')).toHaveLength(baseline + 1)
  401. } finally {
  402. vi.useRealTimers()
  403. }
  404. })
  405. it('marks a loaded parent row expandable only for a direct subagent publication', async () => {
  406. const api = new FakeApiClient()
  407. const root = 'fk-root' as SessionId
  408. api.onSubagentList = () => Promise.resolve(ok({
  409. entries: [
  410. {
  411. kind: 'child', id: S1, mode: 'continuable', label: 'parent',
  412. activity: 'inactive', hasChildren: false,
  413. },
  414. {
  415. kind: 'child', id: S2, mode: 'continuable', label: 'ordinary parent',
  416. activity: 'inactive', hasChildren: false,
  417. },
  418. ] as never[],
  419. parentAvailable: true,
  420. }))
  421. const manager = new SessionManager(api, fakeRemote())
  422. await manager.refreshSubagents(root)
  423. manager.handleHostEnvelope({
  424. rpcId: 'nested-subagent' as never,
  425. payload: {
  426. type: 'host/session-added', sessionId: 'fk-grandchild' as SessionId,
  427. parentSessionId: S1, origin: 'subagent', blank: false,
  428. },
  429. })
  430. manager.handleHostEnvelope({
  431. rpcId: 'ordinary-fork' as never,
  432. payload: {
  433. type: 'host/session-added', sessionId: 'fk-fork' as SessionId,
  434. parentSessionId: S2, blank: false,
  435. },
  436. })
  437. expect(manager.getListSnapshot().subagentsByParent[root]?.entries).toMatchObject([
  438. { kind: 'child', id: S1, hasChildren: true },
  439. { kind: 'child', id: S2, hasChildren: false },
  440. ])
  441. })
  442. it('preserves a live expandability hint across only the older in-flight catalog response', async () => {
  443. const api = new FakeApiClient()
  444. const root = 'fk-root' as SessionId
  445. const response = deferred<Awaited<ReturnType<FakeApiClient['onSubagentList']>>>()
  446. api.onSubagentList = () => response.promise
  447. const manager = new SessionManager(api, fakeRemote())
  448. const refresh = manager.refreshSubagents(root)
  449. manager.handleHostEnvelope({
  450. rpcId: 'nested-subagent' as never,
  451. payload: {
  452. type: 'host/session-added', sessionId: 'fk-grandchild' as SessionId,
  453. parentSessionId: S1, origin: 'subagent', blank: false,
  454. },
  455. })
  456. response.resolve(ok({
  457. entries: [{
  458. kind: 'child', id: S1, mode: 'continuable', label: 'parent',
  459. activity: 'inactive', hasChildren: false,
  460. }] as never[],
  461. parentAvailable: true,
  462. }))
  463. await refresh
  464. expect(manager.getListSnapshot().subagentsByParent[root]?.entries).toMatchObject([
  465. { kind: 'child', id: S1, hasChildren: true },
  466. ])
  467. api.onSubagentList = () => Promise.resolve(ok({
  468. entries: [{
  469. kind: 'child', id: S1, mode: 'continuable', label: 'parent',
  470. activity: 'inactive', hasChildren: false,
  471. }] as never[],
  472. parentAvailable: true,
  473. }))
  474. await manager.refreshSubagents(root)
  475. expect(manager.getListSnapshot().subagentsByParent[root]?.entries).toMatchObject([
  476. { kind: 'child', id: S1, hasChildren: false },
  477. ])
  478. })
  479. it('replays status frames over an older in-flight catalog response', async () => {
  480. const api = new FakeApiClient()
  481. const root = 'fk-root' as SessionId
  482. const response = deferred<Awaited<ReturnType<FakeApiClient['onSubagentList']>>>()
  483. api.onSubagentList = () => response.promise
  484. const manager = new SessionManager(api, fakeRemote())
  485. const refresh = manager.refreshSubagents(root)
  486. manager.handleHostEnvelope({
  487. rpcId: 'child-stopped' as never,
  488. payload: { type: 'host/session-status', sessionId: S1, running: false },
  489. })
  490. manager.handleHostEnvelope({
  491. rpcId: 'child-started' as never,
  492. payload: { type: 'host/session-status', sessionId: S2, running: true },
  493. })
  494. response.resolve(ok({
  495. entries: [
  496. {
  497. kind: 'child', id: S1, mode: 'continuable', label: 'stopped',
  498. activity: 'running', hasChildren: false,
  499. },
  500. {
  501. kind: 'child', id: S2, mode: 'continuable', label: 'started',
  502. activity: 'inactive', hasChildren: false,
  503. },
  504. ] as never[],
  505. parentAvailable: true,
  506. }))
  507. await refresh
  508. expect(manager.getListSnapshot().subagentsByParent[root]?.entries).toMatchObject([
  509. { kind: 'child', id: S1, activity: 'inactive' },
  510. { kind: 'child', id: S2, activity: 'running' },
  511. ])
  512. })
  513. it('marks a detached catalog child inactive without requiring a selected address', async () => {
  514. const api = new FakeApiClient()
  515. api.onSubagentList = () => Promise.resolve(ok({
  516. entries: [{
  517. kind: 'child', id: S2, mode: 'continuable', label: 'worker',
  518. activity: 'running', hasChildren: false,
  519. }] as never[],
  520. parentAvailable: true,
  521. }))
  522. const manager = new SessionManager(api, fakeRemote())
  523. await manager.refreshSubagents(S1)
  524. manager.handleHostEnvelope({
  525. rpcId: 'child-detached' as never,
  526. payload: { type: 'host/session-removed', sessionId: S2 },
  527. })
  528. expect(manager.getListSnapshot().subagentsByParent[S1]?.entries).toMatchObject([
  529. { kind: 'child', id: S2, activity: 'inactive' },
  530. ])
  531. })
  532. it('coalesces overlapping catalog reads without scheduling a trailing pull', async () => {
  533. const api = new FakeApiClient()
  534. const root = 'fk-root' as SessionId
  535. const first = deferred<Awaited<ReturnType<FakeApiClient['onSubagentList']>>>()
  536. api.onSubagentList = () => first.promise
  537. const manager = new SessionManager(api, fakeRemote())
  538. const refresh = manager.refreshSubagents(root)
  539. expect(manager.refreshSubagents(root)).toBe(refresh)
  540. api.onSubagentList = () => Promise.resolve(ok({ entries: [], parentAvailable: true }))
  541. first.resolve(ok({ entries: [], parentAvailable: true }))
  542. await refresh
  543. expect(api.callsOf('subagent.list')).toHaveLength(1)
  544. })
  545. it('runs one trailing catalog refresh for a membership change coalesced into an in-flight pull', async () => {
  546. vi.useFakeTimers()
  547. try {
  548. const api = new FakeApiClient()
  549. const root = 'fk-root' as SessionId
  550. const first = deferred<Awaited<ReturnType<FakeApiClient['onSubagentList']>>>()
  551. const second = deferred<Awaited<ReturnType<FakeApiClient['onSubagentList']>>>()
  552. api.onSubagentList = () => first.promise
  553. const manager = new SessionManager(api, fakeRemote(), root)
  554. const refresh = manager.refreshSubagents(root)
  555. // A membership frame arrives while the pull is in flight; the debounced
  556. // refresh it schedules fires 50ms later and is coalesced into the pull —
  557. // which was requested before the new child existed. The stale mark must
  558. // queue one trailing pull carrying the change.
  559. manager.handleHostEnvelope({
  560. rpcId: 'child-added' as never,
  561. payload: {
  562. type: 'host/session-added', sessionId: S2, parentSessionId: root, blank: false,
  563. },
  564. })
  565. await vi.advanceTimersByTimeAsync(50)
  566. api.onSubagentList = () => second.promise
  567. first.resolve(ok({
  568. entries: [{
  569. kind: 'child', id: S1, mode: 'continuable', label: 'older',
  570. activity: 'inactive', hasChildren: false,
  571. }] as never[],
  572. parentAvailable: true,
  573. }))
  574. await refresh
  575. // The trailing pull is already in flight (kicked synchronously in finally).
  576. second.resolve(ok({
  577. entries: [
  578. {
  579. kind: 'child', id: S1, mode: 'continuable', label: 'older',
  580. activity: 'inactive', hasChildren: false,
  581. },
  582. {
  583. kind: 'child', id: S2, mode: 'continuable', label: 'new child',
  584. activity: 'inactive', hasChildren: false,
  585. },
  586. ] as never[],
  587. parentAvailable: true,
  588. }))
  589. await second.promise
  590. expect(api.callsOf('subagent.list')).toHaveLength(2)
  591. expect(manager.getListSnapshot().subagentsByParent[root]?.entries).toMatchObject([
  592. { kind: 'child', id: S1, label: 'older' },
  593. { kind: 'child', id: S2, label: 'new child' },
  594. ])
  595. } finally {
  596. vi.useRealTimers()
  597. }
  598. })
  599. it('keeps removal invalidation across a stale success and failed trailing pull', async () => {
  600. const api = new FakeApiClient()
  601. const root = 'fk-root' as SessionId
  602. const child = () => ({
  603. kind: 'child' as const, id: S2, mode: 'continuable' as const, label: 'worker',
  604. activity: 'inactive' as const, hasChildren: false,
  605. })
  606. const first = deferred<Awaited<ReturnType<FakeApiClient['onSubagentList']>>>()
  607. api.onSubagentList = () => first.promise
  608. const manager = new SessionManager(api, fakeRemote())
  609. const refresh = manager.refreshSubagents(root)
  610. first.resolve(ok({ entries: [child()] as never[], parentAvailable: true }))
  611. await refresh
  612. manager.selectSubagent({ parentSessionId: root, childSessionId: S2, mode: 'continuable' })
  613. // The removal lands while a second pull is in flight: the invalidation
  614. // must survive the pre-removal ok response, so one trailing pull runs.
  615. const mid = deferred<Awaited<ReturnType<FakeApiClient['onSubagentList']>>>()
  616. api.onSubagentList = () => mid.promise
  617. const midRefresh = manager.refreshSubagents(root)
  618. manager.handleHostEnvelope({
  619. rpcId: 'parent-removed-mid-pull' as never,
  620. payload: { type: 'host/session-removed', sessionId: root },
  621. })
  622. const trailing = deferred<Awaited<ReturnType<FakeApiClient['onSubagentList']>>>()
  623. api.onSubagentList = () => trailing.promise
  624. mid.resolve(ok({ entries: [child()] as never[], parentAvailable: true }))
  625. await midRefresh
  626. expect(manager.getListSnapshot().subagentsByParent[root]?.parentAvailable).toBe(false)
  627. expect(manager.get(S2).getSnapshot().subagent).toMatchObject({ parentAvailable: false })
  628. trailing.resolve(err({ code: 'internal', message: 'trailing pull failed', details: {} }))
  629. await vi.waitFor(() => {
  630. expect(manager.getListSnapshot().subagentsByParent[root]).toMatchObject({
  631. state: 'error',
  632. parentAvailable: false,
  633. })
  634. })
  635. const rootCalls = api.callsOf('subagent.list')
  636. .filter(call => (call as { parentSessionId: SessionId }).parentSessionId === root)
  637. expect(rootCalls).toHaveLength(3)
  638. expect(manager.getListSnapshot().subagentsByParent[root]?.parentAvailable).toBe(false)
  639. expect(manager.get(S2).getSnapshot().subagent).toMatchObject({ parentAvailable: false })
  640. })
  641. it('invalidates catalog availability when the owning parent is removed', async () => {
  642. const api = new FakeApiClient()
  643. const root = 'fk-root' as SessionId
  644. api.onSubagentList = () => Promise.resolve(ok({
  645. entries: [{
  646. kind: 'child', id: S2, mode: 'continuable', label: 'worker',
  647. activity: 'inactive', hasChildren: false,
  648. }] as never[],
  649. parentAvailable: true,
  650. }))
  651. const manager = new SessionManager(api, fakeRemote())
  652. await manager.refreshSubagents(root)
  653. manager.selectSubagent({ parentSessionId: root, childSessionId: S2, mode: 'continuable' })
  654. expect(manager.get(S2).getSnapshot().subagent).toMatchObject({ parentAvailable: true })
  655. manager.handleHostEnvelope({
  656. rpcId: 'parent-removed' as never,
  657. payload: { type: 'host/session-removed', sessionId: root },
  658. })
  659. expect(manager.getListSnapshot().subagentsByParent[root]?.parentAvailable).toBe(false)
  660. expect(manager.get(S2).getSnapshot().subagent).toMatchObject({ parentAvailable: false })
  661. })
  662. })
  663. describe('remaining branches', () => {
  664. it('refreshList folds a transport throw into the error state', async () => {
  665. const api = new FakeApiClient()
  666. api.onList = () => Promise.reject(new Error('list wire down'))
  667. const manager = new SessionManager(api, fakeRemote())
  668. await manager.refreshList()
  669. expect(manager.getListSnapshot()).toMatchObject({ state: 'error', error: { code: 'internal', message: 'list wire down' } })
  670. })
  671. it('refreshList pushes running bits down to already-instantiated sessions', async () => {
  672. const api = new FakeApiClient()
  673. const manager = new SessionManager(api, fakeRemote())
  674. const session = manager.get(S1)
  675. api.onList = () => Promise.resolve(ok({ items: [summary(S1, { running: true })] as never[] }))
  676. await manager.refreshList()
  677. expect(session.getSnapshot().running).toBe(true)
  678. })
  679. it('create passes cwd and a preallocated id, folds transport throws, and deduplicates the echo', async () => {
  680. const api = new FakeApiClient()
  681. api.onCreate = () => Promise.resolve(ok({ sessionId: S1 }))
  682. const manager = new SessionManager(api, fakeRemote())
  683. await manager.create({ cwd: '/tmp/w', sessionId: S1 })
  684. expect(api.callsOf('session.create')).toEqual([{ cwd: '/tmp/w', sessionId: S1 }])
  685. expect(manager.getListSnapshot().items[0]).toMatchObject({ sessionId: S1, cwd: '/tmp/w' })
  686. await manager.create({ cwd: '/tmp/w' }) // same id returned: no duplicate row
  687. expect(manager.getListSnapshot().items).toHaveLength(1)
  688. api.onCreate = () => Promise.reject(new Error('create wire down'))
  689. expect(await manager.create()).toMatchObject({ ok: false, error: { code: 'internal' } })
  690. // Business error passes through untouched.
  691. api.onCreate = () => Promise.resolve(err({ code: 'internal', message: 'no', details: {} }))
  692. expect(await manager.create()).toMatchObject({ ok: false })
  693. })
  694. it('publishes a real Ungrouped summary from workspace-attach-failed', async () => {
  695. const api = new FakeApiClient()
  696. api.onCreate = () => Promise.resolve(err({
  697. code: 'workspace-attach-failed',
  698. message: 'published but unattached',
  699. details: { sessionId: S1, workspaceId: 'w1' },
  700. } as never))
  701. const manager = new SessionManager(api, fakeRemote())
  702. const result = await manager.create({ workspaceId: 'w1' as never, sessionId: S1 })
  703. expect(result).toMatchObject({ ok: false, error: { code: 'workspace-attach-failed' } })
  704. expect(manager.getListSnapshot().items).toEqual([expect.objectContaining({ sessionId: S1 })])
  705. expect(manager.getListSnapshot().items[0]).not.toHaveProperty('cwd')
  706. })
  707. it('reconciles a fork child published before workspace attachment fails', async () => {
  708. const api = new FakeApiClient()
  709. api.onFork = () => Promise.resolve(err({
  710. code: 'workspace-attach-failed',
  711. message: 'forked but unattached',
  712. details: { sessionId: S2, workspaceId: 'w1' },
  713. } as never))
  714. const manager = new SessionManager(api, fakeRemote())
  715. const result = await manager.fork({ sessionId: S1 })
  716. expect(result).toMatchObject({ ok: false, error: { code: 'workspace-attach-failed' } })
  717. expect(manager.getListSnapshot().items).toEqual([expect.objectContaining({
  718. sessionId: S2,
  719. parentSessionId: S1,
  720. blank: false,
  721. })])
  722. })
  723. it('reconciles a preallocated id after an ordinary transport failure', async () => {
  724. const api = new FakeApiClient()
  725. api.onCreate = () => Promise.reject(new Error('response lost'))
  726. const manager = new SessionManager(api, fakeRemote())
  727. const failed = await manager.create({ workspaceId: 'w1' as never, sessionId: S1 })
  728. expect(failed).toMatchObject({ ok: false, error: { message: 'response lost' } })
  729. expect(manager.getListSnapshot().items).toEqual([])
  730. manager.handleHostEnvelope({
  731. rpcId: 'published-later' as never,
  732. payload: { type: 'host/session-added', blank: true, sessionId: S1, cwd: '/w/one' },
  733. })
  734. expect(manager.getListSnapshot().items).toEqual([
  735. expect.objectContaining({ sessionId: S1, cwd: '/w/one' }),
  736. ])
  737. manager.handleHostEnvelope({
  738. rpcId: 'duplicate-frame' as never,
  739. payload: { type: 'host/session-added', blank: true, sessionId: S1, cwd: '/w/one' },
  740. })
  741. expect(manager.getListSnapshot().items).toHaveLength(1)
  742. })
  743. it('subscribe notifies on list changes and stops after unsubscribe', async () => {
  744. const api = new FakeApiClient()
  745. const manager = new SessionManager(api, fakeRemote())
  746. let notified = 0
  747. const unsubscribe = manager.subscribe(() => { notified++ })
  748. await manager.refreshList()
  749. await new Promise(resolve => setTimeout(resolve, 0))
  750. expect(notified).toBeGreaterThan(0)
  751. const seen = notified
  752. unsubscribe()
  753. manager.handleHostEnvelope({ rpcId: 'h' as never, payload: { type: 'host/session-added', blank: true, sessionId: S1 } })
  754. await new Promise(resolve => setTimeout(resolve, 0))
  755. expect(notified).toBe(seen)
  756. })
  757. it('routes stream/error and unknown frames to the documented drops, and dispatches to instantiated sessions', () => {
  758. const api = new FakeApiClient()
  759. const manager = new SessionManager(api, fakeRemote())
  760. manager.handleMuxEnvelope({ rpcId: 'e' as never, payload: { type: 'stream/error', error: { code: 'internal', message: 'x', details: {} } } })
  761. manager.handleHostEnvelope({ rpcId: 'e2' as never, payload: { type: 'stream/error', error: { code: 'internal', message: 'x', details: {} } } })
  762. manager.handleHostEnvelope({ rpcId: 'e3' as never, payload: { type: 'future/host-frame' } as never })
  763. const session = manager.get(S1)
  764. manager.handleMuxEnvelope({ rpcId: 'q1' as never, payload: { type: 'question/requested', sessionId: S1, questions: [] } })
  765. expect(session.getSnapshot().pending).toMatchObject([{ kind: 'question' }])
  766. // status flip for an unknown session only touches summaries (no crash).
  767. manager.handleHostEnvelope({ rpcId: 'h9' as never, payload: { type: 'host/session-status', sessionId: S2, running: true } })
  768. manager.handleHostEnvelope({ rpcId: 'ha' as never, payload: { type: 'host/agent-error', sessionId: S2, message: '无实例' } })
  769. })
  770. it('keeps list-entry identity for unchanged rows across an unrelated list change', async () => {
  771. const api = new FakeApiClient()
  772. api.onList = () => Promise.resolve(ok({ items: [summary(S1), summary(S2, { updatedAt: 200 })] as never[] }))
  773. const manager = new SessionManager(api, fakeRemote())
  774. await manager.refreshList()
  775. const before = manager.getListSnapshot()
  776. manager.handleHostEnvelope({ rpcId: 'h' as never, payload: { type: 'host/session-status', sessionId: S2, running: true } })
  777. const after = manager.getListSnapshot()
  778. expect(after.items).not.toBe(before.items)
  779. const beforeS1 = before.items.find(e => e.sessionId === S1)
  780. const afterS1 = after.items.find(e => e.sessionId === S1)
  781. expect(afterS1).toBe(beforeS1) // untouched entry keeps identity (entryCache)
  782. // Same-order same-entries snapshot reuses the items array.
  783. manager.handleHostEnvelope({ rpcId: 'h2' as never, payload: { type: 'host/agent-error', sessionId: S1, message: 'x' } })
  784. expect(manager.getListSnapshot().items).toBe(after.items)
  785. })
  786. it('carries parentSessionId from host/session-added into the lineage row', () => {
  787. const api = new FakeApiClient()
  788. const manager = new SessionManager(api, fakeRemote())
  789. manager.handleHostEnvelope({ rpcId: 'h1' as never, payload: { type: 'host/session-added', blank: true, sessionId: S1 } })
  790. manager.handleHostEnvelope({
  791. rpcId: 'h2' as never,
  792. payload: {
  793. type: 'host/session-added', blank: true, sessionId: S2,
  794. parentSessionId: S1, origin: 'subagent',
  795. },
  796. })
  797. const items = manager.getListSnapshot().items
  798. expect(items.find(e => e.sessionId === S2)).toMatchObject({
  799. parentSessionId: S1, origin: 'subagent', depth: 1,
  800. })
  801. })
  802. })
  803. describe('connected generation', () => {
  804. it('refreshes the list and resyncs only opened instances', async () => {
  805. const api = new FakeApiClient()
  806. api.onHistory = () => Promise.resolve(ok({
  807. events: entries(plainTurn(0, 0, 'a', 'b')) as never[],
  808. hasMore: false,
  809. modelSelection: { provider: 'deepseek-official', model: 'deepseek-chat' },
  810. }))
  811. const manager = new SessionManager(api, fakeRemote())
  812. const openedSession = manager.get(S1)
  813. await openedSession.open()
  814. manager.get(S2) // instantiated but never opened
  815. const historyCallsBefore = api.callsOf('session.history').length
  816. manager.handleConnected()
  817. await vi.waitFor(() => {
  818. expect(api.callsOf('session.list').length).toBe(1)
  819. // Only the opened instance repulls history; the cold one stays silent.
  820. expect(api.callsOf('session.history').length).toBe(historyCallsBefore + 1)
  821. })
  822. })
  823. it('reloads the durable parent address for a restored child selection', async () => {
  824. const api = new FakeApiClient()
  825. const address = {
  826. parentSessionId: S1, childSessionId: S2, mode: 'continuable' as const,
  827. }
  828. const manager = new SessionManager(api, fakeRemote(), S2, address)
  829. manager.handleConnected()
  830. await vi.waitFor(() => {
  831. expect(api.callsOf('subagent.list')).toContainEqual({ parentSessionId: S1 })
  832. })
  833. expect(manager.getListSnapshot().currentAddress).toEqual(address)
  834. })
  835. })
  836. describe('pending-interaction list status', () => {
  837. it('tracks approval requests through replay and resolution without instantiation', () => {
  838. const manager = new SessionManager(new FakeApiClient(), fakeRemote())
  839. manager.handleHostEnvelope({ rpcId: 'h1' as never, payload: { type: 'host/session-added', sessionId: S1, blank: false } })
  840. expect(manager.getListSnapshot().items[0]?.pendingInteraction).toBeUndefined()
  841. manager.handleMuxEnvelope({ rpcId: 'ra' as never, payload: { type: 'approval/requested', sessionId: S1, approvalId: 'ap1' as never, toolName: 'rm' } })
  842. expect(manager.getListSnapshot().items[0]?.pendingInteraction).toBe('approval')
  843. // Mux-open replay of the same question (same approvalId) is idempotent.
  844. manager.handleMuxEnvelope({ rpcId: 'ra' as never, payload: { type: 'approval/requested', sessionId: S1, approvalId: 'ap1' as never, toolName: 'rm' } })
  845. expect(manager.getListSnapshot().items[0]?.pendingInteraction).toBe('approval')
  846. manager.handleMuxEnvelope({ rpcId: 'rx' as never, payload: { type: 'approval/resolved', sessionId: S1, approvalId: 'ap1' as never, outcome: 'allowed-once' as never } })
  847. expect(manager.getListSnapshot().items[0]?.pendingInteraction).toBeUndefined()
  848. })
  849. it('classifies ordinary questions and renderable plan reviews, then clears by question rpcId', () => {
  850. const manager = new SessionManager(new FakeApiClient(), fakeRemote())
  851. manager.handleHostEnvelope({ rpcId: 'h1' as never, payload: { type: 'host/session-added', sessionId: S1, blank: false } })
  852. manager.handleMuxEnvelope({
  853. rpcId: 'q1' as never,
  854. payload: { type: 'question/requested', sessionId: S1, questions: [{ id: 'name', question: 'Name?' }] },
  855. })
  856. expect(manager.getListSnapshot().items[0]?.pendingInteraction).toBe('question')
  857. manager.handleMuxEnvelope({ rpcId: 'qx' as never, payload: { type: 'question/resolved', sessionId: S1, questionRpcId: 'q1' as never, outcome: 'answered' } })
  858. expect(manager.getListSnapshot().items[0]?.pendingInteraction).toBeUndefined()
  859. manager.handleMuxEnvelope({
  860. rpcId: 'q2' as never,
  861. payload: {
  862. type: 'question/requested',
  863. sessionId: S1,
  864. questions: [{
  865. id: 'plan', question: 'Approve?', detail: '# Plan',
  866. options: [{ label: 'Approve' }, { label: 'Refuse' }],
  867. intent: { kind: 'plan-review', approve: 'Approve' },
  868. }],
  869. },
  870. })
  871. expect(manager.getListSnapshot().items[0]?.pendingInteraction).toBe('plan-review')
  872. manager.handleMuxEnvelope({ rpcId: 'qy' as never, payload: { type: 'question/resolved', sessionId: S1, questionRpcId: 'q2' as never, outcome: 'cancelled' } })
  873. expect(manager.getListSnapshot().items[0]?.pendingInteraction).toBeUndefined()
  874. })
  875. it.each([
  876. ['missing detail', {}],
  877. ['multi-select', { detail: '# Plan', multiSelect: true }],
  878. ['more than two options', { detail: '# Plan', options: [{ label: 'Approve' }, { label: 'Refuse' }, { label: 'Revise' }] }],
  879. ['missing approve option', { detail: '# Plan', options: [{ label: 'Refuse' }] }],
  880. ])('keeps an unrenderable %s plan intent on the ordinary question flow', (_name, over) => {
  881. const manager = new SessionManager(new FakeApiClient(), fakeRemote())
  882. manager.handleHostEnvelope({ rpcId: 'h1' as never, payload: { type: 'host/session-added', sessionId: S1, blank: false } })
  883. manager.handleMuxEnvelope({
  884. rpcId: 'q-plan' as never,
  885. payload: {
  886. type: 'question/requested', sessionId: S1,
  887. questions: [{
  888. id: 'plan', question: 'Approve?', options: [{ label: 'Approve' }],
  889. intent: { kind: 'plan-review', approve: 'Approve' },
  890. ...over,
  891. }],
  892. },
  893. })
  894. expect(manager.getListSnapshot().items[0]?.pendingInteraction).toBe('question')
  895. })
  896. it('the first question outranks sibling approvals and resolving it reveals the remaining wait', () => {
  897. const manager = new SessionManager(new FakeApiClient(), fakeRemote())
  898. manager.handleHostEnvelope({ rpcId: 'h1' as never, payload: { type: 'host/session-added', sessionId: S1, blank: false } })
  899. manager.handleMuxEnvelope({ rpcId: 'r1' as never, payload: { type: 'approval/requested', sessionId: S1, approvalId: 'a1' as never, toolName: 'rm' } })
  900. manager.handleMuxEnvelope({
  901. rpcId: 'q1' as never,
  902. payload: { type: 'question/requested', sessionId: S1, questions: [{ id: 'name', question: 'Name?' }] },
  903. })
  904. expect(manager.getListSnapshot().items[0]?.pendingInteraction).toBe('question')
  905. manager.handleMuxEnvelope({ rpcId: 'qy' as never, payload: { type: 'question/resolved', sessionId: S1, questionRpcId: 'q1' as never, outcome: 'answered' } })
  906. expect(manager.getListSnapshot().items[0]?.pendingInteraction).toBe('approval')
  907. manager.handleMuxEnvelope({ rpcId: 'rx' as never, payload: { type: 'approval/resolved', sessionId: S1, approvalId: 'a1' as never, outcome: 'rejected' as never } })
  908. expect(manager.getListSnapshot().items[0]?.pendingInteraction).toBeUndefined()
  909. manager.handleMuxEnvelope({ rpcId: 'r2' as never, payload: { type: 'approval/requested', sessionId: S1, approvalId: 'a2' as never, toolName: 'rm' } })
  910. manager.handleHostEnvelope({ rpcId: 'h2' as never, payload: { type: 'host/session-removed', sessionId: S1 } })
  911. expect(manager.getListSnapshot().items).toHaveLength(0)
  912. })
  913. it('drops stale status at generation death before replay re-adds live interactions', () => {
  914. const manager = new SessionManager(new FakeApiClient(), fakeRemote())
  915. manager.handleHostEnvelope({ rpcId: 'h1' as never, payload: { type: 'host/session-added', sessionId: S1, blank: false } })
  916. manager.handleMuxEnvelope({ rpcId: 'ra' as never, payload: { type: 'approval/requested', sessionId: S1, approvalId: 'ap1' as never, toolName: 'rm' } })
  917. expect(manager.getListSnapshot().items[0]?.pendingInteraction).toBe('approval')
  918. // Generation death clears (resolved-while-disconnected questions send no frame)…
  919. manager.handleDisconnected()
  920. expect(manager.getListSnapshot().items[0]?.pendingInteraction).toBeUndefined()
  921. // …and a replayed frame arriving before onConnected (stream open precedes
  922. // the readiness handshake) survives the later handleConnected untouched.
  923. manager.handleMuxEnvelope({ rpcId: 'ra' as never, payload: { type: 'approval/requested', sessionId: S1, approvalId: 'ap1' as never, toolName: 'rm' } })
  924. manager.handleConnected()
  925. expect(manager.getListSnapshot().items[0]?.pendingInteraction).toBe('approval')
  926. })
  927. it('generation death drops buffered answerable frames (a dead generation cannot be answered)', () => {
  928. const manager = new SessionManager(new FakeApiClient(), fakeRemote())
  929. manager.handleHostEnvelope({ rpcId: 'h1' as never, payload: { type: 'host/session-added', sessionId: S1, blank: false } })
  930. // Buffered pre-instantiation: an approval pair and a queued row.
  931. manager.handleMuxEnvelope({ rpcId: 'ra' as never, payload: { type: 'approval/requested', sessionId: S1, approvalId: 'ap1' as never, toolName: 'rm' } })
  932. manager.handleMuxEnvelope({ rpcId: 'q1' as never, payload: { type: 'question/requested', sessionId: S1, questions: [] } })
  933. manager.handleDisconnected()
  934. // Instantiate after the death sweep: no zombie interaction replays (the
  935. // pendingBuffers held only dead-generation rpcIds), so the session mints
  936. // no pending waits.
  937. const session = manager.get(S1)
  938. expect(session.getSnapshot().pending).toEqual([])
  939. })
  940. })
  941. describe('completed reminder', () => {
  942. const status = (rpcId: string, sessionId: SessionId, running: boolean) => ({
  943. rpcId: rpcId as never,
  944. payload: { type: 'host/session-status' as const, sessionId, running },
  945. })
  946. const added = (rpcId: string, sessionId: SessionId) => ({
  947. rpcId: rpcId as never,
  948. payload: { type: 'host/session-added' as const, sessionId, blank: false },
  949. })
  950. const entry = (manager: SessionManager, sessionId: SessionId) =>
  951. manager.getListSnapshot().items.find(item => item.sessionId === sessionId)
  952. it('arms on a running→idle flip of a non-selected session and clears on select', () => {
  953. const manager = new SessionManager(new FakeApiClient(), fakeRemote())
  954. manager.handleHostEnvelope(added('h1', S1))
  955. manager.handleHostEnvelope(added('h2', S2))
  956. manager.select(S1)
  957. expect(entry(manager, S2)?.completed).toBe(false)
  958. manager.handleHostEnvelope(status('s1', S2, true))
  959. manager.handleHostEnvelope(status('s2', S2, false))
  960. expect(entry(manager, S2)?.completed).toBe(true)
  961. // Opening the session consumes the reminder.
  962. manager.select(S2)
  963. expect(entry(manager, S2)?.completed).toBe(false)
  964. })
  965. it('never arms for the session being watched and re-arms after a switch-away re-run', () => {
  966. const manager = new SessionManager(new FakeApiClient(), fakeRemote())
  967. manager.handleHostEnvelope(added('h1', S1))
  968. manager.handleHostEnvelope(added('h2', S2))
  969. manager.select(S2)
  970. manager.handleHostEnvelope(status('s1', S2, true))
  971. manager.handleHostEnvelope(status('s2', S2, false))
  972. expect(entry(manager, S2)?.completed).toBe(false) // watched to completion: no reminder
  973. // Switch away; a fresh run completing again arms the reminder.
  974. manager.select(S1)
  975. manager.handleHostEnvelope(status('s3', S2, true))
  976. manager.handleHostEnvelope(status('s4', S2, false))
  977. expect(entry(manager, S2)?.completed).toBe(true)
  978. })
  979. it('a re-run disarms the reminder while running and re-arms on its completion', () => {
  980. const manager = new SessionManager(new FakeApiClient(), fakeRemote())
  981. manager.handleHostEnvelope(added('h1', S1))
  982. manager.handleHostEnvelope(added('h2', S2))
  983. manager.select(S1)
  984. manager.handleHostEnvelope(status('s1', S2, true))
  985. manager.handleHostEnvelope(status('s2', S2, false))
  986. expect(entry(manager, S2)?.completed).toBe(true)
  987. // The user starts a new run without opening the session: running wins.
  988. manager.handleHostEnvelope(status('s3', S2, true))
  989. expect(entry(manager, S2)?.completed).toBe(false)
  990. manager.handleHostEnvelope(status('s4', S2, false))
  991. expect(entry(manager, S2)?.completed).toBe(true)
  992. })
  993. it('session-removed drops the reminder and a re-add starts clean', () => {
  994. const manager = new SessionManager(new FakeApiClient(), fakeRemote())
  995. manager.handleHostEnvelope(added('h1', S1))
  996. manager.handleHostEnvelope(added('h2', S2))
  997. manager.select(S1)
  998. manager.handleHostEnvelope(status('s1', S2, true))
  999. manager.handleHostEnvelope(status('s2', S2, false))
  1000. expect(entry(manager, S2)?.completed).toBe(true)
  1001. manager.handleHostEnvelope({ rpcId: 'rm' as never, payload: { type: 'host/session-removed', sessionId: S2 } })
  1002. expect(manager.getListSnapshot().items.find(item => item.sessionId === S2)).toBeUndefined()
  1003. manager.handleHostEnvelope(added('h3', S2))
  1004. expect(entry(manager, S2)?.completed).toBe(false)
  1005. })
  1006. it('a list refresh carrying the running→idle transition arms the reminder', async () => {
  1007. const api = new FakeApiClient()
  1008. api.onList = () => Promise.resolve(ok({ items: [summary(S1), summary(S2, { updatedAt: 200, running: true })] as never[] }))
  1009. const manager = new SessionManager(api, fakeRemote())
  1010. await manager.refreshList()
  1011. manager.select(S1)
  1012. expect(entry(manager, S2)?.completed).toBe(false)
  1013. api.onList = () => Promise.resolve(ok({ items: [summary(S1), summary(S2, { updatedAt: 200, running: false })] as never[] }))
  1014. await manager.refreshList()
  1015. expect(entry(manager, S2)?.completed).toBe(true)
  1016. })
  1017. it('never arms for sessions already idle at first observation', async () => {
  1018. const api = new FakeApiClient()
  1019. api.onList = () => Promise.resolve(ok({ items: [summary(S1), summary(S2, { updatedAt: 200 })] as never[] }))
  1020. const manager = new SessionManager(api, fakeRemote())
  1021. await manager.refreshList()
  1022. manager.select(S1)
  1023. expect(entry(manager, S2)?.completed).toBe(false)
  1024. api.onList = () => Promise.resolve(ok({ items: [summary(S1), summary(S2, { updatedAt: 201 })] as never[] }))
  1025. await manager.refreshList()
  1026. expect(entry(manager, S2)?.completed).toBe(false)
  1027. })
  1028. it('arms a completion that happened during an in-flight first pull (baseline running, replayed idle)', async () => {
  1029. const api = new FakeApiClient()
  1030. const gate = deferred<Awaited<ReturnType<FakeApiClient['onList']>>>()
  1031. api.onList = () => gate.promise
  1032. const manager = new SessionManager(api, fakeRemote())
  1033. const refresh = manager.refreshList()
  1034. // The session finishes while the first pull is still in flight; the pull
  1035. // response recorded it as running at pull time.
  1036. manager.handleHostEnvelope(status('s-mid', S2, false))
  1037. gate.resolve(ok({ items: [summary(S1), summary(S2, { updatedAt: 200, running: true })] as never[] }))
  1038. await refresh
  1039. expect(entry(manager, S2)?.completed).toBe(true)
  1040. })
  1041. it('arms when a session ran and completed entirely between in-flight mutations (baseline idle)', async () => {
  1042. const api = new FakeApiClient()
  1043. const gate = deferred<Awaited<ReturnType<FakeApiClient['onList']>>>()
  1044. api.onList = () => gate.promise
  1045. const manager = new SessionManager(api, fakeRemote())
  1046. const refresh = manager.refreshList()
  1047. // The unknown session starts and finishes while the first pull is in
  1048. // flight; the pull-time baseline recorded it idle, so the running→idle
  1049. // edge lives entirely inside the replayed mutations.
  1050. manager.handleHostEnvelope(status('s-start', S2, true))
  1051. manager.handleHostEnvelope(status('s-finish', S2, false))
  1052. gate.resolve(ok({ items: [summary(S1), summary(S2, { updatedAt: 200 })] as never[] }))
  1053. await refresh
  1054. expect(entry(manager, S2)?.completed).toBe(true)
  1055. })
  1056. })
  1057. describe('background-job mirror', () => {
  1058. const view = (over: Partial<{ id: string; status: string; label: string }> = {}) => ({
  1059. id: 'bash-1', kind: 'bash', label: 'pnpm run build', status: 'running', startedAt: 5, ...over,
  1060. })
  1061. const tasksFrame = (sessionId: SessionId, jobs: unknown[]) =>
  1062. ({ rpcId: 't' as never, payload: { type: 'session/jobs', sessionId, jobs } as never })
  1063. it('mirrors the whole set last-wins, keyed per session, with no Session instance needed', () => {
  1064. const manager = new SessionManager(new FakeApiClient(), fakeRemote())
  1065. manager.handleMuxEnvelope(tasksFrame(S1, [view()]))
  1066. manager.handleMuxEnvelope(tasksFrame(S2, [view({ id: 'pwsh-1', label: 'other' })]))
  1067. const first = manager.getListSnapshot().jobsBySession
  1068. expect(first[S1]).toEqual([view()])
  1069. expect(first[S2]?.[0]?.label).toBe('other')
  1070. // Last-wins: the newer whole set replaces, it does not merge.
  1071. manager.handleMuxEnvelope(tasksFrame(S1, [view({ status: 'completed' })]))
  1072. expect(manager.getListSnapshot().jobsBySession[S1]).toEqual([view({ status: 'completed' })])
  1073. })
  1074. it('stores an emptied set as an absent key so absence and [] read alike', () => {
  1075. const manager = new SessionManager(new FakeApiClient(), fakeRemote())
  1076. manager.handleMuxEnvelope(tasksFrame(S1, [view()]))
  1077. expect(S1 in manager.getListSnapshot().jobsBySession).toBe(true)
  1078. manager.handleMuxEnvelope(tasksFrame(S1, []))
  1079. expect(S1 in manager.getListSnapshot().jobsBySession).toBe(false)
  1080. })
  1081. it('clears the mirror on re-subscribe, because a task-free generation sends no baseline', () => {
  1082. const manager = new SessionManager(new FakeApiClient(), fakeRemote())
  1083. manager.handleMuxEnvelope(tasksFrame(S1, [view()]))
  1084. manager.handleMuxEnvelope({
  1085. rpcId: 's' as never,
  1086. payload: { type: 'session/subscribed', sessionId: S1, lastSeq: 3 },
  1087. })
  1088. expect(S1 in manager.getListSnapshot().jobsBySession).toBe(false)
  1089. })
  1090. it('drops the rows when the session is removed, whichever stream lands first', () => {
  1091. const manager = new SessionManager(new FakeApiClient(), fakeRemote())
  1092. manager.handleHostEnvelope({ rpcId: 'a' as never, payload: { type: 'host/session-added', blank: true, sessionId: S1 } })
  1093. manager.handleMuxEnvelope(tasksFrame(S1, [view()]))
  1094. manager.handleHostEnvelope({ rpcId: 'r' as never, payload: { type: 'host/session-removed', sessionId: S1 } })
  1095. expect(S1 in manager.getListSnapshot().jobsBySession).toBe(false)
  1096. })
  1097. it('notifies list subscribers so an open header re-renders without a poll', async () => {
  1098. const manager = new SessionManager(new FakeApiClient(), fakeRemote())
  1099. const seen = vi.fn()
  1100. manager.subscribe(seen)
  1101. manager.handleMuxEnvelope(tasksFrame(S1, [view()]))
  1102. // The notifier batches on a microtask; the frame itself is already applied.
  1103. await Promise.resolve()
  1104. expect(seen).toHaveBeenCalled()
  1105. })
  1106. })