loop.ts 31 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740
  1. /**
  2. * Drives one agent across queued durable turns. Turn failures are contained so
  3. * later work can run; the session log, not this driver, owns conversation state.
  4. * See .agents/notes/implemented/architecture/2026-06-18-agent-lifecycle-and-ownership-seams.md.
  5. * @module dsh-agent-loop/loop
  6. */
  7. import type { Context } from 'cordis'
  8. import type { ContentBlock, FinishReason, GenerateOptions, LlmCallConfig, Message } from '@deepseek-ai/dsh-llm'
  9. import { isDeepStrictEqual } from 'node:util'
  10. import { BlockAssembler, HarnessError, assertNever, deepFreeze, errorChain, isLlmAdapterFailure } from '@deepseek-ai/dsh-llm'
  11. import { agentEvents, assembleContextFor } from '@deepseek-ai/dsh-agent'
  12. import type { AgentEventDispatch, ContinuationDecision, HookContext, PromptDecision, RequestError, RequestErrorDecision } from '@deepseek-ai/dsh-agent'
  13. import { canonicalHeader } from '@deepseek-ai/dsh-session'
  14. import type { Session, TurnEndReason, TurnTrigger } from '@deepseek-ai/dsh-session'
  15. import { createTransmissionLog, recordRequestHeader } from './request-log.ts'
  16. import type { TransmissionLog } from './request-log.ts'
  17. import { renderPrompt } from '@deepseek-ai/dsh-system-prompt'
  18. import type { PromptAssembly } from '@deepseek-ai/dsh-system-prompt'
  19. import type {} from '@deepseek-ai/dsh-tools'
  20. import { executeToolCalls } from './tool-calls.ts'
  21. import type { Inbox } from './inbox.ts'
  22. /** Normalize thrown values while preserving an existing error code. */
  23. function toError(error: unknown): RequestError {
  24. return error instanceof Error ? error : new HarnessError(String(error), 'UNKNOWN', { cause: error })
  25. }
  26. /** Distinguishes final model-request failures from failures in later step processing. */
  27. class TerminalModelRequestFailure extends Error {
  28. constructor(readonly requestError: RequestError) {
  29. super(requestError.message, { cause: requestError })
  30. this.name = 'TerminalModelRequestFailure'
  31. }
  32. }
  33. /** Convert terminal failure finishes into step errors; unknown extensible finishes remain successful. */
  34. function finishError(finish: FinishReason): RequestError | undefined {
  35. switch (finish.kind) {
  36. case 'error': {
  37. const error: RequestError = new Error(finish.message)
  38. if (finish.code !== undefined) error.code = finish.code
  39. return error
  40. }
  41. case 'aborted': {
  42. const error: RequestError = new Error('model stream aborted')
  43. error.code = 'ABORTED'
  44. return error
  45. }
  46. // stop / tool-calls / max-tokens / plugin-added kinds → not a failure.
  47. default:
  48. return undefined
  49. }
  50. }
  51. /**
  52. * Build the `{ message, code? }` part of an error payload, omitting the
  53. * `code` key entirely when absent (exactOptionalPropertyTypes-correct).
  54. * The durable message renders the full cause chain: `turn/end` is the single
  55. * durable record of an in-turn failure, so a wrapper message alone (e.g.
  56. * `fetch failed`) would lose the diagnosis the session log exists to keep.
  57. */
  58. function errorData(err: RequestError): { message: string; code?: string } {
  59. return { message: errorChain(err), ...typeof err.code === 'string' ? { code: err.code } : {} }
  60. }
  61. /** Map a successful max-token finish onto the turn reason; other successful finishes add nothing. */
  62. function stepFinishReason(finish: FinishReason): TurnEndReason | undefined {
  63. switch (finish.kind) {
  64. case 'max-tokens':
  65. return { kind: 'max-tokens' }
  66. // stop / tool-calls / plugin-added kinds → no turn-end contribution
  67. // beyond the default `completed`. FinishReason is merge-extensible, so a
  68. // default (not assertNever) handles unknown kinds as ordinary success.
  69. default:
  70. return undefined
  71. }
  72. }
  73. /** Mutable agent controls supplied to the loop driver. */
  74. export interface LoopHandle {
  75. /** Native-private agent inbox handed to the driver only at internal startup. */
  76. readonly inbox: Inbox
  77. /** Maximum parallel-safe calls allowed in one step. */
  78. readonly maxParallelToolCalls: number
  79. setStatus(status: 'idle' | 'running'): void
  80. setAbort(controller: AbortController | undefined): void
  81. /** Resolves when the agent is disposed — unblocks the idle wait. */
  82. disposed: Promise<void>
  83. isDisposed(): boolean
  84. /** Whether cancellation is pending for the current loop iteration. */
  85. isCancelled(): boolean
  86. /** Resolved pending-cancellation reason; meaningful only while {@link isCancelled} is true. */
  87. cancelReason(): string
  88. /** Clear the cancel marker (called once per iteration after the turn returns). */
  89. clearCancel(): void
  90. /** Settle idle waiters before pre-running cancellation publishes idle. */
  91. settleIdle(): void
  92. /** Run an active tool-call batch, accepting post-tool context into the FIFO drained before settlement. */
  93. readonly withToolBatch: <T>(run: (acceptContext: (context: HookContext) => void) => Promise<T>) => Promise<T>
  94. }
  95. /**
  96. * Drive queued messages as independent durable turns until disposal. Plugin
  97. * failures end the current turn without terminating the driver. The caller
  98. * establishes the `ctx.agents.withInitiator()` boundary before entry; package-private
  99. * orchestration recovers that exact Agent and captures its Session locally.
  100. * @param ctx - the plugin context the loop reaches its initiating Agent,
  101. * events (agent/…, session/flush), and services (systemPrompt, llm, tools)
  102. * through.
  103. * @param handle - the bridge to the agent's mutable state: status/abort setters plus the disposal and cancel-marker reads.
  104. * @throws when no initiating Agent is active.
  105. */
  106. export async function runLoop(ctx: Context, handle: LoopHandle): Promise<void> {
  107. const agent = ctx.agents.requireInitiator()
  108. // Per-instance prefix and request-header state; conversation history remains in the session log.
  109. const transmission = createTransmissionLog()
  110. const { session } = agent
  111. // Fused subject and scope carrier for every agent event below.
  112. const events = agentEvents(ctx, agent)
  113. while (!handle.isDisposed()) {
  114. // An idle listener can enqueue and cancel replacement work before the next
  115. // wait is installed. Consume that empty marker before parking the driver.
  116. if (handle.isCancelled()) {
  117. handle.clearCancel()
  118. if (!handle.inbox.hasQueued) {
  119. handle.settleIdle()
  120. handle.setStatus('idle')
  121. continue
  122. }
  123. }
  124. await handle.inbox.waitForQueued(handle.disposed)
  125. if (handle.isDisposed()) break
  126. // Cancellation between wake and `running` skips only the cancelled work;
  127. // a replacement prompt still runs before the eventual idle transition.
  128. if (handle.isCancelled()) {
  129. handle.clearCancel()
  130. if (!handle.inbox.hasQueued) {
  131. // Settle before publishing idle: the already-idle path has no status
  132. // transition, while an idle listener can register waiters for new work.
  133. handle.settleIdle()
  134. handle.setStatus('idle')
  135. continue
  136. }
  137. }
  138. handle.setStatus('running')
  139. if (handle.isDisposed()) break
  140. // A synchronous `running` listener can cancel before `runTurn`; balance the
  141. // status only when no replacement prompt was queued by that listener.
  142. if (handle.isCancelled()) {
  143. handle.clearCancel()
  144. if (!handle.inbox.hasQueued) {
  145. handle.setStatus('idle')
  146. continue
  147. }
  148. }
  149. // Idle injection can add a turn, so derive the next number from the log.
  150. const turn = lastTurnNumber(session) + 1
  151. let terminalStopped = false
  152. try {
  153. terminalStopped = await runTurn(ctx, events, handle, turn, transmission)
  154. } catch (error: unknown) {
  155. // Pre-turn failure has no durable boundary to close; report it without appending outside a turn.
  156. const err = toError(error)
  157. ctx.logger.warn(`agent "${agent.id}": turn ${turn} failed before it started: ${errorChain(err)}`)
  158. try {
  159. events.emit('agent/error', turn, 0, err)
  160. } catch { /* contained: a throwing agent/error listener must not kill the driver */ }
  161. }
  162. // Reset per iteration, including when a prompt arrives during the flush window.
  163. handle.clearCancel()
  164. // Late steering becomes queued input unless terminal policy stopped the turn.
  165. for (const message of handle.inbox.drainSteering()) {
  166. if (!terminalStopped) handle.inbox.enqueue(message)
  167. }
  168. if (!handle.inbox.hasQueued) handle.setStatus('idle')
  169. }
  170. }
  171. async function runTurn(
  172. ctx: Context, events: AgentEventDispatch, handle: LoopHandle, turn: number, transmission: TransmissionLog,
  173. ): Promise<boolean> {
  174. const agent = ctx.agents.requireInitiator()
  175. const { session } = agent
  176. const drainSteering = (): boolean => {
  177. const messages = handle.inbox.drainSteering()
  178. for (const message of messages) {
  179. session.append('steering/message', { turn, content: message.content, source: message.source }, { surfaceOp: 'append' })
  180. }
  181. return messages.length > 0
  182. }
  183. // Claim one queued message before opening its turn, but append it only after `turn/start`.
  184. const message = handle.inbox.dequeueQueued()
  185. /* v8 ignore next 3 -- invariant guard: runLoop only calls runTurn when hasQueued */
  186. if (!message) throw new Error('runTurn invariant violated: no queued message at turn start')
  187. const trigger: TurnTrigger = { kind: 'message', source: message.source }
  188. let reason: TurnEndReason = { kind: 'completed' }
  189. let step = 0
  190. let requestRetryAttempt = 0
  191. let stepOpen = false
  192. let errorReported = false
  193. let terminalStopped = false
  194. // Close the committed step once; pre-commit validation failure still escapes.
  195. const closeStep = (): void => {
  196. if (!stepOpen) return
  197. session.append('step/end', { turn, step })
  198. stepOpen = false
  199. }
  200. // Record the durable turn failure once and contain the live error notification.
  201. const failTurn = (err: RequestError): void => {
  202. if (errorReported) return
  203. errorReported = true
  204. reason = { kind: 'error', step, ...errorData(err) }
  205. try {
  206. events.emit('agent/error', turn, step, err)
  207. } catch {
  208. // contained: the error is already captured on `reason`; a throwing
  209. // agent/error listener must not prevent the turn from closing.
  210. }
  211. }
  212. // Pre-commit validation failure escapes rather than masquerading as a committed boundary.
  213. const closeTurn = (): void => {
  214. session.append('turn/end', { turn, reason })
  215. }
  216. try {
  217. // --- Turn boundary. Once turn/start is appended, a turn/end is owed no
  218. // matter what throws below; the catch + closeTurn guarantee it. A pre-commit
  219. // veto leaves no turn/start in the log and therefore owes no turn/end.
  220. session.append('turn/start', { turn, trigger })
  221. // The claimed message runs the `agent/prompt-submit` waterfall before it
  222. // becomes a `user/message` — a hook can rewrite the prompt or block it.
  223. // Recorded INSIDE the turn (after turn/start) so every event is turn-enclosed;
  224. // turn/end is now owed, so a throwing prompt-submit listener (the waterfall
  225. // throws) is caught below and the turn still closes.
  226. const promptDecision = await events.waterfall(
  227. 'agent/prompt-submit', message.content, message.source,
  228. () => Promise.resolve<PromptDecision>({ kind: 'allow' }),
  229. )
  230. if (promptDecision.kind === 'block') {
  231. session.append('prompt/blocked', { content: message.content, source: message.source, reason: promptDecision.reason })
  232. reason = { kind: 'rejected', reason: promptDecision.reason }
  233. } else {
  234. // `allow.content` REPLACES the prompt bytes (a rewrite); absent keeps them.
  235. const content = promptDecision.content ?? message.content
  236. session.append('user/message', { content, source: message.source }, { surfaceOp: 'append' })
  237. // Every `allow.additionalContexts` entry is a separate context/message the
  238. // next request also sees. The turn is open, so inject() appends each one
  239. // into THIS turn without flattening provenance or metadata.
  240. for (const context of promptDecision.additionalContexts ?? []) {
  241. agent.inject(context.content, {
  242. source: context.source,
  243. ...context.meta !== undefined ? { meta: context.meta } : {},
  244. })
  245. }
  246. }
  247. while (true) {
  248. // A blocked prompt closes its zero-step turn as rejected.
  249. if (promptDecision.kind === 'block') break
  250. step += 1
  251. // Steering from the previous round's continuation listeners joins before
  252. // the request.
  253. drainSteering()
  254. // The step's AbortController exists BEFORE any async pre-step work so a
  255. // dispose() or cancel() — in a synchronous turn-start listener or an
  256. // async listener whose effect fires before we block — always has an armed
  257. // abort to cancel against. isDisposed below covers disposal, which does
  258. // NOT set the cancel marker. Cleared on every exit path below.
  259. const abort = new AbortController()
  260. handle.setAbort(abort)
  261. // Assemble once before pre-step so listener work and the request share one prompt value.
  262. const assembly = await ctx.systemPrompt.assemble(assembleContextFor(agent))
  263. const fullSystemPrompt = renderPrompt(assembly)
  264. // Cancellation or disposal during assembly ends the turn before any step opens.
  265. if (handle.isCancelled() || handle.isDisposed()) {
  266. handle.setAbort(undefined)
  267. reason = handle.isDisposed() ? { kind: 'disposed' } : { kind: 'aborted', reason: handle.cancelReason() }
  268. break
  269. }
  270. // Compose the request-only prefix once per loop instance before the first
  271. // request boundary. It precedes all derived history and is recorded only
  272. // in the request header, not as session history.
  273. if (transmission.sessionPrefix === undefined) {
  274. const emptyPrefix: Message[] = deepFreeze([])
  275. const composed = await events.waterfall(
  276. 'agent/session-prefix', emptyPrefix, abort.signal,
  277. () => Promise.resolve(emptyPrefix),
  278. )
  279. // Never cache an interrupted composition; the next turn recomposes it.
  280. if (handle.isCancelled() || handle.isDisposed()) {
  281. handle.setAbort(undefined)
  282. reason = handle.isDisposed() ? { kind: 'disposed' } : { kind: 'aborted', reason: handle.cancelReason() }
  283. break
  284. }
  285. transmission.sessionPrefix = deepFreeze(structuredClone(composed))
  286. }
  287. // Await surface mutations outside the step before snapshotting history.
  288. await events.serial('agent/pre-step', turn, step, abort.signal)
  289. // Interruption landing during the pre-step seam: do not open an empty step.
  290. if (handle.isCancelled() || handle.isDisposed()) {
  291. handle.setAbort(undefined)
  292. reason = handle.isDisposed() ? { kind: 'disposed' } : { kind: 'aborted', reason: handle.cancelReason() }
  293. break
  294. }
  295. // Snapshot the exact log prefix before step/start: the reconstruction
  296. // boundary. Appends after this synchronous snapshot join the next request.
  297. const boundaryMessages = session.deriveMessages()
  298. session.append('step/start', { turn, step })
  299. // Only a committed step/start creates a balancing obligation. A
  300. // pre-commit veto throws before this assignment; post-commit observers
  301. // are contained inside Session.append().
  302. stepOpen = true
  303. // Cancel landing in the step-start window: a synchronous `session/event`
  304. // step/start listener can cancel after the step is already open. Check
  305. // AFTER the step/start append and before `runStep`: drop the step, end the
  306. // turn accordingly. closeStep balances the already-appended step/start.
  307. if (handle.isCancelled() || handle.isDisposed()) {
  308. handle.setAbort(undefined)
  309. reason = handle.isDisposed() ? { kind: 'disposed' } : { kind: 'aborted', reason: handle.cancelReason() }
  310. closeStep()
  311. break
  312. }
  313. let stepOutcome:
  314. | { hadToolCalls: boolean; finish: FinishReason }
  315. | { requestError: RequestError }
  316. | { error: RequestError }
  317. try {
  318. stepOutcome = await runStep(
  319. ctx, events, handle, turn, step, assembly, fullSystemPrompt, boundaryMessages, transmission, abort.signal)
  320. } catch (error: unknown) {
  321. if (error instanceof TerminalModelRequestFailure) {
  322. stepOutcome = { requestError: error.requestError }
  323. } else {
  324. stepOutcome = { error: toError(error) }
  325. }
  326. }
  327. if ('requestError' in stepOutcome) {
  328. // Recovery observes a balanced failed step and the original provider
  329. // error while the failed step's signal remains the active owner.
  330. closeStep()
  331. if (handle.isDisposed() || abort.signal.aborted) {
  332. handle.setAbort(undefined)
  333. reason = handle.isDisposed()
  334. ? { kind: 'disposed' }
  335. : { kind: 'aborted', reason: String(abort.signal.reason) }
  336. break
  337. }
  338. const defaultDecision: RequestErrorDecision = { action: 'fail' }
  339. let recoveryDecision: RequestErrorDecision = defaultDecision
  340. try {
  341. recoveryDecision = await events.waterfall(
  342. 'agent/request-error', turn, step, stepOutcome.requestError,
  343. requestRetryAttempt, abort.signal,
  344. () => Promise.resolve(defaultDecision),
  345. )
  346. } catch (recoveryError: unknown) {
  347. ctx.logger.warn(
  348. `agent "${agent.id}": request recovery failed at turn ${turn}, step ${step}: ${errorChain(recoveryError)}`,
  349. )
  350. }
  351. handle.setAbort(undefined)
  352. // Cancellation and disposal always win over either a recovery decision
  353. // or a recovery-listener failure.
  354. // eslint-disable-next-line @typescript-eslint/no-unnecessary-condition
  355. if (handle.isDisposed() || abort.signal.aborted) {
  356. reason = handle.isDisposed()
  357. ? { kind: 'disposed' }
  358. : { kind: 'aborted', reason: String(abort.signal.reason) }
  359. break
  360. }
  361. switch (recoveryDecision.action) {
  362. case 'retry':
  363. requestRetryAttempt += 1
  364. continue
  365. case 'fail':
  366. failTurn(stepOutcome.requestError)
  367. break
  368. /* v8 ignore next -- closed-union exhaustiveness guard */
  369. default:
  370. assertNever(recoveryDecision, 'agent request-error decision')
  371. }
  372. break
  373. }
  374. if ('error' in stepOutcome) {
  375. // Steering that arrived during the failed step stays in the inbox —
  376. // runLoop re-enqueues it as a queued message, so an abort-then-steer
  377. // starts a fresh turn instead of being silently consumed.
  378. closeStep()
  379. handle.setAbort(undefined)
  380. const { error } = stepOutcome
  381. /* v8 ignore next -- narrow race: disposal while non-request step work throws. */
  382. if (handle.isDisposed()) {
  383. reason = { kind: 'disposed' }
  384. } else if (abort.signal.aborted) {
  385. /* v8 ignore next -- signal.reason always set: cancel()/disposal provide a default */
  386. reason = { kind: 'aborted', reason: String(abort.signal.reason ?? 'aborted') }
  387. } else {
  388. failTurn(error)
  389. }
  390. break
  391. }
  392. requestRetryAttempt = 0
  393. // Preserve max-token completion unless a later disposal, abort, or error wins.
  394. const stepReason = stepFinishReason(stepOutcome.finish)
  395. if (stepReason) reason = stepReason
  396. // Steering that arrived during streaming/tool execution.
  397. const steered = drainSteering()
  398. try {
  399. await events.serial('agent/post-step', turn, step, abort.signal)
  400. } catch (error: unknown) {
  401. stepOutcome = { error: toError(error) }
  402. }
  403. if ('error' in stepOutcome) {
  404. closeStep()
  405. handle.setAbort(undefined)
  406. /* v8 ignore next -- narrow race: disposal while a post-step listener throws. */
  407. if (handle.isDisposed()) {
  408. reason = { kind: 'disposed' }
  409. } else if (abort.signal.aborted) {
  410. /* v8 ignore next -- signal.reason always set by cancellation or disposal. */
  411. reason = { kind: 'aborted', reason: String(abort.signal.reason ?? 'aborted') }
  412. } else {
  413. failTurn(stepOutcome.error)
  414. }
  415. break
  416. }
  417. if (handle.isDisposed() || abort.signal.aborted) {
  418. reason = handle.isDisposed()
  419. ? { kind: 'disposed' }
  420. : { kind: 'aborted', reason: String(abort.signal.reason) }
  421. closeStep()
  422. handle.setAbort(undefined)
  423. break
  424. }
  425. closeStep()
  426. handle.setAbort(undefined)
  427. const defaultDecision: ContinuationDecision = { action: stepOutcome.hadToolCalls || steered ? 'continue' : 'stop' }
  428. let decision: ContinuationDecision
  429. try {
  430. decision = await events.waterfall(
  431. 'agent/turn-continuation', turn, defaultDecision,
  432. () => Promise.resolve(defaultDecision),
  433. )
  434. } catch (error: unknown) {
  435. // A broken continuation plugin ends the turn, not the loop.
  436. failTurn(toError(error))
  437. break
  438. }
  439. // A continuation reason becomes next-step steering.
  440. if (decision.action === 'continue' && decision.reason) {
  441. handle.inbox.steer({ content: decision.reason.content, source: decision.reason.source })
  442. }
  443. let shouldContinue = decision.action === 'continue'
  444. // Pending steering overrides an ordinary stop.
  445. if (!shouldContinue && handle.inbox.hasSteering) shouldContinue = true
  446. // Terminal policy is monotonic and runs after ordinary continuation folding.
  447. let terminalStop = false
  448. try {
  449. const stop = await events.serial('agent/turn-stop', turn)
  450. terminalStop = stop !== undefined
  451. } catch (error: unknown) {
  452. // A broken terminal policy is an ordinary continuation failure: fail
  453. // this turn closed while leaving the driver alive for later turns.
  454. failTurn(toError(error))
  455. break
  456. }
  457. if (terminalStop) {
  458. terminalStopped = true
  459. // Terminal stop discards steering but preserves ordinary queued prompts.
  460. handle.inbox.drainSteering()
  461. shouldContinue = false
  462. }
  463. // The marker catches cancellation after the step controller was cleared.
  464. if (handle.isCancelled()) {
  465. reason = { kind: 'aborted', reason: handle.cancelReason() }
  466. break
  467. }
  468. if (!shouldContinue || handle.isDisposed()) {
  469. /* v8 ignore next -- disposal during continuation-decision window is a narrow race; error-path disposal is covered elsewhere */
  470. if (handle.isDisposed()) reason = { kind: 'disposed' }
  471. break
  472. }
  473. }
  474. // Normal / inline-error loop exit: close the turn.
  475. closeTurn()
  476. } catch (error: unknown) {
  477. // Close only a turn whose start committed to the log.
  478. const turnStartLogged = session.events.some(e => e.type === 'turn/start' && e.data.turn === turn)
  479. if (!turnStartLogged) throw error
  480. closeStep()
  481. // Preserve an established disposal reason; otherwise report the failure.
  482. if (handle.isDisposed() && !errorReported) { // eslint-disable-line @typescript-eslint/no-unnecessary-condition
  483. reason = { kind: 'disposed' }
  484. } else {
  485. failTurn(toError(error))
  486. }
  487. closeTurn()
  488. }
  489. // Flush through the store-owned durability checkpoint without killing the driver on failure.
  490. try {
  491. await ctx.sessions.flush(session)
  492. } catch (error: unknown) {
  493. // The turn is closed, so report the failed flush live rather than append outside a turn.
  494. const err = toError(error)
  495. ctx.logger.warn(`agent "${agent.id}": session/flush failed at turn ${turn}: ${errorChain(err)}`)
  496. try {
  497. events.emit('agent/error', turn, step, err)
  498. } catch {
  499. // contained: a throwing agent/error listener must not escape the loop.
  500. }
  501. }
  502. return terminalStopped
  503. }
  504. /**
  505. * Run one committed step: transform call config, log the request header, build
  506. * the request from the cached prefix plus the step-boundary snapshot, stream and
  507. * record the response, then execute tools. The caller has already assembled the
  508. * prompt, run `agent/pre-step`, snapshotted history, and opened the step.
  509. */
  510. async function runStep(
  511. ctx: Context,
  512. events: AgentEventDispatch,
  513. handle: LoopHandle,
  514. turn: number,
  515. step: number,
  516. assembly: PromptAssembly,
  517. system: string,
  518. boundaryMessages: Message[],
  519. transmission: TransmissionLog,
  520. signal: AbortSignal,
  521. ): Promise<{ hadToolCalls: boolean; finish: FinishReason }> {
  522. const agent = ctx.agents.requireInitiator()
  523. const { session, options } = agent
  524. // Seed the first request from agent options and later requests from the logged header;
  525. // detach and freeze so listeners must return an attributable replacement.
  526. const seedConfig: LlmCallConfig = deepFreeze(structuredClone(transmission.loggedHeader
  527. // eslint-disable-next-line @typescript-eslint/no-non-null-assertion -- loggedHeader ⟹ a snapshot is in the log
  528. ? session.requestHeader()!.config
  529. : { provider: options.provider ?? '', model: options.model ?? '' }))
  530. // Listener replacements are recorded in the request header before dispatch.
  531. const config = await events.waterfall('agent/request', turn, step, seedConfig, () => Promise.resolve(seedConfig))
  532. if (!config.provider || !config.model) {
  533. throw new Error(`agent "${agent.id}" has no provider/model: set AgentOptions.provider and AgentOptions.model or supply both via the agent/request waterfall`)
  534. }
  535. // eslint-disable-next-line @typescript-eslint/no-non-null-assertion -- runTurn composes the prefix before every runStep call
  536. const sessionPrefix = transmission.sessionPrefix!
  537. // Record the canonical header, including the otherwise-unlogged prefix, before dispatch.
  538. const header = canonicalHeader({
  539. config,
  540. ...system ? { system } : {},
  541. ...assembly.tools.length > 0 ? { tools: assembly.tools } : {},
  542. ...sessionPrefix.length > 0 ? { messagePrefix: sessionPrefix } : {},
  543. })
  544. recordRequestHeader(session, transmission, header)
  545. // Freeze the logged header plus boundary snapshot; the prefix precedes derived history.
  546. const request: GenerateOptions = deepFreeze({
  547. provider: header.config.provider,
  548. model: header.config.model,
  549. messages: [...header.messagePrefix ?? [], ...boundaryMessages],
  550. ...header.system !== undefined ? { system: header.system } : {},
  551. ...header.tools !== undefined ? { tools: header.tools } : {},
  552. ...header.config.temperature !== undefined ? { temperature: header.config.temperature } : {},
  553. ...header.config.maxTokens !== undefined ? { maxTokens: header.config.maxTokens } : {},
  554. ...header.config.stop !== undefined ? { stop: header.config.stop } : {},
  555. sessionId: session.id,
  556. signal,
  557. })
  558. // --- Model call (streaming-first; raw chunks are the replay record) ---
  559. const assembler = new BlockAssembler()
  560. const chunkSeqs: number[] = []
  561. const stream = ctx.llm.stream(request)
  562. try {
  563. for await (const chunk of stream) {
  564. /* v8 ignore next -- signal.reason always set: cancel()/disposal provide a default */
  565. if (signal.aborted) throw new Error(String(signal.reason ?? 'aborted'))
  566. const chunkEvent = session.append('assistant/chunk', { turn, step, chunk })
  567. chunkSeqs.push(chunkEvent.seq)
  568. assembler.push(chunk)
  569. }
  570. } catch (error: unknown) {
  571. if (isLlmAdapterFailure(stream, error)) throw new TerminalModelRequestFailure(error)
  572. throw error
  573. }
  574. // Normalize failure finish chunks into the same path as thrown stream errors.
  575. const stepError = finishError(assembler.finish)
  576. if (stepError) throw new TerminalModelRequestFailure(stepError)
  577. const recordAssistantMessage = (
  578. assembledContent: ContentBlock[],
  579. message: Message,
  580. preserveReplayState = true,
  581. ): void => {
  582. session.append(
  583. 'assistant/message',
  584. {
  585. turn,
  586. step,
  587. content: message.content,
  588. provenance: assistantProvenance(
  589. header.config,
  590. assembler.replayState,
  591. preserveReplayState && isDeepStrictEqual(message.content, assembledContent),
  592. ),
  593. ...assembler.usage === undefined ? {} : { usage: assembler.usage },
  594. },
  595. { surfaceOp: 'append', sourceEventSeqs: chunkSeqs },
  596. )
  597. }
  598. // A rejected result still records the successful provider call without retaining rejected output.
  599. const processStepResult = async (assembledContent: ContentBlock[], message: Message): Promise<Message> => {
  600. try {
  601. return await events.waterfall(
  602. 'agent/step-result', turn, step, message, () => Promise.resolve(message),
  603. )
  604. } catch (error: unknown) {
  605. recordAssistantMessage(assembledContent, { ...message, content: [] }, false)
  606. throw error
  607. }
  608. }
  609. if (assembler.finish.kind === 'max-tokens') {
  610. const assembled = assembler.message()
  611. const assembledContent = structuredClone(assembled.content)
  612. let message: Message = withoutToolCalls(assembled)
  613. message = withoutToolCalls(await processStepResult(assembledContent, message))
  614. // Preserve usage even when max-token truncation produced no content.
  615. recordAssistantMessage(assembledContent, message)
  616. return { hadToolCalls: false, finish: assembler.finish }
  617. }
  618. // Record the post-waterfall message that tool dispatch uses.
  619. const assembled = assembler.message()
  620. const assembledContent = structuredClone(assembled.content)
  621. let message: Message = assembled
  622. message = await processStepResult(assembledContent, message)
  623. // Every successful call records its completion anchor, including explicit
  624. // empty chunk provenance for a contentless, usage-less provider response.
  625. recordAssistantMessage(assembledContent, message)
  626. // Dispatch may overlap; policy, durable results, and result context stay model-ordered.
  627. const toolCalls = message.content.filter(block => block.type === 'tool-call')
  628. if (toolCalls.length === 0) return { hadToolCalls: false, finish: assembler.finish }
  629. return handle.withToolBatch(async (acceptContext) => {
  630. await executeToolCalls(
  631. ctx, turn, step, toolCalls, signal, handle.maxParallelToolCalls, acceptContext,
  632. )
  633. return { hadToolCalls: true, finish: assembler.finish }
  634. })
  635. }
  636. /** Build durable assistant provenance, dropping replay state after any content rewrite. */
  637. function assistantProvenance(config: LlmCallConfig, replayState: unknown, contentUnchanged: boolean): NonNullable<Message['provenance']> {
  638. return {
  639. provider: config.provider,
  640. model: config.model,
  641. ...contentUnchanged && replayState !== undefined ? { replayState } : {},
  642. }
  643. }
  644. function withoutToolCalls(message: Message): Message {
  645. return { ...message, content: message.content.filter(block => block.type !== 'tool-call') }
  646. }
  647. /**
  648. * The last turn number in a (possibly seeded) session log, or 0.
  649. * @param session - the session whose log is scanned for the latest `turn/start`.
  650. * @returns the latest `turn/start`'s turn number, or 0 when the log has none (the next turn is this plus one).
  651. */
  652. export function lastTurnNumber(session: Session): number {
  653. const lastStart = session.events.findLast(event => event.type === 'turn/start')
  654. return lastStart?.data.turn ?? 0
  655. }
  656. /**
  657. * Whether the session log has an unmatched `turn/start`. Agent status is not
  658. * sufficient during pre-start and post-end windows.
  659. * @param session - the session whose log is inspected.
  660. * @returns true when the log's last turn boundary is a `turn/start` with no matching `turn/end` yet.
  661. */
  662. export function isTurnOpen(session: Session): boolean {
  663. const last = session.events.findLast(e => e.type === 'turn/start' || e.type === 'turn/end')
  664. return last?.type === 'turn/start'
  665. }