diff --git a/packages/core/agent-loop/src/agent.ts b/packages/core/agent-loop/src/agent.ts index 3498e69d9a..56e95e8d6a 100644 --- a/packages/core/agent-loop/src/agent.ts +++ b/packages/core/agent-loop/src/agent.ts @@ -1,9 +1,9 @@ /** * The concrete Agent, in the naive-agent shape: the agent IS the machine. * Two inboxes — `queued` (prompts, one turn each) and `outbox` (steering + - * injected context, taken whole at every step boundary). `kick()` admits and - * records one queued prompt; `start()` then steps until the model owes no - * response. + * injected or admitted input, taken at turn start and every step boundary). + * `kick()` claims one queued prompt and resolves admission before `run()` owns + * the turn boundaries, step loop, settlement, and idle handoff. * * The session log IS the transcript: every take appends, every step re-derives * (`session.deriveMessages()`), so editing history between steps is naturally @@ -38,7 +38,7 @@ import { import type { ContentBlock, GenerateOptions, LlmCallConfig, LlmFailure, Message, MessageSource, } from '@deepseek-ai/dsh-llm' -import { canonicalHeader, headerEquals, snapshotJsonValue } from '@deepseek-ai/dsh-session' +import { canonicalHeader, headerEquals } from '@deepseek-ai/dsh-session' import type { PromptMessageData, Session, SessionId, TurnEndReason, TurnTrigger } from '@deepseek-ai/dsh-session' import { renderPrompt } from '@deepseek-ai/dsh-system-prompt' import type {} from '@deepseek-ai/dsh-tools' @@ -53,19 +53,13 @@ interface QueuedMessage { wakeup: boolean } -/** Input awaiting the next step boundary. */ -type OutboxItem = - | ({ kind: 'steering' } & QueuedMessage) - | { kind: 'context'; context: HookContext } - -/** Mutable settlement facts shared by one turn's intake and step loop. */ -interface TurnState { - turn: number - step: number - reason: TurnEndReason +/** One model-facing input awaiting the next step boundary. */ +interface OutboxItem { + data: PromptMessageData + steering?: QueuedMessage } -/** Build one live inbox event payload from an accepted message. */ +/** Build one live inbox event payload from a queued message. */ function inboxMessage(message: QueuedMessage, steering: boolean): AgentMessage { return { id: message.id, @@ -139,7 +133,7 @@ function withoutToolCalls(message: Message): Message { /** * The concrete {@link Agent}: the classic naive agent loop — whole derived * history in, one assistant message out, loop until a reply owes no tool call. - * One `run()` drains the work queue, one turn per unit. + * One `run()` owns one complete turn. */ export class ReactLoopAgent extends Agent { /** Prompts awaiting a turn of their own: one dequeued per turn, FIFO. */ @@ -149,9 +143,9 @@ export class ReactLoopAgent extends Agent { /** Whether observers see one running drain interval; queued turns share it. */ private busy = false - /** The active turn's abort owner; rotated per turn, aborted by {@link cancel}. */ + /** The claimed activity's abort owner, spanning prompt admission and its turn. */ private turnAbort: AbortController | undefined - /** Resolves when the current `run()` has fully exited (quiescence for waiters and teardown). */ + /** Resolves when the current admission-plus-turn activity fully exits. */ done: Promise = Promise.resolve() /** The agent-scoped registration boundary; the lifecycle owner unwinds it after {@link done}. */ @@ -187,15 +181,6 @@ export class ReactLoopAgent extends Agent { return this.busy ? 'running' : 'idle' } - /** Detach and freeze one public payload; rejects non-lossless-JSON input synchronously. */ - private accept(value: T): T { - const accepted = snapshotJsonValue(value) - if (accepted === undefined) { - throw new TypeError('agent message content, source, and contexts must be losslessly JSON-serializable') - } - return deepFreeze(accepted) - } - // ------------------------------------------------------------------------- // Public driving verbs. // ------------------------------------------------------------------------- @@ -206,49 +191,49 @@ export class ReactLoopAgent extends Agent { const target = options.target ?? 'next-turn' const wakeup = options.wakeup ?? true if (target === 'next-step' && !wakeup) { - this.injectContext(content, options) + const context: HookContext = { + content, + source: options.source ?? { kind: 'plugin', plugin: '' }, + } + if (this.turnAbort !== undefined) { + this.outbox.push({ data: { content: context.content, source: context.source } }) + return id + } + + const turn = ++this.lastTurn + let opened = false + try { + this.session.append('turn/start', { turn, trigger: { kind: 'injection', source: context.source } }) + opened = true + this.session.append('user/message', context, { surfaceOp: 'append' }) + } finally { + if (opened) this.session.append('turn/end', { turn, reason: { kind: 'completed' } }) + } return id } const steering = target === 'next-step' && this.turnAbort !== undefined - const accepted = this.accept({ + const message: QueuedMessage = { id, content, source: options.source ?? { kind: 'user' }, contexts: options.contexts ?? [], wakeup, - }) - if (steering) this.outbox.push({ kind: 'steering', ...accepted }) - else this.queued.push(accepted) - emitAgentEvent(this.loopCtx, this, 'agent/inbox/enqueue', inboxMessage(accepted, steering)) + } + if (steering) { + const prepared = preparePromptMessage(message.content, message.source, message.contexts) + this.outbox.push({ data: prepared.data, steering: message }) + for (const context of prepared.separateContexts) { + this.outbox.push({ data: { content: context.content, source: context.source } }) + } + } else { + this.queued.push(message) + } + emitAgentEvent(this.loopCtx, this, 'agent/inbox/enqueue', inboxMessage(message, steering)) if (!steering && wakeup) this.kick() return id } - /** Stage non-waking context at the next step boundary, or in an idle one-shot turn. */ - private injectContext(content: ContentBlock[], options: SendOptions): void { - const context = this.accept({ - content, - source: options.source ?? { kind: 'plugin', plugin: '' }, - }) - if (this.turnAbort !== undefined) { - this.outbox.push({ kind: 'context', context }) - return - } - // Idle: wrap the injection in a one-shot turn so every event stays - // turn-enclosed (the durability/replay boundary is the turn). - const turn = ++this.lastTurn - let opened = false - try { - this.session.append('turn/start', { turn, trigger: { kind: 'injection', source: context.source } }) - opened = true - this.session.append('user/message', context, { surfaceOp: 'append' }) - } finally { - // Close only a turn whose start committed; a pre-commit veto escapes. - if (opened) this.session.append('turn/end', { turn, reason: { kind: 'completed' } }) - } - } - /** * Clear all pending work and abort the active turn; the first cause wins. * The cause is signal payload for observers and the durable turn/end @@ -263,17 +248,17 @@ export class ReactLoopAgent extends Agent { if (cause.kind !== 'disposed') emitAgentEvent(this.loopCtx, this, 'agent/cancel-requested', cause) } if (!options.keepInbox) { - const steering = this.outbox.flatMap(item => item.kind === 'steering' ? [item] : []) const discarded = [ ...this.queued.map(message => inboxMessage(message, false)), - ...steering.map(message => inboxMessage(message, true)), + ...this.outbox.flatMap(item => item.steering === undefined ? [] : [inboxMessage(item.steering, true)]), ] // Clear before abort observers run: replacement work belongs to the next turn. this.queued.length = 0 this.outbox.length = 0 if (discarded.length > 0) emitAgentEvent(this.loopCtx, this, 'agent/inbox/discard', discarded) } - this.turnAbort?.abort(Object.freeze({ kind: cause.kind })) + const reason = Object.freeze({ kind: cause.kind }) + this.turnAbort?.abort(reason) } /** @@ -284,27 +269,36 @@ export class ReactLoopAgent extends Agent { */ retry(): void { if (this.turnAbort !== undefined) throw new Error(`agent "${this.id}" cannot retry while busy`) - this.launch({ kind: 'retry' }, (state, signal) => this.start(state, signal)) + this.done = this.loopCtx.agents.withInitiator(this, () => this.run({ kind: 'retry' })) } /** Resolve at idle quiescence: no run driving and no waking prompt waiting. */ async whenIdle(): Promise { // `done` is replaced per run, so re-reading it each lap follows chained // turns; a run failure still counts as quiescence for the waiter. - while (this.turnAbort !== undefined || this.queued.some(message => message.wakeup)) await this.done.catch(() => undefined) + while (this.turnAbort !== undefined || this.queued.some(message => message.wakeup)) { + await this.done.catch(() => undefined) + } } // ------------------------------------------------------------------------- // The machine. // ------------------------------------------------------------------------- - /** Claim, admit, and record the next queued prompt before starting its step loop. */ + /** Claim and admit the next queued prompt, then start its turn. */ private kick(): void { if (this.turnAbort !== undefined || !this.queued.some(message => message.wakeup)) return const message = this.queued.shift() - if (message !== undefined) { - emitAgentEvent(this.loopCtx, this, 'agent/inbox/dequeue', inboxMessage(message, false)) - this.launch({ kind: 'message', source: message.source }, async (state, signal) => { + if (message === undefined) return + + emitAgentEvent(this.loopCtx, this, 'agent/inbox/dequeue', inboxMessage(message, false)) + const controller = new AbortController() + this.turnAbort = controller + this.done = this.loopCtx.agents.withInitiator(this, async () => { + const signal = controller.signal + const trigger: TurnTrigger = { kind: 'message', source: message.source } + try { + signal.throwIfAborted() const decision = await this.loopCtx.waterfall( agentCarrier(this), 'agent/prompt-submit', this, message.content, message.source, signal, () => Promise.resolve({ @@ -315,94 +309,123 @@ export class ReactLoopAgent extends Agent { signal.throwIfAborted() if (decision.kind === 'block') { - this.session.append('prompt/blocked', { - content: message.content, - source: message.source, - reason: decision.reason, - }) - state.reason = { kind: 'rejected', reason: decision.reason } + this.rejectPrompt(trigger, message, decision.reason, controller) return + } else { + const prepared = preparePromptMessage( + decision.content ?? message.content, + message.source, + decision.additionalContexts ?? [], + ) + this.outbox.push({ data: prepared.data }) + for (const context of prepared.separateContexts) { + this.outbox.push({ data: { content: context.content, source: context.source } }) + } } + } catch (error: unknown) { + this.failAdmission(trigger, error, controller) + return + } + if (this.turnAbort === controller) this.turnAbort = undefined + await this.run(trigger) + }) + } - const prepared = preparePromptMessage( - decision.content ?? message.content, - message.source, - decision.additionalContexts ?? [], - ) - this.session.append('user/message', prepared.data, { surfaceOp: 'append' }) - for (const context of prepared.separateContexts) { - this.outbox.push({ kind: 'context', context: this.accept(context) }) - } - await this.start(state, signal) - }, true) + /** Record a policy-blocked prompt as a zero-step turn. */ + private rejectPrompt( + trigger: TurnTrigger, + message: QueuedMessage, + rejection: string, + controller: AbortController, + ): void { + if (!this.busy) { + this.busy = true + emitAgentEvent(this.loopCtx, this, 'agent/status', 'running') + } + const signal = controller.signal + const turn = ++this.lastTurn + let reason: TurnEndReason = { kind: 'rejected', reason: rejection } + let idle: IdleReason = { kind: 'completed' } + + try { + signal.throwIfAborted() + this.session.append('turn/start', { turn, trigger }) + this.turnOpen = true + signal.throwIfAborted() + this.drainOutbox(turn) + this.session.append('prompt/blocked', { + content: message.content, + source: message.source, + reason: rejection, + }) + } catch (error: unknown) { + ({ reason, idle } = this.settle(turn, 0, error, signal)) + } finally { + this.finishTurn(controller, turn, 0, reason, idle) } } - /** Own one turn from its durable opening through settlement and idle handoff. */ - private launch( - trigger: TurnTrigger, - work: (state: TurnState, signal: AbortSignal) => Promise, - deferOpen = false, - ): void { + /** Settle a prompt-admission failure without entering the step loop. */ + private failAdmission(trigger: TurnTrigger, failure: unknown, controller: AbortController): void { + if (!this.busy) { + this.busy = true + emitAgentEvent(this.loopCtx, this, 'agent/status', 'running') + } + const signal = controller.signal + const turn = ++this.lastTurn + let reason: TurnEndReason = { kind: 'completed' } + let idle: IdleReason = { kind: 'completed' } + + try { + signal.throwIfAborted() + this.session.append('turn/start', { turn, trigger }) + this.turnOpen = true + signal.throwIfAborted() + this.drainOutbox(turn) + throw failure + } catch (error: unknown) { + ({ reason, idle } = this.settle(turn, 0, error, signal)) + } finally { + this.finishTurn(controller, turn, 0, reason, idle) + } + } + + /** Own one complete turn over input already admitted by {@link kick}, or retry history as-is. */ + private async run(trigger: TurnTrigger): Promise { + if (this.turnAbort !== undefined) throw new Error(`agent "${this.id}" is already running`) const controller = new AbortController() this.turnAbort = controller if (!this.busy) { this.busy = true emitAgentEvent(this.loopCtx, this, 'agent/status', 'running') } - // The whole run inherits this agent as its process-local initiator so - // tools, the llm service, and nested factories can attribute their work. - this.done = this.loopCtx.agents.withInitiator(this, async () => { - const signal = controller.signal - const state: TurnState = { - turn: ++this.lastTurn, - step: 0, - reason: { kind: 'completed' }, - } - let idle: IdleReason = { kind: 'completed' } + const signal = controller.signal + const turn = ++this.lastTurn + let step = 0 + let reason: TurnEndReason = { kind: 'completed' } + let idle: IdleReason = { kind: 'completed' } - try { - // A queued claim keeps its established pre-turn cancellation window: - // send() returns before the durable turn opens, while retry starts now. - if (deferOpen) await Promise.resolve() - signal.throwIfAborted() - this.session.append('turn/start', { turn: state.turn, trigger }) - this.turnOpen = true - signal.throwIfAborted() - await work(state, signal) - } catch (error: unknown) { - ({ reason: state.reason, idle } = this.settle(state.turn, state.step, error, signal)) - } finally { - if (this.turnAbort === controller) this.turnAbort = undefined - try { - this.closeTurn(state.turn, state.step, state.reason) - } catch (error: unknown) { - // A rejected boundary append (a pre-commit validation veto) must not - // kill the machine or strand its running interval: report and move on — the - // idle tail below still runs and the next turn still opens. - const err = toError(error) - this.loopCtx.logger.warn(`agent "${this.id}": closing turn ${state.turn} failed: ${errorChain(err)}`) - emitAgentEvent(this.loopCtx, this, 'agent/error', state.turn, state.step, err) - } - this.idle(state.turn, idle) - } - }) - } - - /** Run the naive step loop after retry or admitted prompt intake has prepared the turn. */ - private async start(state: TurnState, signal: AbortSignal): Promise { - while (true) { - state.step += 1 - const { owes, maxTokens } = await this.step(state.turn, state.step, signal) - if (maxTokens) state.reason = { kind: 'max-tokens' } - // The naive rule, data-driven: run another step while the model is - // owed a response. On a would-stop boundary, `agent/stopping` gives - // listeners one chance to object — by steering, not by voting — and - // the outbox is re-read: data decides, so listener order cannot. - if (owes || this.outbox.some(item => item.kind === 'steering')) continue - await this.loopCtx.serial(agentCarrier(this), 'agent/stopping', this, state.turn, signal) + try { signal.throwIfAborted() - if (!this.outbox.some(item => item.kind === 'steering')) break + this.session.append('turn/start', { turn, trigger }) + this.turnOpen = true + signal.throwIfAborted() + + this.drainOutbox(turn) + + while (true) { + step += 1 + const { owes, maxTokens } = await this.step(turn, step, signal) + if (maxTokens) reason = { kind: 'max-tokens' } + if (owes || this.outbox.some(item => item.steering !== undefined)) continue + await this.loopCtx.serial(agentCarrier(this), 'agent/stopping', this, turn, signal) + signal.throwIfAborted() + if (!this.drainOutbox(turn)) break + } + } catch (error: unknown) { + ({ reason, idle } = this.settle(turn, step, error, signal)) + } finally { + this.finishTurn(controller, turn, step, reason, idle) } } @@ -492,7 +515,7 @@ export class ReactLoopAgent extends Agent { if (toolCalls.length > 0) { ({ concluded } = await executeToolCalls( this.loopCtx, turn, step, toolCalls, signal, - context => this.outbox.push({ kind: 'context', context: this.accept(context) }), + context => this.outbox.push({ data: { content: context.content, source: context.source } }), )) } @@ -565,26 +588,18 @@ export class ReactLoopAgent extends Agent { })) } - /** Take the outbox whole into the log: committed from here. Returns whether steering was taken. */ + /** Commit the outbox whole and report whether it contained steering. */ private drainOutbox(turn: number): boolean { let steered = false for (const item of this.outbox.splice(0)) { - if (item.kind === 'context') { - const { content, source } = item.context - this.session.append('user/message', { content, source }, { surfaceOp: 'append' }) + const message = item.steering + if (message === undefined) { + this.session.append('user/message', item.data, { surfaceOp: 'append' }) continue } steered = true - emitAgentEvent(this.loopCtx, this, 'agent/inbox/dequeue', inboxMessage(item, true)) - const prepared = preparePromptMessage(item.content, item.source, item.contexts) - this.session.append('steering/message', { - turn, - ...prepared.data, - }, { surfaceOp: 'append' }) - for (const context of prepared.separateContexts) { - const { content, source } = context - this.session.append('user/message', { content, source }, { surfaceOp: 'append' }) - } + emitAgentEvent(this.loopCtx, this, 'agent/inbox/dequeue', inboxMessage(message, true)) + this.session.append('steering/message', { turn, ...item.data }, { surfaceOp: 'append' }) } return steered } @@ -629,6 +644,26 @@ export class ReactLoopAgent extends Agent { } } + /** Close one claimed turn and hand the machine back to the idle boundary. */ + private finishTurn( + controller: AbortController, + turn: number, + step: number, + reason: TurnEndReason, + idle: IdleReason, + ): void { + try { + this.closeTurn(turn, step, reason) + } catch (error: unknown) { + // A rejected boundary append must not strand the running interval. + const err = toError(error) + this.loopCtx.logger.warn(`agent "${this.id}": closing turn ${turn} failed: ${errorChain(err)}`) + emitAgentEvent(this.loopCtx, this, 'agent/error', turn, step, err) + } + if (this.turnAbort === controller) this.turnAbort = undefined + this.idle(turn, idle) + } + /** * The turn boundary's tail (naive `idle()`): no turn owner remains, * the idle report fires (a listener may synchronously `retry()` or `send()` @@ -640,11 +675,11 @@ export class ReactLoopAgent extends Agent { // Requeue BEFORE the idle report so earlier-arrived steering keeps its // FIFO position ahead of anything a listener send()s synchronously. for (const item of this.outbox.splice(0)) { - if (item.kind === 'context') { + const message = item.steering + if (message === undefined) { this.outbox.push(item) continue } - const { kind: _kind, ...message } = item this.queued.push(message) } emitAgentEvent(this.loopCtx, this, 'agent/idle', turn, idle)