queue-store.spec.ts 11 KB

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