refactor agent pre-step inbox lifecycle
This commit is contained in:
@@ -4,6 +4,7 @@
|
||||
* @module @deepseek-ai/dsh-agent/inbox
|
||||
*/
|
||||
|
||||
import type { MessageId } from '@deepseek-ai/dsh-llm'
|
||||
import type { Session, SessionEventMap, UserMessage } from '@deepseek-ai/dsh-session'
|
||||
|
||||
/** One of the two ordered pending-message lists owned by an agent. */
|
||||
@@ -12,11 +13,22 @@ export type InboxTarget = 'next-turn' | 'next-step'
|
||||
/** Mutable state privately owned by an {@link Inbox}. */
|
||||
type InboxState = Record<InboxTarget, UserMessage[]>
|
||||
|
||||
/** Live notifications committed by inbox mutations. */
|
||||
export interface InboxNotifications {
|
||||
/** Publish one inserted message. */
|
||||
inserted(message: UserMessage): void
|
||||
/** Publish one discarded message. */
|
||||
discarded(message: UserMessage): void
|
||||
}
|
||||
|
||||
/** A replay-once projection that incrementally consumes later inbox splices. */
|
||||
export class Inbox {
|
||||
private readonly state: InboxState = { 'next-turn': [], 'next-step': [] }
|
||||
|
||||
constructor(private readonly session: Session) {
|
||||
constructor(
|
||||
private readonly session: Session,
|
||||
private readonly notifications: InboxNotifications,
|
||||
) {
|
||||
for (const event of session.events.slice(session.header.seedLength ?? 0)) {
|
||||
if (event.type !== 'agent/inbox/spliced') continue
|
||||
try {
|
||||
@@ -32,7 +44,7 @@ export class Inbox {
|
||||
return this.state['next-turn']
|
||||
}
|
||||
|
||||
/** Input awaiting admission at a step boundary. */
|
||||
/** Input awaiting the next step boundary. */
|
||||
get nextStep(): readonly UserMessage[] {
|
||||
return this.state['next-step']
|
||||
}
|
||||
@@ -42,6 +54,68 @@ export class Inbox {
|
||||
return this.nextTurn.length > 0 || this.nextStep.length > 0
|
||||
}
|
||||
|
||||
/**
|
||||
* Remove and return the complete batch proposed for one step. The durable
|
||||
* splices are pure deletions; the caller publishes claimed notifications.
|
||||
* @param target - whether this boundary also consumes one queued turn.
|
||||
* @returns next-step input followed by the queued turn, when requested.
|
||||
*/
|
||||
claim(target: InboxTarget): UserMessage[] {
|
||||
const claimed = this.mutate('next-step', 0, this.nextStep.length, [], false)
|
||||
if (target === 'next-turn') {
|
||||
claimed.push(...this.mutate('next-turn', 0, 1, [], false))
|
||||
}
|
||||
return claimed
|
||||
}
|
||||
|
||||
/**
|
||||
* Append one message to a pending list and durably record the insertion.
|
||||
* @param target - pending list to extend.
|
||||
* @param message - message to append.
|
||||
* @throws if the message identity is already pending.
|
||||
*/
|
||||
append(target: InboxTarget, message: UserMessage): void {
|
||||
this.splice(target, this.state[target].length, 0, [message])
|
||||
}
|
||||
|
||||
/**
|
||||
* Prepend one message to a pending list and durably record the insertion.
|
||||
* @param target - pending list to extend.
|
||||
* @param message - message to prepend.
|
||||
* @throws if the message identity is already pending.
|
||||
*/
|
||||
prepend(target: InboxTarget, message: UserMessage): void {
|
||||
this.splice(target, 0, 0, [message])
|
||||
}
|
||||
|
||||
/**
|
||||
* Replace one pending message in place and durably record the mutation.
|
||||
* @param target - pending list containing the message.
|
||||
* @param messageId - identity of the message to replace.
|
||||
* @param newMessage - replacement message.
|
||||
* @returns whether the message was still pending.
|
||||
* @throws if the replacement duplicates another pending message identity.
|
||||
*/
|
||||
update(target: InboxTarget, messageId: MessageId, newMessage: UserMessage): boolean {
|
||||
const index = this.state[target].findIndex(message => message.id === messageId)
|
||||
if (index < 0) return false
|
||||
this.splice(target, index, 1, [newMessage])
|
||||
return true
|
||||
}
|
||||
|
||||
/**
|
||||
* Remove one pending message and durably record its cancellation.
|
||||
* @param target - pending list containing the message.
|
||||
* @param messageId - identity of the message to remove.
|
||||
* @returns whether the message was still pending.
|
||||
*/
|
||||
remove(target: InboxTarget, messageId: MessageId): boolean {
|
||||
const index = this.state[target].findIndex(message => message.id === messageId)
|
||||
if (index < 0) return false
|
||||
this.splice(target, index, 1, [])
|
||||
return true
|
||||
}
|
||||
|
||||
/**
|
||||
* Apply standard splice semantics and durably record the normalized result.
|
||||
* The durable event commits before the live projection mutates, so synchronous
|
||||
@@ -51,7 +125,6 @@ export class Inbox {
|
||||
* @param start - splice position.
|
||||
* @param deleteCount - maximum number of messages to remove.
|
||||
* @param inserted - messages to insert at the resolved position.
|
||||
* @param outcome - terminal disposition of removed messages.
|
||||
* @returns messages removed by the splice.
|
||||
*/
|
||||
splice(
|
||||
@@ -59,7 +132,17 @@ export class Inbox {
|
||||
start: number,
|
||||
deleteCount: number,
|
||||
inserted: UserMessage[],
|
||||
outcome?: 'admitted' | 'canceled',
|
||||
): UserMessage[] {
|
||||
return this.mutate(target, start, deleteCount, inserted, true)
|
||||
}
|
||||
|
||||
/** Commit one normalized mutation and publish its live notifications. */
|
||||
private mutate(
|
||||
target: InboxTarget,
|
||||
start: number,
|
||||
deleteCount: number,
|
||||
inserted: UserMessage[],
|
||||
discardRemoved: boolean,
|
||||
): UserMessage[] {
|
||||
const inbox = this.state[target]
|
||||
const truncatedStart = Math.trunc(start)
|
||||
@@ -73,17 +156,22 @@ export class Inbox {
|
||||
inbox.length - actualStart,
|
||||
)
|
||||
if (actualDeleteCount === 0 && inserted.length === 0) return []
|
||||
const resolvedOutcome = outcome ?? (actualDeleteCount > 0 ? 'canceled' : undefined)
|
||||
const outcome = discardRemoved && actualDeleteCount > 0 ? 'canceled' : undefined
|
||||
const splice = {
|
||||
target,
|
||||
start: actualStart,
|
||||
...(actualDeleteCount === 0 ? {} : { removedCount: actualDeleteCount }),
|
||||
inserted,
|
||||
...(resolvedOutcome === undefined ? {} : { outcome: resolvedOutcome }),
|
||||
...(outcome === undefined ? {} : { outcome }),
|
||||
}
|
||||
this.validate(splice)
|
||||
const event = this.session.append('agent/inbox/spliced', splice)
|
||||
return inbox.splice(actualStart, actualDeleteCount, ...event.data.inserted)
|
||||
const removed = inbox.splice(actualStart, actualDeleteCount, ...event.data.inserted)
|
||||
if (discardRemoved) {
|
||||
for (const message of removed) this.notifications.discarded(message)
|
||||
}
|
||||
for (const message of event.data.inserted) this.notifications.inserted(message)
|
||||
return removed
|
||||
}
|
||||
|
||||
/** Apply one normalized durable splice to the projection. */
|
||||
|
||||
@@ -42,21 +42,26 @@ export interface CancelOptions {
|
||||
/**
|
||||
* An agent's lifecycle state, emitted on every transition as `agent/status`:
|
||||
* `idle` means no driver is scheduled or active; `running` begins when a
|
||||
* cancellable admission is scheduled and lasts while the driver drains,
|
||||
* cancellable pre-step processing is scheduled and lasts while the driver drains,
|
||||
* closes, or checkpoints turns. Disposal removes the agent from its registry;
|
||||
* it is not a third observable status.
|
||||
*/
|
||||
export type AgentStatus = 'idle' | 'running'
|
||||
|
||||
/**
|
||||
* Prompt interception result. An allowed batch replaces the submitted
|
||||
* messages; a listener wrapping `next()` preserves that batch unless it
|
||||
* intentionally replaces it. A blocked batch explicitly chooses whether to
|
||||
* discard the claimed messages; unclaimed work remains pending.
|
||||
*/
|
||||
export type PromptDecision =
|
||||
| { kind: 'allow'; messages: UserMessage[] }
|
||||
| { kind: 'block'; reason: string; discardClaimed: boolean }
|
||||
/** Coordinates and cancellation for a proposed step. */
|
||||
export interface PreStepContext {
|
||||
/** Turn that will own the step. */
|
||||
readonly turn: number
|
||||
/** Step proposed by the loop. */
|
||||
readonly step: number
|
||||
/** Current turn cancellation signal. */
|
||||
readonly signal: AbortSignal
|
||||
}
|
||||
|
||||
/** Whether and with which messages the loop enters a proposed step. */
|
||||
export type PreStepDecision =
|
||||
| { kind: 'reject' }
|
||||
| { kind: 'enter'; messages: UserMessage[] }
|
||||
|
||||
/** One failed model-request attempt presented to recovery listeners. */
|
||||
export interface RequestFailureContext {
|
||||
@@ -135,11 +140,11 @@ export interface Agent {
|
||||
steer(message: UserMessage): void
|
||||
|
||||
/**
|
||||
* Append model-facing context without running the model. Admission or an
|
||||
* open turn stages it at the next safe log position; outside that window it
|
||||
* appends immediately without opening a turn. If admission closes without a
|
||||
* turn, a context-only boundary appends immediately; context staged beside
|
||||
* steering remains pending with it.
|
||||
* Queue model-facing context for the next pre-step without waking the
|
||||
* driver. Collecting and running drivers claim it at the nearest later
|
||||
* step boundary; idle drivers leave it pending until follow-up or steering
|
||||
* wakes them. It may miss a request whose pre-step already claimed its
|
||||
* batch. Cancellation or disposal may discard pending context.
|
||||
* @param message - identified injected context and its producer provenance.
|
||||
*/
|
||||
inject(message: UserMessage): void
|
||||
@@ -178,6 +183,30 @@ declare module 'cordis' {
|
||||
* @mode emit
|
||||
*/
|
||||
'agent/status'(this: Scoped<Agent>, agent: Agent, status: AgentStatus): void
|
||||
/**
|
||||
* One message entered the live inbox.
|
||||
* @param agent - the agent whose inbox changed.
|
||||
* @param event - the inserted message.
|
||||
* Scope-filtered dispatch (`@deepseek-ai/dsh-scope`): agent-scoped listeners receive only that agent.
|
||||
* @mode emit
|
||||
*/
|
||||
'agent/inbox/inserted'(this: Scoped<Agent>, agent: Agent, event: { message: UserMessage }): void
|
||||
/**
|
||||
* One message left the inbox for a turn.
|
||||
* @param agent - the agent whose inbox changed.
|
||||
* @param event - the claimed message and owning turn.
|
||||
* Scope-filtered dispatch (`@deepseek-ai/dsh-scope`): agent-scoped listeners receive only that agent.
|
||||
* @mode emit
|
||||
*/
|
||||
'agent/inbox/claimed'(this: Scoped<Agent>, agent: Agent, event: { message: UserMessage; turn: number }): void
|
||||
/**
|
||||
* One message was discarded from the live inbox.
|
||||
* @param agent - the agent whose inbox changed.
|
||||
* @param event - the discarded message.
|
||||
* Scope-filtered dispatch (`@deepseek-ai/dsh-scope`): agent-scoped listeners receive only that agent.
|
||||
* @mode emit
|
||||
*/
|
||||
'agent/inbox/discarded'(this: Scoped<Agent>, agent: Agent, event: { message: UserMessage }): void
|
||||
// ---- session lifecycle (emit) ----
|
||||
/**
|
||||
* The session lifecycle began, once before the first turn. Use
|
||||
@@ -193,30 +222,15 @@ declare module 'cordis' {
|
||||
|
||||
// ---- the machine's extension seams ----
|
||||
/**
|
||||
* Allow, rewrite, or block one claimed inbox batch before it becomes
|
||||
* model-visible or opens a turn. Call `next()` for the unchanged default. The
|
||||
* signal controls only this admission attempt; listeners may cooperate with
|
||||
* it but must not retain it for a later attempt or turn.
|
||||
* @param agent - the agent whose driver claimed the batch.
|
||||
* @param messages - the claimed messages.
|
||||
* @param signal - the current turn's explicit abort signal.
|
||||
* Reject a proposed step or replace the messages that enter it. Calling
|
||||
* `next()` preserves the current messages.
|
||||
* @param agent - the agent proposing the step.
|
||||
* @param messages - messages removed from the inbox for this step.
|
||||
* @param context - proposed turn and step coordinates plus cancellation.
|
||||
* Scope-filtered dispatch (`@deepseek-ai/dsh-scope`): agent-scoped listeners receive only that agent.
|
||||
* @mode waterfall
|
||||
*/
|
||||
'agent/prompt-submit'(this: Scoped<Agent>, agent: Agent, messages: UserMessage[], signal: AbortSignal, next: () => Promise<PromptDecision>): Promise<PromptDecision>
|
||||
/**
|
||||
* Awaited serial checkpoint before EVERY request of a turn is built (the
|
||||
* first as well as each post-tools continuation). The single "between
|
||||
* steps" extension point: inject context, steer, or edit the session log
|
||||
* here — the request's history derives from the log right after this settles.
|
||||
* @param agent - the agent about to send a request.
|
||||
* @param turn - the open turn number.
|
||||
* @param step - the step number about to open.
|
||||
* @param signal - the turn abort signal.
|
||||
* Scope-filtered dispatch (`@deepseek-ai/dsh-scope`): agent-scoped listeners receive only that agent.
|
||||
* @mode serial
|
||||
*/
|
||||
'agent/step'(this: Scoped<Agent>, agent: Agent, turn: number, step: number, signal: AbortSignal): Promise<void> | void
|
||||
'agent/pre-step'(this: Scoped<Agent>, agent: Agent, messages: UserMessage[], context: PreStepContext, next: () => Promise<PreStepDecision>): Promise<PreStepDecision>
|
||||
/**
|
||||
* Replace the frozen call configuration. `await next()` yields the config
|
||||
* the machine would use (agent options on the first request, the logged
|
||||
@@ -284,7 +298,7 @@ declare module '@deepseek-ai/dsh-session' {
|
||||
start: number
|
||||
removedCount?: number
|
||||
inserted: UserMessage[]
|
||||
outcome?: 'admitted' | 'canceled'
|
||||
outcome?: 'canceled'
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user