1
/**2
* Driver-owned durable agent inbox projection and command facade.3
*4
* @module @deepseek-ai/dsh-agent-loop/inbox5
*/7
import type { MessageId } from '@deepseek-ai/dsh-llm'8
import type { ProjectionDefinition } from '@deepseek-ai/dsh-session-projection'9
import type SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection'10
import type { Session, SessionEventMap, UserMessage } from '@deepseek-ai/dsh-session'11
import type {12
AgentEventDispatch,13
Inbox as InboxContract,14
InboxState,15
InboxTarget,16
InboxWireState,17
} from '@deepseek-ai/dsh-agent'18
import { z } from 'zod'20
/** Wire validation for pending agent input reconstructed from durable inbox splices. */21
export const inboxProjectionSchema = z.object({22
'next-turn': z.array(z.custom<UserMessage>()).readonly(),23
'next-step': z.array(z.custom<UserMessage>()).readonly(),24
}).readonly()26
/** Standard fold that reconstructs pending input and rejects invalid durable splice history. */27
export 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 state33
const splice = event.data34
try {35
const inbox = state[splice.target]36
const removedCount = splice.removedCount ?? 037
if (!Number.isSafeInteger(splice.start) || splice.start < 0 || splice.start > inbox.length38
|| !Number.isSafeInteger(removedCount) || removedCount < 039
|| 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 already59
// round-trips the session log as lossless JSON. Only the static type60
// 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>67
/**68
* Driver-owned durable Inbox implementation used by ReactLoopAgent and focused69
* 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
*/74
export class ReactLoopInbox implements InboxContract {75
constructor(76
private readonly projections: SessionProjectionRegistry,77
private readonly session: Session,78
private readonly dispatch: AgentEventDispatch,79
) {}81
/** Prompts awaiting individual turns. */82
get nextTurn(): readonly UserMessage[] {83
return this.current()['next-turn']84
}86
/** Input awaiting the next step boundary. */87
get nextStep(): readonly UserMessage[] {88
return this.current()['next-step']89
}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 > 095
}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
}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 claimed114
}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
}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
}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 false143
this.splice(location.target, location.index, 1, [newMessage])144
return true145
}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 false155
this.splice(location.target, location.index, 1, [])156
return true157
}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
}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 undefined184
}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 state195
}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 : truncatedStart209
const actualStart = offset < 0210
? 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 : undefined227
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 removed243
}244
}