sessions-service.client.spec.ts 43 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073
  1. /** Client catalog projection, explicitly retained scopes, streams, and Host operations. */
  2. import { Context } from '@deepseek-ai/cordis'
  3. import { describe, expect, vi } from 'vitest'
  4. import type { SessionId } from '@deepseek-ai/dsh-api-remotes/client'
  5. import type { SessionReference } from '../src/client/contract/sessions.ts'
  6. import { RemoteError } from '@deepseek-ai/dsh-typert-protocol'
  7. import { LlmAttemptId } from '@deepseek-ai/dsh-llm'
  8. import { RemoteStreamCarrierError } from '@deepseek-ai/dsh-api-gateway/client'
  9. import { SESSION_FORMAT_VERSION, SessionSeq } from '@deepseek-ai/dsh-session/types'
  10. import { ok, type RemoteMock } from '@deepseek-ai/dsh-remote-mock'
  11. import { createClientTest, webApp } from '@deepseek-ai/dsh-client-test-runtime/src/assembly/index.ts'
  12. import { ClientSessions, SessionCreateError, SessionForkError } from '../src/client/sessions/service.ts'
  13. import { scopeOf } from '../src/client/scope.ts'
  14. import type {
  15. SessionAssistantStreamBaseline, SessionFollowFrame, SessionFollowRequest,
  16. } from '../src/types.ts'
  17. import { FOLLOW, err, followScript, sessionWorld } from './remote/session.client.ts'
  18. const sid = (s: string): SessionId => s as SessionId
  19. /** ClientSessions uses the Gateway client for stream supervision and the native Remote mocks for responses. */
  20. const API_ROSTER = webApp.closure(['@deepseek-ai/dsh-api-gateway'])
  21. /** The first client boot pays the cold module transform of the api cone. */
  22. const COLD_BOOT_TIMEOUT_MS = 60_000
  23. interface Bench {
  24. ctx: Context
  25. mock: RemoteMock
  26. svc: ClientSessions
  27. unblock: Array<() => void>
  28. }
  29. type BenchFactory = () => Bench
  30. const it = createClientTest({ roster: API_ROSTER }).extend<{ bench: BenchFactory }>({
  31. bench: async ({ mock, start }, use) => {
  32. mock.load(sessionWorld)
  33. const client = await start()
  34. const benches: Bench[] = []
  35. try {
  36. await use(() => {
  37. const ctx = new Context()
  38. const svc = new ClientSessions(ctx, client.ctx.remote)
  39. const b = { ctx, mock, svc, unblock: [] as Array<() => void> }
  40. benches.push(b)
  41. return b
  42. })
  43. } finally {
  44. for (const b of benches) for (const finish of b.unblock) finish()
  45. await Promise.all(benches.map(b => b.ctx.fiber.dispose()))
  46. }
  47. },
  48. })
  49. /** Refresh the manager list from programmable rows and flush the microtask batch. */
  50. type FeedRow = {
  51. id: string
  52. cwd?: string
  53. parentId?: string
  54. origin?: 'subagent'
  55. running?: boolean
  56. blank?: boolean
  57. projections?: Record<string, unknown>
  58. }
  59. async function feedList(b: Bench, rows: FeedRow[]): Promise<void> {
  60. b.mock.remote.session.list.mockResolvedValue(ok({
  61. items: rows.map(r => ({
  62. sessionId: sid(r.id), updatedAt: 1, running: r.running ?? false, blank: r.blank ?? false,
  63. ...(r.cwd !== undefined ? { cwd: r.cwd } : {}),
  64. ...(r.parentId !== undefined ? { parentSessionId: sid(r.parentId) } : {}),
  65. ...(r.origin !== undefined ? { origin: r.origin } : {}),
  66. ...(r.projections === undefined
  67. ? {}
  68. : { projections: { asOfSeq: 0, values: r.projections } }),
  69. })),
  70. }) as never)
  71. await b.svc.refresh()
  72. await Promise.resolve() // manager notifier flush
  73. }
  74. describe('list store projection', () => {
  75. it('projects durable titles separately from cwd/id display fallbacks and parent links', async ({ bench }) => {
  76. const b = bench()
  77. b.svc.handleControlFrame({
  78. type: 'projection', sessionId: sid('s1'), key: 'title', value: 'Durable title', seq: 2,
  79. })
  80. await feedList(b, [
  81. { id: 's1', cwd: '/home/u/proj-a/' },
  82. { id: 's2', parentId: 's1', origin: 'subagent', running: true },
  83. ])
  84. const state = b.svc.list.getSnapshot()
  85. expect(state.ids).toEqual(['s1', 's2'])
  86. expect(state.byId[sid('s1')]).toMatchObject({ title: 'Durable title', displayTitle: 'Durable title', cwd: '/home/u/proj-a/' })
  87. expect(state.byId[sid('s2')]).toMatchObject({
  88. displayTitle: 's2', parentId: 's1', origin: 'subagent', running: true,
  89. })
  90. expect(state.byId[sid('s2')]?.title).toBeUndefined()
  91. }, COLD_BOOT_TIMEOUT_MS)
  92. it('reprojects a blank session from the generic agent-preset projection', async ({ bench }) => {
  93. const b = bench()
  94. await feedList(b, [{ id: 's1', blank: true, projections: { agentPreset: 'standard' } }])
  95. expect(b.svc.list.getSnapshot().byId[sid('s1')]?.projectionValues?.agentPreset).toBe('standard')
  96. b.svc.handleControlFrame({
  97. type: 'projection', sessionId: sid('s1'), key: 'agentPreset', value: 'minimal', seq: 1,
  98. })
  99. await Promise.resolve()
  100. expect(b.svc.list.getSnapshot().byId[sid('s1')]?.projectionValues?.agentPreset).toBe('minimal')
  101. })
  102. it('reflects live increments (host stream via manager) into the store', async ({ bench }) => {
  103. const b = bench()
  104. await feedList(b, [{ id: 's1' }])
  105. b.svc.handleSessionAdded({
  106. sessionId: sid('s2'), updatedAt: 2, running: false, blank: true,
  107. })
  108. await Promise.resolve()
  109. expect(b.svc.list.getSnapshot().ids).toContain('s2')
  110. })
  111. })
  112. describe('search', () => {
  113. it('delegates transient content search without changing the list snapshot', async ({ bench }) => {
  114. const b = bench()
  115. await feedList(b, [{ id: 's1' }])
  116. const before = b.svc.list.getSnapshot()
  117. b.mock.remote.session.search.mockResolvedValue(ok({
  118. items: [{ sessionId: sid('s1'), snippet: 'matching excerpt' }],
  119. hasMore: false,
  120. }))
  121. const signal = new AbortController().signal
  122. const call = vi.spyOn(b.mock.rpc, 'call')
  123. await expect(b.svc.search('needle', signal)).resolves.toEqual({
  124. ok: true,
  125. value: {
  126. items: [{ sessionId: 's1', snippet: 'matching excerpt' }],
  127. hasMore: false,
  128. },
  129. })
  130. expect(call.mock.calls.find(([, endpoint]) => endpoint === 'session/search')?.[3]).toBe(signal)
  131. expect(b.svc.list.getSnapshot()).toBe(before)
  132. })
  133. })
  134. describe('scope tree', () => {
  135. it('publishes transient Assistant chunks and the named durable v2 settlement through one event source', async ({ bench }) => {
  136. const b = bench()
  137. await feedList(b, [{ id: 's1' }])
  138. using _reference = b.svc.retain(sid('s1'), { source: 'controllerOperation' })
  139. await _reference.ready
  140. const binding = b.svc.binding(sid('s1'))
  141. if (binding === undefined) throw new Error('expected Session binding')
  142. await vi.waitFor(() => {
  143. expect(binding.session.getSnapshot().openState).toBe('open')
  144. })
  145. const attemptId = LlmAttemptId('web-live-attempt')
  146. const durableMessage = {
  147. type: 'event' as const,
  148. event: {
  149. type: 'assistant/message', seq: 0, time: 2,
  150. data: {
  151. turn: 1,
  152. step: 1,
  153. message: {
  154. role: 'assistant',
  155. content: [{ type: 'text', text: 'live' }],
  156. source: { kind: 'model', provider: 'p', model: 'm' },
  157. id: 'message-1',
  158. },
  159. stream: [{ type: 'text-chunks', time0: 1, index: 0, dt: [], texts: ['live'] }],
  160. },
  161. surfaceOp: 'append' as const,
  162. },
  163. }
  164. const publications: string[][] = []
  165. const dispose = binding.eventSource.subscribe(() => {
  166. publications.push(binding.eventSource.getSnapshot().entries.map(entry => entry.event.type))
  167. })
  168. b.mock.streams.push(FOLLOW, {
  169. type: 'assistant-stream',
  170. frame: {
  171. type: 'start', attemptId, revision: 1, startedAfterSeq: -1,
  172. turn: 1, step: 1,
  173. },
  174. })
  175. await b.mock.streams.drained(FOLLOW)
  176. b.mock.streams.push(FOLLOW, {
  177. type: 'assistant-stream',
  178. frame: {
  179. type: 'chunk', attemptId, revision: 2, index: 0,
  180. time: 1,
  181. chunk: { type: 'text-delta', index: 0, text: 'live' },
  182. },
  183. })
  184. await b.mock.streams.drained(FOLLOW)
  185. await vi.waitFor(() => {
  186. expect(binding.eventSource.getSnapshot().entries).toHaveLength(1)
  187. })
  188. b.mock.streams.push(FOLLOW, durableMessage)
  189. await b.mock.streams.drained(FOLLOW)
  190. await Promise.resolve()
  191. expect(binding.eventSource.getSnapshot().entries).toHaveLength(1)
  192. b.mock.streams.push(FOLLOW, {
  193. type: 'assistant-stream',
  194. frame: {
  195. type: 'end', attemptId, revision: 3, index: 1,
  196. outcome: { kind: 'committed', eventType: 'assistant/message', seq: 0 },
  197. },
  198. })
  199. await b.mock.streams.drained(FOLLOW)
  200. await vi.waitFor(() => {
  201. expect(binding.eventSource.getSnapshot().entries).toHaveLength(1)
  202. })
  203. expect(publications).toEqual([
  204. ['assistant/live-chunk'],
  205. ['assistant/message'],
  206. ])
  207. dispose()
  208. })
  209. it('replaces an active assistant baseline on reconnect without duplicate chunks', async ({ bench }) => {
  210. const b = bench()
  211. const attemptId = LlmAttemptId('reconnect-attempt')
  212. let records: never[] = []
  213. let assistantStreamBaseline: SessionAssistantStreamBaseline = {
  214. revision: 2,
  215. activeAttempt: {
  216. attemptId, startedAfterSeq: -1, turn: 1, step: 1,
  217. nextIndex: 1,
  218. stream: [{ type: 'text-chunks', time0: 1, index: 0, dt: [], texts: ['a'] }],
  219. },
  220. }
  221. b.mock.stream(FOLLOW, followScript(
  222. () => ok({ records, hasMore: false }),
  223. { assistantStream: () => assistantStreamBaseline },
  224. ))
  225. await feedList(b, [{ id: 's1' }])
  226. using _reference = b.svc.retain(sid('s1'), { source: 'controllerOperation' })
  227. await _reference.ready
  228. const binding = b.svc.binding(sid('s1'))
  229. if (binding === undefined) throw new Error('expected Session binding')
  230. await vi.waitFor(() => {
  231. expect(binding.eventSource.getSnapshot().entries).toHaveLength(1)
  232. })
  233. records = []
  234. assistantStreamBaseline = {
  235. revision: 3,
  236. activeAttempt: {
  237. attemptId, startedAfterSeq: -1, turn: 1, step: 1,
  238. nextIndex: 2,
  239. stream: [{ type: 'text-chunks', time0: 1, index: 0, dt: [1], texts: ['a', 'b'] }],
  240. },
  241. }
  242. b.mock.streams.fail(FOLLOW, new RemoteStreamCarrierError('lost'))
  243. await vi.waitFor(() => {
  244. expect(b.mock.log.requests(FOLLOW)).toHaveLength(2)
  245. expect(binding.eventSource.getSnapshot().entries).toHaveLength(2)
  246. })
  247. expect(binding.eventSource.getSnapshot().entries.map(entry => (
  248. entry.event.type === 'assistant/live-chunk' && entry.event.data.chunk.type === 'text-delta'
  249. ? entry.event.data.chunk.text
  250. : undefined
  251. ))).toEqual(['a', 'b'])
  252. })
  253. it('stages a post-opening assistant settlement behind its exact active attempt', async ({ bench }) => {
  254. const b = bench()
  255. const attemptId = LlmAttemptId('reconnect-settlement-attempt')
  256. const priorMessage = {
  257. type: 'event' as const,
  258. event: {
  259. type: 'assistant/message', seq: 0, time: 30,
  260. data: {
  261. turn: 1,
  262. step: 1,
  263. message: {
  264. role: 'assistant',
  265. content: [{ type: 'text', text: 'retry ' }],
  266. source: { kind: 'model', provider: 'p', model: 'm' },
  267. id: 'prior-attempt-message',
  268. },
  269. stream: [{ type: 'text-chunks', time0: 10, index: 0, dt: [], texts: ['retry '] }],
  270. },
  271. surfaceOp: 'append' as const,
  272. },
  273. }
  274. const currentMessage = {
  275. type: 'event' as const,
  276. event: {
  277. type: 'assistant/message', seq: 1, time: 19,
  278. data: {
  279. turn: 1,
  280. step: 1,
  281. message: {
  282. role: 'assistant',
  283. content: [{ type: 'text', text: 'settled' }],
  284. source: { kind: 'model', provider: 'p', model: 'm' },
  285. id: 'current-attempt-message',
  286. },
  287. stream: [{ type: 'text-chunks', time0: 20, index: 0, dt: [], texts: ['settled'] }],
  288. },
  289. surfaceOp: 'append' as const,
  290. },
  291. }
  292. const history = ok({
  293. records: [priorMessage] as never[],
  294. hasMore: false,
  295. })
  296. const assistantStreamBaseline: SessionAssistantStreamBaseline = {
  297. revision: 2,
  298. activeAttempt: {
  299. attemptId,
  300. startedAfterSeq: SessionSeq(0),
  301. turn: 1,
  302. step: 1,
  303. nextIndex: 1,
  304. stream: currentMessage.event.data.stream,
  305. },
  306. }
  307. b.mock.stream(FOLLOW, followScript(history, { assistantStream: assistantStreamBaseline }))
  308. await feedList(b, [{ id: 's1' }])
  309. using _reference = b.svc.retain(sid('s1'), { source: 'controllerOperation' })
  310. await _reference.ready
  311. const binding = b.svc.binding(sid('s1'))
  312. if (binding === undefined) throw new Error('expected Session binding')
  313. await vi.waitFor(() => {
  314. expect(binding.session.getSnapshot().openState).toBe('open')
  315. })
  316. expect(binding.eventSource.getSnapshot().entries.map(entry => entry.event.type))
  317. .toEqual(['assistant/message', 'assistant/live-chunk'])
  318. expect(binding.eventSource.getSnapshot().entries[0]?.event).toBe(priorMessage.event)
  319. b.mock.streams.push(FOLLOW, currentMessage)
  320. await b.mock.streams.drained(FOLLOW)
  321. await Promise.resolve()
  322. expect(binding.eventSource.getSnapshot().entries.map(entry => entry.event.type))
  323. .toEqual(['assistant/message', 'assistant/live-chunk'])
  324. b.mock.streams.push(FOLLOW, {
  325. type: 'assistant-stream',
  326. frame: {
  327. type: 'end', attemptId, revision: 3, index: 1,
  328. outcome: { kind: 'committed', eventType: 'assistant/message', seq: 1 },
  329. },
  330. })
  331. await b.mock.streams.drained(FOLLOW)
  332. await vi.waitFor(() => {
  333. expect(binding.eventSource.getSnapshot().entries.map(entry => entry.event.type))
  334. .toEqual(['assistant/message', 'assistant/message'])
  335. })
  336. expect(binding.eventSource.getSnapshot().change).toEqual({
  337. kind: 'settle-assistant', attemptId: String(attemptId), entry: currentMessage,
  338. })
  339. })
  340. it('replaces an invalid settlement with the authoritative post-end baseline', async ({ bench }) => {
  341. const b = bench()
  342. const attemptId = LlmAttemptId('reconnect-end-index-attempt')
  343. const prior = {
  344. type: 'event' as const,
  345. event: { type: 'turn/start', seq: 0, time: 19, data: { turn: 1 } },
  346. }
  347. const message = {
  348. type: 'event' as const,
  349. event: {
  350. type: 'assistant/message', seq: 1, time: 21,
  351. data: {
  352. turn: 1,
  353. step: 1,
  354. message: {
  355. role: 'assistant',
  356. content: [{ type: 'text', text: 'settled' }],
  357. source: { kind: 'model', provider: 'p', model: 'm' },
  358. id: 'current-attempt-message',
  359. },
  360. stream: [{ type: 'text-chunks', time0: 20, index: 0, dt: [], texts: ['settled'] }],
  361. },
  362. surfaceOp: 'append' as const,
  363. },
  364. }
  365. let records = [prior] as never[]
  366. let assistantStreamBaseline: SessionAssistantStreamBaseline = {
  367. revision: 2,
  368. activeAttempt: {
  369. attemptId,
  370. startedAfterSeq: SessionSeq(0),
  371. turn: 1,
  372. step: 1,
  373. nextIndex: 1,
  374. stream: message.event.data.stream,
  375. },
  376. }
  377. b.mock.stream(FOLLOW, followScript(
  378. () => ok({ records, hasMore: false }),
  379. { assistantStream: () => assistantStreamBaseline },
  380. ))
  381. await feedList(b, [{ id: 's1' }])
  382. using _reference = b.svc.retain(sid('s1'), { source: 'controllerOperation' })
  383. await _reference.ready
  384. const binding = b.svc.binding(sid('s1'))
  385. if (binding === undefined) throw new Error('expected Session binding')
  386. await vi.waitFor(() => {
  387. expect(binding.eventSource.getSnapshot().entries.map(entry => entry.event.type))
  388. .toEqual(['turn/start', 'assistant/live-chunk'])
  389. })
  390. const openingRevision = binding.eventSource.getSnapshot().revision
  391. b.mock.streams.push(FOLLOW, message)
  392. await b.mock.streams.drained(FOLLOW)
  393. await Promise.resolve()
  394. expect(binding.eventSource.getSnapshot().entries.map(entry => entry.event.type))
  395. .toEqual(['turn/start', 'assistant/live-chunk'])
  396. records = [prior, message] as never[]
  397. assistantStreamBaseline = { revision: 3 }
  398. b.mock.streams.push(FOLLOW, {
  399. type: 'assistant-stream',
  400. frame: {
  401. type: 'end', attemptId, revision: 3, index: 0,
  402. outcome: { kind: 'committed', eventType: 'assistant/message', seq: 0 },
  403. },
  404. })
  405. await b.mock.streams.drained(FOLLOW)
  406. await vi.waitFor(() => {
  407. expect(b.mock.log.requests(FOLLOW)).toHaveLength(2)
  408. expect(b.mock.log.streams(FOLLOW).filter(stream => stream.state === 'open')).toHaveLength(1)
  409. expect(binding.eventSource.getSnapshot().revision).toBeGreaterThan(openingRevision)
  410. expect(binding.eventSource.getSnapshot().entries.map(entry => entry.event.type))
  411. .toEqual(['turn/start', 'assistant/message'])
  412. })
  413. })
  414. it('holds a Host-addressed Context across an empty catalog baseline', async ({ bench }) => {
  415. const b = bench()
  416. using reference = b.svc.retainAgentScope(sid('s-early'))
  417. const scoped = reference.binding.ctx
  418. expect(scopeOf(scoped)).toBe('s-early')
  419. b.svc.handleControlFrame({ type: 'baseline', value: { jobs: {}, projections: {} } })
  420. await feedList(b, [])
  421. expect(b.svc.scope(sid('s-early'))).toBe(scoped)
  422. reference.release()
  423. expect(b.svc.scope(sid('s-early'))).toBeUndefined()
  424. })
  425. it('borrows only retained bindings and preserves them while the catalog changes', async ({ bench }) => {
  426. const b = bench()
  427. await feedList(b, [{ id: 's1' }])
  428. expect(b.svc.scope(sid('s1'))).toBeUndefined()
  429. using reference = b.svc.retain(sid('s1'), { source: 'controllerOperation' })
  430. await reference.ready
  431. const binding = reference.binding
  432. expect(b.svc.scope(sid('s1'))).toBe(binding.ctx)
  433. expect(b.svc.sessionOf(binding.ctx)).toBe(binding.session)
  434. await feedList(b, [])
  435. expect(b.svc.binding(sid('s1'))).toBe(binding)
  436. await feedList(b, [{ id: 's1', running: true }])
  437. expect(b.svc.binding(sid('s1'))).toBe(binding)
  438. reference.release()
  439. expect(b.svc.binding(sid('s1'))).toBeUndefined()
  440. })
  441. it('closes an opened journal when its removed scope drops', async ({ bench }) => {
  442. const b = bench()
  443. const follows = () => b.mock.log.streams(FOLLOW).filter(({ args }) => {
  444. const request = args[0] as SessionFollowRequest
  445. return request.address.kind === 'session' && request.address.sessionId === sid('s1')
  446. })
  447. await feedList(b, [{ id: 's1' }])
  448. using reference = b.svc.retain(sid('s1'), { source: 'controllerOperation' })
  449. await reference.ready
  450. const session = b.svc.binding(sid('s1'))?.session
  451. if (session === undefined) throw new Error('expected the selected Session binding')
  452. await vi.waitFor(() => { expect(follows().filter(stream => stream.state === 'open')).toHaveLength(1) })
  453. const notified = vi.fn()
  454. session.subscribe(notified)
  455. await feedList(b, [])
  456. expect(b.svc.binding(sid('s1'))?.session).toBe(session)
  457. reference.release()
  458. await vi.waitFor(() => { expect(follows().filter(stream => stream.state === 'open')).toHaveLength(0) })
  459. notified.mockClear()
  460. expect(b.mock.streams.push(FOLLOW, {
  461. type: 'event',
  462. event: { seq: 0, timestamp: 0, type: 'turn/start', data: { turn: 0 } } as never,
  463. }, ([request]) => {
  464. const address = (request as SessionFollowRequest).address
  465. return address.kind === 'session' && address.sessionId === sid('s1')
  466. })).toBe(0)
  467. await b.mock.streams.drained(FOLLOW)
  468. await Promise.resolve()
  469. expect(follows()).toHaveLength(1)
  470. expect(notified).not.toHaveBeenCalled()
  471. })
  472. })
  473. describe('Agent scope disposal lifecycle', () => {
  474. it('root disposal runs Agent scope effects', async ({ bench }) => {
  475. const b = bench()
  476. const readiness = b.ctx.plugin(() => undefined)
  477. await readiness
  478. b.svc.handleSessionAdded({
  479. sessionId: sid('live'), updatedAt: 1, running: false, blank: true,
  480. })
  481. await Promise.resolve()
  482. using reference = b.svc.retainAgentScope(sid('live'))
  483. const scoped = reference.binding.ctx
  484. if (scoped === undefined) throw new Error('fixture Agent Context was not minted')
  485. await scoped.fiber.await()
  486. const scopeDisposed = vi.fn()
  487. scoped.effect(() => scopeDisposed, 'fixture Agent scope effect')
  488. await b.ctx.fiber.dispose()
  489. expect(scopeDisposed).toHaveBeenCalledOnce()
  490. expect(b.svc.sessionOf(scoped)).toBeUndefined()
  491. })
  492. it('root disposal waits for an opened Session source to finish closing', async ({ bench }) => {
  493. const closeGate = Promise.withResolvers<undefined>()
  494. const abortObserved = vi.fn()
  495. let followSignal: AbortSignal | undefined
  496. const b = bench()
  497. b.unblock.push(() => { closeGate.resolve(undefined) })
  498. b.mock.remote.session.follow.mockImplementation((request, signal) => {
  499. if (signal === undefined) throw new Error('fixture requires a signal')
  500. followSignal = signal
  501. let opened = false
  502. return {
  503. [Symbol.asyncIterator]: () => ({
  504. next: () => {
  505. if (!opened) {
  506. opened = true
  507. return Promise.resolve({
  508. done: false,
  509. value: {
  510. type: 'snapshot',
  511. header: {
  512. version: SESSION_FORMAT_VERSION,
  513. id: request.address.kind === 'session'
  514. ? request.address.sessionId
  515. : request.address.childSessionId,
  516. createdAt: 0,
  517. isSeeded: false,
  518. },
  519. cursor: -1,
  520. records: [],
  521. hasMore: false,
  522. projections: { asOfSeq: -1, values: {} },
  523. assistantStream: { revision: 0 },
  524. } as const,
  525. })
  526. }
  527. return new Promise((_resolve, reject) => {
  528. signal.addEventListener('abort', () => {
  529. abortObserved()
  530. void closeGate.promise.then(() => {
  531. reject(signal.reason instanceof Error
  532. ? signal.reason
  533. : new Error(String(signal.reason)))
  534. })
  535. }, { once: true })
  536. })
  537. },
  538. }),
  539. }
  540. })
  541. const readiness = b.ctx.plugin(() => undefined)
  542. await readiness
  543. await feedList(b, [{ id: 's1' }])
  544. using _reference = b.svc.retain(sid('s1'), { source: 'controllerOperation' })
  545. await _reference.ready
  546. await vi.waitFor(() => {
  547. expect(b.svc.binding(sid('s1'))?.session.getSnapshot().openState).toBe('open')
  548. })
  549. const disposal = b.ctx.fiber.dispose()
  550. const settled = vi.fn()
  551. const observed = disposal.then(settled)
  552. await vi.waitFor(() => { expect(abortObserved).toHaveBeenCalledOnce() })
  553. expect(followSignal?.aborted).toBe(true)
  554. expect(settled).not.toHaveBeenCalled()
  555. closeGate.resolve(undefined)
  556. await observed
  557. expect(settled).toHaveBeenCalledOnce()
  558. })
  559. it('root disposal joins every Session drop already started by final release under load', async ({ bench }) => {
  560. const closeGates = new Map<SessionId, PromiseWithResolvers<undefined>>()
  561. const aborted = new Set<SessionId>()
  562. const b = bench()
  563. b.unblock.push(() => { for (const gate of closeGates.values()) gate.resolve(undefined) })
  564. b.mock.remote.session.follow.mockImplementation((request, signal) => {
  565. if (signal === undefined) throw new Error('fixture requires a signal')
  566. const sessionId = request.address.kind === 'session'
  567. ? request.address.sessionId
  568. : request.address.childSessionId
  569. const closeGate = Promise.withResolvers<undefined>()
  570. closeGates.set(sessionId, closeGate)
  571. let opened = false
  572. return {
  573. [Symbol.asyncIterator]: () => ({
  574. next: () => {
  575. if (!opened) {
  576. opened = true
  577. return Promise.resolve({
  578. done: false,
  579. value: {
  580. type: 'snapshot',
  581. header: { version: SESSION_FORMAT_VERSION, id: sessionId, createdAt: 0, isSeeded: false },
  582. cursor: -1,
  583. records: [],
  584. hasMore: false,
  585. projections: { asOfSeq: -1, values: {} },
  586. assistantStream: { revision: 0 },
  587. } as const,
  588. })
  589. }
  590. return new Promise<IteratorResult<SessionFollowFrame>>((_resolve, reject) => {
  591. signal.addEventListener('abort', () => {
  592. aborted.add(sessionId)
  593. void closeGate.promise.then(() => {
  594. reject(signal.reason instanceof Error
  595. ? signal.reason
  596. : new Error(String(signal.reason)))
  597. })
  598. }, { once: true })
  599. })
  600. },
  601. }),
  602. }
  603. })
  604. const readiness = b.ctx.plugin(() => undefined)
  605. await readiness
  606. const sessionIds = Array.from({ length: 24 }, (_, index) => sid(`load-${String(index)}`))
  607. const retained = sessionIds.at(-1)
  608. const held = sessionIds[0]
  609. if (retained === undefined || held === undefined) throw new Error('fixture requires sessions')
  610. await feedList(b, sessionIds.map(id => ({ id })))
  611. const references = new Map<SessionId, SessionReference>()
  612. for (const id of sessionIds) {
  613. const reference = b.svc.retain(id, { source: 'controllerOperation' })
  614. await reference.ready
  615. references.set(id, reference)
  616. }
  617. await vi.waitFor(() => {
  618. for (const id of sessionIds) {
  619. expect(b.svc.binding(id)?.session.getSnapshot().openState).toBe('open')
  620. }
  621. })
  622. const pruned = sessionIds.slice(0, -1)
  623. await feedList(b, [{ id: retained }])
  624. for (const id of pruned) references.get(id)?.release()
  625. await vi.waitFor(() => { expect(aborted.size).toBe(pruned.length) })
  626. for (const id of pruned) expect(b.svc.scope(id)).toBeUndefined()
  627. const disposal = b.ctx.fiber.dispose()
  628. const settled = vi.fn()
  629. const observed = disposal.then(settled)
  630. await vi.waitFor(() => { expect(aborted.size).toBe(sessionIds.length) })
  631. const otherClosures: Promise<void>[] = []
  632. for (const [id, gate] of closeGates) {
  633. if (id === held) continue
  634. gate.resolve(undefined)
  635. otherClosures.push(gate.promise)
  636. }
  637. await Promise.all(otherClosures)
  638. expect(settled).not.toHaveBeenCalled()
  639. closeGates.get(held)?.resolve(undefined)
  640. await observed
  641. expect(settled).toHaveBeenCalledOnce()
  642. })
  643. })
  644. describe('borrow-only bindings', () => {
  645. it('keeps catalog discovery separate from history opening and ownership', async ({ bench }) => {
  646. const b = bench()
  647. await feedList(b, [{ id: 's1' }, { id: 's2' }])
  648. expect(b.svc.binding(sid('s1'))).toBeUndefined()
  649. expect(b.svc.scope(sid('s2'))).toBeUndefined()
  650. expect(b.mock.log.requests(FOLLOW)).toHaveLength(0)
  651. using reference = b.svc.retain(sid('s1'), { source: 'controllerOperation' })
  652. await reference.ready
  653. using second = b.svc.retain(sid('s1'), { source: 'controllerOperation' })
  654. await second.ready
  655. expect(second.binding).toBe(reference.binding)
  656. expect(b.mock.log.requests(FOLLOW)).toHaveLength(1)
  657. })
  658. })
  659. describe('catalog-addressed navigation', () => {
  660. it('opens an explicit catalog without retaining it and keeps one-shot labels optional', async ({ bench }) => {
  661. const b = bench()
  662. b.svc.handleControlFrame({
  663. type: 'projection', sessionId: sid('one-shot'), key: 'title', value: 'Investigate startup', seq: 2,
  664. })
  665. b.mock.remote.subagents.list.mockResolvedValue(ok({
  666. entries: [
  667. { kind: 'child', id: sid('one-shot'), mode: 'one-shot', activity: 'inactive', hasChildren: false },
  668. { kind: 'diagnostic', id: sid('missing'), reason: 'unavailable' },
  669. ],
  670. parentAvailable: true,
  671. }))
  672. b.svc.setSubagentCatalogOpen(sid('root'), true)
  673. await b.svc.refreshSubagents(sid('root'))
  674. b.svc.setSubagentCatalogOpen(sid('root'), false)
  675. expect(b.svc.list.getSnapshot().byId[sid('one-shot')]).toMatchObject({
  676. title: 'Investigate startup',
  677. displayTitle: 'Investigate startup',
  678. projectionValues: { title: 'Investigate startup' },
  679. })
  680. expect(b.svc.list.getSnapshot().byId[sid('missing')]).toBeUndefined()
  681. expect(b.svc.binding(sid('one-shot'))).toBeUndefined()
  682. expect(b.mock.remote.subagents.list).toHaveBeenCalledOnce()
  683. })
  684. it('keeps projected titles in standard list rows for an addressed route', async ({ bench }) => {
  685. const b = bench()
  686. b.mock.remote.subagents.list.mockImplementation((payload) => {
  687. const parentSessionId = payload
  688. if (parentSessionId === sid('root')) {
  689. return Promise.resolve(ok({
  690. entries: [{
  691. kind: 'child', id: sid('child'), mode: 'continuable', label: 'Child',
  692. activity: 'inactive', hasChildren: true,
  693. }] as never[],
  694. parentAvailable: true,
  695. }))
  696. }
  697. if (parentSessionId === sid('child')) {
  698. return Promise.resolve(ok({
  699. entries: [{
  700. kind: 'child', id: sid('grandchild'), mode: 'continuable', label: 'Grandchild',
  701. activity: 'inactive', hasChildren: false,
  702. }] as never[],
  703. parentAvailable: false,
  704. }))
  705. }
  706. return Promise.resolve(ok({ entries: [], parentAvailable: false }))
  707. })
  708. await feedList(b, [
  709. { id: 'root' },
  710. {
  711. id: 'child', cwd: '/summary-child', parentId: 'root', origin: 'subagent',
  712. projections: { title: 'Child session title' },
  713. },
  714. { id: 'grandchild', cwd: '/summary-grandchild', parentId: 'child', origin: 'subagent' },
  715. ])
  716. await b.svc.refreshSubagents(sid('root'))
  717. await b.svc.refreshSubagents(sid('child'))
  718. using _reference = b.svc.retain({
  719. parentSessionId: sid('child'), childSessionId: sid('grandchild'), mode: 'continuable',
  720. }, { source: 'controllerOperation' })
  721. await _reference.ready
  722. expect(b.svc.list.getSnapshot().byId[sid('child')]).toMatchObject({
  723. title: 'Child session title',
  724. displayTitle: 'Child session title',
  725. projectionValues: { title: 'Child session title' },
  726. })
  727. expect(b.svc.list.getSnapshot().byId[sid('grandchild')]?.displayTitle).toBe('Grandchild')
  728. })
  729. it('projects a retained descendant and discovers ancestor addresses without retaining ancestor scopes', async ({ bench }) => {
  730. const b = bench()
  731. b.mock.remote.subagents.list.mockImplementation((payload) => {
  732. const parentSessionId = payload
  733. if (parentSessionId === sid('root')) {
  734. return Promise.resolve(ok({
  735. entries: [{
  736. kind: 'child', id: sid('child'), mode: 'continuable', label: 'Child',
  737. activity: 'inactive', hasChildren: true,
  738. }] as never[],
  739. parentAvailable: true,
  740. }))
  741. }
  742. if (parentSessionId === sid('child')) {
  743. return Promise.resolve(ok({
  744. entries: [{
  745. kind: 'child', id: sid('grandchild'), mode: 'continuable', label: 'Grandchild',
  746. activity: 'inactive', hasChildren: false,
  747. }] as never[],
  748. parentAvailable: false,
  749. }))
  750. }
  751. return Promise.resolve(ok({ entries: [], parentAvailable: false }))
  752. })
  753. await feedList(b, [{ id: 'root' }])
  754. await b.svc.refreshSubagents(sid('root'))
  755. await b.svc.refreshSubagents(sid('child'))
  756. using reference = b.svc.retain({
  757. parentSessionId: sid('child'), childSessionId: sid('grandchild'), mode: 'continuable',
  758. }, { source: 'controllerOperation' })
  759. await reference.ready
  760. const list = b.svc.list.getSnapshot()
  761. expect(list.ids).toEqual([sid('root')])
  762. expect(list.byId[sid('child')]).toMatchObject({ parentId: sid('root'), origin: 'subagent' })
  763. expect(list.byId[sid('grandchild')]).toMatchObject({ parentId: sid('child'), origin: 'subagent' })
  764. expect(b.svc.binding(sid('child'))).toBeUndefined()
  765. expect(b.svc.subagentAddress(sid('child'))).toEqual({
  766. parentSessionId: sid('root'), childSessionId: sid('child'), mode: 'continuable',
  767. })
  768. using child = b.svc.retain(sid('child'), { source: 'controllerOperation' })
  769. await child.ready
  770. expect(b.svc.binding(sid('child'))).toBe(child.binding)
  771. expect(b.svc.binding(sid('grandchild'))).toBe(reference.binding)
  772. })
  773. })
  774. describe('create', () => {
  775. it('passes a preallocated id and preserves it on ordinary failure', async ({ bench }) => {
  776. const b = bench()
  777. b.mock.remote.session.create.mockResolvedValue(ok({ sessionId: sid('fresh') }))
  778. await expect(b.svc.create({ cwd: '/w', sessionId: sid('fresh') })).resolves.toBe('fresh')
  779. expect(b.mock.remote.session.create).toHaveBeenCalledExactlyOnceWith({ cwd: '/w', sessionId: 'fresh' })
  780. b.mock.remote.session.create.mockResolvedValue(err(new RemoteError('gateway/internal', '爆了', {})))
  781. const failure = await b.svc.create({ sessionId: sid('candidate') }).catch((error: unknown) => error)
  782. expect(failure).toBeInstanceOf(SessionCreateError)
  783. expect(failure).toMatchObject({
  784. requestedSessionId: 'candidate',
  785. rpcError: { code: 'gateway/internal', message: '爆了' },
  786. })
  787. })
  788. it('publishes a created identity without implicitly retaining its binding', async ({ bench }) => {
  789. const b = bench()
  790. b.mock.remote.session.create.mockResolvedValue(ok({ sessionId: sid('born') }))
  791. const born = await b.svc.create({ workspaceId: 'ws' as never })
  792. // Synchronously after resolution — the draft hand-off contract: the
  793. // create echo IS the entity entering the client's view (blank row +
  794. // catalog publication), no notifier flush in between.
  795. expect(b.svc.list.getSnapshot().byId[born]).toMatchObject({ id: 'born', blank: true })
  796. expect(b.svc.binding(born)).toBeUndefined()
  797. expect(b.svc.scope(born)).toBeUndefined()
  798. using reference = b.svc.retain(born, { source: 'controllerOperation' })
  799. await reference.ready
  800. expect(b.svc.binding(born)).toBe(reference.binding)
  801. })
  802. it('lists the published id after Workspace attachment fails (publication precedes attachment)', async ({ bench }) => {
  803. const b = bench()
  804. b.mock.remote.session.create.mockResolvedValue(err(new RemoteError(
  805. 'session/workspace-attach-failed',
  806. 'ledger unavailable',
  807. { sessionId: sid('published'), workspaceId: 'ws' },
  808. )))
  809. const failure = await b.svc.create({
  810. workspaceId: 'ws' as never,
  811. sessionId: sid('published'),
  812. }).catch((error: unknown) => error)
  813. await Promise.resolve()
  814. expect(failure).toBeInstanceOf(SessionCreateError)
  815. expect(failure).toMatchObject({
  816. requestedSessionId: 'published',
  817. rpcError: { code: 'session/workspace-attach-failed' },
  818. })
  819. expect(b.svc.list.getSnapshot().byId[sid('published')]).toMatchObject({ id: 'published', blank: true })
  820. })
  821. })
  822. describe('fork', () => {
  823. it('propagates a failed fork without creating or retaining a child', async ({ bench }) => {
  824. const b = bench()
  825. const error = new RemoteError('session/not-found', 'source missing', { sessionId: sid('source') })
  826. b.mock.remote.session.fork.mockResolvedValue(err(error))
  827. const failure = await b.svc.fork({ sessionId: sid('source') }).catch((cause: unknown) => cause)
  828. expect(failure).toBeInstanceOf(SessionForkError)
  829. expect(failure).toMatchObject({ sourceSessionId: 'source', rpcError: error })
  830. expect(b.svc.list.getSnapshot().ids).toEqual([])
  831. expect(b.mock.log.requests(FOLLOW)).toHaveLength(0)
  832. })
  833. it.for([
  834. ['Roadmap', 'Roadmap (1)'],
  835. ['Roadmap (1)', 'Roadmap (2)'],
  836. ['计划(1)', '计划(2)'],
  837. ['计划 (9)', '计划 (10)'],
  838. ] as const)('increments the durable title %j after the child is published', async ([sourceTitle, childTitle], { bench }) => {
  839. const b = bench()
  840. b.svc.handleControlFrame({
  841. type: 'projection', sessionId: sid('source'), key: 'title', value: sourceTitle, seq: 2,
  842. })
  843. await feedList(b, [{ id: 'source', cwd: '/work' }])
  844. b.mock.remote.session.fork.mockResolvedValue(ok({ sessionId: sid('child') }))
  845. b.mock.remote.session.rename.mockImplementation((payload) => {
  846. const { title } = payload as { title: string }
  847. return Promise.resolve(ok({ title, seq: 3 }))
  848. })
  849. await expect(b.svc.fork({
  850. sessionId: sid('source'), atSeq: 7, increaseTitle: true,
  851. })).resolves.toBe('child')
  852. expect(b.mock.remote.session.fork).toHaveBeenCalledExactlyOnceWith({ sessionId: 'source', atSeq: 7 })
  853. expect(b.mock.remote.session.rename).toHaveBeenCalledExactlyOnceWith({ sessionId: 'child', title: childTitle })
  854. await Promise.resolve()
  855. expect(b.svc.list.getSnapshot().byId[sid('child')]).toMatchObject({
  856. title: childTitle,
  857. displayTitle: childTitle,
  858. parentId: 'source',
  859. })
  860. })
  861. it('floors a fractional anchor to the real event seq the wire accepts', async ({ bench }) => {
  862. const b = bench()
  863. await feedList(b, [{ id: 'source', cwd: '/work' }])
  864. b.mock.remote.session.fork.mockResolvedValue(ok({ sessionId: sid('child') }))
  865. // The frozen node of an interrupted turn carries turnEnd.seq - 0.9.
  866. await expect(b.svc.fork({ sessionId: sid('source'), atSeq: 41.1 })).resolves.toBe('child')
  867. expect(b.mock.remote.session.fork).toHaveBeenCalledExactlyOnceWith({ sessionId: 'source', atSeq: 41 })
  868. })
  869. it('does not rename without the title policy or a durable source title', async ({ bench }) => {
  870. const b = bench()
  871. await feedList(b, [{ id: 'source', cwd: '/work' }])
  872. b.mock.remote.session.fork.mockResolvedValue(ok({ sessionId: sid('child') }))
  873. await expect(b.svc.fork({ sessionId: sid('source'), increaseTitle: true })).resolves.toBe('child')
  874. expect(b.mock.remote.session.rename).not.toHaveBeenCalled()
  875. b.mock.remote.session.fork.mockResolvedValue(ok({ sessionId: sid('child-2') }))
  876. await expect(b.svc.fork({ sessionId: sid('source') })).resolves.toBe('child-2')
  877. expect(b.mock.remote.session.rename).not.toHaveBeenCalled()
  878. })
  879. it('rejects child rename failure while preserving its catalog row and releasing the operation', async ({ bench }) => {
  880. const b = bench()
  881. b.svc.handleControlFrame({
  882. type: 'projection', sessionId: sid('source'), key: 'title', value: 'Roadmap', seq: 2,
  883. })
  884. await feedList(b, [{ id: 'source' }])
  885. b.mock.remote.session.fork.mockResolvedValue(ok({ sessionId: sid('child') }))
  886. b.mock.remote.session.rename.mockResolvedValue(err(new RemoteError('session/title-invalid', 'rejected', { sessionId: sid('child') })))
  887. await expect(b.svc.fork({ sessionId: sid('source'), increaseTitle: true }))
  888. .rejects.toThrow('fork child rename failed: session/title-invalid: rejected')
  889. expect(b.svc.list.getSnapshot().byId[sid('child')]).toBeDefined()
  890. expect(b.svc.binding(sid('child'))).toBeUndefined()
  891. expect(b.svc.retainInfo(sid('child')).getSnapshot().referenceCount).toBe(0)
  892. })
  893. })
  894. describe('catalog arrival', () => {
  895. it('keeps a retained generation indexed after its Host row is removed', async ({ bench }) => {
  896. const b = bench()
  897. b.svc.handleSessionAdded({ sessionId: sid('s-new'), updatedAt: 1, running: false, blank: true })
  898. await Promise.resolve()
  899. expect(b.svc.binding(sid('s-new'))).toBeUndefined()
  900. const reference = b.svc.retain(sid('s-new'), { source: 'controllerOperation' })
  901. await reference.ready
  902. const binding = reference.binding
  903. b.svc.handleSessionRemoved(sid('s-new'))
  904. await Promise.resolve()
  905. expect(b.svc.list.getSnapshot().ids).not.toContain(sid('s-new'))
  906. expect(b.svc.list.getSnapshot().byId[sid('s-new')]?.retainedBy).toEqual({ controllerOperation: 1 })
  907. expect(b.svc.binding(sid('s-new'))).toBe(binding)
  908. reference.release()
  909. expect(b.svc.list.getSnapshot().byId[sid('s-new')]).toBeUndefined()
  910. })
  911. })
  912. describe('blank mirror', () => {
  913. it('flips blank=false from the running:true status frame (cross-client conversion)', async ({ bench }) => {
  914. const b = bench()
  915. await feedList(b, [{ id: 's1', blank: true }])
  916. expect(b.svc.list.getSnapshot().byId[sid('s1')]).toMatchObject({ blank: true })
  917. using _reference = b.svc.retain(sid('s1'), { source: 'controllerOperation' })
  918. await _reference.ready
  919. b.svc.handleSessionStatus(sid('s1'), true)
  920. await Promise.resolve()
  921. expect(b.svc.list.getSnapshot().byId[sid('s1')]).toMatchObject({ blank: false, running: true })
  922. // The instantiated Session mirrors the same flip.
  923. expect(b.svc.binding(sid('s1'))?.session.getSnapshot().blank).toBe(false)
  924. })
  925. it('flips blank=false on prompt ACCEPTANCE, not on the attempt', async ({ bench }) => {
  926. const b = bench()
  927. await feedList(b, [{ id: 's1', blank: true, cwd: '/w/a' }])
  928. using reference = b.svc.retain(sid('s1'), { source: 'controllerOperation' })
  929. await reference.ready
  930. const session = reference.binding.session
  931. expect(session.getSnapshot().blank).toBe(true)
  932. const gate = Promise.withResolvers<Awaited<ReturnType<typeof b.mock.remote.session.prompt>>>()
  933. b.mock.remote.session.prompt.mockReturnValue(gate.promise)
  934. const send = session.prompt([{ type: 'text', text: 'hi' }], 'queue')
  935. // In flight: still blank (the flip point is the success response, which
  936. // proves the user message reached the host log).
  937. expect(session.getSnapshot().blank).toBe(true)
  938. gate.resolve(ok({ accepted: true as const }))
  939. await send
  940. expect(session.getSnapshot().blank).toBe(false)
  941. await Promise.resolve()
  942. expect(b.svc.list.getSnapshot().byId[sid('s1')]).toMatchObject({ blank: false })
  943. })
  944. it('keeps a rejected first prompt blank: hidden and still reusable', async ({ bench }) => {
  945. const b = bench()
  946. await feedList(b, [{ id: 's1', blank: true, cwd: '/w/a' }])
  947. using reference = b.svc.retain(sid('s1'), { source: 'controllerOperation' })
  948. await reference.ready
  949. const session = reference.binding.session
  950. b.mock.remote.session.prompt.mockResolvedValue(err(new RemoteError('gateway/internal', 'agent busy', {})))
  951. const result = await session.prompt([{ type: 'text', text: 'hi' }], 'queue')
  952. expect(result.ok).toBe(false)
  953. // No flip on failure: local stays aligned with the host authority
  954. // (events.length still 0), so the session stays hidden and reusable.
  955. expect(session.getSnapshot().blank).toBe(true)
  956. await Promise.resolve()
  957. expect(b.svc.list.getSnapshot().byId[sid('s1')]).toMatchObject({ blank: true })
  958. })
  959. it('takes session-added blank=true as the hidden birth and list blank as reconnect authority', async ({ bench }) => {
  960. const b = bench()
  961. await feedList(b, [])
  962. b.svc.handleSessionAdded({
  963. sessionId: sid('s-new'), updatedAt: 2, running: false, blank: true, cwd: '/w/a',
  964. })
  965. await Promise.resolve()
  966. expect(b.svc.list.getSnapshot().byId[sid('s-new')]).toMatchObject({ blank: true })
  967. // Reconnect re-pull: the summary's blank=false wins (authoritative alignment).
  968. await feedList(b, [{ id: 's-new', blank: false, cwd: '/w/a' }])
  969. expect(b.svc.list.getSnapshot().byId[sid('s-new')]).toMatchObject({ blank: false })
  970. })
  971. it('never re-blanks: a stale blank=true summary cannot hide an engaged session', async ({ bench }) => {
  972. const b = bench()
  973. await feedList(b, [{ id: 's1', blank: true }])
  974. using reference = b.svc.retain(sid('s1'), { source: 'controllerOperation' })
  975. await reference.ready
  976. const session = reference.binding.session
  977. await session.prompt([{ type: 'text', text: 'hi' }], 'queue')
  978. await Promise.resolve()
  979. expect(b.svc.list.getSnapshot().byId[sid('s1')]).toMatchObject({ blank: false })
  980. // The next list pull still claims blank (host hasn't logged the message yet).
  981. await feedList(b, [{ id: 's1', blank: true }])
  982. expect(b.svc.binding(sid('s1'))?.session.getSnapshot().blank).toBe(false)
  983. })
  984. })
  985. describe('coverage tails (branch duals)', () => {
  986. it('displayTitleOf falls back to the id for empty and separator-only cwd', async ({ bench }) => {
  987. const b = bench()
  988. await feedList(b, [{ id: 'no-base', cwd: '///' }, { id: 'empty-cwd', cwd: '' }])
  989. const { byId } = b.svc.list.getSnapshot()
  990. expect(byId[sid('no-base')]?.displayTitle).toBe('no-base')
  991. expect(byId[sid('empty-cwd')]?.displayTitle).toBe('empty-cwd')
  992. expect(byId[sid('no-base')]?.title).toBeUndefined()
  993. })
  994. it('reading an unknown binding leaves an existing reference unchanged', async ({ bench }) => {
  995. const b = bench()
  996. await feedList(b, [{ id: 's1' }])
  997. using reference = b.svc.retain(sid('s1'), { source: 'controllerOperation' })
  998. await reference.ready
  999. expect(b.svc.binding(sid('ghost'))).toBeUndefined()
  1000. expect(b.svc.binding(sid('s1'))).toBe(reference.binding)
  1001. })
  1002. })