queue-store.client.spec.ts 11 KB

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