queue-store.spec.ts 10.0 KB

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