返回源码地图

packages/experimental/agent-team/src/mailbox.ts

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

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

1/** Durable Team mailbox admission, target-local dispatch, acknowledgement, and recovery. */
2
3import { randomUUID } from 'node:crypto'
4import type { Context } from '@deepseek-ai/cordis'
5import { brandString } from '@deepseek-ai/dsh-brand'
6import type { Agent } from '@deepseek-ai/dsh-agent'
7import { createUserMessage } from '@deepseek-ai/dsh-llm'
8import type { ContentBlock } from '@deepseek-ai/dsh-llm'
9import { SessionId } from '@deepseek-ai/dsh-session'
10import type { Session, SessionEvent } from '@deepseek-ai/dsh-session'
11import { steerHostSubagentPrompt } from '@deepseek-ai/dsh-subagent/internal'
12import { errorMessage, TeamError } from './error.ts'
13import type { TeamJournal } from './journal.ts'
14import type { TeamRuntimeLifecycle } from './lifecycle.ts'
15import { readPersistedSession } from './persisted.ts'
16import type { TeamRoster } from './roster.ts'
17import { resolveActiveMember } from './roster.ts'
18import { messageAccepted } from './session-message.ts'
19import { TeamId, TeamMessageId } from './types.ts'
20import type {
21 SendTeamMessageRequest,
22 SendTeamMessageResult,
23 TeamMessageSnapshot,
24} from './types.ts'
25
26/** Owns every process-local state transition for the durable Team mailbox. */
27export class TeamMailbox {
28 private readonly dispatchTails = new Map<SessionId, Promise<void>>()
29 private readonly inFlightMessages = new Set<TeamMessageId>()
30 private readonly inFlightDispatches = new Set<Promise<unknown>>()
31
32 /**
33 * @param ctx - Team service context with Agent, Session, persistence, and subagent services.
34 * @param journal - authoritative Lead-log transaction owner.
35 * @param roster - Team membership and member-name resolver.
36 * @param lifecycle - shared Team runtime admission cutoff.
37 * @param maxPendingMessagesPerMember - per-target queued-minus-delivered limit.
38 * @param maxMessageBytes - maximum complete sender-framed delivery size.
39 */
40 constructor(
41 private readonly ctx: Context,
42 private readonly journal: TeamJournal,
43 private readonly roster: TeamRoster,
44 private readonly lifecycle: TeamRuntimeLifecycle,
45 private readonly maxPendingMessagesPerMember: number,
46 private readonly maxMessageBytes: number,
47 ) {}
48
49 /**
50 * Queue one durable peer message, then attempt immediate delivery.
51 * @param caller - exact live sending Team member.
52 * @param request - target name, content, and pre-queue cancellation.
53 * @returns durable message identity and immediate-delivery observation.
54 */
55 async send(caller: Agent, request: SendTeamMessageRequest): Promise<SendTeamMessageResult> {
56 if (this.lifecycle.disposed) throw new TeamError('Agent Teams service is disposing', 'TEAM_DISPOSED')
57 const operation = this.sendAdmitted(caller, {
58 ...request,
59 signal: AbortSignal.any([request.signal, this.lifecycle.signal]),
60 })
61 return await this.trackDispatch(operation)
62 }
63
64 /**
65 * Observe target-side durable receipts and checkpoint their Lead-log acknowledgement.
66 * @param session - exact target Session receiving the event.
67 * @param event - newly appended Session event.
68 */
69 observeSessionEvent(session: Session, event: SessionEvent): void {
70 if (this.lifecycle.disposed || event.type !== 'user/message' || event.data.source.kind !== 'team-message') return
71 const source = event.data.source
72 const acknowledgement = Promise.resolve().then(async () => {
73 const root = this.ctx.agents.get(brandString<SessionId>(source.teamId))
74 if (root !== undefined) await this.checkpointDelivered(root, session, source.messageId)
75 }).catch((error: unknown) => {
76 this.ctx.logger.warn(`Team message "${source.messageId}" acknowledgement failed: ${errorMessage(error)}`)
77 })
78 void this.trackDispatch(acknowledgement)
79 }
80
81 /**
82 * Retry durable pending messages relevant to one started Team member.
83 * @param agent - newly started exact live Agent.
84 * @param signal - shared runtime cancellation.
85 */
86 async recoverFor(agent: Agent, signal: AbortSignal): Promise<void> {
87 signal.throwIfAborted()
88 const membership = this.roster.tryMembership(agent)
89 if (membership === undefined) return
90 const state = this.journal.state(membership.root)
91 const messages = state.messages.filter(message =>
92 !state.delivered.includes(message.id)
93 && (membership.role === 'lead' || message.targetId === agent.id))
94 for (const message of messages) {
95 signal.throwIfAborted()
96 await this.tryDispatch(membership.root, message, signal)
97 }
98 }
99
100 /**
101 * Return admitted dispatch and acknowledgement operations captured for disposal.
102 * @returns detached snapshot ordered only by Set insertion.
103 */
104 pendingDispatches(): readonly Promise<unknown>[] {
105 return [...this.inFlightDispatches]
106 }
107
108 /** Queue and dispatch one mailbox item admitted before the disposal cutoff. */
109 private async sendAdmitted(
110 caller: Agent,
111 request: SendTeamMessageRequest,
112 ): Promise<SendTeamMessageResult> {
113 const membership = this.roster.membership(caller)
114 request.signal.throwIfAborted()
115 const root = membership.root
116 const content = structuredClone(request.content)
117 const queued = await this.journal.transact(root.id, async () => {
118 request.signal.throwIfAborted()
119 const state = this.journal.state(root)
120 const target = resolveActiveMember(root, state, request.target)
121 if (target.id === caller.id) throw new TeamError('a Team member cannot message itself', 'TEAM_SELF_MESSAGE')
122 const pendingForTarget = state.messages.filter(candidate =>
123 candidate.targetId === target.id && !state.delivered.includes(candidate.id)).length
124 if (pendingForTarget >= this.maxPendingMessagesPerMember) {
125 throw new TeamError(
126 `teammate "${target.name}" has ${pendingForTarget} pending messages`,
127 'TEAM_MAILBOX_FULL',
128 )
129 }
130 const queued: TeamMessageSnapshot = {
131 id: TeamMessageId(`team-message-${randomUUID()}`),
132 senderId: caller.id,
133 senderName: membership.name,
134 targetId: target.id,
135 content,
136 }
137 if (Buffer.byteLength(JSON.stringify(this.deliveryContent(queued)), 'utf8') > this.maxMessageBytes) {
138 throw new TeamError(`team message exceeds ${this.maxMessageBytes} bytes`, 'TEAM_MESSAGE_TOO_LARGE')
139 }
140 await this.journal.appendAndFlush(root, 'team/message/queued', {
141 version: 2,
142 teamId: TeamId(root.id),
143 message: queued,
144 })
145 // Register dispatch before releasing the root transaction so concurrent
146 // senders enter the target-local queue in durable mailbox order.
147 return { message: queued, dispatch: this.tryDispatch(root, queued, request.signal) }
148 })
149 const accepted = await queued.dispatch
150 return { messageId: queued.message.id, status: accepted ? 'accepted' : 'queued' }
151 }
152
153 /** Attempt one queued message exactly once in this process at a time. */
154 private tryDispatch(root: Agent, message: TeamMessageSnapshot, signal: AbortSignal): Promise<boolean> {
155 if (this.lifecycle.disposed) return Promise.resolve(false)
156 if (this.inFlightMessages.has(message.id)) return Promise.resolve(false)
157 this.inFlightMessages.add(message.id)
158 const operation = this.trackDispatch(
159 this.tryDispatchAdmitted(
160 root,
161 message,
162 AbortSignal.any([signal, this.lifecycle.signal]),
163 ),
164 )
165 const forget = (): void => {
166 this.inFlightMessages.delete(message.id)
167 }
168 void operation.then(forget, forget)
169 return operation
170 }
171
172 /** Track one dispatch transaction through delivery admission or contained failure. */
173 private trackDispatch<T>(operation: Promise<T>): Promise<T> {
174 this.inFlightDispatches.add(operation)
175 void operation.then(() => {
176 this.inFlightDispatches.delete(operation)
177 }, () => {
178 this.inFlightDispatches.delete(operation)
179 })
180 return operation
181 }
182
183 /** Attempt one queued message admitted before the service lifecycle cutoff. */
184 private async tryDispatchAdmitted(
185 root: Agent,
186 message: TeamMessageSnapshot,
187 signal: AbortSignal,
188 ): Promise<boolean> {
189 return await this.serializeDispatch(message, () => this.dispatchThrough(root, message, signal))
190 }
191
192 /** Serialize delivery admission for one durable target in queued order. */
193 private async serializeDispatch(
194 message: TeamMessageSnapshot,
195 operation: () => Promise<boolean>,
196 ): Promise<boolean> {
197 const targetId = message.targetId
198 const prior = this.dispatchTails.get(targetId) ?? Promise.resolve()
199 /* v8 ignore next -- dispatch tails absorb rejection, so the recovery callback is a fail-safe backstop. */
200 const run = prior.then(operation, operation)
201 /* v8 ignore next -- dispatchOnce contains delivery failures and serializeDispatch itself does not throw. */
202 const tail = run.then(() => undefined, () => undefined)
203 this.dispatchTails.set(targetId, tail)
204 try {
205 return await run
206 } finally {
207 if (this.dispatchTails.get(targetId) === tail) this.dispatchTails.delete(targetId)
208 }
209 }
210
211 /** Deliver every pending target message through `message` in durable queue order. */
212 private async dispatchThrough(
213 root: Agent,
214 message: TeamMessageSnapshot,
215 signal: AbortSignal,
216 ): Promise<boolean> {
217 const state = this.journal.state(root)
218 const pending = state.messages.filter(candidate =>
219 candidate.targetId === message.targetId && !state.delivered.includes(candidate.id))
220 const requested = pending.findIndex(candidate => candidate.id === message.id)
221 if (requested < 0) return state.delivered.includes(message.id)
222 for (const candidate of pending.slice(0, requested + 1)) {
223 const ownsInFlight = !this.inFlightMessages.has(candidate.id)
224 if (ownsInFlight) this.inFlightMessages.add(candidate.id)
225 try {
226 if (!await this.dispatchOnce(root, candidate, signal)) return false
227 } finally {
228 if (ownsInFlight) this.inFlightMessages.delete(candidate.id)
229 }
230 }
231 return true
232 }
233
234 /** Attempt one queued delivery after target-local ordering admits it. */
235 private async dispatchOnce(root: Agent, message: TeamMessageSnapshot, signal: AbortSignal): Promise<boolean> {
236 try {
237 const target = message.targetId === root.id ? root : this.ctx.agents.get(message.targetId)
238 if (target !== undefined && this.targetRecorded(target.session, message.id)) {
239 return await this.checkpointDelivered(root, target.session, message.id)
240 }
241 const source = {
242 kind: 'team-message' as const,
243 teamId: TeamId(root.id),
244 messageId: message.id,
245 senderId: message.senderId,
246 senderName: message.senderName,
247 }
248 const content = this.deliveryContent(message)
249 if (message.targetId === root.id) {
250 const input = createUserMessage({ content, source })
251 root.steer(input)
252 return await this.checkpointDelivered(root, root.session, message.id)
253 }
254 if (target === undefined) {
255 const recorded = await this.persistedTargetRecorded(message.targetId, message.id, signal)
256 if (recorded === undefined) return false
257 if (recorded) {
258 await this.markDelivered(root, message.id, message.targetId)
259 return true
260 }
261 }
262 await steerHostSubagentPrompt(this.ctx.subagents, root, message.targetId, content, source, signal)
263 return target === undefined
264 ? true
265 : await this.checkpointDelivered(root, target.session, message.id)
266 } catch (error: unknown) {
267 this.ctx.logger.warn(`team message "${message.id}" remains queued: ${errorMessage(error)}`)
268 return false
269 }
270 }
271
272 /** Flush one live target receipt before the Lead records its delivered edge. */
273 private async checkpointDelivered(
274 root: Agent,
275 target: Session,
276 messageId: TeamMessageId,
277 ): Promise<boolean> {
278 await this.ctx.sessions.flush(target)
279 if (!this.targetRecorded(target, messageId)) return false
280 await this.markDelivered(root, messageId, target.id)
281 return true
282 }
283
284 /** Record delivery unless the acknowledgement already exists. */
285 private async markDelivered(root: Agent, messageId: TeamMessageId, targetId: SessionId): Promise<void> {
286 await this.journal.transact(root.id, async () => {
287 const state = this.journal.state(root)
288 if (state.delivered.includes(messageId)) return
289 const queued = state.messages.find(message => message.id === messageId)
290 if (queued === undefined || queued.targetId !== targetId) return
291 await this.journal.appendAndFlush(root, 'team/message/delivered', {
292 version: 2,
293 teamId: TeamId(root.id),
294 messageId,
295 targetId,
296 })
297 })
298 }
299
300 /** Whether a target Session already contains the durable message identity. */
301 private targetRecorded(session: Session, messageId: TeamMessageId): boolean {
302 // oxlint-disable-next-line typescript/no-deprecated -- Existing Session history read; migration deferred.
303 const suffix = session.snapshotEvents(session.inheritedEventCount)
304 return messageAccepted(suffix, message => message.source.kind === 'team-message'
305 && message.source.messageId === messageId)
306 }
307
308 /** Frame peer content with stable sender and message identity for the receiving model. */
309 private deliveryContent(message: TeamMessageSnapshot): ContentBlock[] {
310 return [
311 { type: 'text', text: `Team message ${message.id} from ${message.senderName}:` },
312 ...structuredClone(message.content),
313 ]
314 }
315
316 /** Read an inactive target's durable log before cold resume; uncertainty keeps the mailbox queued. */
317 private async persistedTargetRecorded(
318 targetId: SessionId,
319 messageId: TeamMessageId,
320 signal: AbortSignal,
321 ): Promise<boolean | undefined> {
322 try {
323 const stored = await readPersistedSession(this.ctx.sessionPersistence, targetId, signal)
324 const suffix = stored.events.slice(stored.inheritedEventCount)
325 return messageAccepted(suffix, message => message.source.kind === 'team-message'
326 && message.source.messageId === messageId)
327 } catch (error: unknown) {
328 this.ctx.logger.warn(`cannot read Team message target "${targetId}": ${errorMessage(error)}`)
329 return undefined
330 }
331 }
332}