control-queue.host.spec.ts 7.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201
  1. import { Context } from '@deepseek-ai/cordis'
  2. import AgentRegistry from '@deepseek-ai/dsh-agent'
  3. import type { Agent, Inbox } from '@deepseek-ai/dsh-agent'
  4. import { createUserMessage } from '@deepseek-ai/dsh-llm'
  5. import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
  6. import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection'
  7. import { describe, expect, it } from 'vitest'
  8. import { SessionControlController } from '../src/control.ts'
  9. import type { SessionControlFrame } from '../src/types.ts'
  10. import { createInboxFixture } from '@deepseek-ai/dsh-agent-loop-testkit'
  11. async function harness(): Promise<{
  12. ctx: Context
  13. control: SessionControlController
  14. agent: Agent
  15. inbox: Inbox
  16. }> {
  17. const ctx = new Context()
  18. await ctx.plugin(SessionStore)
  19. await ctx.plugin(SessionProjectionRegistry)
  20. await ctx.plugin(AgentRegistry)
  21. const session = ctx.sessions.create(SessionId('queue-session'))
  22. const { inbox } = createInboxFixture(ctx.sessionProjections, session)
  23. const agent: Agent = {
  24. id: session.id, options: {}, session, inbox, status: 'running', ctx,
  25. send: () => {}, followup: () => {}, steer: () => {}, inject: () => {}, cancel: () => {},
  26. runMaintenance: task => task(new AbortController().signal), whenIdle: () => Promise.resolve(),
  27. }
  28. ctx.agents.register(agent)
  29. return { ctx, control: new SessionControlController(ctx), agent, inbox }
  30. }
  31. function message(text: string, source: 'user' | 'plugin' = 'user') {
  32. return createUserMessage({
  33. content: [{ type: 'text', text }],
  34. source: source === 'user' ? { kind: 'user' } : { kind: 'plugin', plugin: 'fixture' },
  35. })
  36. }
  37. describe('Session control queue projection', () => {
  38. /** Consume frames until the next queue replacement (inbox projection frames interleave). */
  39. async function nextQueueFrame(
  40. iterator: AsyncIterator<SessionControlFrame>,
  41. ): Promise<Extract<SessionControlFrame, { type: 'queue' }>> {
  42. for (;;) {
  43. const next = await iterator.next()
  44. if (next.done) throw new Error('stream ended before a queue frame')
  45. if (next.value.type === 'queue') return next.value
  46. }
  47. }
  48. it('projects both pending lists in baselines and live replacement frames', async () => {
  49. const { control, inbox } = await harness()
  50. const queued = message('queued')
  51. const steering = message('steering')
  52. const context = message('context', 'plugin')
  53. inbox.append('next-turn', queued)
  54. inbox.append('next-step', steering)
  55. inbox.append('next-step', context)
  56. const abort = new AbortController()
  57. const iterator = control.control(abort.signal)[Symbol.asyncIterator]()
  58. const opened = await iterator.next()
  59. expect(opened.value).toMatchObject({
  60. type: 'baseline',
  61. value: {
  62. queues: {
  63. 'queue-session': [
  64. { id: queued.id, placement: 'queued' },
  65. { id: steering.id, placement: 'steering' },
  66. { id: context.id, placement: 'context' },
  67. ],
  68. },
  69. },
  70. })
  71. const replacement = message('replacement')
  72. inbox.append('next-turn', replacement)
  73. const replaced = await nextQueueFrame(iterator)
  74. expect(replaced.items.map(item => item.id)).toContain(replacement.id)
  75. inbox.remove(steering.id)
  76. const removed = await nextQueueFrame(iterator)
  77. expect(removed.items.map(item => item.id)).not.toContain(steering.id)
  78. abort.abort()
  79. await iterator.next()
  80. })
  81. it('derives queue replacements from the completed projection regardless of registration order', async () => {
  82. const ctx = new Context()
  83. await ctx.plugin(SessionStore)
  84. await ctx.plugin(SessionProjectionRegistry)
  85. await ctx.plugin(AgentRegistry)
  86. const control = new SessionControlController(ctx)
  87. const session = ctx.sessions.create(SessionId('late-projection-queue'))
  88. const { inbox } = createInboxFixture(ctx.sessionProjections, session)
  89. const agent: Agent = {
  90. id: session.id, options: {}, session, inbox, status: 'running', ctx,
  91. send: () => {}, followup: () => {}, steer: () => {}, inject: () => {}, cancel: () => {},
  92. runMaintenance: task => task(new AbortController().signal), whenIdle: () => Promise.resolve(),
  93. }
  94. ctx.agents.register(agent)
  95. const abort = new AbortController()
  96. const iterator = control.control(abort.signal)[Symbol.asyncIterator]()
  97. await iterator.next()
  98. const pending = message('late projection')
  99. inbox.append('next-turn', pending)
  100. await expect(iterator.next()).resolves.toMatchObject({
  101. value: {
  102. type: 'projection',
  103. key: 'inbox',
  104. value: { 'next-turn': [{ id: pending.id }], 'next-step': [] },
  105. },
  106. })
  107. await expect(nextQueueFrame(iterator)).resolves.toMatchObject({
  108. items: [{ id: pending.id, placement: 'queued' }],
  109. })
  110. abort.abort()
  111. await iterator.next()
  112. })
  113. it('projects the prompt rpcId from a user-rpc source and omits it elsewhere', async () => {
  114. const { control, inbox } = await harness()
  115. const identified = createUserMessage({
  116. content: [{ type: 'text', text: 'browser prompt' }],
  117. source: { kind: 'user', rpcId: 'req-42' as never },
  118. })
  119. inbox.append('next-turn', identified)
  120. inbox.append('next-step', message('plain steering'))
  121. const abort = new AbortController()
  122. const iterator = control.control(abort.signal)[Symbol.asyncIterator]()
  123. const opened = await iterator.next()
  124. if (opened.done || opened.value.type !== 'baseline') throw new Error('missing baseline')
  125. const items = opened.value.value.queues['queue-session' as SessionId] ?? []
  126. expect(items.map(item => ({ id: item.id, placement: item.placement, rpcId: item.rpcId }))).toEqual([
  127. { id: identified.id, placement: 'queued', rpcId: 'req-42' },
  128. { id: items[1]?.id, placement: 'steering', rpcId: undefined },
  129. ])
  130. expect('rpcId' in (items[1] ?? {})).toBe(false)
  131. abort.abort()
  132. await iterator.next()
  133. })
  134. it('ignores inbox events without the exact live Agent session', async () => {
  135. const { ctx, control, agent, inbox } = await harness()
  136. const abort = new AbortController()
  137. const iterator = control.control(abort.signal)[Symbol.asyncIterator]()
  138. await iterator.next()
  139. const unrelated = ctx.sessions.create(SessionId('unrelated-queue'))
  140. unrelated.append('agent/inbox/spliced', {
  141. target: 'next-turn',
  142. start: 0,
  143. inserted: [message('unrelated')],
  144. })
  145. const replacement = ctx.sessions.create(SessionId('replacement-session'))
  146. Object.defineProperty(agent, 'session', { configurable: true, value: replacement })
  147. inbox.append('next-turn', message('wrong-session'))
  148. abort.abort()
  149. await iterator.next()
  150. })
  151. it('drops broadcasts after cancellation has ended its queue', async () => {
  152. const { control, inbox } = await harness()
  153. const abort = new AbortController()
  154. const iterator = control.control(abort.signal)[Symbol.asyncIterator]()
  155. await iterator.next()
  156. const waiting = iterator.next()
  157. await Promise.resolve()
  158. abort.abort()
  159. inbox.append('next-turn', message('late'))
  160. await expect(waiting).resolves.toMatchObject({ done: true })
  161. })
  162. it('ends active streams on context disposal after flushing buffered frames', async () => {
  163. const { ctx, control, inbox } = await harness()
  164. const iterator = control.control(new AbortController().signal)[Symbol.asyncIterator]()
  165. await iterator.next()
  166. const first = message('first')
  167. const second = message('second')
  168. inbox.append('next-turn', first)
  169. inbox.append('next-turn', second)
  170. const queues: Extract<SessionControlFrame, { type: 'queue' }>[] = []
  171. await ctx.fiber.dispose()
  172. for (;;) {
  173. const next = await iterator.next()
  174. if (next.done) break
  175. if (next.value.type === 'queue') queues.push(next.value)
  176. }
  177. expect(queues.map(queue => queue.items.map(item => item.id))).toEqual([[first.id], [first.id, second.id]])
  178. })
  179. })