inbox.ts 9.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244
  1. /**
  2. * Driver-owned durable agent inbox projection and command facade.
  3. *
  4. * @module @deepseek-ai/dsh-agent-loop/inbox
  5. */
  6. import type { MessageId } from '@deepseek-ai/dsh-llm'
  7. import type { ProjectionDefinition } from '@deepseek-ai/dsh-session-projection'
  8. import type SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection'
  9. import type { Session, SessionEventMap, UserMessage } from '@deepseek-ai/dsh-session'
  10. import type {
  11. AgentEventDispatch,
  12. Inbox as InboxContract,
  13. InboxState,
  14. InboxTarget,
  15. InboxWireState,
  16. } from '@deepseek-ai/dsh-agent'
  17. import { z } from 'zod'
  18. /** Wire validation for pending agent input reconstructed from durable inbox splices. */
  19. export const inboxProjectionSchema = z.object({
  20. 'next-turn': z.array(z.custom<UserMessage>()).readonly(),
  21. 'next-step': z.array(z.custom<UserMessage>()).readonly(),
  22. }).readonly()
  23. /** Standard fold that reconstructs pending input and rejects invalid durable splice history. */
  24. export const inboxProjectionDefinition = {
  25. key: 'inbox',
  26. stateSchema: inboxProjectionSchema,
  27. init: (): InboxState => ({ 'next-turn': [], 'next-step': [] }),
  28. apply(state: InboxState, event) {
  29. if (event.type !== 'agent/inbox/spliced') return state
  30. const splice = event.data
  31. try {
  32. const inbox = state[splice.target]
  33. const removedCount = splice.removedCount ?? 0
  34. if (!Number.isSafeInteger(splice.start) || splice.start < 0 || splice.start > inbox.length
  35. || !Number.isSafeInteger(removedCount) || removedCount < 0
  36. || splice.start + removedCount > inbox.length) {
  37. throw new Error('invalid inbox splice')
  38. }
  39. const next = inbox.toSpliced(splice.start, removedCount, ...splice.inserted)
  40. const ids = new Set<string>()
  41. for (const message of splice.target === 'next-turn'
  42. ? [...next, ...state['next-step']]
  43. : [...state['next-turn'], ...next]) {
  44. if (ids.has(message.id)) throw new Error(`message "${message.id}" is already pending`)
  45. ids.add(message.id)
  46. }
  47. return splice.target === 'next-turn'
  48. ? { 'next-turn': next, 'next-step': state['next-step'] }
  49. : { 'next-turn': state['next-turn'], 'next-step': next }
  50. } catch (error: unknown) {
  51. throw new Error(`invalid persisted inbox splice at session seq ${event.seq}`, { cause: error })
  52. }
  53. },
  54. wire: {
  55. // The wire value is the fold state itself: every pending message already
  56. // round-trips the session log as lossless JSON. Only the static type
  57. // narrows to the JSON-safe projection table entry.
  58. viewSchema: inboxProjectionSchema as unknown as z.ZodType<InboxWireState>,
  59. view: (state: InboxState) => state as unknown as InboxWireState,
  60. },
  61. stateVersion: 1,
  62. } satisfies ProjectionDefinition<'inbox', InboxState>
  63. /**
  64. * Driver-owned durable Inbox implementation used by ReactLoopAgent and focused
  65. * provider tests.
  66. * @param projections - registry with the standard Inbox projection registered by AgentLoop.
  67. * @param session - session whose durable events store pending input.
  68. * @param dispatch - agent-scoped notifications for Inbox lifecycle events.
  69. */
  70. export class ReactLoopInbox implements InboxContract {
  71. constructor(
  72. private readonly projections: SessionProjectionRegistry,
  73. private readonly session: Session,
  74. private readonly dispatch: AgentEventDispatch,
  75. ) {}
  76. /** Prompts awaiting individual turns. */
  77. get nextTurn(): readonly UserMessage[] {
  78. return this.current()['next-turn']
  79. }
  80. /** Input awaiting the next step boundary. */
  81. get nextStep(): readonly UserMessage[] {
  82. return this.current()['next-step']
  83. }
  84. /** Whether either pending-message list contains work. */
  85. get hasPending(): boolean {
  86. const state = this.current()
  87. return state['next-turn'].length > 0 || state['next-step'].length > 0
  88. }
  89. /** Durably cancel all pending input, clearing next-step before next-turn. */
  90. clear(): void {
  91. this.splice('next-step', 0, this.nextStep.length, [])
  92. this.splice('next-turn', 0, this.nextTurn.length, [])
  93. }
  94. /**
  95. * Remove and return the complete batch proposed for one step.
  96. * @param target - whether this boundary also consumes one queued turn.
  97. * @param turn - turn that will own the claimed batch.
  98. * @returns next-step input followed by the queued turn, when requested.
  99. */
  100. claim(target: InboxTarget, turn: number): UserMessage[] {
  101. const claimed = this.mutate('next-step', 0, this.nextStep.length, [], false)
  102. if (target === 'next-turn') claimed.push(...this.mutate('next-turn', 0, 1, [], false))
  103. for (const message of claimed) this.dispatch.emit('agent/inbox/claimed', { message, turn })
  104. return claimed
  105. }
  106. /**
  107. * Append one message to a pending list.
  108. * @param target - pending list to extend.
  109. * @param message - message to append.
  110. */
  111. append(target: InboxTarget, message: UserMessage): void {
  112. this.splice(target, this.current()[target].length, 0, [message])
  113. }
  114. /**
  115. * Prepend one message to a pending list.
  116. * @param target - pending list to extend.
  117. * @param message - message to prepend.
  118. */
  119. prepend(target: InboxTarget, message: UserMessage): void {
  120. this.splice(target, 0, 0, [message])
  121. }
  122. /**
  123. * Replace one pending message in place.
  124. * @param messageId - identity of the pending message to replace.
  125. * @param newMessage - replacement message.
  126. * @returns whether the message was still pending.
  127. */
  128. replace(messageId: MessageId, newMessage: UserMessage): boolean {
  129. const location = this.locate(messageId)
  130. if (location === undefined) return false
  131. this.splice(location.target, location.index, 1, [newMessage])
  132. return true
  133. }
  134. /**
  135. * Remove one pending message.
  136. * @param messageId - identity of the pending message to remove.
  137. * @returns whether the message was still pending.
  138. */
  139. remove(messageId: MessageId): boolean {
  140. const location = this.locate(messageId)
  141. if (location === undefined) return false
  142. this.splice(location.target, location.index, 1, [])
  143. return true
  144. }
  145. /**
  146. * Apply standard splice semantics and durably record the normalized result.
  147. * @param target - pending list to mutate.
  148. * @param start - splice position.
  149. * @param deleteCount - maximum number of messages to remove.
  150. * @param inserted - messages to insert at the resolved position.
  151. * @returns messages removed by the splice.
  152. */
  153. splice(
  154. target: InboxTarget,
  155. start: number,
  156. deleteCount: number,
  157. inserted: UserMessage[],
  158. ): UserMessage[] {
  159. return this.mutate(target, start, deleteCount, inserted, true)
  160. }
  161. /** Locate one pending identity across both owned lists. */
  162. private locate(messageId: MessageId): { target: InboxTarget; index: number } | undefined {
  163. const state = this.current()
  164. for (const target of ['next-turn', 'next-step'] as const) {
  165. const index = state[target].findIndex(message => message.id === messageId)
  166. if (index >= 0) return { target, index }
  167. }
  168. return undefined
  169. }
  170. /** Read the current durable projection state. */
  171. private current(): InboxState {
  172. const state = this.projections.stateOf(this.session, 'inbox')
  173. if (state === undefined) {
  174. throw new Error(
  175. `agent "${this.session.id}" cannot read inbox state: its projection registration is not active`,
  176. )
  177. }
  178. return state
  179. }
  180. /** Commit one normalized mutation and publish its live events. */
  181. private mutate(
  182. target: InboxTarget,
  183. start: number,
  184. deleteCount: number,
  185. inserted: UserMessage[],
  186. discardRemoved: boolean,
  187. ): UserMessage[] {
  188. const state = this.current()
  189. const inbox = state[target]
  190. const truncatedStart = Math.trunc(start)
  191. const offset = Number.isNaN(truncatedStart) ? 0 : truncatedStart
  192. const actualStart = offset < 0
  193. ? Math.max(inbox.length + offset, 0)
  194. : Math.min(offset, inbox.length)
  195. const truncatedDeleteCount = Math.trunc(deleteCount)
  196. const actualDeleteCount = Math.min(
  197. Math.max(Number.isNaN(truncatedDeleteCount) ? 0 : truncatedDeleteCount, 0),
  198. inbox.length - actualStart,
  199. )
  200. if (actualDeleteCount === 0 && inserted.length === 0) return []
  201. const candidate = inbox.toSpliced(actualStart, actualDeleteCount, ...inserted)
  202. const ids = new Set<string>()
  203. for (const message of target === 'next-turn'
  204. ? [...candidate, ...state['next-step']]
  205. : [...state['next-turn'], ...candidate]) {
  206. if (ids.has(message.id)) throw new Error(`message "${message.id}" is already pending`)
  207. ids.add(message.id)
  208. }
  209. const outcome = discardRemoved && actualDeleteCount > 0 ? 'canceled' as const : undefined
  210. const splice: SessionEventMap['agent/inbox/spliced'] = {
  211. target,
  212. start: actualStart,
  213. ...(actualDeleteCount === 0 ? {} : { removedCount: actualDeleteCount }),
  214. inserted,
  215. ...(outcome === undefined ? {} : { outcome }),
  216. }
  217. const removed = inbox.slice(actualStart, actualStart + actualDeleteCount)
  218. const event = this.session.append('agent/inbox/spliced', splice)
  219. if (discardRemoved) {
  220. for (const message of removed) this.dispatch.emit('agent/inbox/discarded', { message })
  221. }
  222. for (const message of event.data.inserted) {
  223. this.dispatch.emit('agent/inbox/inserted', { message })
  224. }
  225. return removed
  226. }
  227. }