agent.ts 29 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730
  1. /**
  2. * Concrete Agent loop over two pending-input lists: queued prompts each open a
  3. * turn that logs its admitted input after `turn/start` commits, while steering
  4. * and injected context enter through the outbox at step boundaries. Every
  5. * request is derived from the session log.
  6. *
  7. * @module dsh-agent-loop/agent
  8. */
  9. import type { Context } from 'cordis'
  10. import { agentCarrier, assembleContextFor, emitAgentEvent } from '@deepseek-ai/dsh-agent'
  11. import { createScope } from '@deepseek-ai/dsh-scope'
  12. import type { Scope } from '@deepseek-ai/dsh-scope'
  13. import type {
  14. Agent,
  15. CancelOptions,
  16. AgentInterruptReason,
  17. InboxPlacement,
  18. AgentOptions,
  19. AgentStatus,
  20. SettleReason,
  21. PromptDecision,
  22. RequestError,
  23. RequestErrorAction,
  24. SendOptions,
  25. } from '@deepseek-ai/dsh-agent'
  26. import {
  27. BlockAssembler,
  28. LlmError,
  29. assertNever,
  30. createAssistantMessage,
  31. deepFreeze,
  32. errorChain,
  33. freezeMessage,
  34. isHarnessError,
  35. llmFailureOf,
  36. llmRetryPolicyOf,
  37. markAgentLoopRequest,
  38. } from '@deepseek-ai/dsh-llm'
  39. import type { GenerateOptions, LlmCallConfig, LlmFailure, Message, PreparedLlmCall, ResolvedRetryPolicy } from '@deepseek-ai/dsh-llm'
  40. import { canonicalHeader, headerEquals } from '@deepseek-ai/dsh-session'
  41. import type { AssistantMessage, Session, SessionId, TurnEndReason, TurnTrigger, UserMessage } from '@deepseek-ai/dsh-session'
  42. import { renderPrompt } from '@deepseek-ai/dsh-system-prompt'
  43. import type {} from '@deepseek-ai/dsh-tools'
  44. import { executeToolCalls } from './tool-calls.ts'
  45. /** One completed step or a final-adapter failure eligible for recovery. */
  46. type StepOutcome =
  47. | { kind: 'completed'; continueTurn: boolean; concluded: boolean; maxTokens: boolean }
  48. | { kind: 'request-failed'; error: RequestError; failure: LlmFailure; retryPolicy: ResolvedRetryPolicy | undefined }
  49. /**
  50. * The concrete {@link Agent}: each `run()` owns one turn and repeats model
  51. * steps while tools or steering require another request.
  52. */
  53. export class ReactLoopAgent implements Agent {
  54. /** Prompts awaiting individual turns. */
  55. private queued: { message: UserMessage; wakeup: boolean }[] = []
  56. /** Input taken into the session log at step boundaries. */
  57. private outbox: { message: UserMessage; steering: boolean }[] = []
  58. /** Whether observers see a running interval; consecutive turns share it. */
  59. private busy = false
  60. /** Whether an idle waking send has deferred driver admission. */
  61. private wakeScheduled = false
  62. /** Whether next-step input belongs to the current admission or open turn. */
  63. acceptsNextStep = false
  64. /** Abort owner for the current admission or turn. */
  65. private abort: AbortController | undefined
  66. /** Resolves when the current admission and turn exit. */
  67. done: Promise<void> = Promise.resolve()
  68. /** The agent-scoped registration boundary; the lifecycle owner unwinds it after {@link done}. */
  69. readonly scope: Scope
  70. /** The agent's scoped composition context ({@link Agent.ctx}). */
  71. readonly ctx: Context
  72. /** Last turn number opened by this loop or present in its seeded log. */
  73. private lastTurn: number
  74. /** Whether the session log is owed a matching turn end event. */
  75. private turnOpen = false
  76. private stepOpen = false
  77. /** Whether this loop instance has appended its initial/resume request anchor. */
  78. private requestHeaderLogged = false
  79. constructor(
  80. private loopCtx: Context,
  81. public readonly id: SessionId,
  82. public readonly options: AgentOptions,
  83. public readonly session: Session,
  84. ) {
  85. this.lastTurn = session.events.findLast(event => event.type === 'turn/start')?.data.turn ?? 0
  86. this.scope = createScope(loopCtx, this)
  87. this.ctx = this.scope.ctx.extend({ agent: this })
  88. }
  89. /** Last activity state published to observers. */
  90. get status(): AgentStatus {
  91. return this.busy ? 'running' : 'idle'
  92. }
  93. /** Accept and route one unified send item. */
  94. send(
  95. message: UserMessage,
  96. options: SendOptions,
  97. ): void {
  98. const { target, wakeup } = options
  99. if (target === 'next-step' && !wakeup) {
  100. if (this.acceptsNextStep) {
  101. this.outbox.push({ message, steering: false })
  102. return
  103. }
  104. this.session.append('user/message', message, { surfaceOp: 'append' })
  105. return
  106. }
  107. const placement: InboxPlacement = target === 'next-step' && this.acceptsNextStep ? 'steering' : 'queued'
  108. if (placement === 'steering') {
  109. this.outbox.push({ message, steering: true })
  110. } else {
  111. this.queued.push({ message, wakeup })
  112. }
  113. // Preserve the routing decision for every send in this synchronous caller
  114. // stack, while installing quiescence ownership before enqueue observers
  115. // can cancel or dispose.
  116. if (placement === 'queued' && wakeup) this.scheduleKick()
  117. emitAgentEvent(this.loopCtx, this, 'agent/inbox/enqueue', message, placement)
  118. }
  119. /** Queue one ordinary prompt turn and wake the driver. */
  120. followup(input: UserMessage): void {
  121. this.send(input, {
  122. target: 'next-turn',
  123. wakeup: true,
  124. })
  125. }
  126. /** Steer the open turn, falling back to a waking prompt while idle. */
  127. steer(input: UserMessage): void {
  128. this.send(input, {
  129. target: 'next-step',
  130. wakeup: true,
  131. })
  132. }
  133. /** Append model-facing context without waking the driver. */
  134. inject(input: UserMessage): void {
  135. this.send(input, {
  136. target: 'next-step',
  137. wakeup: false,
  138. })
  139. }
  140. /**
  141. * Clear all pending work and abort the active turn; the first cause wins.
  142. * The cause is signal payload for observers and the durable turn/end
  143. * classification — it selects no machine behavior. Teardown is just
  144. * `cancel({kind:'disposed'})` + await {@link done} + {@link scope} dispose,
  145. * all owned by the factory.
  146. */
  147. cancel(cause: AgentInterruptReason, options: CancelOptions = {}): void {
  148. // Effective only when it aborts the active turn or actually discards
  149. // pending work: a keepInbox call with no active turn is a documented
  150. // no-op, so it must not emit cancel-requested for consumers to misread.
  151. const discards = !options.keepInbox && (this.queued.length > 0 || this.outbox.length > 0)
  152. if (this.abort !== undefined || discards) {
  153. // Observe-only: coordination consumers update their state before the
  154. // inboxes clear; listener failures are contained by the dispatcher.
  155. if (cause.kind !== 'disposed') emitAgentEvent(this.loopCtx, this, 'agent/cancel-requested', cause)
  156. }
  157. if (!options.keepInbox) {
  158. const discarded = this.queued.map(item => item.message)
  159. for (const item of this.outbox) {
  160. if (item.steering) discarded.push(item.message)
  161. }
  162. // Clear before abort observers run: replacement work belongs to the next turn.
  163. this.queued.length = 0
  164. this.outbox.length = 0
  165. if (discarded.length > 0) emitAgentEvent(this.loopCtx, this, 'agent/inbox/discard', discarded)
  166. }
  167. const reason = Object.freeze({ kind: cause.kind })
  168. this.abort?.abort(reason)
  169. }
  170. /** Resolve at idle quiescence: no run driving and no waking prompt waiting. */
  171. async whenIdle(): Promise<void> {
  172. // `done` is replaced per activity, so re-reading it follows chained turns.
  173. // Every driver failure today is contained before it can reject `done`,
  174. // but the waiter must not gamble quiescence on that: a future escape
  175. // still counts as settled activity.
  176. /* v8 ignore next 3 -- the catch arm backstops rejection paths that are all currently contained */
  177. while (this.busy || this.wakeScheduled || this.abort !== undefined || this.queued.some(item => item.wakeup)) {
  178. await this.done.catch(() => undefined)
  179. }
  180. }
  181. /** Defer idle admission while keeping {@link done} as its quiescence owner. */
  182. private scheduleKick(): void {
  183. if (this.abort !== undefined || this.wakeScheduled) return
  184. this.wakeScheduled = true
  185. const pending = Promise.withResolvers<void>()
  186. const scheduled = pending.promise
  187. queueMicrotask(() => {
  188. this.wakeScheduled = false
  189. this.kick()
  190. const activity = this.done
  191. if (activity === scheduled) {
  192. pending.resolve()
  193. } else {
  194. void activity.then(
  195. () => { pending.resolve() },
  196. () => { pending.resolve() },
  197. )
  198. }
  199. })
  200. this.done = scheduled
  201. }
  202. /** Claim and admit the next queued prompt, then start its turn. */
  203. private kick(): void {
  204. if (this.abort !== undefined || !this.queued.some(item => item.wakeup)) return
  205. // The some() guard above proves the queue is non-empty; the non-null
  206. // assertion expresses that invariant.
  207. // eslint-disable-next-line @typescript-eslint/no-non-null-assertion
  208. const { message } = this.queued.shift()!
  209. const inheritedOutboxLength = this.outbox.length
  210. const admission = new AbortController()
  211. this.abort = admission
  212. this.acceptsNextStep = true
  213. // Claimed admission is part of the running interval: it is cancellable
  214. // activity, so observers (and their cancel routing) must see it.
  215. if (!this.busy) {
  216. this.busy = true
  217. emitAgentEvent(this.loopCtx, this, 'agent/status', 'running')
  218. }
  219. // The admission body runs synchronously up to the prompt-submit
  220. // waterfall's first await, so the waterfall snapshots its listeners
  221. // before a disposal initiated by the running-status emit above can
  222. // unregister a vetoing plugin.
  223. this.done = this.loopCtx.agents.withInitiator(this, async () => {
  224. const signal = admission.signal
  225. const trigger: TurnTrigger = { kind: 'message', source: message.source }
  226. // Admitted input stays on the stack until its turn/start commits: the
  227. // turn owns it only once the turn exists in the log.
  228. let admitted: UserMessage[] | undefined
  229. try {
  230. signal.throwIfAborted()
  231. const decision = await this.loopCtx.waterfall(
  232. agentCarrier(this), 'agent/prompt-submit', this, message, signal,
  233. () => Promise.resolve<PromptDecision>({ kind: 'allow' }),
  234. )
  235. signal.throwIfAborted()
  236. if (decision.kind === 'allow') {
  237. admitted = [decision.content === undefined
  238. ? message
  239. : freezeMessage({ ...message, content: decision.content })]
  240. for (const context of decision.additionalContexts ?? []) {
  241. admitted.push(freezeMessage(context))
  242. }
  243. }
  244. } catch (error: unknown) {
  245. if (!signal.aborted) {
  246. this.loopCtx.logger.warn(`agent "${this.id}": prompt admission failed: ${errorChain(error)}`)
  247. }
  248. }
  249. // cancel() aborts but never clears the slot, and kick()/run()
  250. // all refuse to install a new owner while one exists, so the admission
  251. // still owns the slot here and releasing it unconditionally is exact.
  252. this.abort = undefined
  253. if (admitted === undefined) {
  254. this.acceptsNextStep = false
  255. try {
  256. this.flushRejectedAdmissionContexts()
  257. } catch (error: unknown) {
  258. // No turn exists for agent/error coordinates. Preserve the
  259. // uncommitted suffix for a later boundary and report locally.
  260. this.loopCtx.logger.warn(
  261. `agent "${this.id}": committing rejected-admission context failed: ${errorChain(error)}`,
  262. )
  263. }
  264. // A synchronously aborted admission would otherwise publish idle
  265. // inside send()'s own synchronous extent, before any post-send
  266. // subscriber could observe the transition.
  267. await Promise.resolve()
  268. this.continueOrIdle()
  269. return
  270. }
  271. await this.run(trigger, admitted, inheritedOutboxLength)
  272. })
  273. // Published only after the abort owner and pending done are installed: a
  274. // dequeue listener that cancels or disposes must find live cancellation
  275. // and quiescence ownership, not the previous activity's settled state.
  276. emitAgentEvent(this.loopCtx, this, 'agent/inbox/dequeue', message, 'queued')
  277. }
  278. /**
  279. * Run one turn and any request-error retry. `admitted` input enters the log
  280. * only after `turn/start` commits; until then it has no owner state to unwind.
  281. */
  282. private async run(
  283. trigger: TurnTrigger,
  284. admitted: UserMessage[] = [],
  285. inheritedOutboxLength = 0,
  286. priorFailures: readonly LlmFailure[] = Object.freeze([]),
  287. ): Promise<void> {
  288. // Both entries hold the invariant: kick() clears the admission slot before
  289. // awaiting run(), and a retry is entered only after the prior run clears it.
  290. /* v8 ignore next -- unreachable guard: every caller clears or checks the abort slot first */
  291. if (this.abort !== undefined) throw new Error(`agent "${this.id}" is already running`)
  292. const controller = new AbortController()
  293. this.abort = controller
  294. this.acceptsNextStep = true
  295. const signal = controller.signal
  296. const turn = this.lastTurn + 1
  297. let step = 0
  298. let opened = false
  299. let reason: TurnEndReason = { kind: 'completed' }
  300. let settleReason: SettleReason = { kind: 'completed' }
  301. let requestFailureHistory = priorFailures
  302. let retryFailures: readonly LlmFailure[] | undefined
  303. const cancelRetry = (): void => { retryFailures = undefined }
  304. signal.addEventListener('abort', cancelRetry, { once: true })
  305. try {
  306. signal.throwIfAborted()
  307. this.session.append('turn/start', { turn, trigger })
  308. // Committed: publish the turn to the machine's own bookkeeping and let
  309. // the admitted input enter the log it now belongs to.
  310. this.turnOpen = true
  311. opened = true
  312. this.lastTurn = turn
  313. // Context or steering retained by an earlier rejected admission happened
  314. // before this prompt and must occupy the same order in durable history.
  315. this.drainOutbox(turn, inheritedOutboxLength)
  316. for (const input of admitted) {
  317. this.session.append('user/message', input, { surfaceOp: 'append' })
  318. }
  319. signal.throwIfAborted()
  320. this.drainOutbox(turn)
  321. steps: while (true) {
  322. step += 1
  323. const outcome = await this.step(turn, step, signal)
  324. switch (outcome.kind) {
  325. case 'completed':
  326. requestFailureHistory = Object.freeze([])
  327. if (outcome.maxTokens) reason = { kind: 'max-tokens' }
  328. // A concluding tool result is terminal: steering already in the
  329. // log waits for the next turn's request instead of reopening this
  330. // one, and the agent/turn-stopping drain below is skipped for the same
  331. // reason.
  332. if (outcome.concluded) break steps
  333. if (outcome.continueTurn || this.outbox.some(item => item.steering)) continue
  334. break
  335. case 'request-failed': {
  336. // step() reports request failures only after step/start commits
  337. // and before its own step/end, so the step is always open here.
  338. this.stepOpen = false
  339. this.session.append('step/end', { turn, step })
  340. if (!signal.aborted) {
  341. try {
  342. const action = await this.loopCtx.waterfall(
  343. agentCarrier(this), 'agent/request-error', this, turn, step, outcome.error,
  344. outcome.failure, requestFailureHistory, outcome.retryPolicy, signal,
  345. () => Promise.resolve<RequestErrorAction>(undefined),
  346. )
  347. // eslint-disable-next-line @typescript-eslint/no-unnecessary-condition -- signal can abort while recovery is awaited.
  348. if (action?.kind === 'retry' && !signal.aborted) {
  349. retryFailures = Object.freeze([...requestFailureHistory, outcome.failure])
  350. }
  351. } catch (recoveryError: unknown) {
  352. this.loopCtx.logger.warn(
  353. `agent "${this.id}": request recovery failed at turn ${turn}, step ${step}: ${errorChain(recoveryError)}`,
  354. )
  355. }
  356. }
  357. const settlement = this.settle(turn, step, outcome.error, signal, outcome.failure)
  358. reason = settlement.reason
  359. settleReason = settlement.settleReason
  360. break steps
  361. }
  362. /* v8 ignore next 2 -- closed-union exhaustiveness guard */
  363. default:
  364. assertNever(outcome)
  365. }
  366. await this.loopCtx.serial(agentCarrier(this), 'agent/turn-stopping', this, turn, signal)
  367. signal.throwIfAborted()
  368. if (!this.drainOutbox(turn)) break
  369. }
  370. } catch (caught: unknown) {
  371. try {
  372. if (this.stepOpen) {
  373. this.stepOpen = false
  374. this.session.append('step/end', { turn, step })
  375. }
  376. } catch (closeError: unknown) {
  377. // Contained like the finally's turn close: a persistently rejecting
  378. // step boundary must not escape run(), or the post-finally tail would
  379. // never publish the terminal status and observers would see a
  380. // permanently running agent whose whenIdle() already resolved.
  381. this.loopCtx.logger.warn(`agent "${this.id}": closing step ${turn}/${step} failed: ${errorChain(closeError)}`)
  382. emitAgentEvent(this.loopCtx, this, 'agent/error', turn, step, closeError)
  383. }
  384. ({ reason, settleReason } = this.settle(turn, step, caught, signal))
  385. } finally {
  386. // Every step-close happens before this point on both success and
  387. // failure paths (step(), the request-failed branch, the catch), so the
  388. // finally owes only the turn boundary.
  389. this.acceptsNextStep = false
  390. try {
  391. if (this.turnOpen) {
  392. // Re-entrant turn/end listeners must route new input to a later turn.
  393. this.turnOpen = false
  394. this.session.append('turn/end', { turn, reason })
  395. }
  396. } catch (error: unknown) {
  397. retryFailures = undefined
  398. this.loopCtx.logger.warn(`agent "${this.id}": closing turn ${turn} failed: ${errorChain(error)}`)
  399. emitAgentEvent(this.loopCtx, this, 'agent/error', turn, step, error)
  400. }
  401. // cancel() aborts but never clears the slot, and no second run can
  402. // install a controller while this one is still unwinding, so the slot
  403. // is still this run's controller here.
  404. this.abort = undefined
  405. signal.removeEventListener('abort', cancelRetry)
  406. }
  407. if (opened) {
  408. try {
  409. await this.loopCtx.sessions.flush(this.session)
  410. } catch (error: unknown) {
  411. this.loopCtx.logger.warn(`agent "${this.id}": session/flush failed at turn ${turn}: ${errorChain(error)}`)
  412. emitAgentEvent(this.loopCtx, this, 'agent/error', turn, step, error)
  413. }
  414. }
  415. if (retryFailures !== undefined) {
  416. await this.run({ kind: 'retry' }, [], 0, retryFailures)
  417. } else {
  418. // agent/settled names only committed turns: a run aborted or rejected
  419. // before turn/start has no durable turn/end for consumers to settle
  420. // against, so it exits without the notification.
  421. if (opened) emitAgentEvent(this.loopCtx, this, 'agent/settled', turn, settleReason)
  422. this.continueOrIdle()
  423. }
  424. }
  425. /**
  426. * Run the `agent/step` extension point, commit pending input, derive one
  427. * request, and execute its tool calls inside one durable step boundary.
  428. */
  429. private async step(
  430. turn: number,
  431. step: number,
  432. signal: AbortSignal,
  433. ): Promise<StepOutcome> {
  434. const { session } = this
  435. // The single between-steps extension point: listeners inject, steer, or
  436. // edit the log here; the request derives from the log after this settles.
  437. await this.loopCtx.serial(agentCarrier(this), 'agent/step', this, turn, step, signal)
  438. signal.throwIfAborted()
  439. // Take the outbox whole — same-boundary steering and context leave in
  440. // this request together.
  441. this.drainOutbox(turn)
  442. // Assemble the system prompt fresh each step (it may depend on log state).
  443. const assembly = await this.loopCtx.systemPrompt.assemble(assembleContextFor(this, signal))
  444. signal.throwIfAborted()
  445. const system = renderPrompt(assembly)
  446. // Snapshot the exact log prefix: the reconstruction boundary. Appends
  447. // after this synchronous snapshot join the next request.
  448. const boundaryMessages = session.deriveMessages()
  449. session.append('step/start', { turn, step })
  450. this.stepOpen = true
  451. signal.throwIfAborted()
  452. const { request, preparedCall } = await this.buildRequest(
  453. turn, step, assembly.tools, system, boundaryMessages, signal,
  454. )
  455. const assembler = new BlockAssembler()
  456. const chunkSeqs: number[] = []
  457. const stream = preparedCall?.stream(request) ?? this.loopCtx.llm.stream(request)
  458. try {
  459. for await (const chunk of stream) {
  460. signal.throwIfAborted()
  461. const chunkEvent = session.append('assistant/chunk', { turn, step, chunk })
  462. chunkSeqs.push(chunkEvent.seq)
  463. assembler.push(chunk)
  464. }
  465. } catch (error: unknown) {
  466. const facts = llmFailureOf(stream, error)
  467. if (facts !== undefined && error instanceof Error) {
  468. return { kind: 'request-failed', error, failure: facts, retryPolicy: llmRetryPolicyOf(stream) }
  469. }
  470. throw error
  471. }
  472. signal.throwIfAborted()
  473. // Failure finish chunks take the same path as thrown stream errors.
  474. const finish = assembler.finish
  475. if (finish.kind === 'error' || finish.kind === 'aborted') {
  476. const error = new LlmError(finish.failure.message, finish.failure.code, finish.failure)
  477. return { kind: 'request-failed', error, failure: finish.failure, retryPolicy: llmRetryPolicyOf(stream) }
  478. }
  479. // Truncated (max-tokens) output cannot owe tool calls.
  480. const assembled = assembler.blocks()
  481. const content = finish.kind === 'max-tokens'
  482. ? assembled.filter(block => block.type !== 'tool-call')
  483. : assembled
  484. const message: AssistantMessage = createAssistantMessage({
  485. content,
  486. source: {
  487. provider: request.provider,
  488. model: request.model,
  489. ...assembler.replayState !== undefined ? { replayState: assembler.replayState } : {},
  490. },
  491. })
  492. session.append(
  493. 'assistant/message',
  494. {
  495. turn,
  496. step,
  497. message,
  498. ...assembler.usage === undefined ? {} : { usage: assembler.usage },
  499. },
  500. { surfaceOp: 'append', sourceEventSeqs: chunkSeqs },
  501. )
  502. const toolCalls = content.filter(block => block.type === 'tool-call')
  503. let concluded = false
  504. if (toolCalls.length > 0) {
  505. ({ concluded } = await executeToolCalls(
  506. this.loopCtx, turn, step, toolCalls, signal,
  507. context => this.outbox.push({ message: freezeMessage(context), steering: false }),
  508. ))
  509. }
  510. // Tool results stay adjacent to their calls; input accepted during the
  511. // request enters the log only after the complete result batch.
  512. const steered = this.drainOutbox(turn)
  513. session.append('step/end', { turn, step })
  514. this.stepOpen = false
  515. return {
  516. kind: 'completed',
  517. continueTurn: (toolCalls.length > 0 && !concluded) || steered,
  518. concluded,
  519. maxTokens: finish.kind === 'max-tokens',
  520. }
  521. }
  522. /**
  523. * Compose one frozen request and bind it to the adapter registration that
  524. * resolved its exact-model defaults.
  525. */
  526. private async buildRequest(
  527. turn: number,
  528. step: number,
  529. tools: GenerateOptions['tools'] & object,
  530. system: string,
  531. boundaryMessages: Message[],
  532. signal: AbortSignal,
  533. ): Promise<{ request: GenerateOptions; preparedCall?: PreparedLlmCall }> {
  534. const { session } = this
  535. // A loop instance starts from its declared route, restoring only an opaque
  536. // effort owned by that exact model. Later steps fold the config it logged.
  537. const persistedConfig = session.requestHeader()?.config
  538. const route = { provider: this.options.provider ?? '', model: this.options.model ?? '' }
  539. const reasoningEffort = persistedConfig?.provider === route.provider
  540. && persistedConfig.model === route.model
  541. ? persistedConfig.reasoningEffort
  542. : undefined
  543. const maxTokens = this.options.maxTokens
  544. const seedConfig = deepFreeze(structuredClone(
  545. this.requestHeaderLogged
  546. // eslint-disable-next-line @typescript-eslint/no-non-null-assertion -- the instance logged the header it now folds
  547. ? persistedConfig!
  548. : {
  549. ...route,
  550. ...reasoningEffort === undefined ? {} : { reasoningEffort },
  551. ...maxTokens === undefined ? {} : { maxTokens },
  552. },
  553. ))
  554. const proposedConfig = await this.loopCtx.waterfall(
  555. agentCarrier(this), 'agent/request', this, turn, step, signal,
  556. () => Promise.resolve(seedConfig),
  557. )
  558. signal.throwIfAborted()
  559. if (!proposedConfig.provider || !proposedConfig.model) {
  560. throw new Error(`agent "${this.id}" has no provider/model: set AgentOptions.provider and AgentOptions.model or supply both via the agent/request waterfall`)
  561. }
  562. let config: LlmCallConfig
  563. let preparedCall: PreparedLlmCall | undefined
  564. try {
  565. preparedCall = await this.loopCtx.llm.prepareCall(proposedConfig, signal)
  566. config = preparedCall.config
  567. } catch (error: unknown) {
  568. // A llm/stream listener may own and short-circuit a route with no
  569. // adapter. Terminal dispatch still raises NO_ADAPTER when none does.
  570. if (!(error instanceof LlmError) || error.code !== 'NO_ADAPTER') throw error
  571. config = proposedConfig
  572. }
  573. signal.throwIfAborted()
  574. const header = canonicalHeader({
  575. config,
  576. ...system ? { system } : {},
  577. ...tools.length > 0 ? { tools } : {},
  578. })
  579. const baseline = session.requestHeader()
  580. if (!this.requestHeaderLogged) {
  581. session.append('request/header', { header, reason: baseline === undefined ? 'initial' : 'resume' })
  582. this.requestHeaderLogged = true
  583. } else if (baseline === undefined || !headerEquals(baseline, header)) {
  584. session.append('request/header', { header, reason: 'change' })
  585. }
  586. const request = markAgentLoopRequest(deepFreeze({
  587. ...header.config,
  588. messages: boundaryMessages,
  589. ...header.system !== undefined ? { system: header.system } : {},
  590. ...header.tools !== undefined ? { tools: header.tools } : {},
  591. sessionId: session.id,
  592. signal,
  593. }))
  594. return { request, ...preparedCall === undefined ? {} : { preparedCall } }
  595. }
  596. /** Commit the outbox and report whether it contained steering. */
  597. private drainOutbox(turn: number, limit = this.outbox.length): boolean {
  598. let steered = false
  599. for (const item of this.outbox.splice(0, limit)) {
  600. if (item.steering) {
  601. steered = true
  602. emitAgentEvent(this.loopCtx, this, 'agent/inbox/dequeue', item.message, 'steering')
  603. this.session.append(
  604. 'steering/message',
  605. { turn, message: item.message },
  606. { surfaceOp: 'append' },
  607. )
  608. } else {
  609. this.session.append('user/message', item.message, { surfaceOp: 'append' })
  610. }
  611. }
  612. return steered
  613. }
  614. /**
  615. * Give context-only input its ordinary idle placement when admission
  616. * produces no turn. Steering keeps the whole boundary staged so context
  617. * accepted beside it cannot split from the request it accompanies.
  618. */
  619. private flushRejectedAdmissionContexts(): void {
  620. if (this.outbox.some(item => item.steering)) return
  621. const contexts = this.outbox.splice(0)
  622. for (let index = 0; index < contexts.length; index += 1) {
  623. const item = contexts[index]
  624. /* v8 ignore next 2 -- the steering precheck proves this batch is context-only */
  625. if (item === undefined || item.steering) throw new Error('rejected-admission context batch changed')
  626. try {
  627. this.session.append('user/message', item.message, { surfaceOp: 'append' })
  628. } catch (error: unknown) {
  629. this.outbox.unshift(...contexts.slice(index))
  630. throw error
  631. }
  632. }
  633. }
  634. /**
  635. * The single settlement funnel: classify one turn failure (interruption
  636. * beats error) into the durable turn/end reason and live settlement report.
  637. */
  638. private settle(
  639. turn: number,
  640. step: number,
  641. error: unknown,
  642. signal: AbortSignal,
  643. failure?: LlmFailure,
  644. ): { reason: TurnEndReason; settleReason: SettleReason } {
  645. if (signal.aborted) {
  646. // Slot invariant, stated rather than re-validated: the turn controller
  647. // is machine-private and cancel() is its only aborter, always with one
  648. // frozen canonical cause as the reason.
  649. const interrupt = signal.reason as AgentInterruptReason
  650. return {
  651. reason: { kind: interrupt.kind === 'disposed' ? 'disposed' : 'aborted' },
  652. settleReason: { kind: 'aborted' },
  653. }
  654. }
  655. if (failure !== undefined) {
  656. emitAgentEvent(this.loopCtx, this, 'agent/error', turn, step, error)
  657. // The durable record renders the full cause chain: turn/end is the one
  658. // durable trace of the failure, so a wrapper message alone would lose
  659. // the transport detail the log exists to keep.
  660. const rendered = errorChain(error)
  661. return {
  662. reason: { kind: 'error', step, failure: { ...failure, ...rendered === '<unrenderable value>' ? {} : { message: rendered } } },
  663. settleReason: { kind: 'error', error, failure },
  664. }
  665. }
  666. emitAgentEvent(this.loopCtx, this, 'agent/error', turn, step, error)
  667. return {
  668. reason: { kind: 'error', step, message: errorChain(error), ...isHarnessError(error) ? { code: error.code } : {} },
  669. settleReason: { kind: 'error', error },
  670. }
  671. }
  672. /** Continue with a waking prompt, or publish the idle status. */
  673. private continueOrIdle(): void {
  674. if (this.queued.some(item => item.wakeup)) {
  675. this.kick()
  676. } else {
  677. // Every caller sits inside an admission or run whose install marked the
  678. // interval busy, so the flag is still set here.
  679. this.busy = false
  680. emitAgentEvent(this.loopCtx, this, 'agent/status', 'idle')
  681. }
  682. }
  683. }