1
/** Durable Team mailbox admission, target-local dispatch, acknowledgement, and recovery. */3
import { randomUUID } from 'node:crypto'4
import type { Context } from '@deepseek-ai/cordis'5
import { brandString } from '@deepseek-ai/dsh-brand'6
import type { Agent } from '@deepseek-ai/dsh-agent'7
import { createUserMessage } from '@deepseek-ai/dsh-llm'8
import type { ContentBlock } from '@deepseek-ai/dsh-llm'9
import { SessionId } from '@deepseek-ai/dsh-session'10
import type { Session, SessionEvent } from '@deepseek-ai/dsh-session'11
import { steerHostSubagentPrompt } from '@deepseek-ai/dsh-subagent/internal'12
import { errorMessage, TeamError } from './error.ts'13
import type { TeamJournal } from './journal.ts'14
import type { TeamRuntimeLifecycle } from './lifecycle.ts'15
import { readPersistedSession } from './persisted.ts'16
import type { TeamRoster } from './roster.ts'17
import { resolveActiveMember } from './roster.ts'18
import { messageAccepted } from './session-message.ts'19
import { TeamId, TeamMessageId } from './types.ts'20
import type {21
SendTeamMessageRequest,22
SendTeamMessageResult,23
TeamMessageSnapshot,24
} from './types.ts'26
/** Owns every process-local state transition for the durable Team mailbox. */27
export 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>>()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
) {}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
}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') return71
const source = event.data.source72
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
}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) return90
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
}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
}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.root116
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)).length124
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 concurrent146
// 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.dispatch150
return { messageId: queued.message.id, status: accepted ? 'accepted' : 'queued' }151
}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 operation170
}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 operation181
}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
}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.targetId198
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 run206
} finally {207
if (this.dispatchTails.get(targetId) === tail) this.dispatchTails.delete(targetId)208
}209
}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 false227
} finally {228
if (ownsInFlight) this.inFlightMessages.delete(candidate.id)229
}230
}231
return true232
}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 false257
if (recorded) {258
await this.markDelivered(root, message.id, message.targetId)259
return true260
}261
}262
await steerHostSubagentPrompt(this.ctx.subagents, root, message.targetId, content, source, signal)263
return target === undefined264
? true265
: 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 false269
}270
}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 false280
await this.markDelivered(root, messageId, target.id)281
return true282
}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)) return289
const queued = state.messages.find(message => message.id === messageId)290
if (queued === undefined || queued.targetId !== targetId) return291
await this.journal.appendAndFlush(root, 'team/message/delivered', {292
version: 2,293
teamId: TeamId(root.id),294
messageId,295
targetId,296
})297
})298
}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
}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
}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 undefined330
}331
}332
}