inbox.spec.ts 5.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139
  1. import { describe, expect, it } from 'vitest'
  2. import { AgentMessageId } from '@deepseek-ai/dsh-agent'
  3. import { Inbox, agentMessage } from '../src/inbox.ts'
  4. function message(text: string) {
  5. return { id: AgentMessageId(text), content: [{ type: 'text' as const, text }], source: { kind: 'user' as const }, contexts: [], wakeup: true }
  6. }
  7. describe('agentMessage', () => {
  8. it('returns a frozen payload so a listener cannot mutate it for later listeners', () => {
  9. const payload = agentMessage(message('m'), false)
  10. expect(Object.isFrozen(payload)).toBe(true)
  11. expect(() => { (payload as { id: string }).id = 'mutated' }).toThrow()
  12. expect(payload.id).toBe(AgentMessageId('m'))
  13. })
  14. })
  15. function resolverPair() {
  16. let r!: () => void
  17. const p = new Promise<void>((resolve) => { r = resolve })
  18. return { promise: p, resolve: r }
  19. }
  20. describe('Inbox', () => {
  21. it('dequeues one queued message at a time in FIFO order', () => {
  22. const inbox = new Inbox()
  23. inbox.enqueue(message('first'))
  24. inbox.enqueue(message('second'))
  25. expect(inbox.hasQueued).toBe(true)
  26. expect(inbox.dequeueQueued()?.content[0]).toMatchObject({ text: 'first' })
  27. expect(inbox.hasQueued).toBe(true)
  28. expect(inbox.dequeueQueued()?.content[0]).toMatchObject({ text: 'second' })
  29. expect(inbox.hasQueued).toBe(false)
  30. expect(inbox.dequeueQueued()).toBeUndefined()
  31. })
  32. it('enqueue(msg, false) queues without waking a parked waiter', async () => {
  33. const inbox = new Inbox()
  34. let woke = false
  35. const waiter = inbox.waitForQueued(new Promise(() => {})).then(() => { woke = true })
  36. inbox.enqueue(message('quiet'), false)
  37. // The item is queued, but the parked waiter was not resolved by it.
  38. expect(inbox.hasQueued).toBe(true)
  39. await Promise.resolve()
  40. expect(woke).toBe(false)
  41. // A later waking enqueue resolves the same waiter.
  42. inbox.enqueue(message('loud'))
  43. await waiter
  44. expect(woke).toBe(true)
  45. })
  46. it('pending() snapshots queued then steering without removing them', () => {
  47. const inbox = new Inbox()
  48. inbox.enqueue(message('q'))
  49. inbox.steer(message('s'))
  50. const pending = inbox.pending()
  51. expect(pending.map(p => p.steering)).toEqual([false, true])
  52. // Snapshot does not drain the FIFOs.
  53. expect(inbox.hasQueued).toBe(true)
  54. expect(inbox.hasSteering).toBe(true)
  55. })
  56. it('pushes and drains steering messages separately from queued', () => {
  57. const inbox = new Inbox()
  58. inbox.steer(message('steer'))
  59. expect(inbox.hasQueued).toBe(false)
  60. expect(inbox.hasSteering).toBe(true)
  61. const steering = inbox.drainSteering()
  62. expect(steering).toHaveLength(1)
  63. expect(inbox.hasSteering).toBe(false)
  64. })
  65. it('waitForQueued returns immediately when a queued message is already present', async () => {
  66. const inbox = new Inbox()
  67. inbox.enqueue(message('ready'))
  68. const started = Date.now()
  69. await inbox.waitForQueued(new Promise(() => {})) // never-resolving cancel
  70. expect(Date.now() - started).toBeLessThan(50)
  71. })
  72. it('waitForQueued resolves when a message is enqueued', async () => {
  73. const inbox = new Inbox()
  74. const waiter = inbox.waitForQueued(new Promise(() => {})) // never-resolving cancel
  75. // enqueue after starting the wait
  76. setTimeout(() => { inbox.enqueue(message('wake')) }, 5)
  77. await waiter
  78. })
  79. it('waitForQueued resolves when the cancel promise resolves', async () => {
  80. const inbox = new Inbox()
  81. const { promise, resolve } = resolverPair()
  82. const waiter = inbox.waitForQueued(promise)
  83. resolve()
  84. await waiter
  85. })
  86. it('waitForQueued overwrites the previous wakeup callback (only the latest waiter is notified)', async () => {
  87. const inbox = new Inbox()
  88. const { promise: p1, resolve: r1 } = resolverPair()
  89. void inbox.waitForQueued(new Promise(() => {})) // first call, never resolved
  90. void inbox.waitForQueued(p1) // second call overwrites wakeup
  91. // Cancelling the latest waiter clears the shared callback; enqueue must neither
  92. // wake the stale waiter nor fail on the cleared callback.
  93. r1()
  94. await p1
  95. inbox.enqueue(message('hey'))
  96. })
  97. it('clears wakeup in finally handler when enqueue resolves', async () => {
  98. const inbox = new Inbox()
  99. void inbox.waitForQueued(new Promise(() => {})) // never-resolving cancel
  100. // The wakeup is set. Now trigger it via enqueue → wakeup() calls resolve,
  101. // promise resolves, finally clears wakeup because wakeup === resolve.
  102. inbox.enqueue(message('wake'))
  103. // No explicit await needed — enqueue is synchronous, and the microtask
  104. // (finally) runs. The key coverage hit is finally with wakeup === resolve.
  105. })
  106. it('finally handler does not clear wakeup when a different waiter overwrote it', async () => {
  107. // A stale waiter's finally must not clear the replacement waiter.
  108. const inbox = new Inbox()
  109. const { promise: c1, resolve: r1 } = resolverPair()
  110. void inbox.waitForQueued(c1) // wakeup = resolve1, c1.then(resolve1)
  111. void inbox.waitForQueued(new Promise(() => {})) // wakeup = resolve2, cancel never resolves
  112. r1()
  113. await c1
  114. // The replacement remains registered and is resolved by enqueue.
  115. inbox.enqueue(message('hey'))
  116. })
  117. })