fixture.spec.ts 59 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108
  1. /**
  2. * Fixture impl semantics: the demo data source must honor the same contract
  3. * shapes as the real host (paging boundaries, rpcId echo, replay lifecycle,
  4. * baseline replay, timing hooks) — this is the vitest-side drift detector for
  5. * the hand-written fixture/host parallel implementations.
  6. */
  7. import { afterEach, describe, expect, it, vi } from 'vitest'
  8. import type { SessionId, WorkspaceId } from '../src/client/api.ts'
  9. import { RpcId } from '../src/client/api.ts'
  10. import type { HostFrame, MuxFrame, RpcMessage, RpcRequest } from '../src/client/api.ts'
  11. import { FixtureApiClient, createFixtureApi } from '../src/client/fixture.ts'
  12. const sid = (id: string): SessionId => id as SessionId
  13. const req = <P>(payload: P): RpcRequest<P> => ({ rpcId: RpcId(`t-${Math.abs(Math.sin(reqCount++)).toString(36).slice(2, 10)}`), payload })
  14. let reqCount = 0
  15. interface TimingHooks {
  16. setHistoryDelay(ms: number): void
  17. failNextHistory(): void
  18. appendUser(id: string, msg: string): void
  19. appendTitle(id: string, title: string): void
  20. startReasoningChunkStorm(id: string, chunkCount: number, chunksPerInterval: number, intervalMs: number): string
  21. reasoningChunkStormState(): {
  22. sessionId: string
  23. chunkCount: number
  24. chunksPerInterval: number
  25. intervalMs: number
  26. emitted: number
  27. marker: string
  28. emitting: boolean
  29. } | null
  30. beginModelRetry(id: string): void
  31. scheduleModelRetry(id: string, retry?: number, delayMs?: number): void
  32. cancelModelRetryDuringBackoff(id: string, delayMs?: number): void
  33. completeModelRetry(id: string): void
  34. appendSilent(id: string, msg: string): void
  35. breakStreams(): void
  36. }
  37. const timing = (): TimingHooks => (globalThis as Record<string, unknown>).__fxTiming as TimingHooks
  38. /** Collect stream frames until the predicate or a soft cap; abort ends the stream. */
  39. async function collect<F>(stream: AsyncIterable<RpcRequest<F>>, abort: AbortController, done: (frames: F[]) => boolean): Promise<F[]> {
  40. const frames: F[] = []
  41. for await (const envelope of stream) {
  42. frames.push(envelope.payload)
  43. if (done(frames) || frames.length > 500) {
  44. abort.abort()
  45. break
  46. }
  47. }
  48. return frames
  49. }
  50. describe('createFixtureApi', () => {
  51. it('serves the session list sorted by updatedAt desc and echoes rpcIds on every unary', async () => {
  52. const api = createFixtureApi()
  53. const request = req({})
  54. const response = await api.sessions.list(request)
  55. expect(response.rpcId).toBe(request.rpcId)
  56. if (!response.result.ok) throw new Error('list failed')
  57. expect(response.result.value.items.map(s => s.sessionId)).toEqual(['fx-alpha', 'fx-beta', 'fx-gamma'])
  58. expect(response.result.value.items[1]?.parentSessionId).toBe('fx-alpha') // lineage material
  59. })
  60. it('searches current message text with literal unicode61-style token phrases', async () => {
  61. const api = createFixtureApi()
  62. const signal = new AbortController().signal
  63. const phrase = await api.sessions.search(req({ query: 'FIXTURE 历史消息' }), signal)
  64. expect(phrase.result).toMatchObject({
  65. ok: true,
  66. value: {
  67. items: [{ sessionId: 'fx-alpha' }],
  68. hasMore: false,
  69. },
  70. })
  71. if (!phrase.result.ok) throw new Error('search failed')
  72. expect(phrase.result.value.items[0]?.snippet).toContain('fixture 历史消息')
  73. timing().appendUser(
  74. 'fx-alpha',
  75. `${'leading context '.repeat(20)}late café token${' trailing context'.repeat(20)}`,
  76. )
  77. const late = await api.sessions.search(req({ query: 'LATE CAFE TOKEN' }), signal)
  78. if (!late.result.ok) throw new Error('late search failed')
  79. const lateSnippet = late.result.value.items[0]?.snippet ?? ''
  80. expect(lateSnippet).toContain('late café token')
  81. expect(lateSnippet.startsWith('…')).toBe(true)
  82. expect(lateSnippet.endsWith('…')).toBe(true)
  83. expect(Array.from(lateSnippet).length).toBeLessThanOrEqual(120)
  84. timing().appendUser('fx-alpha', 'Greek final sigma: ος')
  85. const finalSigma = await api.sessions.search(req({ query: 'ΟΣ' }), signal)
  86. if (!finalSigma.result.ok) throw new Error('final sigma search failed')
  87. expect(finalSigma.result.value.items[0]?.snippet).toContain('ος')
  88. const substring = await api.sessions.search(req({ query: 'ixtur' }), signal)
  89. expect(substring.result).toEqual({
  90. ok: true,
  91. value: { items: [], hasMore: false },
  92. })
  93. const punctuationOnly = await api.sessions.search(req({ query: '*' }), signal)
  94. expect(punctuationOnly.result).toEqual({
  95. ok: true,
  96. value: { items: [], hasMore: false },
  97. })
  98. const reasoningOnly = await api.sessions.search(req({ query: '思考过程' }), signal)
  99. expect(reasoningOnly.result).toEqual({
  100. ok: true,
  101. value: { items: [], hasMore: false },
  102. })
  103. const aborted = new AbortController()
  104. aborted.abort()
  105. await expect(api.sessions.search(req({ query: 'fixture' }), aborted.signal))
  106. .resolves.toMatchObject({ result: { ok: false, error: { code: 'cancelled' } } })
  107. })
  108. it('pages history backwards on message-boundary cuts with seq-contiguous stitching', async () => {
  109. const api = createFixtureApi()
  110. const tail = await api.sessions.history(req({ sessionId: sid('fx-alpha'), maxMessages: 10 }))
  111. if (!tail.result.ok) throw new Error('history failed')
  112. const tailPage = tail.result.value
  113. expect(tailPage.hasMore).toBe(true)
  114. expect(tailPage.events[0]?.event.type).toBe('turn/start') // cut lands on a turn boundary
  115. const boundary = tailPage.events[0]?.event.seq ?? 0
  116. expect(boundary).toBeGreaterThan(0)
  117. const older = await api.sessions.history(req({ sessionId: sid('fx-alpha'), beforeSeq: boundary, maxMessages: 10 }))
  118. if (!older.result.ok) throw new Error('older failed')
  119. const olderTail = older.result.value.events.at(-1)?.event
  120. expect((olderTail?.seq ?? -1) + 1).toBe(boundary) // pages stitch with no hole/overlap
  121. // Out-of-range beforeSeq clamps instead of exploding.
  122. const clamped = await api.sessions.history(req({ sessionId: sid('fx-alpha'), beforeSeq: -5, maxMessages: 10 }))
  123. if (!clamped.result.ok) throw new Error('clamped failed')
  124. expect(clamped.result.value.events).toEqual([])
  125. // Unknown session: empty page, not an error (history of a bare id). The
  126. // tail block still rides it — empty-log cut at -1, the host convention.
  127. const empty = await api.sessions.history(req({ sessionId: sid('no-such'), maxMessages: 10 }))
  128. if (!empty.result.ok) throw new Error('empty failed')
  129. // Fixture composes the todos + plan units (host parallel when tool-todo
  130. // and plan-mode are mounted): the empty-log values.
  131. expect(empty.result.value).toEqual({
  132. events: [], hasMore: false, projections: { asOfSeq: -1, values: {
  133. todos: null,
  134. // Permission unit composed: the composition-default select.
  135. permissions: {
  136. options: [
  137. { value: 'workspace-write', name: 'workspace-write', description: 'Write inside the workspace and permitted temporary directories; wider retries require approval.' },
  138. { value: 'danger-full-access', name: 'danger-full-access', description: 'Full file access without approval prompts.' },
  139. ],
  140. currentValue: 'workspace-write',
  141. },
  142. plan: { active: false, pending: false },
  143. goal: null,
  144. tokenUsage: {
  145. uncachedInputTokens: 0,
  146. outputTokens: 0,
  147. cacheReadTokens: 0,
  148. cacheWriteTokens: 0,
  149. },
  150. // No request ran, so neither pressure nor capacity is known yet.
  151. contextPressure: {},
  152. contextBreakdown: {
  153. systemTokens: 0,
  154. toolsTokens: 0,
  155. messageTokens: 0,
  156. },
  157. } },
  158. })
  159. })
  160. it('serves grouped models and keeps a selection for later history and fixture requests', async () => {
  161. const api = createFixtureApi()
  162. const sessionId = sid('fx-alpha')
  163. const catalog = await api.sessions.models(req({ sessionId }))
  164. if (!catalog.result.ok) throw new Error('models failed')
  165. expect(catalog.result.value.groups.map(group => group.name)).toEqual(['DeepSeek', 'OpenAI'])
  166. expect(catalog.result.value.groups[0]?.models.map(model => model.id))
  167. .toEqual(['deepseek-v4-flash', 'deepseek-v4-pro'])
  168. const selected = await api.sessions.selectModel(req({
  169. sessionId,
  170. provider: 'openai',
  171. model: 'gpt-5',
  172. }))
  173. if (!selected.result.ok) throw new Error('selection failed')
  174. expect(selected.result.value.selected).toEqual({ provider: 'openai', model: 'gpt-5' })
  175. const history = await api.sessions.history(req({ sessionId }))
  176. if (!history.result.ok) throw new Error('history failed')
  177. const prompt = await api.sessions.prompt(req({
  178. sessionId,
  179. mode: 'queue',
  180. content: [{ type: 'text', text: 'report model' }],
  181. }))
  182. expect(prompt.result.ok).toBe(true)
  183. await new Promise(resolve => setTimeout(resolve, 600))
  184. const after = await api.sessions.history(req({ sessionId }))
  185. if (!after.result.ok) throw new Error('history failed')
  186. expect(JSON.stringify(after.result.value.events)).toContain('openai/gpt-5')
  187. })
  188. it('serves configured DeepSeek readiness and keeps credential values write-only', async () => {
  189. const api = createFixtureApi()
  190. const settings = await api.settings.describe(req({}))
  191. if (!settings.result.ok) throw new Error('settings describe failed')
  192. expect(settings.result.value.namespaces).toMatchObject([{
  193. ns: 'llm-deepseek',
  194. value: { apiKeyEnv: 'DEEPSEEK_API_KEY' },
  195. secrets: [{ path: ['apiKey'], set: false }],
  196. }])
  197. const initial = await api.credentials.describe(req({ refs: ['DEEPSEEK_API_KEY', 'TEST_API_KEY'] }))
  198. if (!initial.result.ok) throw new Error('credential describe failed')
  199. expect(initial.result.value.credentials).toEqual({
  200. DEEPSEEK_API_KEY: { configured: true, source: 'file', writable: true },
  201. TEST_API_KEY: { configured: false, writable: true },
  202. })
  203. await api.credentials.set(req({ ref: 'TEST_API_KEY', value: 'write-only-fixture-secret' }))
  204. const configured = await api.credentials.describe(req({ refs: ['TEST_API_KEY'] }))
  205. if (!configured.result.ok) throw new Error('credential describe failed')
  206. expect(configured.result.value.credentials.TEST_API_KEY).toEqual({
  207. configured: true,
  208. source: 'file',
  209. writable: true,
  210. })
  211. await api.credentials.unset(req({ ref: 'TEST_API_KEY' }))
  212. const cleared = await api.credentials.describe(req({ refs: ['TEST_API_KEY'] }))
  213. if (!cleared.result.ok) throw new Error('credential describe failed')
  214. expect(cleared.result.value.credentials.TEST_API_KEY).toEqual({ configured: false, writable: true })
  215. })
  216. it('emits the todo/write snapshot at the real tool boundary: between tool/call and tool/result, timestamps monotonic', async () => {
  217. const api = createFixtureApi()
  218. const tail = await api.sessions.history(req({ sessionId: sid('fx-alpha'), maxMessages: 10 }))
  219. if (!tail.result.ok) throw new Error('history failed')
  220. const events = tail.result.value.events.map(e => e.event)
  221. const todoAt = events.findIndex(e => e.type === 'todo/write')
  222. expect(todoAt).toBeGreaterThan(0)
  223. // Production ordering (the tool appends mid-execution): call → snapshot → result.
  224. expect(events[todoAt - 1]?.type).toBe('tool/call')
  225. expect(events[todoAt + 1]?.type).toBe('tool/result')
  226. const times = events.slice(todoAt - 1, todoAt + 2).map(e => e.time)
  227. expect(times[0]).toBeLessThanOrEqual(times[1] ?? 0)
  228. expect(times[1]).toBeLessThanOrEqual(times[2] ?? 0)
  229. // The sample is a parallel plan: this fixture chooses the parallel policy,
  230. // so the surfaces fed from here face more than one active item.
  231. const snapshot = events[todoAt] as { data: { todos: { status: string }[] } }
  232. expect(snapshot.data.todos.filter(t => t.status === 'in_progress')).toHaveLength(2)
  233. })
  234. it('create adds a session and pushes host/session-added to open host streams', async () => {
  235. const api = createFixtureApi()
  236. const abort = new AbortController()
  237. const seen: HostFrame[] = []
  238. const consuming = (async () => {
  239. for await (const envelope of api.events.host(req({}), abort.signal)) {
  240. seen.push(envelope.payload)
  241. if (seen.length >= 1) abort.abort()
  242. }
  243. })()
  244. await new Promise(resolve => setTimeout(resolve, 10)) // let the stream register
  245. const created = await api.sessions.create(req({}))
  246. if (!created.result.ok) throw new Error('create failed')
  247. await consuming
  248. if (!created.result.ok) throw new Error('create failed')
  249. const createdId = created.result.value.sessionId
  250. expect(seen).toEqual([{ type: 'host/session-added', sessionId: createdId, blank: true, cwd: '/tmp/fixture' }])
  251. const list = await api.sessions.list(req({}))
  252. if (!list.result.ok) throw new Error('list failed')
  253. expect(list.result.value.items.some(s => s.sessionId === createdId)).toBe(true)
  254. })
  255. it('prompt replays a full streamed turn and cancel mid-replay freezes with (已中断)', async () => {
  256. const api = createFixtureApi()
  257. const created = await api.sessions.create(req({}))
  258. if (!created.result.ok) throw new Error('create failed')
  259. const id = created.result.value.sessionId
  260. const abort = new AbortController()
  261. const frames: MuxFrame[] = []
  262. const consuming = (async () => {
  263. for await (const envelope of api.events.mux(req({}), abort.signal)) {
  264. frames.push(envelope.payload)
  265. const last = envelope.payload
  266. if (last.type === 'session/event' && last.event.type === 'turn/end') {
  267. abort.abort()
  268. }
  269. }
  270. })()
  271. await new Promise(resolve => setTimeout(resolve, 10))
  272. // Unknown session → session-not-found with the id echoed in details.
  273. const missing = await api.sessions.prompt(req({ sessionId: sid('ghost'), mode: 'queue' as const, content: [{ type: 'text' as const, text: 'x' }] }))
  274. expect(missing.result).toMatchObject({ ok: false, error: { code: 'session-not-found', details: { sessionId: 'ghost' } } })
  275. // Real prompt: replay starts (running flips true), cancel freezes it.
  276. const accepted = await api.sessions.prompt(req({ sessionId: id, mode: 'queue' as const, content: [{ type: 'text' as const, text: 'render markdown' }] }))
  277. expect(accepted.result).toMatchObject({ ok: true, value: { accepted: true } })
  278. await new Promise(resolve => setTimeout(resolve, 120)) // a couple of typewriter ticks
  279. await api.sessions.cancel(req({ sessionId: id }))
  280. await consuming
  281. const types = frames.filter((f): f is Extract<MuxFrame, { type: 'session/event' }> => f.type === 'session/event').map(f => f.event.type)
  282. expect(types).toContain('turn/start')
  283. expect(types).toContain('user/message')
  284. expect(types).toContain('assistant/chunk')
  285. expect(types).toContain('assistant/message')
  286. expect(types.at(-1)).toBe('turn/end')
  287. // Capacity is durable log state, not a transient frame: the prompt path
  288. // records request/context and the projection carries it to the client.
  289. expect(types).toContain('request/context')
  290. expect(frames.some(frame =>
  291. frame.type === 'session/projection'
  292. && frame.key === 'tokenUsage'
  293. && (frame.value as { outputTokens?: number }).outputTokens === 8)).toBe(true)
  294. expect(frames.some(frame =>
  295. frame.type === 'session/projection'
  296. && frame.key === 'contextPressure'
  297. && (frame.value as { contextWindow?: number }).contextWindow === 128_000)).toBe(true)
  298. expect(frames.some(frame =>
  299. frame.type === 'session/projection'
  300. && frame.key === 'contextBreakdown'
  301. && (frame.value as { messageTokens?: number }).messageTokens! > 0)).toBe(true)
  302. const finalize = frames.find((f): f is Extract<MuxFrame, { type: 'session/event' }> => f.type === 'session/event' && f.event.type === 'assistant/message')
  303. expect(JSON.stringify(finalize?.event.data)).toContain('(已中断)')
  304. // Idle cancel: no replay in flight, must not explode; running flips false.
  305. const idleCancel = await api.sessions.cancel(req({ sessionId: id }))
  306. expect(idleCancel.result).toMatchObject({ ok: true })
  307. })
  308. it('steer during a replay lands a user/message inside the current turn and the replay continues', async () => {
  309. const api = createFixtureApi()
  310. const created = await api.sessions.create(req({}))
  311. if (!created.result.ok) throw new Error('create failed')
  312. const id = created.result.value.sessionId
  313. const abort = new AbortController()
  314. const framesPromise = collect<MuxFrame>(api.events.mux(req({}), abort.signal), abort,
  315. frames => frames.some(f => f.type === 'session/event' && f.event.type === 'turn/end'))
  316. await new Promise(resolve => setTimeout(resolve, 10))
  317. await api.sessions.prompt(req({ sessionId: id, mode: 'queue' as const, content: [{ type: 'text' as const, text: '短' }] }))
  318. await api.sessions.prompt(req({ sessionId: id, mode: 'steer' as const, content: [{ type: 'text' as const, text: '插话' }] }))
  319. const frames = await framesPromise
  320. const types = frames.filter((f): f is Extract<MuxFrame, { type: 'session/event' }> => f.type === 'session/event').map(f => f.event.type)
  321. expect(JSON.stringify(frames)).toContain('插话')
  322. expect(types.at(-1)).toBe('turn/end') // steer did not restart the turn
  323. })
  324. it('mux open replays subscribed sessions and resident interactions with stable rpcIds', async () => {
  325. const api = createFixtureApi()
  326. const openOnce = async (): Promise<RpcRequest<MuxFrame>[]> => {
  327. const abort = new AbortController()
  328. const envelopes: RpcRequest<MuxFrame>[] = []
  329. for await (const envelope of api.events.mux(req({}), abort.signal)) {
  330. envelopes.push(envelope)
  331. if (envelopes.length >= 11) abort.abort()
  332. }
  333. return envelopes
  334. }
  335. const first = await openOnce()
  336. const second = await openOnce()
  337. expect(first[0]?.payload).toMatchObject({ type: 'session/subscribed', sessionId: 'fx-alpha' })
  338. expect((first[0]?.payload as { lastSeq: number }).lastSeq).toBeGreaterThan(0)
  339. // Projection baseline frames follow subscribed (domain units + token usage).
  340. expect(first[1]?.payload).toMatchObject({ type: 'session/projection', sessionId: 'fx-alpha', key: 'title', value: 'Fixture 历史会话' })
  341. expect(first[2]?.payload).toMatchObject({ type: 'session/projection', sessionId: 'fx-alpha', key: 'todos' })
  342. expect(first[3]?.payload).toMatchObject({ type: 'session/projection', sessionId: 'fx-alpha', key: 'permissions' })
  343. expect(first[4]?.payload).toMatchObject({ type: 'session/projection', sessionId: 'fx-alpha', key: 'plan', value: { active: false, pending: false } })
  344. expect(first[5]?.payload).toMatchObject({ type: 'session/projection', sessionId: 'fx-alpha', key: 'goal', value: null })
  345. expect(first[6]?.payload).toMatchObject({ type: 'session/projection', sessionId: 'fx-alpha', key: 'tokenUsage' })
  346. expect(first[7]?.payload).toMatchObject({ type: 'session/projection', sessionId: 'fx-alpha', key: 'contextPressure' })
  347. expect(first[8]?.payload).toMatchObject({
  348. type: 'session/projection', sessionId: 'fx-alpha', key: 'contextBreakdown',
  349. value: { systemTokens: 0, toolsTokens: 0 },
  350. })
  351. expect((first[8]?.payload as { value: { messageTokens: number } }).value.messageTokens).toBeGreaterThan(0)
  352. expect(first[9]?.payload).toMatchObject({ type: 'approval/requested', toolName: 'dangerous_tool' })
  353. expect(second[9]?.rpcId).toBe(first[9]?.rpcId) // stable rpcId across replays (host replay semantics)
  354. expect(first[10]?.payload).toMatchObject({ type: 'question/requested', sessionId: 'fx-alpha' })
  355. expect(second[10]?.rpcId).toBe(first[10]?.rpcId)
  356. })
  357. it('steer with no replay in flight falls through to a fresh queued turn; non-text blocks stringify empty', async () => {
  358. const api = createFixtureApi()
  359. const abort = new AbortController()
  360. const framesPromise = collect<MuxFrame>(api.events.mux(req({}), abort.signal), abort,
  361. frames => frames.some(f => f.type === 'session/event' && f.event.type === 'turn/end'))
  362. await new Promise(resolve => setTimeout(resolve, 10))
  363. const created = await api.sessions.create(req({}))
  364. if (!created.result.ok) throw new Error('create failed')
  365. // steer while idle + a non-text content block (covers the '' arm of the text join).
  366. await api.sessions.prompt(req({
  367. sessionId: created.result.value.sessionId, mode: 'steer' as const,
  368. content: [{ type: 'text' as const, text: '短' }, { type: 'image', data: 'x' } as never],
  369. }))
  370. const frames = await framesPromise
  371. const types = frames.filter((f): f is Extract<MuxFrame, { type: 'session/event' }> => f.type === 'session/event').map(f => f.event.type)
  372. expect(types[0]).toBe('turn/start') // idle steer degraded to a queued turn, not an in-turn insert
  373. })
  374. it('gamma interval flip emits host/session-status and a running log-less session subscribes at lastSeq -1', async () => {
  375. vi.useFakeTimers()
  376. try {
  377. const api = createFixtureApi()
  378. const abort = new AbortController()
  379. const hostSeen: HostFrame[] = []
  380. const consuming = (async () => {
  381. for await (const envelope of api.events.host(req({}), abort.signal)) hostSeen.push(envelope.payload)
  382. })()
  383. await vi.advanceTimersByTimeAsync(5001) // interval fires: fx-gamma flips running=true (no log exists)
  384. expect(hostSeen).toContainEqual({ type: 'host/session-status', sessionId: sid('fx-gamma'), running: true })
  385. // A mux stream opened now sees gamma in the baseline with lastSeq = -1 (empty log arm).
  386. const mabort = new AbortController()
  387. const baseline: MuxFrame[] = []
  388. const muxConsuming = (async () => {
  389. for await (const envelope of api.events.mux(req({}), mabort.signal)) {
  390. baseline.push(envelope.payload)
  391. if (baseline.length >= 3) mabort.abort()
  392. }
  393. })()
  394. await vi.advanceTimersByTimeAsync(10)
  395. mabort.abort()
  396. await muxConsuming
  397. expect(baseline).toContainEqual({ type: 'session/subscribed', sessionId: sid('fx-gamma'), lastSeq: -1 })
  398. abort.abort()
  399. await vi.advanceTimersByTimeAsync(10)
  400. await consuming
  401. } finally {
  402. vi.useRealTimers()
  403. }
  404. })
  405. it('respond resolves the resident question once and rejects duplicate or unrelated ids', async () => {
  406. const api = createFixtureApi()
  407. expect(await api.respond({ type: 'client-response', rpcId: RpcId('x'), result: { ok: true, value: {} } })).toEqual({ accepted: false, reason: 'not-pending' })
  408. const abort = new AbortController()
  409. let question: RpcRequest<MuxFrame> | undefined
  410. for await (const envelope of api.events.mux(req({}), abort.signal)) {
  411. if (envelope.payload.type !== 'question/requested') continue
  412. question = envelope
  413. abort.abort()
  414. }
  415. if (question === undefined) throw new Error('fixture question missing')
  416. const response = { type: 'client-response' as const, rpcId: question.rpcId, result: { ok: true as const, value: {} } }
  417. expect(await api.respond(response)).toEqual({ accepted: true })
  418. expect(await api.respond(response)).toEqual({ accepted: false, reason: 'not-pending' })
  419. const replayAbort = new AbortController()
  420. const replayed = await collect(api.events.mux(req({}), replayAbort.signal), replayAbort, frames => frames.length === 2)
  421. expect(replayed.every(frame => frame.type !== 'question/requested')).toBe(true)
  422. const cancelledApi = createFixtureApi()
  423. const cancelAbort = new AbortController()
  424. let cancelQuestion: RpcRequest<MuxFrame> | undefined
  425. for await (const envelope of cancelledApi.events.mux(req({}), cancelAbort.signal)) {
  426. if (envelope.payload.type !== 'question/requested') continue
  427. cancelQuestion = envelope
  428. cancelAbort.abort()
  429. }
  430. if (cancelQuestion === undefined) throw new Error('fixture cancellation question missing')
  431. expect(await cancelledApi.respond({
  432. type: 'client-response', rpcId: cancelQuestion.rpcId,
  433. result: { ok: false, error: { code: 'cancelled', message: 'skip', details: {} } },
  434. })).toEqual({ accepted: true })
  435. })
  436. it('respond answers the resident approval once: routing, validation, resolved broadcast, then not-pending', async () => {
  437. const api = createFixtureApi()
  438. // Discover the resident approval's stable rpcId from the mux baseline.
  439. const abort = new AbortController()
  440. const seen: { rpcId: string; frame: MuxFrame }[] = []
  441. const consuming = (async () => {
  442. for await (const envelope of api.events.mux(req({}), abort.signal)) seen.push({ rpcId: envelope.rpcId, frame: envelope.payload })
  443. })()
  444. await vi.waitFor(() => {
  445. expect(seen.some(s => s.frame.type === 'approval/requested')).toBe(true)
  446. })
  447. const requested = seen.find(s => s.frame.type === 'approval/requested')
  448. if (requested === undefined || requested.frame.type !== 'approval/requested') throw new Error('unreachable')
  449. const approvalId = requested.frame.approvalId
  450. // Routed but malformed answers.
  451. expect(await api.respond({ type: 'client-response', rpcId: RpcId(requested.rpcId), result: { ok: false, error: { code: 'internal', message: 'x', details: {} } } }))
  452. .toEqual({ accepted: false, reason: 'bad-response' })
  453. expect(await api.respond({ type: 'client-response', rpcId: RpcId(requested.rpcId), result: { ok: true, value: { approvalId: 'wrong', outcome: 'rejected' } } }))
  454. .toEqual({ accepted: false, reason: 'bad-response' })
  455. expect(await api.respond({ type: 'client-response', rpcId: RpcId(requested.rpcId), result: { ok: true, value: { approvalId, outcome: 'maybe' } } }))
  456. .toEqual({ accepted: false, reason: 'bad-response' })
  457. // The real answer settles the question and broadcasts resolved.
  458. expect(await api.respond({ type: 'client-response', rpcId: RpcId(requested.rpcId), result: { ok: true, value: { sessionId: sid('fx-alpha'), approvalId, outcome: 'allowed-once' } } }))
  459. .toEqual({ accepted: true })
  460. await vi.waitFor(() => {
  461. expect(seen.some(s => s.frame.type === 'approval/resolved' && s.frame.outcome === 'allowed-once')).toBe(true)
  462. })
  463. // Settled: a duplicate answer is late, and a fresh mux open replays nothing.
  464. expect(await api.respond({ type: 'client-response', rpcId: RpcId(requested.rpcId), result: { ok: true, value: { sessionId: sid('fx-alpha'), approvalId, outcome: 'rejected' } } }))
  465. .toEqual({ accepted: false, reason: 'not-pending' })
  466. abort.abort()
  467. await consuming
  468. const abort2 = new AbortController()
  469. const replayed = await collect(api.events.mux(req({}), abort2.signal), abort2, frames => frames.length === 2)
  470. expect(replayed.some(f => f.type === 'approval/requested')).toBe(false)
  471. })
  472. it('describe answers the fixture identity', async () => {
  473. const api = createFixtureApi()
  474. const response = await api.host.describe(req({}))
  475. expect(response.result).toMatchObject({ ok: true, value: { version: '0.0.0-fixture', attachedSessions: 1 } })
  476. const empty = await createFixtureApi({ empty: true }).host.describe(req({}))
  477. expect(empty.result).toMatchObject({ ok: true, value: { attachedSessions: 0 } })
  478. })
  479. it('createDirectory under the root mints /name whose listing and crumbs share the identity', async () => {
  480. const api = createFixtureApi()
  481. const created = await api.host.createDirectory(req({ path: '/', name: 'srv' }))
  482. if (!created.result.ok) throw new Error('create failed')
  483. expect(created.result.value.path).toBe('/srv')
  484. const listed = await api.host.listDirectory(req({ path: '/srv' }), new AbortController().signal)
  485. if (!listed.result.ok) throw new Error('list failed')
  486. expect(listed.result.value.crumbs).toEqual([
  487. { name: '/', path: '/', hidden: false },
  488. { name: 'srv', path: '/srv', hidden: false },
  489. ])
  490. const root = await api.host.listDirectory(req({ path: '/' }), new AbortController().signal)
  491. if (!root.result.ok) throw new Error('root list failed')
  492. expect(root.result.value.entries).toContainEqual({ name: 'srv', path: '/srv', hidden: false })
  493. })
  494. it('workspace.list serves the resident account and create reuses on path collision', async () => {
  495. const api = createFixtureApi()
  496. const listed = await api.workspace.list(req({}))
  497. if (!listed.result.ok) throw new Error('list failed')
  498. expect(listed.result.value.items).toEqual([expect.objectContaining({
  499. workspaceId: 'fx-ws-fixture', path: '/tmp/fixture', title: 'fixture',
  500. sessionIds: ['fx-alpha', 'fx-beta', 'fx-gamma'],
  501. })])
  502. // path collision → the existing entity comes back, created:false, no frame.
  503. const reused = await api.workspace.create(req({ path: '/tmp/fixture' }))
  504. if (!reused.result.ok) throw new Error('reuse failed')
  505. expect(reused.result.value).toMatchObject({ created: false, workspace: { workspaceId: 'fx-ws-fixture' } })
  506. })
  507. it('workspace.create on a fresh path mints a new entity and pushes host/workspace-changed', async () => {
  508. const api = createFixtureApi()
  509. const abort = new AbortController()
  510. const seen: HostFrame[] = []
  511. const consuming = (async () => {
  512. for await (const envelope of api.events.host(req({}), abort.signal)) {
  513. seen.push(envelope.payload)
  514. abort.abort()
  515. }
  516. })()
  517. await new Promise(resolve => setTimeout(resolve, 10))
  518. const created = await api.workspace.create(req({ path: '/tmp/fixture-workspaces/nova' }))
  519. if (!created.result.ok) throw new Error('create failed')
  520. expect(created.result.value.created).toBe(true)
  521. expect(created.result.value.workspace).toMatchObject({
  522. path: '/tmp/fixture-workspaces/nova', title: 'nova', sessionIds: [],
  523. })
  524. await consuming
  525. expect(seen).toEqual([{ type: 'host/workspace-changed', workspace: created.result.value.workspace }])
  526. // A basename-less path serves as its own title.
  527. const rootPath = await api.workspace.create(req({ path: '/' }))
  528. if (!rootPath.result.ok) throw new Error('rootPath failed')
  529. expect(rootPath.result.value.workspace.title).toBe('/')
  530. })
  531. it('workspace.rename covers not-found, conflict, no-op, and the changed frame', async () => {
  532. const api = createFixtureApi()
  533. const abort = new AbortController()
  534. const seen: HostFrame[] = []
  535. const consuming = (async () => {
  536. for await (const envelope of api.events.host(req({}), abort.signal)) {
  537. seen.push(envelope.payload)
  538. if (seen.length >= 2) abort.abort()
  539. }
  540. })()
  541. await new Promise(resolve => setTimeout(resolve, 10))
  542. const wsid = 'fx-ws-fixture' as WorkspaceId
  543. const missing = await api.workspace.rename(req({ workspaceId: 'fx-ws-void' as WorkspaceId, title: 'x' }))
  544. expect(missing.result).toMatchObject({ ok: false, error: { code: 'workspace-not-found', details: { workspaceId: 'fx-ws-void' } } })
  545. await api.workspace.create(req({ path: '/tmp/fixture-workspaces/occupied' }))
  546. const conflict = await api.workspace.rename(req({ workspaceId: wsid, title: ' occupied ' }))
  547. expect(conflict.result).toMatchObject({ ok: false, error: { code: 'workspace-name-conflict', details: { name: 'occupied' } } })
  548. const noop = await api.workspace.rename(req({ workspaceId: wsid, title: ' fixture ' }))
  549. if (!noop.result.ok) throw new Error('no-op rename failed')
  550. expect(noop.result.value.workspace.title).toBe('fixture')
  551. const renamed = await api.workspace.rename(req({ workspaceId: wsid, title: 'renamed' }))
  552. if (!renamed.result.ok) throw new Error('rename failed')
  553. expect(renamed.result.value.workspace.title).toBe('renamed')
  554. await consuming
  555. // Only the create and the effective rename emit frames; the no-op stays silent.
  556. expect(seen.map(f => f.type)).toEqual(['host/workspace-changed', 'host/workspace-changed'])
  557. })
  558. it('session.rename covers not-found, blank title, and the accepted append + title frame', async () => {
  559. const api = createFixtureApi()
  560. const abort = new AbortController()
  561. const framesPromise = (async () => {
  562. const frames: MuxFrame[] = []
  563. for await (const envelope of api.events.mux(req({}), abort.signal)) {
  564. frames.push(envelope.payload)
  565. if (frames.some(f => f.type === 'session/projection' && f.key === 'title' && f.value === '重命名')) abort.abort()
  566. }
  567. return frames
  568. })()
  569. await new Promise(resolve => setTimeout(resolve, 10))
  570. const missing = await api.sessions.rename(req({ sessionId: sid('fx-void'), title: 'x' }))
  571. expect(missing.result).toMatchObject({ ok: false, error: { code: 'session-not-found', details: { sessionId: 'fx-void' } } })
  572. const blank = await api.sessions.rename(req({ sessionId: sid('fx-alpha'), title: ' ' }))
  573. expect(blank.result).toMatchObject({ ok: false, error: { code: 'title-invalid', details: { sessionId: 'fx-alpha' } } })
  574. const renamed = await api.sessions.rename(req({ sessionId: sid('fx-alpha'), title: ' 重命名 ' }))
  575. if (!renamed.result.ok) throw new Error('rename failed')
  576. expect(renamed.result.value.title).toBe('重命名')
  577. const acceptedSeq = renamed.result.value.seq
  578. // The response seq addresses the appended title event (the client plane
  579. // has no session/title in its event union — titles ride the projection —
  580. // so the event is located by seq and its payload checked structurally).
  581. const history = await api.sessions.history(req({ sessionId: sid('fx-alpha'), maxMessages: 100 }))
  582. if (!history.result.ok) throw new Error('history failed')
  583. const appended = history.result.value.events.find(entry => entry.event.seq === acceptedSeq)
  584. expect(appended?.event).toMatchObject({
  585. type: 'session/title',
  586. data: { title: '重命名', messageSeqs: [], source: { kind: 'user' } },
  587. })
  588. // Beyond the subscribe-time baseline replay, the append emitted exactly
  589. // one title projection frame carrying the new value at the response seq.
  590. const frames = await framesPromise
  591. const titleFrames = frames.filter(f => f.type === 'session/projection' && f.key === 'title' && f.sessionId === sid('fx-alpha') && f.value === '重命名')
  592. expect(titleFrames).toHaveLength(1)
  593. expect(titleFrames[0]).toMatchObject({ seq: acceptedSeq })
  594. })
  595. it('workspace.insertSessionBefore moves, appends, no-ops, and rejects invalid ids', async () => {
  596. const api = createFixtureApi()
  597. const wsid = 'fx-ws-fixture' as WorkspaceId
  598. const missing = await api.workspace.insertSessionBefore(req({ workspaceId: 'fx-ws-void' as WorkspaceId, sessionId: sid('fx-alpha') }))
  599. expect(missing.result).toMatchObject({ ok: false, error: { code: 'workspace-not-found' } })
  600. const ghost = await api.workspace.insertSessionBefore(req({ workspaceId: wsid, sessionId: sid('fx-ghost') }))
  601. expect(ghost.result).toMatchObject({ ok: false, error: { code: 'workspace-move-invalid', details: { sessionId: 'fx-ghost' } } })
  602. const badAnchor = await api.workspace.insertSessionBefore(req({ workspaceId: wsid, sessionId: sid('fx-alpha'), beforeSessionId: sid('fx-ghost') }))
  603. expect(badAnchor.result).toMatchObject({ ok: false, error: { code: 'workspace-move-invalid', details: { beforeSessionId: 'fx-ghost' } } })
  604. const moved = await api.workspace.insertSessionBefore(req({ workspaceId: wsid, sessionId: sid('fx-gamma'), beforeSessionId: sid('fx-beta') }))
  605. if (!moved.result.ok) throw new Error('move failed')
  606. expect(moved.result.value.workspace.sessionIds).toEqual(['fx-alpha', 'fx-gamma', 'fx-beta'])
  607. const appended = await api.workspace.insertSessionBefore(req({ workspaceId: wsid, sessionId: sid('fx-alpha') }))
  608. if (!appended.result.ok) throw new Error('append failed')
  609. expect(appended.result.value.workspace.sessionIds).toEqual(['fx-gamma', 'fx-beta', 'fx-alpha'])
  610. const before = appended.result.value.workspace.updatedAt
  611. const noop = await api.workspace.insertSessionBefore(req({ workspaceId: wsid, sessionId: sid('fx-alpha') }))
  612. if (!noop.result.ok) throw new Error('no-op move failed')
  613. expect(noop.result.value.workspace.sessionIds).toEqual(['fx-gamma', 'fx-beta', 'fx-alpha'])
  614. expect(noop.result.value.workspace.updatedAt).toBe(before)
  615. })
  616. it('workspace.delete removes only the Workspace row and emits the removal frame', async () => {
  617. const api = createFixtureApi()
  618. const abort = new AbortController()
  619. const seen: HostFrame[] = []
  620. const consuming = (async () => {
  621. for await (const envelope of api.events.host(req({}), abort.signal)) {
  622. seen.push(envelope.payload)
  623. abort.abort()
  624. }
  625. })()
  626. await new Promise(resolve => setTimeout(resolve, 10))
  627. const missing = await api.workspace.delete(req({ workspaceId: 'fx-ws-void' as WorkspaceId }))
  628. expect(missing.result).toMatchObject({ ok: false, error: { code: 'workspace-not-found' } })
  629. const deleted = await api.workspace.delete(req({ workspaceId: 'fx-ws-fixture' as WorkspaceId }))
  630. expect(deleted.result).toEqual({ ok: true, value: { deleted: true } })
  631. await consuming
  632. expect(seen).toEqual([{ type: 'host/workspace-removed', workspaceId: 'fx-ws-fixture' }])
  633. const list = await api.workspace.list(req({}))
  634. if (!list.result.ok) throw new Error('workspace list failed')
  635. expect(list.result.value.items.some(workspace => workspace.workspaceId === 'fx-ws-fixture')).toBe(false)
  636. const sessions = await api.sessions.list(req({}))
  637. if (!sessions.result.ok) throw new Error('session list failed')
  638. expect(sessions.result.value.items.map(session => session.sessionId)).toContain('fx-alpha')
  639. })
  640. it('session.create({workspaceId}) lands on the account and unknown ids error', async () => {
  641. const api = createFixtureApi()
  642. const abort = new AbortController()
  643. const seen: HostFrame[] = []
  644. const consuming = (async () => {
  645. for await (const envelope of api.events.host(req({}), abort.signal)) {
  646. seen.push(envelope.payload)
  647. if (seen.length >= 2) abort.abort()
  648. }
  649. })()
  650. await new Promise(resolve => setTimeout(resolve, 10))
  651. const missing = await api.sessions.create(req({ workspaceId: 'fx-ws-void' as WorkspaceId }))
  652. expect(missing.result).toMatchObject({ ok: false, error: { code: 'workspace-not-found', details: { workspaceId: 'fx-ws-void' } } })
  653. const created = await api.sessions.create(req({ workspaceId: 'fx-ws-fixture' as WorkspaceId }))
  654. if (!created.result.ok) throw new Error('create failed')
  655. const id = created.result.value.sessionId
  656. await consuming
  657. // The session lands with the workspace's path as cwd, and the account
  658. // write pushes the fresh workspace snapshot after session-added.
  659. expect(seen[0]).toEqual({ type: 'host/session-added', sessionId: id, blank: true, cwd: '/tmp/fixture' })
  660. expect(seen[1]).toMatchObject({
  661. type: 'host/workspace-changed',
  662. workspace: { workspaceId: 'fx-ws-fixture', sessionIds: [id, 'fx-alpha', 'fx-beta', 'fx-gamma'] },
  663. })
  664. })
  665. it('supports an empty baseline, preallocated ids, workspace-first frames, and idempotent retry', async () => {
  666. const api = createFixtureApi({ empty: true, createFrameOrder: 'workspace-first' })
  667. const initialSessions = await api.sessions.list(req({}))
  668. const initialWorkspaces = await api.workspace.list(req({}))
  669. expect(initialSessions.result).toMatchObject({ ok: true, value: { items: [] } })
  670. expect(initialWorkspaces.result).toMatchObject({ ok: true, value: { items: [] } })
  671. const made = await api.workspace.create(req({ path: '/tmp/fixture-workspaces/nova' }))
  672. if (!made.result.ok) throw new Error('workspace create failed')
  673. const abort = new AbortController()
  674. const framesPromise = collect(api.events.host(req({}), abort.signal), abort, frames => frames.length === 2)
  675. await new Promise(resolve => setTimeout(resolve, 10))
  676. const preallocated = sid('fx-preallocated')
  677. const created = await api.sessions.create(req({
  678. workspaceId: made.result.value.workspace.workspaceId,
  679. sessionId: preallocated,
  680. }))
  681. expect(created.result).toEqual({ ok: true, value: { sessionId: preallocated } })
  682. const frames = await framesPromise
  683. expect(frames[0]).toMatchObject({
  684. type: 'host/workspace-changed', workspace: { sessionIds: [preallocated] },
  685. })
  686. expect(frames[1]).toEqual({ type: 'host/session-added', sessionId: preallocated, blank: true, cwd: made.result.value.workspace.path })
  687. const retried = await api.sessions.create(req({
  688. workspaceId: made.result.value.workspace.workspaceId,
  689. sessionId: preallocated,
  690. }))
  691. expect(retried.result).toEqual({ ok: true, value: { sessionId: preallocated } })
  692. const listed = await api.sessions.list(req({}))
  693. if (!listed.result.ok) throw new Error('session list failed')
  694. expect(listed.result.value.items.filter(item => item.sessionId === preallocated)).toHaveLength(1)
  695. const conflict = await api.sessions.create(req({ sessionId: preallocated, cwd: '/elsewhere' }))
  696. expect(conflict.result).toMatchObject({
  697. ok: false,
  698. error: { code: 'session-conflict', details: { sessionId: preallocated, requestedCwd: '/elsewhere' } },
  699. })
  700. })
  701. it('attaches an existing ungrouped Session to a matching Workspace', async () => {
  702. const api = createFixtureApi()
  703. const sessionId = sid('fx-existing-ungrouped')
  704. await expect(api.sessions.create(req({ sessionId, cwd: '/tmp/fixture' }))).resolves.toMatchObject({
  705. result: { ok: true, value: { sessionId } },
  706. })
  707. await expect(api.sessions.create(req({
  708. sessionId,
  709. workspaceId: 'fx-ws-fixture' as WorkspaceId,
  710. }))).resolves.toMatchObject({ result: { ok: true, value: { sessionId } } })
  711. const workspaces = await api.workspace.list(req({}))
  712. if (!workspaces.result.ok) throw new Error('workspace list failed')
  713. expect(workspaces.result.value.items[0]?.sessionIds).toContain(sessionId)
  714. })
  715. it('reports a conflict without an existing cwd detail for an unrecorded cwd', async () => {
  716. const api = createFixtureApi()
  717. const listed = await api.sessions.list(req({}))
  718. if (!listed.result.ok) throw new Error('session list failed')
  719. const existing = listed.result.value.items.find(item => item.sessionId === sid('fx-alpha'))
  720. if (existing === undefined) throw new Error('fixture Session missing')
  721. delete existing.cwd
  722. const conflict = await api.sessions.create(req({ sessionId: existing.sessionId }))
  723. expect(conflict.result).toEqual({
  724. ok: false,
  725. error: {
  726. code: 'session-conflict',
  727. message: `session ${existing.sessionId} already uses no cwd`,
  728. details: { sessionId: existing.sessionId, requestedCwd: '/tmp/fixture' },
  729. },
  730. })
  731. })
  732. it('publishes an ungrouped Session when Workspace attachment fails', async () => {
  733. const api = createFixtureApi({ failWorkspaceAttach: true })
  734. const sessionId = sid('fx-partial')
  735. const created = await api.sessions.create(req({
  736. workspaceId: 'fx-ws-fixture' as WorkspaceId,
  737. sessionId,
  738. }))
  739. expect(created.result).toMatchObject({
  740. ok: false,
  741. error: { code: 'workspace-attach-failed', details: { sessionId, workspaceId: 'fx-ws-fixture' } },
  742. })
  743. const listed = await api.sessions.list(req({}))
  744. const workspaces = await api.workspace.list(req({}))
  745. if (!listed.result.ok || !workspaces.result.ok) throw new Error('list failed')
  746. expect(listed.result.value.items.filter(item => item.sessionId === sessionId)).toHaveLength(1)
  747. expect(workspaces.result.value.items[0]?.sessionIds).not.toContain(sessionId)
  748. const retried = await api.sessions.create(req({
  749. workspaceId: 'fx-ws-fixture' as WorkspaceId,
  750. sessionId,
  751. }))
  752. expect(retried.result).toMatchObject({ ok: false, error: { code: 'workspace-attach-failed' } })
  753. const afterRetry = await api.sessions.list(req({}))
  754. if (!afterRetry.result.ok) throw new Error('list failed')
  755. expect(afterRetry.result.value.items.filter(item => item.sessionId === sessionId)).toHaveLength(1)
  756. })
  757. it('reconciles a dropped create response and can reject a prompt before acceptance', async () => {
  758. const sessionId = sid('fx-lost-response')
  759. const dropped = createFixtureApi({ dropSessionCreateResponse: true })
  760. await expect(Promise.resolve().then(() => dropped.sessions.create(req({
  761. workspaceId: 'fx-ws-fixture' as WorkspaceId,
  762. sessionId,
  763. })))).rejects.toThrow(/dropped session\.create response/)
  764. const listed = await dropped.sessions.list(req({}))
  765. const workspaces = await dropped.workspace.list(req({}))
  766. if (!listed.result.ok || !workspaces.result.ok) throw new Error('list failed')
  767. expect(listed.result.value.items.some(item => item.sessionId === sessionId)).toBe(true)
  768. expect(workspaces.result.value.items[0]?.sessionIds).toContain(sessionId)
  769. await expect(dropped.sessions.create(req({
  770. workspaceId: 'fx-ws-fixture' as WorkspaceId,
  771. sessionId,
  772. }))).resolves.toMatchObject({ result: { ok: true, value: { sessionId } } })
  773. const rejecting = createFixtureApi({ empty: true, rejectPrompt: true })
  774. const real = await rejecting.sessions.create(req({ sessionId: sid('fx-rejected') }))
  775. if (!real.result.ok) throw new Error('session create failed')
  776. const prompt = await rejecting.sessions.prompt(req({
  777. sessionId: real.result.value.sessionId,
  778. mode: 'queue' as const,
  779. content: [{ type: 'text' as const, text: 'keep me' }],
  780. }))
  781. expect(prompt.result).toMatchObject({ ok: false, error: { code: 'agent-busy' } })
  782. })
  783. it('timing hooks: history delay + one-shot failure, silent append, and breakStreams end open generators', async () => {
  784. const api = createFixtureApi()
  785. const hooks = timing()
  786. // One-shot transport failure after transit delay.
  787. hooks.setHistoryDelay(5)
  788. hooks.failNextHistory()
  789. await expect(api.sessions.history(req({ sessionId: sid('fx-alpha'), maxMessages: 5 }))).rejects.toThrow(/simulated history transport failure/)
  790. hooks.setHistoryDelay(0)
  791. // The failure was one-shot: the next call succeeds.
  792. const ok = await api.sessions.history(req({ sessionId: sid('fx-alpha'), maxMessages: 5 }))
  793. expect(ok.result.ok).toBe(true)
  794. // appendUser emits on the mux stream; appendSilent only lands in the log (lost frame).
  795. const abort = new AbortController()
  796. const seen: MuxFrame[] = []
  797. const consuming = (async () => {
  798. for await (const envelope of api.events.mux(req({}), abort.signal)) seen.push(envelope.payload)
  799. })()
  800. await new Promise(resolve => setTimeout(resolve, 10))
  801. hooks.appendSilent('fx-alpha', '静默丢帧')
  802. hooks.appendUser('fx-alpha', '正常直播')
  803. hooks.appendTitle('fx-alpha', 'Fixture 修订标题')
  804. hooks.beginModelRetry('fx-alpha')
  805. hooks.scheduleModelRetry('fx-alpha')
  806. hooks.completeModelRetry('fx-alpha')
  807. hooks.beginModelRetry('fx-alpha')
  808. hooks.cancelModelRetryDuringBackoff('fx-alpha')
  809. await vi.waitFor(() => {
  810. expect(seen.some(f => f.type === 'session/event' && JSON.stringify(f.event.data).includes('正常直播'))).toBe(true)
  811. expect(seen.some(f => f.type === 'session/event' && (f.event as { type: string }).type === 'llm/retry')).toBe(true)
  812. expect(seen.some(f => f.type === 'session/event' && JSON.stringify(f.event.data).includes('重试后的完整回复'))).toBe(true)
  813. expect(seen.some(f => f.type === 'session/event'
  814. && f.event.type === 'turn/end'
  815. && f.event.data.reason.kind === 'aborted')).toBe(true)
  816. expect(seen.some(f => f.type === 'session/projection' && f.key === 'title' && f.value === 'Fixture 修订标题')).toBe(true)
  817. })
  818. expect(seen.some(f => f.type === 'session/event' && JSON.stringify(f.event.data).includes('静默丢帧'))).toBe(false)
  819. const rawTitleIndex = seen.findIndex(f => f.type === 'session/event' && (f.event as { type: string }).type === 'session/title')
  820. const titleControlIndex = seen.findIndex(f => f.type === 'session/projection' && f.key === 'title' && f.value === 'Fixture 修订标题')
  821. expect(titleControlIndex).toBe(rawTitleIndex + 1)
  822. // But history serves the silent event (the client's repull finds it).
  823. const repull = await api.sessions.history(req({ sessionId: sid('fx-alpha'), maxMessages: 5 }))
  824. if (!repull.result.ok) throw new Error('repull failed')
  825. expect(JSON.stringify(repull.result.value.events)).toContain('静默丢帧')
  826. // breakStreams force-ends BOTH stream kinds without the client abort.
  827. const habort = new AbortController()
  828. const hostConsuming = (async () => {
  829. for await (const _ of api.events.host(req({}), habort.signal)) { /* drain */ }
  830. })()
  831. await new Promise(resolve => setTimeout(resolve, 10))
  832. hooks.breakStreams()
  833. await consuming // returns because the stream broke, not because we aborted
  834. await hostConsuming
  835. expect(abort.signal.aborted).toBe(false)
  836. expect(habort.signal.aborted).toBe(false)
  837. })
  838. it('paces the opt-in reasoning stress hook from an external interval', async () => {
  839. vi.useFakeTimers()
  840. vi.setSystemTime(0)
  841. const api = createFixtureApi()
  842. const hooks = timing()
  843. expect(hooks.reasoningChunkStormState()).toBeNull()
  844. expect(() => hooks.startReasoningChunkStorm('fx-alpha', 0, 1, 16)).toThrow(/chunk count/)
  845. expect(() => hooks.startReasoningChunkStorm('fx-alpha', 1, 0, 16)).toThrow(/chunks per interval/)
  846. expect(() => hooks.startReasoningChunkStorm('fx-alpha', 1, 1, 0)).toThrow(/reasoning interval/)
  847. const abort = new AbortController()
  848. try {
  849. const streamed = collect(api.events.mux(req({}), abort.signal), abort, frames => frames.some(frame => (
  850. frame.type === 'session/event'
  851. && frame.event.type === 'assistant/chunk'
  852. && frame.event.data.chunk.type === 'reasoning-delta'
  853. && frame.event.data.chunk.text.includes('REASONING_STRESS_COMPLETE')
  854. )))
  855. const marker = hooks.startReasoningChunkStorm('fx-alpha', 3, 2, 16)
  856. expect(() => hooks.startReasoningChunkStorm('fx-alpha', 1, 1, 16)).toThrow(/already running/)
  857. expect(hooks.reasoningChunkStormState()).toMatchObject({ emitted: 0, emitting: true, marker })
  858. await vi.advanceTimersByTimeAsync(0)
  859. expect(hooks.reasoningChunkStormState()).toMatchObject({ emitted: 2, emitting: true })
  860. await vi.advanceTimersByTimeAsync(16)
  861. expect(hooks.reasoningChunkStormState()).toEqual({
  862. sessionId: 'fx-alpha', chunkCount: 3, chunksPerInterval: 2, intervalMs: 16,
  863. emitted: 3, marker, emitting: false,
  864. })
  865. const frames = await streamed
  866. const deltas = frames.flatMap(frame => (
  867. frame.type === 'session/event'
  868. && frame.event.type === 'assistant/chunk'
  869. && frame.event.data.chunk.type === 'reasoning-delta'
  870. ? [frame.event.data.chunk.text]
  871. : []
  872. ))
  873. expect(deltas).toEqual(['推理', '推理', `\n${marker}`])
  874. } finally {
  875. abort.abort()
  876. vi.useRealTimers()
  877. }
  878. })
  879. })
  880. describe('FixtureApiClient (protocol-level fake carrier)', () => {
  881. afterEach(() => {
  882. vi.restoreAllMocks()
  883. vi.unstubAllGlobals()
  884. })
  885. it('doFetch is an unreachable tripwire (all protocol paths overridden)', () => {
  886. const client = new FixtureApiClient()
  887. // Protected at compile time only; reach it directly to pin the tripwire message.
  888. expect(() => (client as unknown as { doFetch(): Promise<Response> }).doFetch()).toThrow(/doFetch must be unreachable/)
  889. })
  890. it('mints request ids, taps all four full forms, and never touches doFetch', async () => {
  891. const client = new FixtureApiClient()
  892. const tapped: RpcMessage[] = []
  893. client.subscribeEnvelopes(batch => tapped.push(...batch))
  894. const response = await client.sessions.list({})
  895. expect(response.result.ok).toBe(true)
  896. await client.respond({ type: 'client-response', rpcId: RpcId('r-x'), result: { ok: true, value: {} } })
  897. await vi.waitFor(() => {
  898. const kinds = tapped.map(m => m.type)
  899. expect(kinds).toContain('client-request')
  900. expect(kinds).toContain('server-response')
  901. expect(kinds).toContain('client-response')
  902. })
  903. const request = tapped.find(m => m.type === 'client-request')
  904. const reply = tapped.find(m => m.type === 'server-response')
  905. expect(request?.rpcId).toBe(reply?.rpcId) // echo discipline holds through the fake carrier
  906. })
  907. it('covers the whole unary dispatch table', async () => {
  908. const client = new FixtureApiClient()
  909. expect((await client.sessions.search(
  910. { query: 'fixture' },
  911. new AbortController().signal,
  912. )).result.ok).toBe(true)
  913. const created = await client.sessions.create({})
  914. if (!created.result.ok) throw new Error('create failed')
  915. const id = created.result.value.sessionId
  916. expect((await client.sessions.history({ sessionId: id })).result.ok).toBe(true)
  917. expect((await client.sessions.prompt({ sessionId: id, mode: 'queue', content: [{ type: 'text', text: '嗨' }] })).result.ok).toBe(true)
  918. expect((await client.sessions.cancel({ sessionId: id })).result.ok).toBe(true)
  919. expect((await client.host.describe({})).result.ok).toBe(true)
  920. expect((await client.workspace.list({})).result.ok).toBe(true)
  921. const workspace = await client.workspace.create({ path: '/tmp/fixture-workspaces/via-client' })
  922. if (!workspace.result.ok) throw new Error('workspace create failed')
  923. expect(workspace.result.value.workspace.title).toBe('via-client')
  924. const wsid = workspace.result.value.workspace.workspaceId
  925. const renamed = await client.workspace.rename({ workspaceId: wsid, title: 'via-client-2' })
  926. if (!renamed.result.ok) throw new Error('workspace rename failed')
  927. expect(renamed.result.value.workspace.title).toBe('via-client-2')
  928. const attached = await client.sessions.create({ workspaceId: wsid })
  929. if (!attached.result.ok) throw new Error('attached create failed')
  930. const moved = await client.workspace.insertSessionBefore({ workspaceId: wsid, sessionId: attached.result.value.sessionId })
  931. if (!moved.result.ok) throw new Error('workspace move failed')
  932. expect(moved.result.value.workspace.sessionIds).toEqual([attached.result.value.sessionId])
  933. // Goal lifecycle over the fixture fold: create → edit → pause → resume → complete → clear;
  934. // every mutation acknowledges with the NEW CAS ref (state rides the projection frames).
  935. const goalCreated = await client.goals.create({ sessionId: id, objective: 'ship it' })
  936. if (!goalCreated.result.ok) throw new Error('goal create failed')
  937. let ref = goalCreated.result.value.ref
  938. expect(ref.revision).toBe(1)
  939. const edited = await client.goals.edit({ sessionId: id, ref, objective: 'ship it v2' })
  940. if (!edited.result.ok) throw new Error('goal edit failed')
  941. ref = edited.result.value.ref
  942. const paused = await client.goals.pause({ sessionId: id, ref })
  943. if (!paused.result.ok) throw new Error('goal pause failed')
  944. ref = paused.result.value.ref
  945. const resumed = await client.goals.resume({ sessionId: id, ref })
  946. if (!resumed.result.ok) throw new Error('goal resume failed')
  947. ref = resumed.result.value.ref
  948. // A stale ref loses the CAS check.
  949. expect((await client.goals.pause({ sessionId: id, ref: { ...ref, revision: 1 } })).result.ok).toBe(false)
  950. const completed = await client.goals.complete({ sessionId: id, ref })
  951. if (!completed.result.ok) throw new Error('goal complete failed')
  952. ref = completed.result.value.ref
  953. // complete → complete is an invalid transition.
  954. expect((await client.goals.complete({ sessionId: id, ref })).result.ok).toBe(false)
  955. expect((await client.goals.clear({ sessionId: id, ref })).result).toEqual({ ok: true, value: { cleared: true } })
  956. const goalHistory = await client.sessions.history({ sessionId: id })
  957. if (!goalHistory.result.ok) throw new Error('goal history failed')
  958. const goalEvents = goalHistory.result.value.events.map(entry => entry.event as unknown as {
  959. type: string
  960. data: {
  961. operation?: string
  962. source?: { kind?: string; round?: number }
  963. }
  964. })
  965. const goalChanges = goalEvents.filter(event => event.type === 'goal/change')
  966. expect(goalChanges.map(event => event.data.operation))
  967. .toEqual(['create', 'edit', 'pause', 'resume', 'complete', 'clear'])
  968. expect(goalEvents.some(event => event.type === 'user/message'
  969. && event.data.source?.kind === 'goal' && event.data.source.round === 0)).toBe(false)
  970. })
  971. it('maps empty, prompt-reject, and workspace-first query scenarios', async () => {
  972. vi.stubGlobal('location', {
  973. search: '?fixture=empty&fixturePrompt=reject&fixtureFrames=workspace-first',
  974. })
  975. const client = new FixtureApiClient()
  976. await expect(client.sessions.list({})).resolves.toMatchObject({ result: { ok: true, value: { items: [] } } })
  977. const made = await client.workspace.create({ path: '/tmp/fixture-workspaces/query-workspace' })
  978. if (!made.result.ok) throw new Error('workspace create failed')
  979. const abort = new AbortController()
  980. const framesPromise = collect(client.events.host({}, abort.signal), abort, frames => frames.length === 2)
  981. await new Promise(resolve => setTimeout(resolve, 10))
  982. const sessionId = sid('fx-query-session')
  983. const created = await client.sessions.create({
  984. workspaceId: made.result.value.workspace.workspaceId,
  985. sessionId,
  986. })
  987. expect(created.result).toMatchObject({ ok: true, value: { sessionId } })
  988. const frames = await framesPromise
  989. expect(frames.map(frame => frame.type)).toEqual(['host/workspace-changed', 'host/session-added'])
  990. const rejected = await client.sessions.prompt({
  991. sessionId,
  992. mode: 'queue',
  993. content: [{ type: 'text', text: 'retain' }],
  994. })
  995. expect(rejected.result).toMatchObject({ ok: false, error: { code: 'agent-busy' } })
  996. })
  997. it('maps attach-failure and dropped-response query scenarios', async () => {
  998. vi.stubGlobal('location', { search: '?fixture&fixtureAttach=fail' })
  999. const partial = new FixtureApiClient()
  1000. const partialResult = await partial.sessions.create({
  1001. workspaceId: 'fx-ws-fixture' as WorkspaceId,
  1002. sessionId: sid('fx-query-partial'),
  1003. })
  1004. expect(partialResult.result).toMatchObject({
  1005. ok: false,
  1006. error: { code: 'workspace-attach-failed', details: { sessionId: 'fx-query-partial' } },
  1007. })
  1008. vi.stubGlobal('location', { search: '?fixture&fixtureSessionCreate=drop-response' })
  1009. const dropped = new FixtureApiClient()
  1010. await expect(dropped.sessions.create({
  1011. workspaceId: 'fx-ws-fixture' as WorkspaceId,
  1012. sessionId: sid('fx-query-dropped'),
  1013. })).rejects.toThrow(/dropped session\.create response/)
  1014. })
  1015. it('fires onOpen at stream-iteration start and taps server-request full forms', async () => {
  1016. const client = new FixtureApiClient()
  1017. const tapped: RpcMessage[] = []
  1018. client.subscribeEnvelopes(batch => tapped.push(...batch))
  1019. const order: string[] = []
  1020. const abort = new AbortController()
  1021. for await (const envelope of client.events.mux({}, abort.signal, () => order.push('open'))) {
  1022. order.push(envelope.payload.type)
  1023. abort.abort()
  1024. }
  1025. expect(order[0]).toBe('open')
  1026. expect(order[1]).toBe('session/subscribed')
  1027. await vi.waitFor(() => {
  1028. expect(tapped.some(m => m.type === 'server-request')).toBe(true)
  1029. })
  1030. // Host stream side of the pair (same tap path).
  1031. const habort = new AbortController()
  1032. const hostOrder: string[] = []
  1033. const hostIterator = client.events.host({}, habort.signal, () => hostOrder.push('open'))[Symbol.asyncIterator]()
  1034. const raced = await Promise.race([hostIterator.next(), new Promise<'idle'>(resolve => setTimeout(() => { resolve('idle') }, 50))])
  1035. expect(hostOrder).toEqual(['open']) // established even though the host stream stays silent
  1036. habort.abort()
  1037. if (raced === 'idle') await hostIterator.return?.(undefined)
  1038. })
  1039. })