inbox.ts 3.2 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697
  1. /**
  2. * Per-agent message inbox: queued and steering FIFOs. Purely an in-memory
  3. * mechanism of the loop driver — the public surface is `Agent.send()` and
  4. * `Agent.steer()`.
  5. *
  6. * @module dsh-agent-loop/inbox
  7. */
  8. import type { ContentBlock, MessageSource } from '@deepseek-ai/dsh-llm'
  9. /** One message waiting in an agent's inbox. */
  10. export interface InboxMessage {
  11. content: ContentBlock[]
  12. source: MessageSource
  13. }
  14. /**
  15. * Per-agent inbox: a queued FIFO (dequeued once per turn start) and a steering FIFO
  16. * (drained between steps of a running turn). Purely an in-memory mechanism of
  17. * the loop — the public surface is `Agent.send()` / `Agent.steer()`.
  18. */
  19. export class Inbox {
  20. private queuedMessages: InboxMessage[] = []
  21. private steeringMessages: InboxMessage[] = []
  22. private wakeup: (() => void) | undefined
  23. /** True while queued messages are pending — read by the idle wait's fast path and the loop's turn-start checks. */
  24. get hasQueued(): boolean {
  25. return this.queuedMessages.length > 0
  26. }
  27. /** True while steering messages are pending — read by `cancel()`'s arm gate and the loop's stop-override check. */
  28. get hasSteering(): boolean {
  29. return this.steeringMessages.length > 0
  30. }
  31. /**
  32. * Add a message to the queued FIFO and wake a parked {@link waitForQueued}.
  33. * @param message - the message to queue for the next turn start.
  34. */
  35. enqueue(message: InboxMessage): void {
  36. this.queuedMessages.push(message)
  37. this.wakeup?.()
  38. }
  39. /**
  40. * Add a message to the steering FIFO. Deliberately no wakeup: steering is
  41. * drained between steps of a running turn, never by the idle wait —
  42. * `Agent.steer()` on an idle agent falls back to `send()` instead.
  43. * @param message - the message to inject between steps of the running turn.
  44. */
  45. steer(message: InboxMessage): void {
  46. this.steeringMessages.push(message)
  47. }
  48. /**
  49. * Remove the oldest queued message for one turn start.
  50. * @returns the oldest message, or `undefined` when the queued FIFO is empty.
  51. */
  52. dequeueQueued(): InboxMessage | undefined {
  53. return this.queuedMessages.shift()
  54. }
  55. /**
  56. * Drain all steering messages (between steps).
  57. * @returns the drained messages in arrival order; the steering FIFO is left empty.
  58. */
  59. drainSteering(): InboxMessage[] {
  60. return this.steeringMessages.splice(0)
  61. }
  62. /**
  63. * Discard all pending messages (queued + steering) without delivering them —
  64. * used by `cancel()`, which drops un-started work rather than draining it into
  65. * a turn. Unlike `dequeueQueued`/`drainSteering`, the messages are thrown away.
  66. */
  67. clear(): void {
  68. this.queuedMessages.length = 0
  69. this.steeringMessages.length = 0
  70. }
  71. /**
  72. * Wait until a queued message arrives or `cancel` resolves.
  73. * @param cancel - a promise whose resolution abandons the wait without a
  74. * message (the driver loop passes the agent's disposed promise so a parked
  75. * loop can exit).
  76. */
  77. waitForQueued(cancel: Promise<void>): Promise<void> {
  78. if (this.hasQueued) return Promise.resolve()
  79. const { promise, resolve } = Promise.withResolvers<void>()
  80. this.wakeup = resolve
  81. void cancel.then(resolve)
  82. return promise.finally(() => {
  83. if (this.wakeup === resolve) this.wakeup = undefined
  84. })
  85. }
  86. }