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.
2026-07-19 22:50:49 +08:00
* See .agents/notes/implemented/architecture/2026-06-18-agent-lifecycle-and-ownership-seams.md.
2026-07-13 23:27:00 +08:00
* @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-20 03:34:19 +08:00
import type { ContentBlock , FinishReason , GenerateOptions , LlmCallConfig , LlmFailure , Message } from '@deepseek-ai/dsh-llm'
2026-07-14 21:57:52 +08:00
import { isDeepStrictEqual } from 'node:util'
2026-07-21 14:09:17 +08:00
import { BlockAssembler , HarnessError , LlmError , assertNever , deepFreeze , errorChain , llmFailureOf , markAgentLoopRequest } from '@deepseek-ai/dsh-llm'
2026-07-16 18:12:34 +08:00
import { agentEvents , agentInterruptReasonOf , 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-07-22 17:34:31 +08:00
import type { PromptMessageData , 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-07-16 18:12:34 +08:00
import type { TurnCancellation } from './cancellation.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 {
2026-07-20 03:34:19 +08:00
constructor (
readonly requestError : RequestError ,
readonly failure : LlmFailure ,
) {
2026-07-20 22:11:26 +08:00
super ( failure . message , { cause : requestError } )
2026-07-15 16:03:52 +08:00
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-20 03:34:19 +08:00
function finishError ( finish : FinishReason ) : { error : RequestError ; failure : LlmFailure } | undefined {
2026-06-13 00:28:29 +08:00
switch ( finish . kind ) {
2026-07-20 03:34:19 +08:00
case 'error' :
2026-06-13 00:28:29 +08:00
case 'aborted' : {
2026-07-20 03:34:19 +08:00
const facts = finish . failure
const error = new LlmError ( facts . message , facts . code , {
. . . facts . status === undefined ? { } : { status : facts.status } ,
2026-07-20 18:38:22 +08:00
. . . facts . providerRetryAfterMs === undefined
? { }
: { providerRetryAfterMs : facts.providerRetryAfterMs } ,
2026-07-20 03:34:19 +08:00
. . . facts . requestId === undefined ? { } : { requestId : facts.requestId } ,
} )
return { error , failure : error.failure }
2026-06-13 00:28:29 +08:00
}
// 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-20 11:17:09 +08:00
* The durable message renders the full cause chain: `turn/end` is the single
* durable record of an in-turn failure, so a wrapper message alone (e.g.
* `fetch failed`) would lose the diagnosis the session log exists to keep.
2026-06-11 14:02:47 +08:00
*/
2026-07-15 16:03:52 +08:00
function errorData ( err : RequestError ) : { message : string ; code? : string } {
2026-07-20 11:17:09 +08:00
return { message : errorChain ( err ) , . . . typeof err . code === 'string' ? { code : err.code } : { } }
2026-06-11 14:02:47 +08:00
}
2026-07-20 22:11:26 +08:00
/** Preserve cause diagnostics, falling back to adapter-normalized prose for a hostile Error. */
function durableFailure ( err : RequestError , failure : LlmFailure ) : LlmFailure {
const message = errorChain ( err )
return { . . . failure , message : message === '<unrenderable value>' ? failure.message : message }
}
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-16 18:12:34 +08:00
/** Internal control-flow sentinel; durable classification comes only from the turn signal. */
const TURN_INTERRUPTED = new Error ( 'turn interrupted' )
2026-07-22 17:34:31 +08:00
const PROMPT_PREFIX_REQUEST_DELIMITER : ContentBlock = {
type : 'text' ,
text : '\n\n## My request:\n' ,
}
interface PreparedPromptMessage {
data : PromptMessageData
separateContexts : HookContext [ ]
}
/** Bake declared prefix contexts into one reconstructable prompt message. */
function preparePromptMessage (
content : ContentBlock [ ] ,
source : PromptMessageData [ 'source' ] ,
contexts : readonly HookContext [ ] ,
) : PreparedPromptMessage {
const prefixContexts = contexts . filter ( context = > context . placement === 'prompt-prefix' )
const separateContexts = contexts . filter ( context = > context . placement !== 'prompt-prefix' )
if ( prefixContexts . length === 0 ) return { data : { content , source } , separateContexts }
return {
data : {
content : [
. . . prefixContexts . flatMap ( context = > context . content ) ,
PROMPT_PREFIX_REQUEST_DELIMITER ,
. . . content ,
] ,
source ,
envelope : {
displayContent : content ,
prefixContexts : prefixContexts.map ( context = > ( {
source : context.source ,
. . . context . meta === undefined ? { } : { meta : context.meta } ,
} ) ) ,
} ,
} ,
separateContexts ,
}
}
2026-07-16 18:12:34 +08:00
/** Stop at an explicit cooperative boundary without stringifying the runtime reason. */
function interruptionCheckpoint ( signal : AbortSignal ) : void {
if ( signal . aborted ) throw TURN_INTERRUPTED
}
/** Classify a supported turn interruption, with lifecycle disposal taking precedence. */
function interruptionTurnEndReason ( handle : LoopHandle , signal : AbortSignal ) : TurnEndReason | undefined {
if ( handle . isDisposed ( ) ) return { kind : 'disposed' }
const reason = agentInterruptReasonOf ( signal )
if ( reason === undefined ) return undefined
switch ( reason . kind ) {
case 'user' :
case 'parent' :
return { kind : 'aborted' }
2026-07-20 21:38:49 +08:00
/* v8 ignore next 2 -- the private holder requests disposed only after lifecycle state flips, which returns above. */
2026-07-16 18:12:34 +08:00
case 'disposed' :
return { kind : 'disposed' }
2026-07-20 21:38:49 +08:00
/* v8 ignore next 2 -- AgentInterruptReason is closed and the public helper filters unsupported reasons. */
2026-07-16 18:12:34 +08:00
default :
return assertNever ( reason , 'AgentInterruptReason' )
}
}
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
2026-07-16 18:12:34 +08:00
/** Install a fresh active-turn owner before the running notification. */
installTurnCancellation ( ) : TurnCancellation
2026-07-21 12:14:53 +08:00
/** Clear only the exact owner whose turn reached its terminal event boundary. */
2026-07-16 18:12:34 +08:00
clearTurnCancellation ( cancellation : TurnCancellation ) : void
2026-06-11 10:54:31 +08:00
/** Resolves when the agent is disposed — unblocks the idle wait. */
disposed : Promise < void >
isDisposed ( ) : boolean
2026-07-16 18:12:34 +08:00
/** Whether queued work was cancelled before an active turn owner existed. */
isPreRunCancelled ( ) : boolean
/** Clear the cause-less pre-run marker without affecting replacement work. */
clearPreRunCancel ( ) : void
2026-07-20 13:15:11 +08:00
/** Settle idle waiters before pre-running cancellation publishes idle. */
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-20 11:52:30 +08:00
* Drive queued messages as independent durable turns until disposal. Plugin
2026-07-20 11:53:49 +08:00
* failures end the current turn without terminating the driver. The caller
* establishes the `ctx.agents.withInitiator()` boundary before entry; package-private
2026-07-19 16:28:31 +08:00
* 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-16 18:12:34 +08:00
* @param handle - the bridge to status, turn cancellation ownership, disposal, and pre-run cancellation state.
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-20 13:15:11 +08:00
// An idle listener can enqueue and cancel replacement work before the next
// wait is installed. Consume that empty marker before parking the driver.
2026-07-20 21:38:49 +08:00
if ( handle . isPreRunCancelled ( ) ) {
handle . clearPreRunCancel ( )
2026-07-20 13:15:11 +08:00
if ( ! handle . inbox . hasQueued ) {
handle . settleIdle ( )
handle . setStatus ( 'idle' )
continue
}
}
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;
2026-07-20 13:20:51 +08:00
// a replacement prompt still runs before the eventual idle transition.
2026-07-16 18:12:34 +08:00
if ( handle . isPreRunCancelled ( ) ) {
handle . clearPreRunCancel ( )
2026-07-11 22:55:26 +08:00
if ( ! handle . inbox . hasQueued ) {
2026-07-20 13:15:11 +08:00
// Settle before publishing idle: the already-idle path has no status
// transition, while an idle listener can register waiters for new work.
2026-06-20 05:10:16 +08:00
handle . settleIdle ( )
2026-07-20 13:15:11 +08:00
handle . setStatus ( 'idle' )
2026-06-20 05:10:16 +08:00
continue
}
2026-06-20 04:51:32 +08:00
}
2026-07-16 18:12:34 +08:00
let cancellation = handle . installTurnCancellation ( )
2026-06-11 10:54:31 +08:00
handle . setStatus ( 'running' )
2026-07-16 18:12:34 +08:00
if ( handle . isDisposed ( ) ) {
handle . clearTurnCancellation ( cancellation )
break
}
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-07-16 18:12:34 +08:00
if ( cancellation . signal . aborted ) {
handle . clearTurnCancellation ( cancellation )
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-07-16 18:12:34 +08:00
cancellation = handle . installTurnCancellation ( )
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-21 12:14:53 +08:00
terminalStopped = await runTurn ( ctx , events , handle , turn , transmission , cancellation )
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 )
2026-07-20 11:17:09 +08:00
ctx . logger . warn ( ` agent " ${ agent . id } ": turn ${ turn } failed before it started: ${ errorChain ( err ) } ` )
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-07-16 18:12:34 +08:00
} finally {
handle . clearTurnCancellation ( cancellation )
2026-06-11 12:18:52 +08:00
}
2026-06-11 10:54:31 +08:00
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-21 12:14:53 +08:00
cancellation : TurnCancellation ,
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-21 12:14:53 +08:00
const { signal } = cancellation
2026-07-19 16:28:31 +08:00
const drainSteering = ( ) : boolean = > {
const messages = handle . inbox . drainSteering ( )
for ( const message of messages ) {
2026-07-22 17:34:31 +08:00
const prepared = preparePromptMessage ( message . content , message . source , message . contexts )
session . append ( 'steering/message' , { turn , . . . prepared . data } , { surfaceOp : 'append' } )
for ( const context of prepared . separateContexts ) {
session . append ( 'context/message' , {
content : context.content ,
source : context.source ,
. . . context . meta === undefined ? { } : { meta : context.meta } ,
} , { surfaceOp : 'append' } )
2026-07-21 16:46:48 +08:00
}
2026-07-19 16:28:31 +08:00
}
return messages . length > 0
}
2026-06-11 10:54:31 +08:00
2026-07-17 17:14:52 +08:00
// Claim one queued message before opening its turn, but append it only after `turn/start`.
const message = handle . inbox . dequeueQueued ( )
2026-06-11 14:58:36 +08:00
/* v8 ignore next 3 -- invariant guard: runLoop only calls runTurn when hasQueued */
2026-07-17 17:14:52 +08:00
if ( ! message ) throw new Error ( 'runTurn invariant violated: no queued message at turn start' )
const trigger : TurnTrigger = { kind : 'message' , source : message.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-20 03:34:19 +08:00
let requestFailureHistory : readonly LlmFailure [ ] = Object . freeze ( [ ] )
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-20 03:34:19 +08:00
const failTurn = ( err : RequestError , failure? : LlmFailure ) : void = > {
2026-06-14 23:53:37 +08:00
if ( errorReported ) return
errorReported = true
2026-07-20 03:34:19 +08:00
reason = failure === undefined
? { kind : 'error' , step , . . . errorData ( err ) }
2026-07-20 22:11:26 +08:00
: { kind : 'error' , step , failure : durableFailure ( err , failure ) }
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-21 12:14:53 +08:00
// Retire cancellation authority before publishing the terminal event. The
// following durability flush is quiescent turn work, but no longer part of
// the cancellable turn lifetime.
2026-07-02 03:26:45 +08:00
const closeTurn = ( ) : void = > {
2026-07-21 12:14:53 +08:00
handle . clearTurnCancellation ( cancellation )
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-07-16 18:12:34 +08:00
interruptionCheckpoint ( signal )
2026-07-17 17:14:52 +08:00
// The claimed message runs the `agent/prompt-submit` waterfall before it
// becomes a `user/message` — a hook can rewrite the prompt or block it.
2026-06-30 17:11:18 +08:00
// 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.
2026-07-17 17:14:52 +08:00
const promptDecision = await events . waterfall (
2026-07-20 21:38:49 +08:00
'agent/prompt-submit' , message . content , message . source , signal ,
2026-07-21 16:46:48 +08:00
( ) = > Promise . resolve < PromptDecision > ( {
kind : 'allow' ,
. . . message . contexts . length === 0 ? { } : { additionalContexts : message.contexts } ,
} ) ,
2026-07-17 17:14:52 +08:00
)
2026-07-20 21:38:49 +08:00
interruptionCheckpoint ( signal )
2026-07-17 17:14:52 +08:00
if ( promptDecision . kind === 'block' ) {
session . append ( 'prompt/blocked' , { content : message.content , source : message.source , reason : promptDecision.reason } )
reason = { kind : 'rejected' , reason : promptDecision.reason }
} else {
2026-06-30 17:11:18 +08:00
// `allow.content` REPLACES the prompt bytes (a rewrite); absent keeps them.
2026-07-17 17:14:52 +08:00
const content = promptDecision . content ? ? message . content
2026-07-22 17:34:31 +08:00
const prepared = preparePromptMessage ( content , message . source , promptDecision . additionalContexts ? ? [ ] )
session . append ( 'user/message' , prepared . data , { surfaceOp : 'append' } )
// Separate contexts still enter THIS turn through inject(). Prefix
// contexts are already baked into the user/message with their durable
// display envelope, so appending them again would duplicate model input.
for ( const context of prepared . separateContexts ) {
2026-07-13 16:31:03 +08:00
agent . inject ( context . content , {
source : context.source ,
. . . 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-17 17:14:52 +08:00
// A blocked prompt closes its zero-step turn as rejected.
if ( promptDecision . kind === 'block' ) 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-07-15 16:50:44 +08:00
// Assemble once before pre-step so listener work and the request share one prompt value.
2026-07-16 18:12:34 +08:00
const assembly = await ctx . systemPrompt . assemble ( assembleContextFor ( agent , signal ) )
interruptionCheckpoint ( signal )
2026-07-05 01:54:46 +08:00
const fullSystemPrompt = renderPrompt ( assembly )
2026-06-26 13:51:01 +08:00
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 (
2026-07-16 18:12:34 +08:00
'agent/session-prefix' , emptyPrefix , 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-16 18:12:34 +08:00
interruptionCheckpoint ( signal )
2026-07-08 21:46:03 +08:00
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.
2026-07-20 21:38:49 +08:00
await events . serial ( 'agent/pre-step' , turn , step , signal )
2026-07-16 18:12:34 +08:00
interruptionCheckpoint ( signal )
2026-06-30 10:56:34 +08:00
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-16 18:12:34 +08:00
// A synchronous step/start observer can cancel after the step opened.
interruptionCheckpoint ( signal )
2026-06-20 04:51:32 +08:00
2026-07-15 16:03:52 +08:00
let stepOutcome :
| { hadToolCalls : boolean ; finish : FinishReason }
2026-07-20 03:34:19 +08:00
| { requestError : RequestError ; failure : LlmFailure }
2026-07-15 16:03:52 +08:00
| { error : RequestError }
2026-06-14 23:53:37 +08:00
try {
2026-07-09 01:17:47 +08:00
stepOutcome = await runStep (
2026-07-20 21:38:49 +08:00
ctx , events , handle , turn , step , assembly , fullSystemPrompt , boundaryMessages , transmission , 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-20 03:34:19 +08:00
stepOutcome = { requestError : error.requestError , failure : error.failure }
2026-07-15 16:03:52 +08:00
} 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 ( )
2026-07-20 21:38:49 +08:00
const interrupted = interruptionTurnEndReason ( handle , signal )
if ( interrupted !== undefined ) {
reason = interrupted
2026-07-15 16:03:52 +08:00
break
}
const defaultDecision : RequestErrorDecision = { action : 'fail' }
let recoveryDecision : RequestErrorDecision = defaultDecision
try {
recoveryDecision = await events . waterfall (
'agent/request-error' , turn , step , stepOutcome . requestError ,
2026-07-21 00:02:17 +08:00
stepOutcome . failure , requestFailureHistory , signal ,
2026-07-15 16:03:52 +08:00
( ) = > Promise . resolve ( defaultDecision ) ,
)
} catch ( recoveryError : unknown ) {
ctx . logger . warn (
2026-07-20 11:17:09 +08:00
` agent " ${ agent . id } ": request recovery failed at turn ${ turn } , step ${ step } : ${ errorChain ( recoveryError ) } ` ,
2026-07-15 16:03:52 +08:00
)
}
// Cancellation and disposal always win over either a recovery decision
// or a recovery-listener failure.
2026-07-20 21:38:49 +08:00
const recoveryInterrupted = interruptionTurnEndReason ( handle , signal )
if ( recoveryInterrupted !== undefined ) {
reason = recoveryInterrupted
2026-07-15 16:03:52 +08:00
break
}
2026-07-19 15:02:11 +08:00
switch ( recoveryDecision . action ) {
case 'retry' :
2026-07-20 03:34:19 +08:00
requestFailureHistory = Object . freeze ( [ . . . requestFailureHistory , stepOutcome . failure ] )
2026-07-19 15:02:11 +08:00
continue
case 'fail' :
2026-07-20 03:34:19 +08:00
failTurn ( stepOutcome . requestError , stepOutcome . failure )
2026-07-19 15:02:11 +08:00
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 ( )
const { error } = stepOutcome
2026-07-20 21:38:49 +08:00
const interrupted = interruptionTurnEndReason ( handle , signal )
if ( interrupted === undefined ) failTurn ( error )
else reason = interrupted
2026-06-14 23:53:37 +08:00
break
}
2026-06-11 12:18:52 +08:00
2026-07-20 03:34:19 +08:00
requestFailureHistory = Object . freeze ( [ ] )
2026-07-15 16:03:52 +08:00
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 {
2026-07-20 21:38:49 +08:00
await events . serial ( 'agent/post-step' , turn , step , signal )
2026-07-15 16:03:52 +08:00
} catch ( error : unknown ) {
stepOutcome = { error : toError ( error ) }
}
if ( 'error' in stepOutcome ) {
closeStep ( )
2026-07-20 21:38:49 +08:00
const interrupted = interruptionTurnEndReason ( handle , signal )
if ( interrupted === undefined ) failTurn ( stepOutcome . error )
else reason = interrupted
2026-07-15 16:03:52 +08:00
break
}
2026-07-20 21:38:49 +08:00
const postStepInterrupted = interruptionTurnEndReason ( handle , signal )
if ( postStepInterrupted !== undefined ) {
reason = postStepInterrupted
2026-07-15 16:03:52 +08:00
closeStep ( )
break
}
2026-07-12 18:57:42 +08:00
closeStep ( )
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 (
2026-07-16 18:12:34 +08:00
'agent/turn-continuation' , turn , defaultDecision , signal ,
2026-06-14 23:53:37 +08:00
( ) = > Promise . resolve ( defaultDecision ) ,
)
2026-07-16 18:12:34 +08:00
interruptionCheckpoint ( signal )
2026-06-14 23:53:37 +08:00
} catch ( error : unknown ) {
2026-07-20 21:38:49 +08:00
const interrupted = interruptionTurnEndReason ( handle , signal )
if ( interrupted === undefined ) failTurn ( toError ( error ) )
else reason = interrupted
2026-06-14 23:53:37 +08:00
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-21 16:46:48 +08:00
handle . inbox . steer ( { content : decision.reason.content , source : decision.reason.source , contexts : [ ] } )
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-16 18:12:34 +08:00
const stop = await events . serial ( 'agent/turn-stop' , turn , signal )
interruptionCheckpoint ( signal )
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.
2026-07-20 21:38:49 +08:00
const interrupted = interruptionTurnEndReason ( handle , signal )
if ( interrupted === undefined ) failTurn ( toError ( error ) )
else reason = interrupted
2026-07-11 22:55:26 +08:00
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-20 21:38:49 +08:00
if ( ! shouldContinue ) 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-20 21:38:49 +08:00
const interrupted = interruptionTurnEndReason ( handle , signal )
if ( interrupted === undefined ) failTurn ( toError ( error ) )
else reason = interrupted
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-07-20 11:17:09 +08:00
ctx . logger . warn ( ` agent " ${ agent . id } ": session/flush failed at turn ${ turn } : ${ errorChain ( err ) } ` )
2026-06-15 20:56:17 +08:00
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-20 21:38:49 +08:00
const config = await events . waterfall (
'agent/request' , turn , step , seedConfig , signal , ( ) = > Promise . resolve ( seedConfig ) ,
)
2026-07-16 18:12:34 +08:00
interruptionCheckpoint ( signal )
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-21 14:09:17 +08:00
const request : GenerateOptions = markAgentLoopRequest ( 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-21 00:25:38 +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 ) {
2026-07-20 21:38:49 +08:00
interruptionCheckpoint ( signal )
2026-07-19 16:14:58 +08:00
const chunkEvent = session . append ( 'assistant/chunk' , { turn , step , chunk } )
chunkSeqs . push ( chunkEvent . seq )
assembler . push ( chunk )
}
} catch ( error : unknown ) {
2026-07-20 03:34:19 +08:00
const failure = llmFailureOf ( stream , error )
if ( failure !== undefined && error instanceof Error ) throw new TerminalModelRequestFailure ( error , failure )
2026-07-19 16:14:58 +08:00
throw error
2026-06-11 10:54:31 +08:00
}
2026-07-16 18:12:34 +08:00
interruptionCheckpoint ( signal )
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-20 03:34:19 +08:00
if ( stepError ) throw new TerminalModelRequestFailure ( stepError . error , stepError . failure )
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 {
2026-07-20 21:38:49 +08:00
const processed = await events . waterfall (
'agent/step-result' , turn , step , message , signal , ( ) = > Promise . resolve ( message ) ,
2026-07-19 16:28:31 +08:00
)
2026-07-20 21:38:49 +08:00
interruptionCheckpoint ( signal )
return processed
2026-07-19 16:28:31 +08:00
} 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'
}