2026-07-13 23:27:00 +08:00
/**
* Drives one agent across queued durable turns. Turn failures are contained so
* later work can run; the session log, not this driver, owns conversation state.
* See docs/rfc/implemented/architecture/2026-06-18-agent-lifecycle-and-ownership-seams.md.
* @module dsh-agent-loop/loop
*/
2026-06-11 12:39:27 +08:00
2026-06-11 10:54:31 +08:00
import type { Context } from 'cordis'
2026-07-15 13:45:19 +08:00
import type { ContentBlock , FinishReason , GenerateOptions , LlmCallConfig , Message } from '@deepseek-ai/dsh-llm'
2026-07-14 21:57:52 +08:00
import { isDeepStrictEqual } from 'node:util'
2026-07-19 15:02:11 +08:00
import { BlockAssembler , HarnessError , assertNever , deepFreeze , isLlmAdapterFailure } from '@deepseek-ai/dsh-llm'
2026-07-09 01:17:47 +08:00
import { agentEvents , assembleContextFor } from '@deepseek-ai/dsh-agent'
2026-07-15 16:03:52 +08:00
import type { AgentEventDispatch , ContinuationDecision , HookContext , PromptDecision , RequestError , RequestErrorDecision } from '@deepseek-ai/dsh-agent'
2026-07-06 03:07:34 +08:00
import { canonicalHeader } from '@deepseek-ai/dsh-session'
2026-06-11 12:18:52 +08:00
import type { Session , TurnEndReason , TurnTrigger } from '@deepseek-ai/dsh-session'
2026-07-06 03:07:34 +08:00
import { createTransmissionLog , recordRequestHeader } from './request-log.ts'
import type { TransmissionLog } from './request-log.ts'
2026-06-11 10:54:31 +08:00
import { renderPrompt } from '@deepseek-ai/dsh-system-prompt'
2026-06-26 13:51:01 +08:00
import type { PromptAssembly } from '@deepseek-ai/dsh-system-prompt'
2026-06-11 10:54:31 +08:00
import type { } from '@deepseek-ai/dsh-tools'
2026-07-16 14:36:16 +08:00
import { executeToolCalls } from './tool-calls.ts'
2026-07-11 22:55:26 +08:00
import type { Inbox } from './inbox.ts'
2026-06-11 10:54:31 +08:00
2026-07-13 16:24:32 +08:00
/** Normalize thrown values while preserving an existing error code. */
2026-07-15 16:03:52 +08:00
function toError ( error : unknown ) : RequestError {
2026-06-14 01:07:28 +08:00
return error instanceof Error ? error : new HarnessError ( String ( error ) , 'UNKNOWN' , { cause : error } )
2026-06-11 14:02:47 +08:00
}
2026-07-19 16:14:58 +08:00
/** Distinguishes final model-request failures from failures in later step processing. */
2026-07-15 16:03:52 +08:00
class TerminalModelRequestFailure extends Error {
constructor ( readonly requestError : RequestError ) {
super ( requestError . message , { cause : requestError } )
this . name = 'TerminalModelRequestFailure'
}
}
2026-07-13 16:24:32 +08:00
/** Convert terminal failure finishes into step errors; unknown extensible finishes remain successful. */
2026-07-15 16:03:52 +08:00
function finishError ( finish : FinishReason ) : RequestError | undefined {
2026-06-13 00:28:29 +08:00
switch ( finish . kind ) {
case 'error' : {
2026-07-15 16:03:52 +08:00
const error : RequestError = new Error ( finish . message )
2026-06-13 00:28:29 +08:00
if ( finish . code !== undefined ) error . code = finish . code
return error
}
case 'aborted' : {
2026-07-15 16:03:52 +08:00
const error : RequestError = new Error ( 'model stream aborted' )
2026-06-13 00:28:29 +08:00
error . code = 'ABORTED'
return error
}
// stop / tool-calls / max-tokens / plugin-added kinds → not a failure.
default :
return undefined
}
}
2026-06-11 14:02:47 +08:00
/**
* Build the `{ message, code? }` part of an error payload, omitting the
* `code` key entirely when absent (exactOptionalPropertyTypes-correct).
*/
2026-07-15 16:03:52 +08:00
function errorData ( err : RequestError ) : { message : string ; code? : string } {
2026-06-11 14:02:47 +08:00
return { message : err.message , . . . typeof err . code === 'string' ? { code : err.code } : { } }
}
2026-07-13 16:24:32 +08:00
/** Map a successful max-token finish onto the turn reason; other successful finishes add nothing. */
2026-06-15 23:53:47 +08:00
function stepFinishReason ( finish : FinishReason ) : TurnEndReason | undefined {
switch ( finish . kind ) {
case 'max-tokens' :
return { kind : 'max-tokens' }
// stop / tool-calls / plugin-added kinds → no turn-end contribution
// beyond the default `completed`. FinishReason is merge-extensible, so a
// default (not assertNever) handles unknown kinds as ordinary success.
default :
return undefined
}
}
2026-07-13 16:24:32 +08:00
/** Mutable agent controls supplied to the loop driver. */
2026-06-11 10:54:31 +08:00
export interface LoopHandle {
2026-07-11 22:55:26 +08:00
/** Native-private agent inbox handed to the driver only at internal startup. */
readonly inbox : Inbox
2026-07-18 14:59:26 +08:00
/** Maximum parallel-safe calls allowed in one step. */
2026-07-16 14:36:16 +08:00
readonly maxParallelToolCalls : number
2026-06-11 10:54:31 +08:00
setStatus ( status : 'idle' | 'running' ) : void
setAbort ( controller : AbortController | undefined ) : void
/** Resolves when the agent is disposed — unblocks the idle wait. */
disposed : Promise < void >
isDisposed ( ) : boolean
2026-07-13 16:24:32 +08:00
/** Whether cancellation is pending for the current loop iteration. */
2026-06-20 04:51:32 +08:00
isCancelled ( ) : boolean
2026-07-13 16:24:32 +08:00
/** Resolved pending-cancellation reason; meaningful only while {@link isCancelled} is true. */
2026-06-20 11:51:32 +08:00
cancelReason ( ) : string
2026-06-20 04:51:32 +08:00
/** Clear the cancel marker (called once per iteration after the turn returns). */
clearCancel ( ) : void
2026-07-13 23:27:00 +08:00
/** Settle idle waiters when pre-running cancellation skips a turn, without emitting `agent/status`. */
2026-06-20 04:51:32 +08:00
settleIdle ( ) : void
2026-07-15 12:49:50 +08:00
/** Run an active tool-call batch, accepting post-tool context into the FIFO drained before settlement. */
readonly withToolBatch : < T > ( run : ( acceptContext : ( context : HookContext ) = > void ) = > Promise < T > ) = > Promise < T >
2026-06-11 10:54:31 +08:00
}
/**
2026-07-13 16:24:32 +08:00
* Drive queued batches as durable turns until disposal. Plugin failures end the
2026-07-19 14:20:27 +08:00
* current turn without terminating the driver. The caller establishes the
2026-07-19 16:28:31 +08:00
* `ctx.agents.withInitiator()` boundary before entry; package-private
* orchestration recovers that exact Agent and captures its Session locally.
2026-07-19 14:06:47 +08:00
* @param ctx - the plugin context the loop reaches its initiating Agent,
* events (agent/…, session/flush), and services (systemPrompt, llm, tools)
* through.
2026-07-06 22:09:30 +08:00
* @param handle - the bridge to the agent's mutable state: status/abort setters plus the disposal and cancel-marker reads.
2026-07-19 14:20:27 +08:00
* @throws when no initiating Agent is active.
2026-06-11 10:54:31 +08:00
*/
2026-07-19 14:06:47 +08:00
export async function runLoop ( ctx : Context , handle : LoopHandle ) : Promise < void > {
const agent = ctx . agents . requireInitiator ( )
2026-07-13 16:24:32 +08:00
// Per-instance prefix and request-header state; conversation history remains in the session log.
2026-07-06 03:07:34 +08:00
const transmission = createTransmissionLog ( )
2026-06-11 10:54:31 +08:00
const { session } = agent
2026-07-13 16:24:32 +08:00
// Fused subject and scope carrier for every agent event below.
2026-07-09 01:17:47 +08:00
const events = agentEvents ( ctx , agent )
2026-06-11 10:54:31 +08:00
while ( ! handle . isDisposed ( ) ) {
2026-07-11 22:55:26 +08:00
await handle . inbox . waitForQueued ( handle . disposed )
2026-06-11 10:54:31 +08:00
if ( handle . isDisposed ( ) ) break
2026-07-13 16:24:32 +08:00
// Cancellation between wake and `running` skips only the cancelled work;
// a replacement prompt still runs and owns the eventual idle transition.
2026-06-20 04:51:32 +08:00
if ( handle . isCancelled ( ) ) {
handle . clearCancel ( )
2026-07-11 22:55:26 +08:00
if ( ! handle . inbox . hasQueued ) {
2026-06-20 05:10:16 +08:00
handle . settleIdle ( )
continue
}
2026-06-20 04:51:32 +08:00
}
2026-06-11 10:54:31 +08:00
handle . setStatus ( 'running' )
2026-06-20 04:51:32 +08:00
2026-07-13 16:24:32 +08:00
// A synchronous `running` listener can cancel before `runTurn`; balance the
// status only when no replacement prompt was queued by that listener.
2026-06-20 04:51:32 +08:00
if ( handle . isCancelled ( ) ) {
handle . clearCancel ( )
2026-07-11 22:55:26 +08:00
if ( ! handle . inbox . hasQueued ) {
2026-06-20 12:57:32 +08:00
handle . setStatus ( 'idle' )
continue
}
2026-06-20 04:51:32 +08:00
}
2026-07-13 16:24:32 +08:00
// Idle injection can add a turn, so derive the next number from the log.
2026-06-15 20:56:17 +08:00
const turn = lastTurnNumber ( session ) + 1
2026-07-11 22:55:26 +08:00
let terminalStopped = false
2026-06-11 12:18:52 +08:00
try {
2026-07-19 14:06:47 +08:00
terminalStopped = await runTurn ( ctx , events , handle , turn , transmission )
2026-06-11 12:39:27 +08:00
} catch ( error : unknown ) {
2026-07-13 16:24:32 +08:00
// Pre-turn failure has no durable boundary to close; report it without appending outside a turn.
2026-06-15 20:56:17 +08:00
const err = toError ( error )
ctx . logger . warn ( ` agent " ${ agent . id } ": turn ${ turn } failed before it started: ${ err . message } ` )
2026-06-11 12:18:52 +08:00
try {
2026-07-09 01:17:47 +08:00
events . emit ( 'agent/error' , turn , 0 , err )
2026-06-15 20:56:17 +08:00
} catch { /* contained: a throwing agent/error listener must not kill the driver */ }
2026-06-11 12:18:52 +08:00
}
2026-06-11 10:54:31 +08:00
2026-07-13 16:24:32 +08:00
// Reset per iteration, including when a prompt arrives during the flush window.
2026-06-20 04:51:32 +08:00
handle . clearCancel ( )
2026-07-13 16:24:32 +08:00
// Late steering becomes queued input unless terminal policy stopped the turn.
2026-07-11 22:55:26 +08:00
for ( const message of handle . inbox . drainSteering ( ) ) {
if ( ! terminalStopped ) handle . inbox . enqueue ( message )
2026-06-11 10:54:31 +08:00
}
2026-07-11 22:55:26 +08:00
if ( ! handle . inbox . hasQueued ) handle . setStatus ( 'idle' )
2026-06-11 12:18:52 +08:00
}
}
2026-06-11 10:54:31 +08:00
2026-07-06 03:07:34 +08:00
async function runTurn (
2026-07-19 14:06:47 +08:00
ctx : Context , events : AgentEventDispatch , handle : LoopHandle , turn : number , transmission : TransmissionLog ,
2026-07-11 22:55:26 +08:00
) : Promise < boolean > {
2026-07-19 14:06:47 +08:00
const agent = ctx . agents . requireInitiator ( )
2026-06-11 12:18:52 +08:00
const { session } = agent
2026-07-19 16:28:31 +08:00
const drainSteering = ( ) : boolean = > {
const messages = handle . inbox . drainSteering ( )
for ( const message of messages ) {
session . append ( 'steering/message' , { turn , content : message.content , source : message.source } , { surfaceOp : 'append' } )
}
return messages . length > 0
}
2026-06-11 10:54:31 +08:00
2026-07-13 16:24:32 +08:00
// Drain before opening the turn, but append only after `turn/start`.
2026-07-11 22:55:26 +08:00
const queued = handle . inbox . drainQueued ( )
2026-06-11 14:17:58 +08:00
const first = queued [ 0 ]
2026-06-11 14:58:36 +08:00
/* v8 ignore next 3 -- invariant guard: runLoop only calls runTurn when hasQueued */
2026-06-11 14:17:58 +08:00
if ( ! first ) throw new Error ( 'runTurn invariant violated: no queued message at turn start' )
const trigger : TurnTrigger = { kind : 'message' , source : first.source }
2026-06-11 10:54:31 +08:00
2026-06-11 12:18:52 +08:00
let reason : TurnEndReason = { kind : 'completed' }
let step = 0
2026-07-15 16:03:52 +08:00
let requestRetryAttempt = 0
2026-06-14 23:53:37 +08:00
let stepOpen = false
let errorReported = false
2026-07-11 22:55:26 +08:00
let terminalStopped = false
2026-06-11 10:54:31 +08:00
2026-07-13 16:24:32 +08:00
// Close the committed step once; pre-commit validation failure still escapes.
2026-07-12 18:57:42 +08:00
const closeStep = ( ) : void = > {
if ( ! stepOpen ) return
session . append ( 'step/end' , { turn , step } )
2026-06-14 23:53:37 +08:00
stepOpen = false
}
2026-06-11 12:18:52 +08:00
2026-07-13 16:24:32 +08:00
// Record the durable turn failure once and contain the live error notification.
2026-07-15 16:03:52 +08:00
const failTurn = ( err : RequestError ) : void = > {
2026-06-14 23:53:37 +08:00
if ( errorReported ) return
errorReported = true
2026-07-02 03:26:45 +08:00
reason = { kind : 'error' , step , . . . errorData ( err ) }
2026-06-14 23:53:37 +08:00
try {
2026-07-09 01:17:47 +08:00
events . emit ( 'agent/error' , turn , step , err )
2026-06-14 23:53:37 +08:00
} catch {
2026-07-02 03:26:45 +08:00
// contained: the error is already captured on `reason`; a throwing
// agent/error listener must not prevent the turn from closing.
2026-06-14 23:53:37 +08:00
}
}
2026-06-11 12:18:52 +08:00
2026-07-13 16:24:32 +08:00
// Pre-commit validation failure escapes rather than masquerading as a committed boundary.
2026-07-02 03:26:45 +08:00
const closeTurn = ( ) : void = > {
2026-07-12 18:57:42 +08:00
session . append ( 'turn/end' , { turn , reason } )
2026-06-14 23:53:37 +08:00
}
2026-06-11 12:18:52 +08:00
2026-06-14 23:53:37 +08:00
try {
// --- Turn boundary. Once turn/start is appended, a turn/end is owed no
2026-07-12 18:57:42 +08:00
// matter what throws below; the catch + closeTurn guarantee it. A pre-commit
// veto leaves no turn/start in the log and therefore owes no turn/end.
2026-06-14 23:53:37 +08:00
session . append ( 'turn/start' , { turn , trigger } )
2026-06-30 17:11:18 +08:00
// Each drained queued message runs the `agent/prompt-submit` waterfall before
// it becomes a `user/message` — a hook can rewrite the prompt or block it.
// Recorded INSIDE the turn (after turn/start) so every event is turn-enclosed;
// turn/end is now owed, so a throwing prompt-submit listener (the waterfall
// throws) is caught below and the turn still closes.
let anyAllowed = false
// Seeded with a floor (only observable if the batch were empty, which
// runTurn never allows — it is called with ≥1 queued message); each `block`
// decision carries a required `reason` and overwrites it, so a fully-blocked
// batch always reports the last vetoing reason.
let lastBlockReason = 'prompt blocked by hook'
2026-06-15 20:56:17 +08:00
for ( const message of queued ) {
2026-07-09 01:17:47 +08:00
const decision = await events . waterfall (
'agent/prompt-submit' , message . content , message . source ,
2026-06-30 17:11:18 +08:00
( ) = > Promise . resolve < PromptDecision > ( { kind : 'allow' } ) ,
)
if ( decision . kind === 'block' ) {
lastBlockReason = decision . reason
2026-07-02 17:03:53 +08:00
// Record the veto durably: `PromptDecision.reason` is the durable record
// of why a prompt was blocked, but a fully-blocked batch's `rejected`
// turn/end only preserves the LAST reason, and a MIXED batch (this prompt
// blocked, another allowed) does not end `rejected` at all — so without
// this append a blocked prompt would vanish from the log whenever any
// sibling prompt is allowed. `prompt/blocked` sits in the open turn in
// place of the `user/message` this prompt would have become.
session . append ( 'prompt/blocked' , { content : message.content , source : message.source , reason : decision.reason } )
2026-06-30 17:11:18 +08:00
continue
}
anyAllowed = true
// `allow.content` REPLACES the prompt bytes (a rewrite); absent keeps them.
const content = decision . content ? ? message . content
session . append ( 'user/message' , { content , source : message.source } , { surfaceOp : 'append' } )
2026-07-13 16:31:03 +08:00
// Every `allow.additionalContexts` entry is a separate context/message the
// next request also sees. The turn is open, so inject() appends each one
// into THIS turn without flattening provenance, framing, or metadata.
for ( const context of decision . additionalContexts ? ? [ ] ) {
agent . inject ( context . content , {
source : context.source ,
. . . context . envelope !== undefined ? { envelope : context.envelope } : { } ,
. . . context . meta !== undefined ? { meta : context.meta } : { } ,
2026-07-10 14:32:44 +08:00
} )
2026-06-30 17:11:18 +08:00
}
2026-06-15 20:56:17 +08:00
}
2026-06-11 10:54:31 +08:00
2026-06-14 23:53:37 +08:00
while ( true ) {
2026-07-13 16:24:32 +08:00
// A fully blocked batch closes its zero-step turn as rejected.
2026-06-30 17:11:18 +08:00
if ( ! anyAllowed ) {
reason = { kind : 'rejected' , reason : lastBlockReason }
break
}
2026-06-14 23:53:37 +08:00
step += 1
2026-06-11 12:18:52 +08:00
2026-07-02 03:26:45 +08:00
// Steering from the previous round's continuation listeners joins before
// the request.
2026-07-19 16:28:31 +08:00
drainSteering ( )
2026-06-14 23:53:37 +08:00
2026-06-29 15:59:52 +08:00
// The step's AbortController exists BEFORE any async pre-step work so a
// dispose() or cancel() — in a synchronous turn-start listener or an
// async listener whose effect fires before we block — always has an armed
// abort to cancel against. isDisposed below covers disposal, which does
// NOT set the cancel marker. Cleared on every exit path below.
const abort = new AbortController ( )
handle . setAbort ( abort )
2026-07-15 16:50:44 +08:00
// Assemble once before pre-step so listener work and the request share one prompt value.
2026-07-09 01:17:47 +08:00
const assembly = await ctx . systemPrompt . assemble ( assembleContextFor ( agent ) )
2026-07-05 01:54:46 +08:00
const fullSystemPrompt = renderPrompt ( assembly )
2026-06-26 13:51:01 +08:00
2026-07-13 16:24:32 +08:00
// Cancellation or disposal during assembly ends the turn before any step opens.
2026-06-29 15:59:52 +08:00
if ( handle . isCancelled ( ) || handle . isDisposed ( ) ) {
2026-06-26 13:51:01 +08:00
handle . setAbort ( undefined )
2026-06-29 15:59:52 +08:00
reason = handle . isDisposed ( ) ? { kind : 'disposed' } : { kind : 'aborted' , reason : handle.cancelReason ( ) }
2026-06-26 13:51:01 +08:00
break
}
2026-07-15 16:50:44 +08:00
// Compose the request-only prefix once per loop instance before the first
// request boundary. It precedes all derived history and is recorded only
// in the request header, not as session history.
2026-07-08 20:30:10 +08:00
if ( transmission . sessionPrefix === undefined ) {
const emptyPrefix : Message [ ] = deepFreeze ( [ ] )
2026-07-09 23:24:42 +08:00
const composed = await events . waterfall (
'agent/session-prefix' , emptyPrefix , abort . signal ,
2026-07-08 20:30:10 +08:00
( ) = > Promise . resolve ( emptyPrefix ) ,
2026-07-08 21:46:03 +08:00
)
2026-07-13 16:24:32 +08:00
// Never cache an interrupted composition; the next turn recomposes it.
2026-07-08 21:46:03 +08:00
if ( handle . isCancelled ( ) || handle . isDisposed ( ) ) {
handle . setAbort ( undefined )
reason = handle . isDisposed ( ) ? { kind : 'disposed' } : { kind : 'aborted' , reason : handle.cancelReason ( ) }
break
}
transmission . sessionPrefix = deepFreeze ( structuredClone ( composed ) )
2026-07-08 20:30:10 +08:00
}
2026-07-15 16:50:44 +08:00
// Await surface mutations outside the step before snapshotting history.
await events . serial ( 'agent/pre-step' , turn , step , abort . signal )
2026-06-26 13:51:01 +08:00
2026-07-02 02:16:07 +08:00
// Interruption landing during the pre-step seam: do not open an empty step.
2026-06-30 10:56:34 +08:00
if ( handle . isCancelled ( ) || handle . isDisposed ( ) ) {
handle . setAbort ( undefined )
reason = handle . isDisposed ( ) ? { kind : 'disposed' } : { kind : 'aborted' , reason : handle.cancelReason ( ) }
break
}
2026-07-13 23:27:00 +08:00
// Snapshot the exact log prefix before step/start: the reconstruction
// boundary. Appends after this synchronous snapshot join the next request.
2026-07-06 03:07:34 +08:00
const boundaryMessages = session . deriveMessages ( )
2026-06-30 12:19:18 +08:00
session . append ( 'step/start' , { turn , step } )
2026-07-12 18:57:42 +08:00
// Only a committed step/start creates a balancing obligation. A
// pre-commit veto throws before this assignment; post-commit observers
// are contained inside Session.append().
stepOpen = true
2026-06-14 23:53:37 +08:00
2026-07-02 02:16:07 +08:00
// Cancel landing in the step-start window: a synchronous `session/event`
// step/start listener can cancel after the step is already open. Check
// AFTER the step/start append and before `runStep`: drop the step, end the
// turn accordingly. closeStep balances the already-appended step/start.
2026-06-29 15:59:52 +08:00
if ( handle . isCancelled ( ) || handle . isDisposed ( ) ) {
2026-06-20 04:51:32 +08:00
handle . setAbort ( undefined )
2026-06-29 15:59:52 +08:00
reason = handle . isDisposed ( ) ? { kind : 'disposed' } : { kind : 'aborted' , reason : handle.cancelReason ( ) }
2026-06-20 04:51:32 +08:00
closeStep ( )
break
}
2026-07-15 16:03:52 +08:00
let stepOutcome :
| { hadToolCalls : boolean ; finish : FinishReason }
| { requestError : RequestError }
| { error : RequestError }
2026-06-14 23:53:37 +08:00
try {
2026-07-09 01:17:47 +08:00
stepOutcome = await runStep (
2026-07-19 14:06:47 +08:00
ctx , events , handle , turn , step , assembly , fullSystemPrompt , boundaryMessages , transmission , abort . signal )
2026-06-14 23:53:37 +08:00
} catch ( error : unknown ) {
2026-07-19 16:14:58 +08:00
if ( error instanceof TerminalModelRequestFailure ) {
2026-07-15 16:03:52 +08:00
stepOutcome = { requestError : error.requestError }
} else {
stepOutcome = { error : toError ( error ) }
}
}
if ( 'requestError' in stepOutcome ) {
// Recovery observes a balanced failed step and the original provider
// error while the failed step's signal remains the active owner.
closeStep ( )
if ( handle . isDisposed ( ) || abort . signal . aborted ) {
handle . setAbort ( undefined )
reason = handle . isDisposed ( )
? { kind : 'disposed' }
: { kind : 'aborted' , reason : String ( abort . signal . reason ) }
break
}
const defaultDecision : RequestErrorDecision = { action : 'fail' }
let recoveryDecision : RequestErrorDecision = defaultDecision
try {
recoveryDecision = await events . waterfall (
'agent/request-error' , turn , step , stepOutcome . requestError ,
requestRetryAttempt , abort . signal ,
( ) = > Promise . resolve ( defaultDecision ) ,
)
} catch ( recoveryError : unknown ) {
ctx . logger . warn (
` agent " ${ agent . id } ": request recovery failed at turn ${ turn } , step ${ step } : ${ toError ( recoveryError ) . message } ` ,
)
}
2026-06-14 23:53:37 +08:00
handle . setAbort ( undefined )
2026-07-15 16:03:52 +08:00
// Cancellation and disposal always win over either a recovery decision
// or a recovery-listener failure.
// eslint-disable-next-line @typescript-eslint/no-unnecessary-condition
if ( handle . isDisposed ( ) || abort . signal . aborted ) {
reason = handle . isDisposed ( )
? { kind : 'disposed' }
: { kind : 'aborted' , reason : String ( abort . signal . reason ) }
break
}
2026-07-19 15:02:11 +08:00
switch ( recoveryDecision . action ) {
case 'retry' :
requestRetryAttempt += 1
continue
case 'fail' :
failTurn ( stepOutcome . requestError )
break
/* v8 ignore next -- closed-union exhaustiveness guard */
default :
assertNever ( recoveryDecision , 'agent request-error decision' )
2026-07-15 16:03:52 +08:00
}
break
2026-06-11 10:54:31 +08:00
}
2026-06-11 12:18:52 +08:00
2026-06-14 23:53:37 +08:00
if ( 'error' in stepOutcome ) {
// Steering that arrived during the failed step stays in the inbox —
// runLoop re-enqueues it as a queued message, so an abort-then-steer
// starts a fresh turn instead of being silently consumed.
closeStep ( )
2026-07-15 16:03:52 +08:00
handle . setAbort ( undefined )
2026-06-14 23:53:37 +08:00
const { error } = stepOutcome
2026-07-15 16:03:52 +08:00
/* v8 ignore next -- narrow race: disposal while non-request step work throws. */
2026-06-14 23:53:37 +08:00
if ( handle . isDisposed ( ) ) {
reason = { kind : 'disposed' }
} else if ( abort . signal . aborted ) {
2026-06-21 05:50:39 +08:00
/* v8 ignore next -- signal.reason always set: cancel()/disposal provide a default */
2026-06-14 23:53:37 +08:00
reason = { kind : 'aborted' , reason : String ( abort . signal . reason ? ? 'aborted' ) }
} else {
failTurn ( error )
}
break
}
2026-06-11 12:18:52 +08:00
2026-07-15 16:03:52 +08:00
requestRetryAttempt = 0
2026-07-13 16:24:32 +08:00
// Preserve max-token completion unless a later disposal, abort, or error wins.
2026-06-15 23:53:47 +08:00
const stepReason = stepFinishReason ( stepOutcome . finish )
if ( stepReason ) reason = stepReason
2026-06-14 23:53:37 +08:00
// Steering that arrived during streaming/tool execution.
2026-07-19 16:28:31 +08:00
const steered = drainSteering ( )
2026-06-11 10:54:31 +08:00
2026-07-15 16:03:52 +08:00
try {
await events . serial ( 'agent/post-step' , turn , step , abort . signal )
} catch ( error : unknown ) {
stepOutcome = { error : toError ( error ) }
}
if ( 'error' in stepOutcome ) {
closeStep ( )
handle . setAbort ( undefined )
/* v8 ignore next -- narrow race: disposal while a post-step listener throws. */
if ( handle . isDisposed ( ) ) {
reason = { kind : 'disposed' }
} else if ( abort . signal . aborted ) {
/* v8 ignore next -- signal.reason always set by cancellation or disposal. */
reason = { kind : 'aborted' , reason : String ( abort . signal . reason ? ? 'aborted' ) }
} else {
failTurn ( stepOutcome . error )
}
break
}
if ( handle . isDisposed ( ) || abort . signal . aborted ) {
reason = handle . isDisposed ( )
? { kind : 'disposed' }
: { kind : 'aborted' , reason : String ( abort . signal . reason ) }
closeStep ( )
handle . setAbort ( undefined )
break
}
2026-07-12 18:57:42 +08:00
closeStep ( )
2026-07-15 16:03:52 +08:00
handle . setAbort ( undefined )
2026-06-14 23:53:37 +08:00
2026-06-30 17:11:18 +08:00
const defaultDecision : ContinuationDecision = { action : stepOutcome.hadToolCalls || steered ? 'continue' : 'stop' }
let decision : ContinuationDecision
2026-06-14 23:53:37 +08:00
try {
2026-07-09 01:17:47 +08:00
decision = await events . waterfall (
'agent/turn-continuation' , turn , defaultDecision ,
2026-06-14 23:53:37 +08:00
( ) = > Promise . resolve ( defaultDecision ) ,
)
} catch ( error : unknown ) {
// A broken continuation plugin ends the turn, not the loop.
failTurn ( toError ( error ) )
break
}
2026-06-11 10:54:31 +08:00
2026-07-13 16:24:32 +08:00
// A continuation reason becomes next-step steering.
2026-06-30 17:11:18 +08:00
if ( decision . action === 'continue' && decision . reason ) {
2026-07-11 22:55:26 +08:00
handle . inbox . steer ( { content : decision.reason.content , source : decision.reason.source } )
2026-06-30 17:11:18 +08:00
}
let shouldContinue = decision . action === 'continue'
2026-07-13 16:24:32 +08:00
// Pending steering overrides an ordinary stop.
2026-07-11 22:55:26 +08:00
if ( ! shouldContinue && handle . inbox . hasSteering ) shouldContinue = true
2026-07-13 16:24:32 +08:00
// Terminal policy is monotonic and runs after ordinary continuation folding.
2026-07-11 22:55:26 +08:00
let terminalStop = false
try {
2026-07-12 22:36:04 +08:00
const stop = await events . serial ( 'agent/turn-stop' , turn )
2026-07-11 22:55:26 +08:00
terminalStop = stop !== undefined
} catch ( error : unknown ) {
// A broken terminal policy is an ordinary continuation failure: fail
// this turn closed while leaving the driver alive for later turns.
failTurn ( toError ( error ) )
break
}
if ( terminalStop ) {
terminalStopped = true
2026-07-13 16:24:32 +08:00
// Terminal stop discards steering but preserves ordinary queued prompts.
2026-07-11 22:55:26 +08:00
handle . inbox . drainSteering ( )
shouldContinue = false
}
2026-06-11 12:18:52 +08:00
2026-07-13 16:24:32 +08:00
// The marker catches cancellation after the step controller was cleared.
2026-06-20 04:51:32 +08:00
if ( handle . isCancelled ( ) ) {
2026-06-20 11:51:32 +08:00
reason = { kind : 'aborted' , reason : handle.cancelReason ( ) }
2026-06-20 04:51:32 +08:00
break
}
2026-06-14 23:53:37 +08:00
if ( ! shouldContinue || handle . isDisposed ( ) ) {
/* v8 ignore next -- disposal during continuation-decision window is a narrow race; error-path disposal is covered elsewhere */
if ( handle . isDisposed ( ) ) reason = { kind : 'disposed' }
break
}
2026-06-11 10:54:31 +08:00
}
2026-06-11 12:18:52 +08:00
2026-07-02 03:26:45 +08:00
// Normal / inline-error loop exit: close the turn.
closeTurn ( )
2026-06-14 23:53:37 +08:00
} catch ( error : unknown ) {
2026-07-13 16:24:32 +08:00
// Close only a turn whose start committed to the log.
2026-06-15 23:44:54 +08:00
const turnStartLogged = session . events . some ( e = > e . type === 'turn/start' && e . data . turn === turn )
if ( ! turnStartLogged ) throw error
2026-06-14 23:53:37 +08:00
closeStep ( )
2026-07-13 16:24:32 +08:00
// Preserve an established disposal reason; otherwise report the failure.
2026-06-14 23:53:37 +08:00
if ( handle . isDisposed ( ) && ! errorReported ) { // eslint-disable-line @typescript-eslint/no-unnecessary-condition
reason = { kind : 'disposed' }
} else {
failTurn ( toError ( error ) )
}
2026-07-02 03:26:45 +08:00
closeTurn ( )
2026-06-14 23:53:37 +08:00
}
2026-06-11 10:54:31 +08:00
2026-07-13 16:24:32 +08:00
// Flush through the store-owned durability checkpoint without killing the driver on failure.
2026-06-11 12:18:52 +08:00
try {
2026-07-09 01:17:47 +08:00
await ctx . sessions . flush ( session )
2026-06-11 12:39:27 +08:00
} catch ( error : unknown ) {
2026-07-13 16:24:32 +08:00
// The turn is closed, so report the failed flush live rather than append outside a turn.
2026-06-11 14:02:47 +08:00
const err = toError ( error )
2026-06-15 20:56:17 +08:00
ctx . logger . warn ( ` agent " ${ agent . id } ": session/flush failed at turn ${ turn } : ${ err . message } ` )
try {
2026-07-09 01:17:47 +08:00
events . emit ( 'agent/error' , turn , step , err )
2026-06-15 20:56:17 +08:00
} catch {
// contained: a throwing agent/error listener must not escape the loop.
}
2026-06-11 12:18:52 +08:00
}
2026-07-11 22:55:26 +08:00
return terminalStopped
2026-06-11 12:18:52 +08:00
}
2026-06-11 10:54:31 +08:00
2026-07-13 23:27:00 +08:00
/**
* Run one committed step: transform call config, log the request header, build
* the request from the cached prefix plus the step-boundary snapshot, stream and
* record the response, then execute tools. The caller has already assembled the
* prompt, run `agent/pre-step`, snapshotted history, and opened the step.
*/
2026-06-11 10:54:31 +08:00
async function runStep (
ctx : Context ,
2026-07-09 01:17:47 +08:00
events : AgentEventDispatch ,
2026-07-15 12:22:39 +08:00
handle : LoopHandle ,
2026-06-11 10:54:31 +08:00
turn : number ,
step : number ,
2026-06-26 13:51:01 +08:00
assembly : PromptAssembly ,
system : string ,
2026-07-06 03:07:34 +08:00
boundaryMessages : Message [ ] ,
transmission : TransmissionLog ,
2026-06-11 10:54:31 +08:00
signal : AbortSignal ,
2026-06-15 23:53:47 +08:00
) : Promise < { hadToolCalls : boolean ; finish : FinishReason } > {
2026-07-19 14:06:47 +08:00
const agent = ctx . agents . requireInitiator ( )
2026-06-11 10:54:31 +08:00
const { session , options } = agent
2026-07-13 16:24:32 +08:00
// Seed the first request from agent options and later requests from the logged header;
// detach and freeze so listeners must return an attributable replacement.
2026-07-06 04:13:51 +08:00
const seedConfig : LlmCallConfig = deepFreeze ( structuredClone ( transmission . loggedHeader
2026-07-06 03:07:34 +08:00
// eslint-disable-next-line @typescript-eslint/no-non-null-assertion -- loggedHeader ⟹ a snapshot is in the log
? session . requestHeader ( ) ! . config
2026-07-14 21:57:52 +08:00
: { provider : options.provider ? ? '' , model : options.model ? ? '' } ) )
2026-07-06 03:07:34 +08:00
2026-07-13 16:24:32 +08:00
// Listener replacements are recorded in the request header before dispatch.
2026-07-09 01:17:47 +08:00
const config = await events . waterfall ( 'agent/request' , turn , step , seedConfig , ( ) = > Promise . resolve ( seedConfig ) )
2026-07-14 21:57:52 +08:00
if ( ! config . provider || ! config . model ) {
throw new Error ( ` agent " ${ agent . id } " has no provider/model: set AgentOptions.provider and AgentOptions.model or supply both via the agent/request waterfall ` )
2026-07-06 03:07:34 +08:00
}
2026-07-08 20:30:10 +08:00
// eslint-disable-next-line @typescript-eslint/no-non-null-assertion -- runTurn composes the prefix before every runStep call
const sessionPrefix = transmission . sessionPrefix !
2026-07-07 19:42:30 +08:00
2026-07-13 16:24:32 +08:00
// Record the canonical header, including the otherwise-unlogged prefix, before dispatch.
2026-07-06 03:07:34 +08:00
const header = canonicalHeader ( {
config ,
2026-06-11 14:02:47 +08:00
. . . system ? { system } : { } ,
. . . assembly . tools . length > 0 ? { tools : assembly.tools } : { } ,
2026-07-08 15:44:30 +08:00
. . . sessionPrefix . length > 0 ? { messagePrefix : sessionPrefix } : { } ,
2026-07-06 03:07:34 +08:00
} )
recordRequestHeader ( session , transmission , header )
2026-07-13 16:24:32 +08:00
// Freeze the logged header plus boundary snapshot; the prefix precedes derived history.
2026-07-06 03:07:34 +08:00
const request : GenerateOptions = deepFreeze ( {
2026-07-14 21:57:52 +08:00
provider : header.config.provider ,
2026-07-06 03:07:34 +08:00
model : header.config.model ,
2026-07-08 15:44:30 +08:00
messages : [ . . . header . messagePrefix ? ? [ ] , . . . boundaryMessages ] ,
2026-07-06 03:07:34 +08:00
. . . header . system !== undefined ? { system : header.system } : { } ,
. . . header . tools !== undefined ? { tools : header.tools } : { } ,
. . . header . config . temperature !== undefined ? { temperature : header.config.temperature } : { } ,
. . . header . config . maxTokens !== undefined ? { maxTokens : header.config.maxTokens } : { } ,
. . . header . config . stop !== undefined ? { stop : header.config.stop } : { } ,
2026-06-22 08:39:36 +08:00
sessionId : session.id ,
2026-06-11 10:54:31 +08:00
signal ,
2026-07-06 03:07:34 +08:00
} )
2026-06-11 10:54:31 +08:00
// --- Model call (streaming-first; raw chunks are the replay record) ---
const assembler = new BlockAssembler ( )
2026-06-17 19:25:29 +08:00
const chunkSeqs : number [ ] = [ ]
2026-07-19 16:14:58 +08:00
const stream = ctx . llm . stream ( request )
try {
for await ( const chunk of stream ) {
/* v8 ignore next -- signal.reason always set: cancel()/disposal provide a default */
if ( signal . aborted ) throw new Error ( String ( signal . reason ? ? 'aborted' ) )
const chunkEvent = session . append ( 'assistant/chunk' , { turn , step , chunk } )
chunkSeqs . push ( chunkEvent . seq )
assembler . push ( chunk )
}
} catch ( error : unknown ) {
if ( isLlmAdapterFailure ( stream , error ) ) throw new TerminalModelRequestFailure ( error )
throw error
2026-06-11 10:54:31 +08:00
}
2026-07-13 16:24:32 +08:00
// Normalize failure finish chunks into the same path as thrown stream errors.
2026-06-13 00:28:29 +08:00
const stepError = finishError ( assembler . finish )
2026-07-15 16:03:52 +08:00
if ( stepError ) throw new TerminalModelRequestFailure ( stepError )
2026-06-13 00:28:29 +08:00
2026-07-19 16:28:31 +08:00
const recordAssistantMessage = (
assembledContent : ContentBlock [ ] ,
message : Message ,
preserveReplayState = true ,
) : void = > {
session . append (
'assistant/message' ,
{
turn ,
step ,
content : message.content ,
provenance : assistantProvenance (
header . config ,
assembler . replayState ,
preserveReplayState && isDeepStrictEqual ( message . content , assembledContent ) ,
) ,
. . . assembler . usage === undefined ? { } : { usage : assembler.usage } ,
} ,
{ surfaceOp : 'append' , sourceEventSeqs : chunkSeqs } ,
)
}
// A rejected result still records the successful provider call without retaining rejected output.
const processStepResult = async ( assembledContent : ContentBlock [ ] , message : Message ) : Promise < Message > = > {
try {
return await events . waterfall (
'agent/step-result' , turn , step , message , ( ) = > Promise . resolve ( message ) ,
)
} catch ( error : unknown ) {
recordAssistantMessage ( assembledContent , { . . . message , content : [ ] } , false )
throw error
}
}
2026-06-18 23:41:14 +08:00
if ( assembler . finish . kind === 'max-tokens' ) {
2026-07-14 21:57:52 +08:00
const assembled = assembler . message ( )
2026-07-15 13:45:19 +08:00
const assembledContent = structuredClone ( assembled . content )
2026-07-14 21:57:52 +08:00
let message : Message = withoutToolCalls ( assembled )
2026-07-19 16:28:31 +08:00
message = withoutToolCalls ( await processStepResult ( assembledContent , message ) )
2026-07-13 16:24:32 +08:00
// Preserve usage even when max-token truncation produced no content.
2026-07-19 16:28:31 +08:00
recordAssistantMessage ( assembledContent , message )
2026-06-18 23:41:14 +08:00
return { hadToolCalls : false , finish : assembler.finish }
}
2026-07-13 16:24:32 +08:00
// Record the post-waterfall message that tool dispatch uses.
2026-07-14 21:57:52 +08:00
const assembled = assembler . message ( )
2026-07-15 13:45:19 +08:00
const assembledContent = structuredClone ( assembled . content )
2026-07-14 21:57:52 +08:00
let message : Message = assembled
2026-07-19 16:28:31 +08:00
message = await processStepResult ( assembledContent , message )
2026-06-11 12:18:52 +08:00
2026-07-17 22:38:46 +08:00
// Every successful call records its completion anchor, including explicit
// empty chunk provenance for a contentless, usage-less provider response.
2026-07-19 16:28:31 +08:00
recordAssistantMessage ( assembledContent , message )
2026-06-11 10:54:31 +08:00
2026-07-18 15:09:39 +08:00
// Dispatch may overlap; policy, durable results, and result context stay model-ordered.
2026-06-11 10:54:31 +08:00
const toolCalls = message . content . filter ( block = > block . type === 'tool-call' )
2026-07-15 12:22:39 +08:00
if ( toolCalls . length === 0 ) return { hadToolCalls : false , finish : assembler.finish }
2026-07-15 12:49:50 +08:00
return handle . withToolBatch ( async ( acceptContext ) = > {
2026-07-18 15:09:39 +08:00
await executeToolCalls (
2026-07-19 14:06:47 +08:00
ctx , turn , step , toolCalls , signal , handle . maxParallelToolCalls , acceptContext ,
2026-07-18 15:09:39 +08:00
)
2026-07-15 12:22:39 +08:00
return { hadToolCalls : true , finish : assembler.finish }
} )
2026-07-15 15:13:57 +08:00
}
2026-07-14 21:57:52 +08:00
/** Build durable assistant provenance, dropping replay state after any content rewrite. */
function assistantProvenance ( config : LlmCallConfig , replayState : unknown , contentUnchanged : boolean ) : NonNullable < Message [ 'provenance' ] > {
return {
provider : config.provider ,
model : config.model ,
. . . contentUnchanged && replayState !== undefined ? { replayState } : { } ,
}
}
2026-06-19 00:37:12 +08:00
function withoutToolCalls ( message : Message ) : Message {
return { . . . message , content : message.content.filter ( block = > block . type !== 'tool-call' ) }
}
2026-07-06 22:09:30 +08:00
/**
* The last turn number in a (possibly seeded) session log, or 0.
* @param session - the session whose log is scanned for the latest `turn/start`.
* @returns the latest `turn/start`'s turn number, or 0 when the log has none (the next turn is this plus one).
*/
2026-06-15 20:56:17 +08:00
export function lastTurnNumber ( session : Session ) : number {
2026-06-11 14:17:58 +08:00
const lastStart = session . events . findLast ( event = > event . type === 'turn/start' )
return lastStart ? . data . turn ? ? 0
2026-06-11 10:54:31 +08:00
}
2026-06-15 20:56:17 +08:00
/**
2026-07-13 16:24:32 +08:00
* Whether the session log has an unmatched `turn/start`. Agent status is not
* sufficient during pre-start and post-end windows.
2026-07-06 22:09:30 +08:00
* @param session - the session whose log is inspected.
* @returns true when the log's last turn boundary is a `turn/start` with no matching `turn/end` yet.
2026-06-15 20:56:17 +08:00
*/
export function isTurnOpen ( session : Session ) : boolean {
const last = session . events . findLast ( e = > e . type === 'turn/start' || e . type === 'turn/end' )
return last ? . type === 'turn/start'
}