fixture.spec.ts 20 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375
  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 } 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. appendSilent(id: string, msg: string): void
  21. breakStreams(): void
  22. }
  23. const timing = (): TimingHooks => (globalThis as Record<string, unknown>).__fxTiming as TimingHooks
  24. /** Collect stream frames until the predicate or a soft cap; abort ends the stream. */
  25. async function collect<F>(stream: AsyncIterable<RpcRequest<F>>, abort: AbortController, done: (frames: F[]) => boolean): Promise<F[]> {
  26. const frames: F[] = []
  27. for await (const envelope of stream) {
  28. frames.push(envelope.payload)
  29. if (done(frames) || frames.length > 500) {
  30. abort.abort()
  31. break
  32. }
  33. }
  34. return frames
  35. }
  36. describe('createFixtureApi', () => {
  37. it('serves the session list sorted by updatedAt desc and echoes rpcIds on every unary', async () => {
  38. const api = createFixtureApi()
  39. const request = req({})
  40. const response = await api.sessions.list(request)
  41. expect(response.rpcId).toBe(request.rpcId)
  42. if (!response.result.ok) throw new Error('list failed')
  43. expect(response.result.value.items.map(s => s.sessionId)).toEqual(['fx-alpha', 'fx-beta', 'fx-gamma'])
  44. expect(response.result.value.items[1]?.parentSessionId).toBe('fx-alpha') // lineage material
  45. })
  46. it('pages history backwards on message-boundary cuts with seq-contiguous stitching', async () => {
  47. const api = createFixtureApi()
  48. const tail = await api.sessions.history(req({ sessionId: sid('fx-alpha'), maxMessages: 10 }))
  49. if (!tail.result.ok) throw new Error('history failed')
  50. const tailPage = tail.result.value
  51. expect(tailPage.hasMore).toBe(true)
  52. expect(tailPage.events[0]?.event.type).toBe('turn/start') // cut lands on a turn boundary
  53. const boundary = tailPage.events[0]?.event.seq ?? 0
  54. expect(boundary).toBeGreaterThan(0)
  55. const older = await api.sessions.history(req({ sessionId: sid('fx-alpha'), beforeSeq: boundary, maxMessages: 10 }))
  56. if (!older.result.ok) throw new Error('older failed')
  57. const olderTail = older.result.value.events.at(-1)?.event
  58. expect((olderTail?.seq ?? -1) + 1).toBe(boundary) // pages stitch with no hole/overlap
  59. // Out-of-range beforeSeq clamps instead of exploding.
  60. const clamped = await api.sessions.history(req({ sessionId: sid('fx-alpha'), beforeSeq: -5, maxMessages: 10 }))
  61. if (!clamped.result.ok) throw new Error('clamped failed')
  62. expect(clamped.result.value.events).toEqual([])
  63. // Unknown session: empty page, not an error (history of a bare id).
  64. const empty = await api.sessions.history(req({ sessionId: sid('no-such'), maxMessages: 10 }))
  65. if (!empty.result.ok) throw new Error('empty failed')
  66. expect(empty.result.value).toEqual({ events: [], hasMore: false })
  67. })
  68. it('create adds a session and pushes host/session-added to open host streams', async () => {
  69. const api = createFixtureApi()
  70. const abort = new AbortController()
  71. const seen: HostFrame[] = []
  72. const consuming = (async () => {
  73. for await (const envelope of api.events.host(req({}), abort.signal)) {
  74. seen.push(envelope.payload)
  75. if (seen.length >= 1) abort.abort()
  76. }
  77. })()
  78. await new Promise(resolve => setTimeout(resolve, 10)) // let the stream register
  79. const created = await api.sessions.create(req({}))
  80. if (!created.result.ok) throw new Error('create failed')
  81. await consuming
  82. if (!created.result.ok) throw new Error('create failed')
  83. const createdId = created.result.value.sessionId
  84. expect(seen).toEqual([{ type: 'host/session-added', sessionId: createdId }])
  85. const list = await api.sessions.list(req({}))
  86. if (!list.result.ok) throw new Error('list failed')
  87. expect(list.result.value.items.some(s => s.sessionId === createdId)).toBe(true)
  88. })
  89. it('prompt replays a full streamed turn and cancel mid-replay freezes with (已中断)', async () => {
  90. const api = createFixtureApi()
  91. const created = await api.sessions.create(req({}))
  92. if (!created.result.ok) throw new Error('create failed')
  93. const id = created.result.value.sessionId
  94. const abort = new AbortController()
  95. const frames: MuxFrame[] = []
  96. const consuming = (async () => {
  97. for await (const envelope of api.events.mux(req({}), abort.signal)) {
  98. frames.push(envelope.payload)
  99. const last = envelope.payload
  100. if (last.type === 'session/event' && last.event.type === 'turn/end') {
  101. abort.abort()
  102. }
  103. }
  104. })()
  105. await new Promise(resolve => setTimeout(resolve, 10))
  106. // Unknown session → session-not-found with the id echoed in details.
  107. const missing = await api.sessions.prompt(req({ sessionId: sid('ghost'), mode: 'queue' as const, content: [{ type: 'text' as const, text: 'x' }] }))
  108. expect(missing.result).toMatchObject({ ok: false, error: { code: 'session-not-found', details: { sessionId: 'ghost' } } })
  109. // Real prompt: replay starts (running flips true), cancel freezes it.
  110. const accepted = await api.sessions.prompt(req({ sessionId: id, mode: 'queue' as const, content: [{ type: 'text' as const, text: 'render markdown' }] }))
  111. expect(accepted.result).toMatchObject({ ok: true, value: { accepted: true } })
  112. await new Promise(resolve => setTimeout(resolve, 120)) // a couple of typewriter ticks
  113. await api.sessions.cancel(req({ sessionId: id }))
  114. await consuming
  115. const types = frames.filter((f): f is Extract<MuxFrame, { type: 'session/event' }> => f.type === 'session/event').map(f => f.event.type)
  116. expect(types).toContain('turn/start')
  117. expect(types).toContain('user/message')
  118. expect(types).toContain('assistant/chunk')
  119. expect(types).toContain('assistant/message')
  120. expect(types.at(-1)).toBe('turn/end')
  121. const finalize = frames.find((f): f is Extract<MuxFrame, { type: 'session/event' }> => f.type === 'session/event' && f.event.type === 'assistant/message')
  122. expect(JSON.stringify(finalize?.event.data)).toContain('(已中断)')
  123. // Idle cancel: no replay in flight, must not explode; running flips false.
  124. const idleCancel = await api.sessions.cancel(req({ sessionId: id }))
  125. expect(idleCancel.result).toMatchObject({ ok: true })
  126. })
  127. it('steer during a replay inserts a steering message and the replay continues to completion', async () => {
  128. const api = createFixtureApi()
  129. const created = await api.sessions.create(req({}))
  130. if (!created.result.ok) throw new Error('create failed')
  131. const id = created.result.value.sessionId
  132. const abort = new AbortController()
  133. const framesPromise = collect<MuxFrame>(api.events.mux(req({}), abort.signal), abort,
  134. frames => frames.some(f => f.type === 'session/event' && f.event.type === 'turn/end'))
  135. await new Promise(resolve => setTimeout(resolve, 10))
  136. await api.sessions.prompt(req({ sessionId: id, mode: 'queue' as const, content: [{ type: 'text' as const, text: '短' }] }))
  137. await api.sessions.prompt(req({ sessionId: id, mode: 'steer' as const, content: [{ type: 'text' as const, text: '插话' }] }))
  138. const frames = await framesPromise
  139. const types = frames.filter((f): f is Extract<MuxFrame, { type: 'session/event' }> => f.type === 'session/event').map(f => f.event.type)
  140. expect(types).toContain('steering/message')
  141. expect(types.at(-1)).toBe('turn/end') // steer did not restart the turn
  142. })
  143. it('mux open replays subscribed sessions and resident interactions with stable rpcIds', async () => {
  144. const api = createFixtureApi()
  145. const openOnce = async (): Promise<RpcRequest<MuxFrame>[]> => {
  146. const abort = new AbortController()
  147. const envelopes: RpcRequest<MuxFrame>[] = []
  148. for await (const envelope of api.events.mux(req({}), abort.signal)) {
  149. envelopes.push(envelope)
  150. if (envelopes.length >= 4) abort.abort()
  151. }
  152. return envelopes
  153. }
  154. const first = await openOnce()
  155. const second = await openOnce()
  156. expect(first[0]?.payload).toMatchObject({ type: 'session/subscribed', sessionId: 'fx-alpha' })
  157. expect((first[0]?.payload as { lastSeq: number }).lastSeq).toBeGreaterThan(0)
  158. expect(first[1]?.payload).toMatchObject({ type: 'session/title', sessionId: 'fx-alpha', title: 'Fixture 历史会话' })
  159. expect(first[2]?.payload).toMatchObject({ type: 'approval/requested', toolName: 'dangerous_tool' })
  160. expect(second[2]?.rpcId).toBe(first[2]?.rpcId) // stable rpcId across replays (host replay semantics)
  161. expect(first[3]?.payload).toMatchObject({ type: 'question/requested', sessionId: 'fx-alpha' })
  162. expect(second[3]?.rpcId).toBe(first[3]?.rpcId)
  163. })
  164. it('steer with no replay in flight falls through to a fresh queued turn; non-text blocks stringify empty', async () => {
  165. const api = createFixtureApi()
  166. const abort = new AbortController()
  167. const framesPromise = collect<MuxFrame>(api.events.mux(req({}), abort.signal), abort,
  168. frames => frames.some(f => f.type === 'session/event' && f.event.type === 'turn/end'))
  169. await new Promise(resolve => setTimeout(resolve, 10))
  170. const created = await api.sessions.create(req({}))
  171. if (!created.result.ok) throw new Error('create failed')
  172. // steer while idle + a non-text content block (covers the '' arm of the text join).
  173. await api.sessions.prompt(req({
  174. sessionId: created.result.value.sessionId, mode: 'steer' as const,
  175. content: [{ type: 'text' as const, text: '短' }, { type: 'image', data: 'x' } as never],
  176. }))
  177. const frames = await framesPromise
  178. const types = frames.filter((f): f is Extract<MuxFrame, { type: 'session/event' }> => f.type === 'session/event').map(f => f.event.type)
  179. expect(types[0]).toBe('turn/start') // idle steer degraded to a queued turn, not a steering insert
  180. })
  181. it('gamma interval flip emits host/session-status and a running log-less session subscribes at lastSeq -1', async () => {
  182. vi.useFakeTimers()
  183. try {
  184. const api = createFixtureApi()
  185. const abort = new AbortController()
  186. const hostSeen: HostFrame[] = []
  187. const consuming = (async () => {
  188. for await (const envelope of api.events.host(req({}), abort.signal)) hostSeen.push(envelope.payload)
  189. })()
  190. await vi.advanceTimersByTimeAsync(5001) // interval fires: fx-gamma flips running=true (no log exists)
  191. expect(hostSeen).toContainEqual({ type: 'host/session-status', sessionId: sid('fx-gamma'), running: true })
  192. // A mux stream opened now sees gamma in the baseline with lastSeq = -1 (empty log arm).
  193. const mabort = new AbortController()
  194. const baseline: MuxFrame[] = []
  195. const muxConsuming = (async () => {
  196. for await (const envelope of api.events.mux(req({}), mabort.signal)) {
  197. baseline.push(envelope.payload)
  198. if (baseline.length >= 3) mabort.abort()
  199. }
  200. })()
  201. await vi.advanceTimersByTimeAsync(10)
  202. mabort.abort()
  203. await muxConsuming
  204. expect(baseline).toContainEqual({ type: 'session/subscribed', sessionId: sid('fx-gamma'), lastSeq: -1 })
  205. abort.abort()
  206. await vi.advanceTimersByTimeAsync(10)
  207. await consuming
  208. } finally {
  209. vi.useRealTimers()
  210. }
  211. })
  212. it('respond resolves the resident question once and rejects duplicate or unrelated ids', async () => {
  213. const api = createFixtureApi()
  214. expect(await api.respond({ type: 'client-response', rpcId: RpcId('x'), result: { ok: true, value: {} } })).toEqual({ accepted: false, reason: 'not-pending' })
  215. const abort = new AbortController()
  216. let question: RpcRequest<MuxFrame> | undefined
  217. for await (const envelope of api.events.mux(req({}), abort.signal)) {
  218. if (envelope.payload.type !== 'question/requested') continue
  219. question = envelope
  220. abort.abort()
  221. }
  222. if (question === undefined) throw new Error('fixture question missing')
  223. const response = { type: 'client-response' as const, rpcId: question.rpcId, result: { ok: true as const, value: {} } }
  224. expect(await api.respond(response)).toEqual({ accepted: true })
  225. expect(await api.respond(response)).toEqual({ accepted: false, reason: 'not-pending' })
  226. const replayAbort = new AbortController()
  227. const replayed = await collect(api.events.mux(req({}), replayAbort.signal), replayAbort, frames => frames.length === 2)
  228. expect(replayed.every(frame => frame.type !== 'question/requested')).toBe(true)
  229. const cancelledApi = createFixtureApi()
  230. const cancelAbort = new AbortController()
  231. let cancelQuestion: RpcRequest<MuxFrame> | undefined
  232. for await (const envelope of cancelledApi.events.mux(req({}), cancelAbort.signal)) {
  233. if (envelope.payload.type !== 'question/requested') continue
  234. cancelQuestion = envelope
  235. cancelAbort.abort()
  236. }
  237. if (cancelQuestion === undefined) throw new Error('fixture cancellation question missing')
  238. expect(await cancelledApi.respond({
  239. type: 'client-response', rpcId: cancelQuestion.rpcId,
  240. result: { ok: false, error: { code: 'cancelled', message: 'skip', details: {} } },
  241. })).toEqual({ accepted: true })
  242. })
  243. it('describe answers the fixture identity', async () => {
  244. const api = createFixtureApi()
  245. const response = await api.host.describe(req({}))
  246. expect(response.result).toMatchObject({ ok: true, value: { version: '0.0.0-fixture', attachedSessions: 1 } })
  247. })
  248. it('timing hooks: history delay + one-shot failure, silent append, and breakStreams end open generators', async () => {
  249. const api = createFixtureApi()
  250. const hooks = timing()
  251. // One-shot transport failure after transit delay.
  252. hooks.setHistoryDelay(5)
  253. hooks.failNextHistory()
  254. await expect(api.sessions.history(req({ sessionId: sid('fx-alpha'), maxMessages: 5 }))).rejects.toThrow(/simulated history transport failure/)
  255. hooks.setHistoryDelay(0)
  256. // The failure was one-shot: the next call succeeds.
  257. const ok = await api.sessions.history(req({ sessionId: sid('fx-alpha'), maxMessages: 5 }))
  258. expect(ok.result.ok).toBe(true)
  259. // appendUser emits on the mux stream; appendSilent only lands in the log (lost frame).
  260. const abort = new AbortController()
  261. const seen: MuxFrame[] = []
  262. const consuming = (async () => {
  263. for await (const envelope of api.events.mux(req({}), abort.signal)) seen.push(envelope.payload)
  264. })()
  265. await new Promise(resolve => setTimeout(resolve, 10))
  266. hooks.appendSilent('fx-alpha', '静默丢帧')
  267. hooks.appendUser('fx-alpha', '正常直播')
  268. hooks.appendTitle('fx-alpha', 'Fixture 修订标题')
  269. await vi.waitFor(() => {
  270. expect(seen.some(f => f.type === 'session/event' && JSON.stringify(f.event.data).includes('正常直播'))).toBe(true)
  271. expect(seen.some(f => f.type === 'session/title' && f.title === 'Fixture 修订标题')).toBe(true)
  272. })
  273. expect(seen.some(f => f.type === 'session/event' && JSON.stringify(f.event.data).includes('静默丢帧'))).toBe(false)
  274. const rawTitleIndex = seen.findIndex(f => f.type === 'session/event' && (f.event as { type: string }).type === 'session/title')
  275. const titleControlIndex = seen.findIndex(f => f.type === 'session/title' && f.title === 'Fixture 修订标题')
  276. expect(titleControlIndex).toBe(rawTitleIndex + 1)
  277. // But history serves the silent event (the client's repull finds it).
  278. const repull = await api.sessions.history(req({ sessionId: sid('fx-alpha'), maxMessages: 5 }))
  279. if (!repull.result.ok) throw new Error('repull failed')
  280. expect(JSON.stringify(repull.result.value.events)).toContain('静默丢帧')
  281. // breakStreams force-ends BOTH stream kinds without the client abort.
  282. const habort = new AbortController()
  283. const hostConsuming = (async () => {
  284. for await (const _ of api.events.host(req({}), habort.signal)) { /* drain */ }
  285. })()
  286. await new Promise(resolve => setTimeout(resolve, 10))
  287. hooks.breakStreams()
  288. await consuming // returns because the stream broke, not because we aborted
  289. await hostConsuming
  290. expect(abort.signal.aborted).toBe(false)
  291. expect(habort.signal.aborted).toBe(false)
  292. })
  293. })
  294. describe('FixtureApiClient (protocol-level fake carrier)', () => {
  295. afterEach(() => {
  296. vi.restoreAllMocks()
  297. })
  298. it('doFetch is an unreachable tripwire (all protocol paths overridden)', () => {
  299. const client = new FixtureApiClient()
  300. // Protected at compile time only; reach it directly to pin the tripwire message.
  301. expect(() => (client as unknown as { doFetch(): Promise<Response> }).doFetch()).toThrow(/doFetch must be unreachable/)
  302. })
  303. it('mints request ids, taps all four full forms, and never touches doFetch', async () => {
  304. const client = new FixtureApiClient()
  305. const tapped: RpcMessage[] = []
  306. client.subscribeEnvelopes(batch => tapped.push(...batch))
  307. const response = await client.sessions.list({})
  308. expect(response.result.ok).toBe(true)
  309. await client.respond({ type: 'client-response', rpcId: RpcId('r-x'), result: { ok: true, value: {} } })
  310. await vi.waitFor(() => {
  311. const kinds = tapped.map(m => m.type)
  312. expect(kinds).toContain('client-request')
  313. expect(kinds).toContain('server-response')
  314. expect(kinds).toContain('client-response')
  315. })
  316. const request = tapped.find(m => m.type === 'client-request')
  317. const reply = tapped.find(m => m.type === 'server-response')
  318. expect(request?.rpcId).toBe(reply?.rpcId) // echo discipline holds through the fake carrier
  319. })
  320. it('covers the whole unary dispatch table', async () => {
  321. const client = new FixtureApiClient()
  322. const created = await client.sessions.create({})
  323. if (!created.result.ok) throw new Error('create failed')
  324. const id = created.result.value.sessionId
  325. expect((await client.sessions.history({ sessionId: id })).result.ok).toBe(true)
  326. expect((await client.sessions.prompt({ sessionId: id, mode: 'queue', content: [{ type: 'text', text: '嗨' }] })).result.ok).toBe(true)
  327. expect((await client.sessions.cancel({ sessionId: id })).result.ok).toBe(true)
  328. expect((await client.host.describe({})).result.ok).toBe(true)
  329. })
  330. it('fires onOpen at stream-iteration start and taps server-request full forms', async () => {
  331. const client = new FixtureApiClient()
  332. const tapped: RpcMessage[] = []
  333. client.subscribeEnvelopes(batch => tapped.push(...batch))
  334. const order: string[] = []
  335. const abort = new AbortController()
  336. for await (const envelope of client.events.mux({}, abort.signal, () => order.push('open'))) {
  337. order.push(envelope.payload.type)
  338. abort.abort()
  339. }
  340. expect(order[0]).toBe('open')
  341. expect(order[1]).toBe('session/subscribed')
  342. await vi.waitFor(() => {
  343. expect(tapped.some(m => m.type === 'server-request')).toBe(true)
  344. })
  345. // Host stream side of the pair (same tap path).
  346. const habort = new AbortController()
  347. const hostOrder: string[] = []
  348. const hostIterator = client.events.host({}, habort.signal, () => hostOrder.push('open'))[Symbol.asyncIterator]()
  349. const raced = await Promise.race([hostIterator.next(), new Promise<'idle'>(resolve => setTimeout(() => { resolve('idle') }, 50))])
  350. expect(hostOrder).toEqual(['open']) // established even though the host stream stays silent
  351. habort.abort()
  352. if (raced === 'idle') await hostIterator.return?.(undefined)
  353. })
  354. })