manager.client.spec.ts 42 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962
  1. /**
  2. * SessionManager orchestration: lazy resident instances, list lifecycle, host
  3. * frame routing, and control baselines for uninstantiated sessions.
  4. */
  5. import { describe, expect, vi } from 'vitest'
  6. import type { SessionId } from '@deepseek-ai/dsh-api-remotes/client'
  7. import { SessionSeq } from '@deepseek-ai/dsh-session/types'
  8. import { RemoteError } from '@deepseek-ai/dsh-typert-protocol'
  9. import type { SessionControlFrame } from '@deepseek-ai/dsh-api-session-controller/types'
  10. import type { SubagentAddress } from '@deepseek-ai/dsh-subagent/client'
  11. import { ok, type RemoteMock } from '@deepseek-ai/dsh-remote-mock'
  12. import {
  13. createClientTest, type ClientTestFixtures, webApp,
  14. } from '@deepseek-ai/dsh-client-test-runtime/src/assembly/index.ts'
  15. import type {} from '@deepseek-ai/dsh-session-title/client'
  16. import { SessionManager } from '../src/client/sessions/manager.ts'
  17. import type { SessionRemotes } from '../src/client/sessions/remotes.ts'
  18. import { entries, plainTurn } from './event-script.client.ts'
  19. import { FOLLOW, err, followScript, sessionWorld } from './remote/session.client.ts'
  20. const S1 = 'fk-m1' as SessionId
  21. const S2 = 'fk-m2' as SessionId
  22. /** A SessionManager's Remote methods use the same native mocks as an assembled client. */
  23. const API_ROSTER = webApp.closure(['@deepseek-ai/dsh-api-gateway'])
  24. const it = createClientTest({ roster: API_ROSTER })
  25. type SummaryOver = Partial<{
  26. updatedAt: number
  27. running: boolean
  28. blank: boolean
  29. cwd: string
  30. parentSessionId: SessionId
  31. origin: 'subagent'
  32. }>
  33. function summary(sessionId: SessionId, over: SummaryOver = {}) {
  34. return { sessionId, updatedAt: 100, running: false, blank: false, ...over }
  35. }
  36. function makeManager(
  37. mock: RemoteMock,
  38. remote: ClientTestFixtures['remote'],
  39. restoredSelection?: SessionId,
  40. restoredAddress?: SubagentAddress,
  41. ): SessionManager {
  42. mock.load(sessionWorld)
  43. // Cases using this helper never open a Session, so they do not need the broader Client Remote's $stream member.
  44. return new SessionManager(remote as unknown as SessionRemotes, restoredSelection, restoredAddress)
  45. }
  46. describe('SessionManager instances', () => {
  47. it('lazily builds one resident instance per id and syncs the running bit from the list', async ({ mock, remote }) => {
  48. remote.session.list.mockResolvedValue(ok({ items: [summary(S1, { running: true })] as never[] }))
  49. const manager = makeManager(mock, remote)
  50. await manager.refreshList()
  51. const session = manager.get(S1)
  52. expect(manager.get(S1)).toBe(session) // resident: same instance forever
  53. expect(session.getSnapshot().running).toBe(true) // list preceded instantiation
  54. })
  55. })
  56. describe('list lifecycle', () => {
  57. it('single-flights refreshList and preserves the Host baseline order', async ({ mock, remote }) => {
  58. const gate = Promise.withResolvers<Awaited<ReturnType<typeof remote.session.list>>>()
  59. remote.session.list.mockReturnValue(gate.promise)
  60. const manager = makeManager(mock, remote)
  61. const first = manager.refreshList()
  62. const second = manager.refreshList()
  63. expect(manager.getListSnapshot().state).toBe('loading')
  64. gate.resolve(ok({ items: [summary(S2, { updatedAt: 200 }), summary(S1)] as never[] }))
  65. await Promise.all([first, second])
  66. expect(remote.session.list).toHaveBeenCalledOnce()
  67. const snapshot = manager.getListSnapshot()
  68. expect(snapshot.state).toBe('idle')
  69. expect(snapshot.items.map(i => i.sessionId)).toEqual([S2, S1])
  70. })
  71. it('replays incremental frames over hydration and never batch-reorders established ids', async ({ mock, remote }) => {
  72. const first = Promise.withResolvers<Awaited<ReturnType<typeof remote.session.list>>>()
  73. remote.session.list.mockReturnValue(first.promise)
  74. const manager = makeManager(mock, remote)
  75. const hydration = manager.refreshList()
  76. manager.handleSessionAdded(summary(S2, { blank: true }))
  77. first.resolve(ok({ items: [summary(S1)] as never[] }))
  78. await hydration
  79. expect(manager.getListSnapshot().items.map(item => item.sessionId)).toEqual([S2, S1])
  80. remote.session.list.mockResolvedValue(ok({
  81. items: [summary(S1, { updatedAt: 900 }), summary(S2, { updatedAt: 800 })] as never[],
  82. }))
  83. await manager.refreshList()
  84. expect(manager.getListSnapshot().items.map(item => item.sessionId)).toEqual([S2, S1])
  85. })
  86. it('advances list activity from the filtered Host notification', async ({ mock, remote }) => {
  87. remote.session.list.mockResolvedValue(ok({ items: [summary(S1)] as never[] }))
  88. const manager = makeManager(mock, remote)
  89. await manager.refreshList()
  90. manager.handleSessionActivity(S1, 500)
  91. expect(manager.getListSnapshot().items[0]?.updatedAt).toBe(500)
  92. })
  93. it('keeps the error in the list snapshot on failure', async ({ mock, remote }) => {
  94. remote.session.list.mockResolvedValue(err(new RemoteError('gateway/internal', 'boom', {})))
  95. const manager = makeManager(mock, remote)
  96. await manager.refreshList()
  97. expect(manager.getListSnapshot()).toMatchObject({ state: 'error', error: { code: 'gateway/internal' } })
  98. // A failed pull does not step the arrival phase: still pending.
  99. expect(manager.getListSnapshot().phase).toBe('pending')
  100. })
  101. it('phase steps pending → ready on the first successful pull and never returns', async ({ mock, remote }) => {
  102. const manager = makeManager(mock, remote)
  103. expect(manager.getListSnapshot().phase).toBe('pending')
  104. await manager.refreshList()
  105. expect(manager.getListSnapshot().phase).toBe('ready')
  106. // Sticky across later failures: the pull-activity axis reports the error,
  107. // the arrival phase holds.
  108. remote.session.list.mockResolvedValue(err(new RemoteError('gateway/internal', 'down', {})))
  109. await manager.refreshList()
  110. expect(manager.getListSnapshot()).toMatchObject({ state: 'error', phase: 'ready' })
  111. // And across an empty re-pull (empty-with-ready = truly no sessions).
  112. remote.session.list.mockResolvedValue(ok({ items: [] as never[] }))
  113. await manager.refreshList()
  114. expect(manager.getListSnapshot()).toMatchObject({ state: 'idle', phase: 'ready' })
  115. expect(manager.getListSnapshot().items).toEqual([])
  116. })
  117. it('merges create into the list immediately without waiting for a refresh', async ({ mock, remote }) => {
  118. remote.session.create.mockResolvedValue(ok({ sessionId: S2 }))
  119. const manager = makeManager(mock, remote)
  120. const result = await manager.create()
  121. expect(result).toMatchObject({ ok: true, value: { sessionId: S2 } })
  122. expect(manager.getListSnapshot().items.map(i => i.sessionId)).toEqual([S2])
  123. })
  124. it('retains title projections before list arrival, keeps last-wins by seq, and clears them on removal', async ({ mock, remote }) => {
  125. const manager = makeManager(mock, remote)
  126. const titleFrame = (title: string, seq: number) => {
  127. manager.handleControlFrame({ type: 'projection', sessionId: S1, key: 'title', value: title, seq })
  128. }
  129. titleFrame('Newest', 4)
  130. titleFrame('Stale', 3)
  131. titleFrame('Equal', 4)
  132. remote.session.list.mockResolvedValue(ok({
  133. items: [summary(S1), summary(S2, { updatedAt: 200 })] as never[],
  134. }))
  135. await manager.refreshList()
  136. const titled = manager.getListSnapshot()
  137. expect(titled.items.map(item => item.sessionId)).toEqual([S1, S2])
  138. expect(titled.items[0]?.title).toBe('Newest')
  139. expect(titled.items[1]?.title).toBeUndefined()
  140. manager.handleSessionRemoved(S1)
  141. manager.handleSessionAdded(summary(S1, { blank: true }))
  142. expect(manager.getListSnapshot().items.find(item => item.sessionId === S1)?.title).toBeUndefined()
  143. })
  144. it('seeds cold titles from the list rows\' projections block under higher-seq-wins', async ({ mock, remote }) => {
  145. const manager = makeManager(mock, remote)
  146. // A push frame landed before the list (S2's title is newer than the block's cut).
  147. manager.handleControlFrame({
  148. type: 'projection', sessionId: S2, key: 'title', value: 'Pushed', seq: 9,
  149. })
  150. remote.session.list.mockResolvedValue(ok({
  151. items: [
  152. { ...summary(S1), projections: { asOfSeq: 4, values: { title: 'Cold cached' } } },
  153. { ...summary(S2, { updatedAt: 200 }), projections: { asOfSeq: 5, values: { title: 'List stale' } } },
  154. ] as never[],
  155. }))
  156. await manager.refreshList()
  157. const items = manager.getListSnapshot().items
  158. // Cold row: title surfaces straight from the list block — no open, no history.
  159. expect(items.find(item => item.sessionId === S1)?.title).toBe('Cold cached')
  160. // The stale list block (seq 5) cannot overwrite the newer push frame (seq 9).
  161. expect(items.find(item => item.sessionId === S2)?.title).toBe('Pushed')
  162. })
  163. it('drops a projection row beyond the subscription baseline before accepting its durable replay', async ({ mock, remote }) => {
  164. remote.session.list.mockResolvedValue(ok({ items: [summary(S1)] as never[] }))
  165. const manager = makeManager(mock, remote)
  166. await manager.refreshList()
  167. const frame = (payload: SessionControlFrame) => { manager.handleControlFrame(payload) }
  168. frame({ type: 'projection', sessionId: S1, key: 'title', value: 'Unflushed', seq: 4 })
  169. // The durable baseline says the host only knows up to seq 2: the phantom
  170. // row rode lost state and must drop, or last-wins pins it forever.
  171. frame({
  172. type: 'baseline',
  173. value: {
  174. queues: {}, jobs: {},
  175. projections: { [S1]: { asOfSeq: 2, values: {} } },
  176. },
  177. })
  178. expect(manager.getListSnapshot().items[0]?.title).toBeUndefined()
  179. frame({ type: 'projection', sessionId: S1, key: 'title', value: 'Durable', seq: 2 })
  180. expect(manager.getListSnapshot().items[0]?.title).toBe('Durable')
  181. // A baseline at or past the row's seq keeps it (nothing phantom to drop).
  182. frame({
  183. type: 'baseline',
  184. value: {
  185. queues: {}, jobs: {},
  186. projections: { [S1]: { asOfSeq: 2, values: { title: 'Durable' } } },
  187. },
  188. })
  189. expect(manager.getListSnapshot().items[0]?.title).toBe('Durable')
  190. })
  191. })
  192. describe('search', () => {
  193. it('returns bounded Host results and forwards the caller signal', async ({ mock, remote }) => {
  194. remote.session.search.mockResolvedValue(ok({
  195. items: [{ sessionId: S1, snippet: 'matching excerpt' }],
  196. hasMore: true,
  197. }))
  198. const manager = makeManager(mock, remote)
  199. const signal = new AbortController().signal
  200. await expect(manager.search('exact phrase', signal)).resolves.toEqual({
  201. ok: true,
  202. value: {
  203. items: [{ sessionId: S1, snippet: 'matching excerpt' }],
  204. hasMore: true,
  205. },
  206. })
  207. expect(remote.session.search).toHaveBeenCalledWith({ query: 'exact phrase' }, signal)
  208. })
  209. it('preserves business errors and propagates a non-Remote throw', async ({ mock, remote }) => {
  210. const manager = makeManager(mock, remote)
  211. remote.session.search.mockResolvedValue(err(new RemoteError('gateway/internal', 'index unavailable', {})))
  212. const signal = new AbortController().signal
  213. await expect(manager.search('first', signal)).resolves.toMatchObject({
  214. ok: false,
  215. error: { code: 'gateway/internal', message: 'index unavailable' },
  216. })
  217. remote.session.search.mockRejectedValue(new Error('wire down'))
  218. await expect(manager.search('second', signal)).rejects.toThrow('wire down')
  219. })
  220. })
  221. describe('Host Remote event routing', () => {
  222. it('adds/removes/flips sessions and keeps removed instances resident', async ({ mock, remote }) => {
  223. const manager = makeManager(mock, remote)
  224. manager.handleSessionAdded(summary(S1, { blank: true }))
  225. manager.handleSessionAdded(summary(S1, { blank: true })) // dup: ignored
  226. expect(manager.getListSnapshot().items).toHaveLength(1)
  227. const session = manager.get(S1)
  228. manager.handleSessionStatus(S1, true)
  229. expect(session.getSnapshot().running).toBe(true)
  230. expect(manager.getListSnapshot().items[0]?.running).toBe(true)
  231. manager.handleSessionError(S1, '炸了')
  232. expect(session.getSnapshot().lastAgentError).toBe('炸了')
  233. manager.handleSessionRemoved(S1)
  234. expect(manager.getListSnapshot().items).toHaveLength(0)
  235. expect(session.getSnapshot().removed).toBe(true)
  236. expect(manager.get(S1)).toBe(session) // resident-instance rule survives removal
  237. })
  238. })
  239. describe('subagent catalogs', () => {
  240. it('keeps a catalog-discovered child address across ordinary selection and status frames', async ({ mock, remote, start }) => {
  241. remote.session.list.mockResolvedValue(ok({ items: [
  242. summary(S1),
  243. summary(S2, { parentSessionId: S1, origin: 'subagent' }),
  244. ] as never[] }))
  245. remote.subagents.list.mockResolvedValue(ok({
  246. entries: [{
  247. kind: 'child', id: S2, mode: 'continuable', label: 'worker',
  248. activity: 'running', hasChildren: false,
  249. }] as never[],
  250. parentAvailable: true,
  251. }))
  252. mock.load(sessionWorld)
  253. const client = await start()
  254. const manager = new SessionManager(client.ctx.remote)
  255. await manager.refreshList()
  256. await manager.refreshSubagents(S1)
  257. manager.selectSubagent({ parentSessionId: S1, childSessionId: S2, mode: 'continuable' })
  258. expect(manager.getListSnapshot().currentAddress).toEqual({
  259. parentSessionId: S1, childSessionId: S2, mode: 'continuable',
  260. })
  261. expect(manager.get(S2).getSnapshot().subagent).toEqual({
  262. address: { parentSessionId: S1, childSessionId: S2, mode: 'continuable' },
  263. parentAvailable: true,
  264. })
  265. // Clicking the same child through an ordinary list-selection path must not
  266. // erase the catalog-derived address and fall back to session.* transport.
  267. manager.select(S2)
  268. expect(manager.getListSnapshot().currentAddress).toEqual({
  269. parentSessionId: S1, childSessionId: S2, mode: 'continuable',
  270. })
  271. expect(manager.get(S2).getSnapshot().subagent).toEqual({
  272. address: { parentSessionId: S1, childSessionId: S2, mode: 'continuable' },
  273. parentAvailable: true,
  274. })
  275. await manager.get(S2).open()
  276. await manager.get(S2).prompt([{ type: 'text', text: 'continue' }], 'queue')
  277. expect(remote.session.follow.mock.calls.map(([request]) => request)).toEqual([
  278. {
  279. address: {
  280. kind: 'subagent', parentSessionId: S1, childSessionId: S2, mode: 'continuable',
  281. },
  282. assistantStream: true,
  283. maxMessages: 50,
  284. },
  285. ])
  286. expect(remote.session.page).not.toHaveBeenCalled()
  287. expect(remote.subagents.prompt.mock.calls.map(([request]) => request)).toEqual([
  288. {
  289. requestId: expect.any(String) as unknown as string,
  290. parentSessionId: S1, childSessionId: S2,
  291. mode: 'continuable',
  292. delivery: 'queue',
  293. content: [{ type: 'text', text: 'continue' }],
  294. clientTimeZone: new Intl.DateTimeFormat().resolvedOptions().timeZone,
  295. },
  296. ])
  297. expect(remote.session.prompt).not.toHaveBeenCalled()
  298. const listCalls = remote.subagents.list.mock.calls.length
  299. manager.handleSessionStatus(S2, false)
  300. expect(manager.getListSnapshot().subagentsByParent[S1]?.entries[0]).toMatchObject({
  301. kind: 'child', id: S2, activity: 'inactive',
  302. })
  303. expect(remote.subagents.list).toHaveBeenCalledTimes(listCalls)
  304. manager.handleSessionRemoved(S2)
  305. expect(manager.getListSnapshot().items.find(item => item.sessionId === S2)).toMatchObject({
  306. origin: 'subagent', parentSessionId: S1, running: false,
  307. })
  308. expect(manager.get(S2).getSnapshot()).toMatchObject({
  309. removed: false,
  310. subagent: {
  311. address: { parentSessionId: S1, childSessionId: S2, mode: 'continuable' },
  312. },
  313. })
  314. })
  315. it('refetches debounced membership only while the parent catalog is open', async ({ mock, remote }) => {
  316. vi.useFakeTimers()
  317. try {
  318. const manager = makeManager(mock, remote)
  319. await manager.refreshSubagents(S1)
  320. manager.setSubagentCatalogOpen(S1, true)
  321. await Promise.resolve()
  322. const baseline = remote.subagents.list.mock.calls.length
  323. manager.handleSessionAdded(summary(S2, { parentSessionId: S1 }))
  324. manager.handleSessionAdded(summary('fk-m3' as SessionId, { parentSessionId: S1 }))
  325. await vi.advanceTimersByTimeAsync(50)
  326. expect(remote.subagents.list).toHaveBeenCalledTimes(baseline + 1)
  327. manager.setSubagentCatalogOpen(S1, false)
  328. manager.handleSessionAdded(summary('fk-m4' as SessionId, { parentSessionId: S1 }))
  329. await vi.advanceTimersByTimeAsync(50)
  330. expect(remote.subagents.list).toHaveBeenCalledTimes(baseline + 1)
  331. } finally {
  332. vi.useRealTimers()
  333. }
  334. })
  335. it('marks a loaded parent row expandable only for a direct subagent publication', async ({ mock, remote }) => {
  336. const root = 'fk-root' as SessionId
  337. remote.subagents.list.mockResolvedValue(ok({
  338. entries: [
  339. {
  340. kind: 'child', id: S1, mode: 'continuable', label: 'parent',
  341. activity: 'inactive', hasChildren: false,
  342. },
  343. {
  344. kind: 'child', id: S2, mode: 'continuable', label: 'ordinary parent',
  345. activity: 'inactive', hasChildren: false,
  346. },
  347. ] as never[],
  348. parentAvailable: true,
  349. }))
  350. const manager = makeManager(mock, remote)
  351. await manager.refreshSubagents(root)
  352. manager.handleSessionAdded(summary('fk-grandchild' as SessionId, {
  353. parentSessionId: S1, origin: 'subagent',
  354. }))
  355. manager.handleSessionAdded(summary('fk-fork' as SessionId, { parentSessionId: S2 }))
  356. expect(manager.getListSnapshot().subagentsByParent[root]?.entries).toMatchObject([
  357. { kind: 'child', id: S1, hasChildren: true },
  358. { kind: 'child', id: S2, hasChildren: false },
  359. ])
  360. })
  361. it('preserves a live expandability hint across only the older in-flight catalog response', async ({ mock, remote }) => {
  362. const root = 'fk-root' as SessionId
  363. const response = Promise.withResolvers<Awaited<ReturnType<typeof remote.subagents.list>>>()
  364. remote.subagents.list.mockReturnValue(response.promise)
  365. const manager = makeManager(mock, remote)
  366. const refresh = manager.refreshSubagents(root)
  367. manager.handleSessionAdded(summary('fk-grandchild' as SessionId, {
  368. parentSessionId: S1, origin: 'subagent',
  369. }))
  370. response.resolve(ok({
  371. entries: [{
  372. kind: 'child', id: S1, mode: 'continuable', label: 'parent',
  373. activity: 'inactive', hasChildren: false,
  374. }] as never[],
  375. parentAvailable: true,
  376. }))
  377. await refresh
  378. expect(manager.getListSnapshot().subagentsByParent[root]?.entries).toMatchObject([
  379. { kind: 'child', id: S1, hasChildren: true },
  380. ])
  381. remote.subagents.list.mockResolvedValue(ok({
  382. entries: [{
  383. kind: 'child', id: S1, mode: 'continuable', label: 'parent',
  384. activity: 'inactive', hasChildren: false,
  385. }] as never[],
  386. parentAvailable: true,
  387. }))
  388. await manager.refreshSubagents(root)
  389. expect(manager.getListSnapshot().subagentsByParent[root]?.entries).toMatchObject([
  390. { kind: 'child', id: S1, hasChildren: false },
  391. ])
  392. })
  393. it('replays status frames over an older in-flight catalog response', async ({ mock, remote }) => {
  394. const root = 'fk-root' as SessionId
  395. const response = Promise.withResolvers<Awaited<ReturnType<typeof remote.subagents.list>>>()
  396. remote.subagents.list.mockReturnValue(response.promise)
  397. const manager = makeManager(mock, remote)
  398. const refresh = manager.refreshSubagents(root)
  399. manager.handleSessionStatus(S1, false)
  400. manager.handleSessionStatus(S2, true)
  401. response.resolve(ok({
  402. entries: [
  403. {
  404. kind: 'child', id: S1, mode: 'continuable', label: 'stopped',
  405. activity: 'running', hasChildren: false,
  406. },
  407. {
  408. kind: 'child', id: S2, mode: 'continuable', label: 'started',
  409. activity: 'inactive', hasChildren: false,
  410. },
  411. ] as never[],
  412. parentAvailable: true,
  413. }))
  414. await refresh
  415. expect(manager.getListSnapshot().subagentsByParent[root]?.entries).toMatchObject([
  416. { kind: 'child', id: S1, activity: 'inactive' },
  417. { kind: 'child', id: S2, activity: 'running' },
  418. ])
  419. })
  420. it('marks a detached catalog child inactive without requiring a selected address', async ({ mock, remote }) => {
  421. remote.subagents.list.mockResolvedValue(ok({
  422. entries: [{
  423. kind: 'child', id: S2, mode: 'continuable', label: 'worker',
  424. activity: 'running', hasChildren: false,
  425. }] as never[],
  426. parentAvailable: true,
  427. }))
  428. const manager = makeManager(mock, remote)
  429. await manager.refreshSubagents(S1)
  430. manager.handleSessionRemoved(S2)
  431. expect(manager.getListSnapshot().subagentsByParent[S1]?.entries).toMatchObject([
  432. { kind: 'child', id: S2, activity: 'inactive' },
  433. ])
  434. })
  435. it('coalesces overlapping catalog reads without scheduling a trailing pull', async ({ mock, remote }) => {
  436. const root = 'fk-root' as SessionId
  437. const first = Promise.withResolvers<Awaited<ReturnType<typeof remote.subagents.list>>>()
  438. remote.subagents.list.mockReturnValue(first.promise)
  439. const manager = makeManager(mock, remote)
  440. const refresh = manager.refreshSubagents(root)
  441. expect(manager.refreshSubagents(root)).toBe(refresh)
  442. remote.subagents.list.mockResolvedValue(ok({ entries: [], parentAvailable: true }))
  443. first.resolve(ok({ entries: [], parentAvailable: true }))
  444. await refresh
  445. expect(remote.subagents.list).toHaveBeenCalledOnce()
  446. })
  447. it('runs one trailing catalog refresh for a membership change coalesced into an in-flight pull', async ({ mock, remote }) => {
  448. vi.useFakeTimers()
  449. try {
  450. const root = 'fk-root' as SessionId
  451. const first = Promise.withResolvers<Awaited<ReturnType<typeof remote.subagents.list>>>()
  452. const second = Promise.withResolvers<Awaited<ReturnType<typeof remote.subagents.list>>>()
  453. remote.subagents.list.mockReturnValue(first.promise)
  454. const manager = makeManager(mock, remote, root)
  455. const refresh = manager.refreshSubagents(root)
  456. manager.setSubagentCatalogOpen(root, true)
  457. // A membership frame arrives while the pull is in flight; the debounced
  458. // refresh it schedules fires 50ms later and is coalesced into the pull —
  459. // which was requested before the new child existed. The stale mark must
  460. // queue one trailing pull carrying the change.
  461. manager.handleSessionAdded(summary(S2, { parentSessionId: root }))
  462. await vi.advanceTimersByTimeAsync(50)
  463. remote.subagents.list.mockReturnValueOnce(second.promise)
  464. first.resolve(ok({
  465. entries: [{
  466. kind: 'child', id: S1, mode: 'continuable', label: 'older',
  467. activity: 'inactive', hasChildren: false,
  468. }] as never[],
  469. parentAvailable: true,
  470. }))
  471. await refresh
  472. // The trailing pull is already in flight (kicked synchronously in finally).
  473. second.resolve(ok({
  474. entries: [
  475. {
  476. kind: 'child', id: S1, mode: 'continuable', label: 'older',
  477. activity: 'inactive', hasChildren: false,
  478. },
  479. {
  480. kind: 'child', id: S2, mode: 'continuable', label: 'new child',
  481. activity: 'inactive', hasChildren: false,
  482. },
  483. ] as never[],
  484. parentAvailable: true,
  485. }))
  486. await second.promise
  487. // The Remote face resolves one microtask after the response settles.
  488. await vi.advanceTimersByTimeAsync(0)
  489. expect(remote.subagents.list).toHaveBeenCalledTimes(2)
  490. expect(manager.getListSnapshot().subagentsByParent[root]?.entries).toMatchObject([
  491. { kind: 'child', id: S1, label: 'older' },
  492. { kind: 'child', id: S2, label: 'new child' },
  493. ])
  494. } finally {
  495. vi.useRealTimers()
  496. }
  497. })
  498. it('keeps removal invalidation across a stale success and failed trailing pull', async ({ mock, remote }) => {
  499. const root = 'fk-root' as SessionId
  500. const child = () => ({
  501. kind: 'child' as const, id: S2, mode: 'continuable' as const, label: 'worker',
  502. activity: 'inactive' as const, hasChildren: false,
  503. })
  504. const first = Promise.withResolvers<Awaited<ReturnType<typeof remote.subagents.list>>>()
  505. remote.subagents.list.mockReturnValue(first.promise)
  506. const manager = makeManager(mock, remote)
  507. const refresh = manager.refreshSubagents(root)
  508. first.resolve(ok({ entries: [child()] as never[], parentAvailable: true }))
  509. await refresh
  510. manager.selectSubagent({ parentSessionId: root, childSessionId: S2, mode: 'continuable' })
  511. // The removal lands while a second pull is in flight: the invalidation
  512. // must survive the pre-removal ok response, so one trailing pull runs.
  513. const mid = Promise.withResolvers<Awaited<ReturnType<typeof remote.subagents.list>>>()
  514. remote.subagents.list.mockReturnValueOnce(mid.promise)
  515. const midRefresh = manager.refreshSubagents(root)
  516. manager.handleSessionRemoved(root)
  517. const trailing = Promise.withResolvers<Awaited<ReturnType<typeof remote.subagents.list>>>()
  518. remote.subagents.list.mockReturnValueOnce(trailing.promise)
  519. mid.resolve(ok({ entries: [child()] as never[], parentAvailable: true }))
  520. await midRefresh
  521. expect(manager.getListSnapshot().subagentsByParent[root]?.parentAvailable).toBe(false)
  522. expect(manager.get(S2).getSnapshot().subagent).toMatchObject({ parentAvailable: false })
  523. trailing.resolve(err(new RemoteError('gateway/internal', 'trailing pull failed', {})))
  524. await vi.waitFor(() => {
  525. expect(manager.getListSnapshot().subagentsByParent[root]).toMatchObject({
  526. state: 'error',
  527. parentAvailable: false,
  528. })
  529. })
  530. const rootCalls = remote.subagents.list.mock.calls.filter(([call]) => call === root)
  531. expect(rootCalls).toHaveLength(3)
  532. expect(manager.getListSnapshot().subagentsByParent[root]?.parentAvailable).toBe(false)
  533. expect(manager.get(S2).getSnapshot().subagent).toMatchObject({ parentAvailable: false })
  534. })
  535. it('invalidates catalog availability when the owning parent is removed', async ({ mock, remote }) => {
  536. const root = 'fk-root' as SessionId
  537. remote.subagents.list.mockResolvedValue(ok({
  538. entries: [{
  539. kind: 'child', id: S2, mode: 'continuable', label: 'worker',
  540. activity: 'inactive', hasChildren: false,
  541. }] as never[],
  542. parentAvailable: true,
  543. }))
  544. const manager = makeManager(mock, remote)
  545. await manager.refreshSubagents(root)
  546. manager.selectSubagent({ parentSessionId: root, childSessionId: S2, mode: 'continuable' })
  547. expect(manager.get(S2).getSnapshot().subagent).toMatchObject({ parentAvailable: true })
  548. manager.handleSessionRemoved(root)
  549. expect(manager.getListSnapshot().subagentsByParent[root]?.parentAvailable).toBe(false)
  550. expect(manager.get(S2).getSnapshot().subagent).toMatchObject({ parentAvailable: false })
  551. })
  552. })
  553. describe('remaining branches', () => {
  554. it('refreshList propagates a non-Remote throw', async ({ mock, remote }) => {
  555. remote.session.list.mockRejectedValue(new Error('list wire down'))
  556. const manager = makeManager(mock, remote)
  557. await expect(manager.refreshList()).rejects.toThrow('list wire down')
  558. })
  559. it('refreshList pushes running bits down to already-instantiated sessions', async ({ mock, remote }) => {
  560. const manager = makeManager(mock, remote)
  561. const session = manager.get(S1)
  562. remote.session.list.mockResolvedValue(ok({ items: [summary(S1, { running: true })] as never[] }))
  563. await manager.refreshList()
  564. expect(session.getSnapshot().running).toBe(true)
  565. })
  566. it('create passes cwd and a preallocated id, folds transport throws, and deduplicates the echo', async ({ mock, remote }) => {
  567. remote.session.create.mockResolvedValue(ok({ sessionId: S1 }))
  568. const manager = makeManager(mock, remote)
  569. await manager.create({ cwd: '/tmp/w', sessionId: S1 })
  570. expect(remote.session.create).toHaveBeenCalledWith({ cwd: '/tmp/w', sessionId: S1 })
  571. expect(manager.getListSnapshot().items[0]).toMatchObject({ sessionId: S1, cwd: '/tmp/w' })
  572. await manager.create({ cwd: '/tmp/w' }) // same id returned: no duplicate row
  573. expect(manager.getListSnapshot().items).toHaveLength(1)
  574. remote.session.create.mockRejectedValue(new Error('create wire down'))
  575. await expect(manager.create()).rejects.toThrow('create wire down')
  576. // Business error passes through untouched.
  577. remote.session.create.mockResolvedValue(err(new RemoteError('gateway/internal', 'no', {})))
  578. expect(await manager.create()).toMatchObject({ ok: false })
  579. })
  580. it('publishes a real Ungrouped summary from workspace-attach-failed', async ({ mock, remote }) => {
  581. remote.session.create.mockResolvedValue(err(new RemoteError('session/workspace-attach-failed', 'published but unattached', {
  582. sessionId: S1, workspaceId: 'w1',
  583. })))
  584. const manager = makeManager(mock, remote)
  585. const result = await manager.create({ workspaceId: 'w1' as never, sessionId: S1 })
  586. expect(result).toMatchObject({ ok: false, error: { code: 'session/workspace-attach-failed' } })
  587. expect(manager.getListSnapshot().items).toEqual([expect.objectContaining({ sessionId: S1 })])
  588. expect(manager.getListSnapshot().items[0]).not.toHaveProperty('cwd')
  589. })
  590. it('reconciles a fork child published before workspace attachment fails', async ({ mock, remote }) => {
  591. remote.session.fork.mockResolvedValue(err(new RemoteError('session/workspace-attach-failed', 'forked but unattached', {
  592. sessionId: S2, workspaceId: 'w1',
  593. })))
  594. const manager = makeManager(mock, remote)
  595. const result = await manager.fork({ sessionId: S1 })
  596. expect(result).toMatchObject({ ok: false, error: { code: 'session/workspace-attach-failed' } })
  597. expect(manager.getListSnapshot().items).toEqual([expect.objectContaining({
  598. sessionId: S2,
  599. parentSessionId: S1,
  600. blank: false,
  601. })])
  602. })
  603. it('reconciles a preallocated id after an ordinary transport failure', async ({ mock, remote }) => {
  604. remote.session.create.mockRejectedValue(new Error('response lost'))
  605. const manager = makeManager(mock, remote)
  606. await expect(manager.create({ workspaceId: 'w1' as never, sessionId: S1 }))
  607. .rejects.toThrow('response lost')
  608. expect(manager.getListSnapshot().items).toEqual([])
  609. manager.handleSessionAdded(summary(S1, { blank: true, cwd: '/w/one' }))
  610. expect(manager.getListSnapshot().items).toEqual([
  611. expect.objectContaining({ sessionId: S1, cwd: '/w/one' }),
  612. ])
  613. manager.handleSessionAdded(summary(S1, { blank: true, cwd: '/w/one' }))
  614. expect(manager.getListSnapshot().items).toHaveLength(1)
  615. })
  616. it('subscribe notifies on list changes and stops after unsubscribe', async ({ mock, remote }) => {
  617. const manager = makeManager(mock, remote)
  618. let notified = 0
  619. const unsubscribe = manager.subscribe(() => { notified++ })
  620. await manager.refreshList()
  621. await new Promise(resolve => setTimeout(resolve, 0))
  622. expect(notified).toBeGreaterThan(0)
  623. const seen = notified
  624. unsubscribe()
  625. manager.handleSessionAdded(summary(S1, { blank: true }))
  626. await new Promise(resolve => setTimeout(resolve, 0))
  627. expect(notified).toBe(seen)
  628. })
  629. it('ignores Host status and error events for sessions without an instance', ({ mock, remote }) => {
  630. const manager = makeManager(mock, remote)
  631. manager.handleSessionStatus(S2, true)
  632. manager.handleSessionError(S2, '无实例')
  633. })
  634. it('keeps list-entry identity for unchanged rows across an unrelated list change', async ({ mock, remote }) => {
  635. remote.session.list.mockResolvedValue(ok({ items: [summary(S1), summary(S2, { updatedAt: 200 })] as never[] }))
  636. const manager = makeManager(mock, remote)
  637. await manager.refreshList()
  638. const before = manager.getListSnapshot()
  639. manager.handleSessionStatus(S2, true)
  640. const after = manager.getListSnapshot()
  641. expect(after.items).not.toBe(before.items)
  642. const beforeS1 = before.items.find(e => e.sessionId === S1)
  643. const afterS1 = after.items.find(e => e.sessionId === S1)
  644. expect(afterS1).toBe(beforeS1) // untouched entry keeps identity (entryCache)
  645. // Same-order same-entries snapshot reuses the items array.
  646. manager.handleSessionError(S1, 'x')
  647. expect(manager.getListSnapshot().items).toBe(after.items)
  648. })
  649. it('carries parentSessionId from the added event into the lineage row', ({ mock, remote }) => {
  650. const manager = makeManager(mock, remote)
  651. manager.handleSessionAdded(summary(S1, { blank: true }))
  652. manager.handleSessionAdded(summary(S2, {
  653. blank: true, parentSessionId: S1, origin: 'subagent',
  654. }))
  655. const items = manager.getListSnapshot().items
  656. expect(items.find(e => e.sessionId === S2)).toMatchObject({
  657. parentSessionId: S1, origin: 'subagent', depth: 1,
  658. })
  659. })
  660. })
  661. describe('connected generation', () => {
  662. it('refreshes query baselines without rebuilding independently resumed Session sources', async ({ mock, remote, start }) => {
  663. mock.stream(FOLLOW, followScript(ok({
  664. records: entries(plainTurn(SessionSeq(0), 0, 'a', 'b')) as never[],
  665. hasMore: false,
  666. modelSelection: { provider: 'deepseek-official', model: 'deepseek-chat' },
  667. })))
  668. const client = await start()
  669. const manager = new SessionManager(client.ctx.remote)
  670. const openedSession = manager.get(S1)
  671. await openedSession.open()
  672. manager.get(S2) // instantiated but never opened
  673. const historyCallsBefore = remote.session.page.mock.calls.length
  674. manager.handleConnected()
  675. await vi.waitFor(() => {
  676. expect(remote.session.list).toHaveBeenCalledOnce()
  677. })
  678. expect(remote.session.follow).toHaveBeenCalledOnce()
  679. expect(remote.session.page).toHaveBeenCalledTimes(historyCallsBefore)
  680. })
  681. it('retains the durable parent address and refreshes its catalogs across reconnect', async ({ mock, remote }) => {
  682. const address = {
  683. parentSessionId: S1, childSessionId: S2, mode: 'continuable' as const,
  684. }
  685. const parent = Promise.withResolvers<Awaited<ReturnType<typeof remote.subagents.list>>>()
  686. const child = Promise.withResolvers<Awaited<ReturnType<typeof remote.subagents.list>>>()
  687. remote.subagents.list.mockImplementation(payload => (payload === S1 ? parent.promise : child.promise))
  688. const manager = makeManager(mock, remote, S2, address)
  689. manager.handleConnected()
  690. expect(manager.get(S2).getSnapshot().subagent).toEqual({ address })
  691. parent.resolve(ok({ entries: [], parentAvailable: true }))
  692. child.resolve(ok({ entries: [], parentAvailable: true }))
  693. await vi.waitFor(() => {
  694. expect(remote.session.list).toHaveBeenCalledOnce()
  695. })
  696. await vi.waitFor(() => {
  697. expect(remote.subagents.list.mock.calls.map(([parentSessionId]) => parentSessionId)).toEqual([S1, S2])
  698. })
  699. expect(manager.get(S2).getSnapshot().subagent).toEqual({
  700. address,
  701. parentAvailable: true,
  702. })
  703. expect(manager.getListSnapshot().currentAddress).toEqual(address)
  704. })
  705. })
  706. describe('completed reminder', () => {
  707. const status = (manager: SessionManager, sessionId: SessionId, running: boolean): void => {
  708. manager.handleSessionStatus(sessionId, running)
  709. }
  710. const added = (manager: SessionManager, sessionId: SessionId): void => {
  711. manager.handleSessionAdded(summary(sessionId))
  712. }
  713. const entry = (manager: SessionManager, sessionId: SessionId) =>
  714. manager.getListSnapshot().items.find(item => item.sessionId === sessionId)
  715. it('arms on a running→idle flip of a non-selected session and clears on select', ({ mock, remote }) => {
  716. const manager = makeManager(mock, remote)
  717. added(manager, S1)
  718. added(manager, S2)
  719. manager.select(S1)
  720. expect(entry(manager, S2)?.completed).toBe(false)
  721. status(manager, S2, true)
  722. status(manager, S2, false)
  723. expect(entry(manager, S2)?.completed).toBe(true)
  724. // Opening the session consumes the reminder.
  725. manager.select(S2)
  726. expect(entry(manager, S2)?.completed).toBe(false)
  727. })
  728. it('never arms for the session being watched and re-arms after a switch-away re-run', ({ mock, remote }) => {
  729. const manager = makeManager(mock, remote)
  730. added(manager, S1)
  731. added(manager, S2)
  732. manager.select(S2)
  733. status(manager, S2, true)
  734. status(manager, S2, false)
  735. expect(entry(manager, S2)?.completed).toBe(false) // watched to completion: no reminder
  736. // Switch away; a fresh run completing again arms the reminder.
  737. manager.select(S1)
  738. status(manager, S2, true)
  739. status(manager, S2, false)
  740. expect(entry(manager, S2)?.completed).toBe(true)
  741. })
  742. it('a re-run disarms the reminder while running and re-arms on its completion', ({ mock, remote }) => {
  743. const manager = makeManager(mock, remote)
  744. added(manager, S1)
  745. added(manager, S2)
  746. manager.select(S1)
  747. status(manager, S2, true)
  748. status(manager, S2, false)
  749. expect(entry(manager, S2)?.completed).toBe(true)
  750. // The user starts a new run without opening the session: running wins.
  751. status(manager, S2, true)
  752. expect(entry(manager, S2)?.completed).toBe(false)
  753. status(manager, S2, false)
  754. expect(entry(manager, S2)?.completed).toBe(true)
  755. })
  756. it('session-removed drops the reminder and a re-add starts clean', ({ mock, remote }) => {
  757. const manager = makeManager(mock, remote)
  758. added(manager, S1)
  759. added(manager, S2)
  760. manager.select(S1)
  761. status(manager, S2, true)
  762. status(manager, S2, false)
  763. expect(entry(manager, S2)?.completed).toBe(true)
  764. manager.handleSessionRemoved(S2)
  765. expect(manager.getListSnapshot().items.find(item => item.sessionId === S2)).toBeUndefined()
  766. added(manager, S2)
  767. expect(entry(manager, S2)?.completed).toBe(false)
  768. })
  769. it('a list refresh carrying the running→idle transition arms the reminder', async ({ mock, remote }) => {
  770. remote.session.list.mockResolvedValue(ok({ items: [summary(S1), summary(S2, { updatedAt: 200, running: true })] as never[] }))
  771. const manager = makeManager(mock, remote)
  772. await manager.refreshList()
  773. manager.select(S1)
  774. expect(entry(manager, S2)?.completed).toBe(false)
  775. remote.session.list.mockResolvedValue(ok({ items: [summary(S1), summary(S2, { updatedAt: 200, running: false })] as never[] }))
  776. await manager.refreshList()
  777. expect(entry(manager, S2)?.completed).toBe(true)
  778. })
  779. it('never arms for sessions already idle at first observation', async ({ mock, remote }) => {
  780. remote.session.list.mockResolvedValue(ok({ items: [summary(S1), summary(S2, { updatedAt: 200 })] as never[] }))
  781. const manager = makeManager(mock, remote)
  782. await manager.refreshList()
  783. manager.select(S1)
  784. expect(entry(manager, S2)?.completed).toBe(false)
  785. remote.session.list.mockResolvedValue(ok({ items: [summary(S1), summary(S2, { updatedAt: 201 })] as never[] }))
  786. await manager.refreshList()
  787. expect(entry(manager, S2)?.completed).toBe(false)
  788. })
  789. it('arms a completion that happened during an in-flight first pull (baseline running, replayed idle)', async ({ mock, remote }) => {
  790. const gate = Promise.withResolvers<Awaited<ReturnType<typeof remote.session.list>>>()
  791. remote.session.list.mockReturnValue(gate.promise)
  792. const manager = makeManager(mock, remote)
  793. const refresh = manager.refreshList()
  794. // The session finishes while the first pull is still in flight; the pull
  795. // response recorded it as running at pull time.
  796. status(manager, S2, false)
  797. gate.resolve(ok({ items: [summary(S1), summary(S2, { updatedAt: 200, running: true })] as never[] }))
  798. await refresh
  799. expect(entry(manager, S2)?.completed).toBe(true)
  800. })
  801. it('arms when a session ran and completed entirely between in-flight mutations (baseline idle)', async ({ mock, remote }) => {
  802. const gate = Promise.withResolvers<Awaited<ReturnType<typeof remote.session.list>>>()
  803. remote.session.list.mockReturnValue(gate.promise)
  804. const manager = makeManager(mock, remote)
  805. const refresh = manager.refreshList()
  806. // The unknown session starts and finishes while the first pull is in
  807. // flight; the pull-time baseline recorded it idle, so the running→idle
  808. // edge lives entirely inside the replayed mutations.
  809. status(manager, S2, true)
  810. status(manager, S2, false)
  811. gate.resolve(ok({ items: [summary(S1), summary(S2, { updatedAt: 200 })] as never[] }))
  812. await refresh
  813. expect(entry(manager, S2)?.completed).toBe(true)
  814. })
  815. })
  816. describe('background-job mirror', () => {
  817. const view = (over: Partial<{ id: string; status: string; label: string }> = {}) => ({
  818. id: 'bash-1', kind: 'bash', label: 'pnpm run build', status: 'running', startedAt: 5, ...over,
  819. })
  820. const tasksFrame = (
  821. sessionId: SessionId,
  822. jobs: unknown[],
  823. ): Extract<SessionControlFrame, { type: 'jobs' }> => ({
  824. type: 'jobs', sessionId, jobs: jobs as never,
  825. })
  826. it('mirrors the whole set last-wins, keyed per session, with no Session instance needed', ({ mock, remote }) => {
  827. const manager = makeManager(mock, remote)
  828. manager.handleControlFrame(tasksFrame(S1, [view()]))
  829. manager.handleControlFrame(tasksFrame(S2, [view({ id: 'pwsh-1', label: 'other' })]))
  830. const first = manager.getListSnapshot().jobsBySession
  831. expect(first[S1]).toEqual([view()])
  832. expect(first[S2]?.[0]?.label).toBe('other')
  833. // Last-wins: the newer whole set replaces, it does not merge.
  834. manager.handleControlFrame(tasksFrame(S1, [view({ status: 'completed' })]))
  835. expect(manager.getListSnapshot().jobsBySession[S1]).toEqual([view({ status: 'completed' })])
  836. })
  837. it('stores an emptied set as an absent key so absence and [] read alike', ({ mock, remote }) => {
  838. const manager = makeManager(mock, remote)
  839. manager.handleControlFrame(tasksFrame(S1, [view()]))
  840. expect(S1 in manager.getListSnapshot().jobsBySession).toBe(true)
  841. manager.handleControlFrame(tasksFrame(S1, []))
  842. expect(S1 in manager.getListSnapshot().jobsBySession).toBe(false)
  843. })
  844. it('clears the mirror when the next control baseline has no jobs', ({ mock, remote }) => {
  845. const manager = makeManager(mock, remote)
  846. manager.handleControlFrame(tasksFrame(S1, [view()]))
  847. manager.handleControlFrame({
  848. type: 'baseline',
  849. value: { queues: {}, jobs: {}, projections: {} },
  850. })
  851. expect(S1 in manager.getListSnapshot().jobsBySession).toBe(false)
  852. })
  853. it('drops the rows when the session is removed, whichever stream lands first', ({ mock, remote }) => {
  854. const manager = makeManager(mock, remote)
  855. manager.handleSessionAdded(summary(S1, { blank: true }))
  856. manager.handleControlFrame(tasksFrame(S1, [view()]))
  857. manager.handleSessionRemoved(S1)
  858. expect(S1 in manager.getListSnapshot().jobsBySession).toBe(false)
  859. })
  860. it('notifies list subscribers so an open header re-renders without a poll', async ({ mock, remote }) => {
  861. const manager = makeManager(mock, remote)
  862. const seen = vi.fn()
  863. manager.subscribe(seen)
  864. manager.handleControlFrame(tasksFrame(S1, [view()]))
  865. // The notifier batches on a microtask; the frame itself is already applied.
  866. await Promise.resolve()
  867. expect(seen).toHaveBeenCalled()
  868. })
  869. })