返回源码地图

packages/core/agent-loop/src/inbox.ts

main snapshot · da00f7f5358f · 正文引用章节 07;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

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