queue-store.client.spec.ts 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302
  1. /**
  2. * Queue snapshot semantics: authoritative replacement after every host-side
  3. * change, reconnect re-baselining, pre-instantiation buffering, editable-text
  4. * projection, and snapshot reference stability.
  5. */
  6. import { describe, expect, it, vi } from 'vitest'
  7. import { createUserMessage } from '@deepseek-ai/dsh-llm'
  8. import type { ContentBlock, UserMessage } from '@deepseek-ai/dsh-llm/types'
  9. import type { SessionEvent } from '@deepseek-ai/dsh-session/types'
  10. import type { MessageId, RpcId, SessionId } from '@deepseek-ai/dsh-api-remotes/client'
  11. import type {
  12. SessionControlFrame,
  13. SessionInteractionId,
  14. } from '@deepseek-ai/dsh-api-session-controller/types'
  15. import { Session } from '../src/client/sessions/session.ts'
  16. import { SessionManager } from '../src/client/sessions/manager.ts'
  17. import { FakeApiClient, fakeRemote } from './fake-api.client.ts'
  18. const SID = 'fk-q1' as SessionId
  19. const text = (value: string): ContentBlock[] => [{ type: 'text', text: value }]
  20. const rid = (id: string): RpcId => id as RpcId
  21. const iid = (id: string): MessageId => id as MessageId
  22. interface QueueFixture {
  23. id: string
  24. body: string
  25. content?: ContentBlock[]
  26. placement?: 'queued' | 'steering'
  27. message?: UserMessage
  28. }
  29. /** Build one authoritative queue snapshot. */
  30. function queueFrame(items: QueueFixture[]): Extract<SessionControlFrame, { type: 'queue' }> {
  31. return {
  32. type: 'queue',
  33. sessionId: SID,
  34. items: items.map(item => ({
  35. id: iid(item.id),
  36. placement: item.placement ?? 'queued',
  37. message: (item.message ?? createUserMessage({
  38. content: item.content ?? text(item.body),
  39. source: { kind: 'user', rpcId: rid(`rpc-${item.id}`) } as never,
  40. })) as never,
  41. })),
  42. }
  43. }
  44. function makeSession(): Session {
  45. return makeBench().session
  46. }
  47. function makeBench(): { api: FakeApiClient; session: Session } {
  48. const api = new FakeApiClient()
  49. return { api, session: new Session(SID, api, fakeRemote(api)) }
  50. }
  51. function makeManager(): SessionManager {
  52. const api = new FakeApiClient()
  53. return new SessionManager(api, fakeRemote(api))
  54. }
  55. function interactionId(value: string): SessionInteractionId {
  56. return value as SessionInteractionId
  57. }
  58. describe('queue snapshot intake', () => {
  59. it('projects stable ids, flat previews, and complete text', () => {
  60. const session = makeSession()
  61. session.handleControlFrame(queueFrame([
  62. { id: 'q-1', body: '第一条 排队\n消息' },
  63. ]))
  64. const queue = session.getSnapshot().queue
  65. expect(typeof queue[0]?.messageId).toBe('string')
  66. expect(queue).toMatchObject([
  67. {
  68. id: 'q-1', placement: 'queued',
  69. content: [{ type: 'text', text: '第一条 排队\n消息' }],
  70. preview: '第一条 排队 消息', text: '第一条 排队\n消息',
  71. },
  72. ])
  73. })
  74. it('marks mixed-content messages non-editable while retaining their preview', () => {
  75. const session = makeSession()
  76. session.handleControlFrame(queueFrame([{
  77. id: 'q-image',
  78. body: '',
  79. content: [{ type: 'text', text: 'hi' }, { type: 'image', data: 'x' } as never],
  80. }]))
  81. const queue = session.getSnapshot().queue
  82. expect(typeof queue[0]?.messageId).toBe('string')
  83. expect(queue).toMatchObject([
  84. {
  85. id: 'q-image', placement: 'queued',
  86. content: [{ type: 'text', text: 'hi' }, { type: 'image', data: 'x' }],
  87. preview: 'hi [image]', text: null,
  88. },
  89. ])
  90. })
  91. it('caps previews at 200 code points and preserves the full editable text', () => {
  92. const session = makeSession()
  93. const body = '长'.repeat(201)
  94. session.handleControlFrame(queueFrame([{ id: 'q-cap', body }]))
  95. const row = session.getSnapshot().queue[0]
  96. expect(Array.from(row?.preview ?? '')).toHaveLength(201)
  97. expect(row?.preview.endsWith('…')).toBe(true)
  98. expect(row?.text).toBe(body)
  99. })
  100. it('replaces content, order, and membership from each authoritative frame', () => {
  101. const session = makeSession()
  102. session.handleControlFrame(queueFrame([
  103. { id: 'q-1', body: 'one' },
  104. { id: 'q-2', body: 'two' },
  105. ]))
  106. session.handleControlFrame(queueFrame([
  107. { id: 'q-2', body: 'two edited' },
  108. ]))
  109. const queue = session.getSnapshot().queue
  110. expect(typeof queue[0]?.messageId).toBe('string')
  111. expect(queue).toMatchObject([
  112. {
  113. id: 'q-2', placement: 'queued',
  114. content: [{ type: 'text', text: 'two edited' }],
  115. preview: 'two edited', text: 'two edited',
  116. },
  117. ])
  118. session.handleControlFrame(queueFrame([]))
  119. expect(session.getSnapshot().queue).toEqual([])
  120. })
  121. it('keeps the queue array reference stable across unrelated snapshot swaps', () => {
  122. const session = makeSession()
  123. session.handleControlFrame(queueFrame([{ id: 'q-stable', body: '稳定' }]))
  124. const before = session.getSnapshot().queue
  125. session.handleAgentError('unrelated')
  126. expect(session.getSnapshot().queue).toBe(before)
  127. })
  128. it('retains steering placement and complete content in the same authoritative snapshot', () => {
  129. const session = makeSession()
  130. session.handleControlFrame(queueFrame([
  131. { id: 'q-next', body: 'later' },
  132. { id: 's-now', body: 'interrupt now', placement: 'steering' },
  133. ]))
  134. expect(session.getSnapshot().queue.map(item => ({
  135. id: item.id, placement: item.placement, content: item.content,
  136. }))).toEqual([
  137. { id: 'q-next', placement: 'queued', content: text('later') },
  138. { id: 's-now', placement: 'steering', content: text('interrupt now') },
  139. ])
  140. })
  141. it('hands off exactly one current occurrence when live steering becomes durable', async () => {
  142. const { api, session } = makeBench()
  143. await session.open()
  144. const message = createUserMessage({
  145. content: text('same message'),
  146. source: { kind: 'user' },
  147. })
  148. session.handleControlFrame(queueFrame([
  149. { id: 's-first', body: '', placement: 'steering', message },
  150. { id: 's-second', body: '', placement: 'steering', message },
  151. ]))
  152. const durable = {
  153. seq: 0,
  154. time: 1_700_000_000_000,
  155. type: 'user/message',
  156. surfaceOp: 'append',
  157. data: message,
  158. } as SessionEvent
  159. await api.pushFollow(SID, { type: 'event', event: durable as never })
  160. await vi.waitFor(() => {
  161. expect(session.getSnapshot().queue.map(item => item.id)).toEqual(['s-second'])
  162. })
  163. session.handleControlFrame(queueFrame([
  164. { id: 's-later', body: '', placement: 'steering', message },
  165. ]))
  166. await api.pushFollow(SID, { type: 'event', event: durable as never })
  167. await vi.waitFor(() => {
  168. expect(session.getSnapshot().queue.map(item => item.id)).toEqual(['s-later'])
  169. })
  170. })
  171. it('hands off live steering when the agent claims it as a user message', async () => {
  172. const { api, session } = makeBench()
  173. await session.open()
  174. const message = createUserMessage({
  175. content: text('claimed steering'),
  176. source: { kind: 'user' },
  177. })
  178. session.handleControlFrame(queueFrame([
  179. { id: 's-claimed', body: '', placement: 'steering', message },
  180. ]))
  181. await api.pushFollow(SID, {
  182. type: 'event',
  183. event: {
  184. seq: 0,
  185. time: 1_700_000_000_000,
  186. type: 'user/message',
  187. surfaceOp: 'append',
  188. data: message,
  189. } as never,
  190. })
  191. await vi.waitFor(() => {
  192. expect(session.getSnapshot().queue).toEqual([])
  193. })
  194. })
  195. })
  196. describe('queue operation transport', () => {
  197. it('addresses the session.updateQueue RPC without optimistic local mutation', async () => {
  198. const api = new FakeApiClient()
  199. const session = new Session(SID, api, fakeRemote(api))
  200. session.handleControlFrame(queueFrame([{ id: 'q-op', body: 'pending' }]))
  201. const before = session.getSnapshot().queue
  202. await expect(session.updateQueue(iid('q-op'), { kind: 'edit', content: text('next') }))
  203. .resolves.toEqual({ ok: true, value: { accepted: true } })
  204. await expect(session.updateQueue(iid('q-op'), { kind: 'steer' }))
  205. .resolves.toEqual({ ok: true, value: { accepted: true } })
  206. expect(api.callsOf('session.updateQueue')).toEqual([
  207. {
  208. sessionId: SID,
  209. itemId: 'q-op',
  210. action: { kind: 'edit', content: text('next') },
  211. },
  212. {
  213. sessionId: SID,
  214. itemId: 'q-op',
  215. action: { kind: 'steer' },
  216. },
  217. ])
  218. expect(session.getSnapshot().queue).toBe(before)
  219. })
  220. })
  221. describe('queue reconnect semantics', () => {
  222. it('a control baseline clears stale state before a fresh update lands', () => {
  223. const session = makeSession()
  224. session.handleControlFrame(queueFrame([{ id: 'q-old', body: '旧连接' }]))
  225. session.replaceControl([], [])
  226. expect(session.getSnapshot().queue).toEqual([])
  227. session.handleControlFrame(queueFrame([{ id: 'q-new', body: '新基线' }]))
  228. expect(session.getSnapshot().queue.map(row => row.id)).toEqual(['q-new'])
  229. })
  230. it('resync does not clear a baseline that raced ahead of the host connection signal', async () => {
  231. const session = makeSession()
  232. await session.open()
  233. session.handleControlFrame(queueFrame([{ id: 'q-fresh', body: '新基线' }]))
  234. await session.resync()
  235. expect(session.getSnapshot().queue.map(row => row.id)).toEqual(['q-fresh'])
  236. })
  237. it('running-status changes never guess at queue retirement', () => {
  238. const session = makeSession()
  239. session.handleControlFrame(queueFrame([{ id: 'q-live', body: '保留' }]))
  240. session.handleRunning(true)
  241. session.handleRunning(false)
  242. expect(session.getSnapshot().queue.map(row => row.id)).toEqual(['q-live'])
  243. })
  244. })
  245. describe('manager buffering of queue snapshots', () => {
  246. it('replays only the latest snapshot for an uninstantiated session', () => {
  247. const manager = makeManager()
  248. manager.handleControlFrame(queueFrame([{ id: 'q-old', body: '旧' }]))
  249. manager.handleControlFrame(queueFrame([{ id: 'q-new', body: '新' }]))
  250. expect(manager.get(SID).getSnapshot().queue.map(row => row.id)).toEqual(['q-new'])
  251. })
  252. it('a control baseline replaces the prior queue and restores answerable interactions', () => {
  253. const manager = makeManager()
  254. manager.handleControlFrame(queueFrame([{ id: 'q-g1', body: '第一代' }]))
  255. const nextQueue = queueFrame([{ id: 'q-g2', body: '第二代' }]).items
  256. manager.handleControlFrame({
  257. type: 'baseline',
  258. value: {
  259. queues: { [SID]: nextQueue },
  260. jobs: {},
  261. approvals: [{
  262. interactionId: interactionId('approval-1'),
  263. sessionId: SID,
  264. approvalId: 'ap-1' as never,
  265. toolName: 'bash',
  266. }],
  267. questions: [],
  268. projections: {},
  269. },
  270. })
  271. const snapshot = manager.get(SID).getSnapshot()
  272. expect(snapshot.queue.map(row => row.id)).toEqual(['q-g2'])
  273. expect(snapshot.pending.map(pending => pending.kind)).toEqual(['approval'])
  274. })
  275. })