agent.ts 42 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015
  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. An idle turn-admission reservation
  6. * can withhold the driver from the queue without touching its contents.
  7. *
  8. * @module dsh-agent-loop/agent
  9. */
  10. import type { Context } from 'cordis'
  11. import { randomUUID } from 'node:crypto'
  12. import { agentCarrier, assembleContextFor, emitAgentEvent, InboxItemId } from '@deepseek-ai/dsh-agent'
  13. import { createScope } from '@deepseek-ai/dsh-scope'
  14. import type { Scope } from '@deepseek-ai/dsh-scope'
  15. import type {
  16. Agent,
  17. CancelOptions,
  18. AgentInterruptReason,
  19. InboxAction,
  20. InboxActionResult,
  21. InboxItem,
  22. InboxItemId as InboxItemIdType,
  23. InboxPlacement,
  24. AgentOptions,
  25. AgentStatus,
  26. SettleReason,
  27. PromptDecision,
  28. RequestError,
  29. RequestErrorAction,
  30. SendOptions,
  31. SteeringOutcome,
  32. SteeringReceipt,
  33. } from '@deepseek-ai/dsh-agent'
  34. import {
  35. BlockAssembler,
  36. LlmError,
  37. assertNever,
  38. createAssistantMessage,
  39. createUserMessage,
  40. deepFreeze,
  41. errorChain,
  42. freezeMessage,
  43. isHarnessError,
  44. llmFailureOf,
  45. llmRetryPolicyOf,
  46. markAgentLoopRequest,
  47. } from '@deepseek-ai/dsh-llm'
  48. import type { GenerateOptions, LlmCallConfig, LlmFailure, Message, PreparedLlmCall, ResolvedRetryPolicy } from '@deepseek-ai/dsh-llm'
  49. import { canonicalHeader, headerEquals } from '@deepseek-ai/dsh-session'
  50. import type { AssistantMessage, EpochHeader, RequestContext, Session, SessionId, TurnEndReason, TurnTrigger, UserMessage } from '@deepseek-ai/dsh-session'
  51. import { renderContextSnapshot, renderPrompt } from '@deepseek-ai/dsh-system-prompt'
  52. import type {} from '@deepseek-ai/dsh-tools'
  53. import { executeToolCalls } from './tool-calls.ts'
  54. /** One completed step or a final-adapter failure eligible for recovery. */
  55. type StepOutcome =
  56. | { kind: 'completed'; continueTurn: boolean; concluded: boolean; maxTokens: boolean }
  57. | { kind: 'request-failed'; error: RequestError; failure: LlmFailure; retryPolicy: ResolvedRetryPolicy | undefined }
  58. /** Internal one-shot controller paired with a public steering receipt. */
  59. interface SteeringDelivery {
  60. readonly receipt: SteeringReceipt
  61. settle(outcome: SteeringOutcome): void
  62. }
  63. /** Create one idempotent steering-admission controller. */
  64. function createSteeringDelivery(): SteeringDelivery {
  65. const { promise, resolve } = Promise.withResolvers<SteeringOutcome>()
  66. let settled = false
  67. return {
  68. receipt: { outcome: promise },
  69. settle(outcome): void {
  70. /* v8 ignore next -- each ownership transfer removes the delivery before another settlement path can reach it. */
  71. if (settled) return
  72. settled = true
  73. resolve(outcome)
  74. },
  75. }
  76. }
  77. const RUNTIME_CONTEXT_SOURCE = '@deepseek-ai/dsh-system-prompt'
  78. /** Clearing marker kept distinct from every prefixed {@link renderContextSnapshot} result. */
  79. const CLEARED_RUNTIME_CONTEXT = 'Current runtime context: none. Earlier runtime-context snapshots no longer apply.'
  80. /** Whether one user message is owned by runtime-context materialization. */
  81. function isRuntimeContextMessage(message: UserMessage): boolean {
  82. return message.source.kind === 'plugin' && message.source.plugin === RUNTIME_CONTEXT_SOURCE
  83. }
  84. /** Latest retained runtime-context snapshot; `found` distinguishes malformed content from absence. */
  85. function retainedRuntimeContext(session: Session): { found: boolean; text: string | undefined } {
  86. const events = session.events
  87. const nodes = session.surface.nodes
  88. for (let index = nodes.length - 1; index >= 0; index -= 1) {
  89. const event = events[nodes[index] as number]
  90. if (event?.type !== 'user/message' || !isRuntimeContextMessage(event.data)) continue
  91. const [block] = event.data.content
  92. return {
  93. found: true,
  94. text: event.data.content.length === 1 && block?.type === 'text' ? block.text : undefined,
  95. }
  96. }
  97. return { found: false, text: undefined }
  98. }
  99. /** Append a full current snapshot only when it changed or compaction removed it. */
  100. function materializeRuntimeContext(session: Session, current: string): void {
  101. const previous = retainedRuntimeContext(session)
  102. if (!previous.found && current.length === 0) {
  103. const compactedPriorSnapshot = session.surface.replaceGeneration > 0
  104. && session.events.some(event => event.type === 'user/message' && isRuntimeContextMessage(event.data))
  105. if (!compactedPriorSnapshot) return
  106. }
  107. const snapshot = current.length === 0 ? CLEARED_RUNTIME_CONTEXT : current
  108. if (previous.text === snapshot) return
  109. session.append('user/message', createUserMessage({
  110. content: [{ type: 'text', text: snapshot }],
  111. source: { kind: 'plugin', plugin: RUNTIME_CONTEXT_SOURCE },
  112. }), { surfaceOp: 'append' })
  113. }
  114. /** Remove adapter-derived values before plugins propose the next request config. */
  115. function requestProposal(header: EpochHeader): LlmCallConfig {
  116. if (header.adapterDefaults === undefined) return header.config
  117. const proposal = { ...header.config }
  118. if (header.adapterDefaults.reasoningEffort === true) delete proposal.reasoningEffort
  119. if (header.adapterDefaults.maxTokens === true) delete proposal.maxTokens
  120. return proposal
  121. }
  122. /**
  123. * The concrete {@link Agent}: each `run()` owns one turn and repeats model
  124. * steps while tools or steering require another request.
  125. */
  126. export class ReactLoopAgent implements Agent {
  127. /** Prompts awaiting individual turns. */
  128. private queued: { item: InboxItem; wakeup: boolean; delivery?: SteeringDelivery }[] = []
  129. /** Input taken into the session log at step boundaries. */
  130. private outbox: { message: UserMessage; steering: boolean; item?: InboxItem; delivery?: SteeringDelivery }[] = []
  131. /** Steering already committed to the log but not yet captured by a request. */
  132. private pendingAdmissions: SteeringDelivery[] = []
  133. /** Whether the active cancellation preserves already committed pending delivery. */
  134. private preservePendingAdmissionsOnAbort = false
  135. /** Whether observers see a running interval; consecutive turns share it. */
  136. private busy = false
  137. /** Whether an idle waking send has deferred driver admission. */
  138. private wakeScheduled = false
  139. /**
  140. * The live idle turn-admission reservation, holding the driver out of the
  141. * queue until its owner releases. It settles idle waiters instead of
  142. * {@link done} so lifecycle teardown never awaits the reserving operation.
  143. */
  144. private admission: { readonly settled: Promise<void>; readonly settle: () => void } | undefined
  145. /** Whether next-step input belongs to the current admission or open turn. */
  146. acceptsNextStep = false
  147. /** Abort owner for the current admission or turn. */
  148. private abort: AbortController | undefined
  149. /** Resolves when the current admission and turn exit. */
  150. done: Promise<void> = Promise.resolve()
  151. /** The agent-scoped registration boundary; the lifecycle owner unwinds it after {@link done}. */
  152. readonly scope: Scope
  153. /** The agent's scoped composition context ({@link Agent.ctx}). */
  154. readonly ctx: Context
  155. /** Last turn number opened by this loop or present in its seeded log. */
  156. private lastTurn: number
  157. /** Whether the session log is owed a matching turn end event. */
  158. private turnOpen = false
  159. private stepOpen = false
  160. /** Whether this loop instance has appended its initial/resume request anchor. */
  161. private requestHeaderLogged = false
  162. constructor(
  163. private loopCtx: Context,
  164. public readonly id: SessionId,
  165. public readonly options: AgentOptions,
  166. public readonly session: Session,
  167. ) {
  168. this.lastTurn = session.events.findLast(event => event.type === 'turn/start')?.data.turn ?? 0
  169. this.scope = createScope(loopCtx, this)
  170. this.ctx = this.scope.ctx.extend({ agent: this })
  171. }
  172. /** Last activity state published to observers. */
  173. get status(): AgentStatus {
  174. return this.busy ? 'running' : 'idle'
  175. }
  176. /** Accept and route one unified send item. */
  177. send(
  178. message: UserMessage,
  179. options: SendOptions,
  180. ): void {
  181. this.route(message, options)
  182. }
  183. /** Route one accepted message, optionally tracking steering admission. */
  184. private route(
  185. message: UserMessage,
  186. options: SendOptions,
  187. delivery?: SteeringDelivery,
  188. ): void {
  189. const { target, wakeup } = options
  190. if (target === 'next-step' && !wakeup) {
  191. if (this.acceptsNextStep) {
  192. this.outbox.push({ message, steering: false })
  193. return
  194. }
  195. this.session.append('user/message', message, { surfaceOp: 'append' })
  196. return
  197. }
  198. const placement: InboxPlacement = target === 'next-step' && this.acceptsNextStep ? 'steering' : 'queued'
  199. const item: InboxItem = Object.freeze({
  200. id: InboxItemId(randomUUID()),
  201. message,
  202. placement,
  203. })
  204. if (placement === 'steering') {
  205. this.outbox.push({ message, steering: true, item, ...delivery === undefined ? {} : { delivery } })
  206. } else {
  207. this.queued.push({ item, wakeup, ...delivery === undefined ? {} : { delivery } })
  208. }
  209. // Preserve the routing decision for every send in this synchronous caller
  210. // stack, while installing quiescence ownership before enqueue observers
  211. // can cancel or dispose.
  212. if (placement === 'queued' && wakeup) this.scheduleKick()
  213. emitAgentEvent(this.loopCtx, this, 'agent/inbox/enqueue', item)
  214. }
  215. /** Apply one synchronous mutation to a still-pending queued occurrence. */
  216. updateInbox(id: InboxItemIdType, action: InboxAction): InboxActionResult {
  217. const queuedIndex = this.queued.findIndex(candidate => candidate.item.id === id)
  218. if (queuedIndex === -1) return 'not-found'
  219. const pending = this.queued[queuedIndex]
  220. /* v8 ignore next -- the index was resolved from this array without an async boundary. */
  221. if (pending === undefined) throw new Error(`agent "${this.id}" queued item disappeared during update`)
  222. /* v8 ignore next -- InboxAction is a closed discriminated union; all variants are covered below. */
  223. switch (action.kind) {
  224. case 'edit': {
  225. const item: InboxItem = Object.freeze({
  226. ...pending.item,
  227. message: freezeMessage({ ...pending.item.message, content: action.content }),
  228. })
  229. this.queued[queuedIndex] = { ...pending, item }
  230. emitAgentEvent(this.loopCtx, this, 'agent/inbox/update', item)
  231. return 'applied'
  232. }
  233. case 'remove': {
  234. this.queued.splice(queuedIndex, 1)
  235. pending.delivery?.settle({ status: 'rejected' })
  236. emitAgentEvent(this.loopCtx, this, 'agent/inbox/discard', [pending.item])
  237. return 'applied'
  238. }
  239. default:
  240. /* v8 ignore next -- InboxAction is a closed discriminated union. */
  241. return assertNever(action)
  242. }
  243. }
  244. /** Queue one ordinary prompt turn and wake the driver. */
  245. followup(input: UserMessage): void {
  246. this.send(input, {
  247. target: 'next-turn',
  248. wakeup: true,
  249. })
  250. }
  251. /** Steer the open turn, falling back to a tracked waking prompt while idle. */
  252. steer(input: UserMessage): SteeringReceipt {
  253. const delivery = createSteeringDelivery()
  254. this.route(input, {
  255. target: 'next-step',
  256. wakeup: true,
  257. }, delivery)
  258. return delivery.receipt
  259. }
  260. /** Append model-facing context without waking the driver. */
  261. inject(input: UserMessage): void {
  262. this.send(input, {
  263. target: 'next-step',
  264. wakeup: false,
  265. })
  266. }
  267. /**
  268. * Hold the idle admission boundary so no queued prompt can open a turn until
  269. * the returned release runs. Later sends keep their ordinary placement and
  270. * `wakeup` facts; only the driver's claim waits.
  271. * @returns the idempotent release, or `undefined` when the driver is active or already committed to waking work.
  272. */
  273. reserveTurnAdmission(): (() => void) | undefined {
  274. // `busy` covers every abort owner: kick() and run() mark the interval
  275. // running before they install one. `wakeScheduled` is the same-tick state
  276. // of an accepted waking prompt whose claim is still a pending microtask.
  277. if (this.busy || this.wakeScheduled || this.admission !== undefined
  278. || this.queued.some(item => item.wakeup)) return undefined
  279. const pending = Promise.withResolvers<void>()
  280. const reservation = { settled: pending.promise, settle: pending.resolve }
  281. this.admission = reservation
  282. return () => {
  283. // Idempotent, and inert once a later reservation owns the boundary.
  284. if (this.admission !== reservation) return
  285. this.admission = undefined
  286. // Re-arm the ordinary path first, so an idle waiter released below
  287. // re-reads live admission activity instead of settled state.
  288. if (this.queued.some(item => item.wakeup)) this.scheduleKick()
  289. reservation.settle()
  290. }
  291. }
  292. /**
  293. * Clear all pending work and abort the active turn; the first cause wins.
  294. * The cause is signal payload for observers and the durable turn/end
  295. * classification — it selects no machine behavior. Teardown is just
  296. * `cancel({kind:'disposed'})` + await {@link done} + {@link scope} dispose,
  297. * all owned by the factory.
  298. */
  299. cancel(cause: AgentInterruptReason, options: CancelOptions = {}): void {
  300. // Effective only when it aborts the active turn or actually discards
  301. // pending work: a keepInbox call with no active turn is a documented
  302. // no-op, so it must not emit cancel-requested for consumers to misread.
  303. const discards = !options.keepInbox && (this.queued.length > 0 || this.outbox.length > 0)
  304. if (this.abort !== undefined || discards) {
  305. // Observe-only: coordination consumers update their state before the
  306. // inboxes clear; listener failures are contained by the dispatcher.
  307. if (cause.kind !== 'disposed') emitAgentEvent(this.loopCtx, this, 'agent/cancel-requested', cause)
  308. }
  309. if (options.keepInbox && this.abort !== undefined) this.preservePendingAdmissionsOnAbort = true
  310. if (!options.keepInbox) {
  311. const discarded = this.queued.map(item => item.item)
  312. for (const item of this.queued) item.delivery?.settle({ status: 'rejected' })
  313. for (const item of this.outbox) {
  314. if (item.steering && item.item !== undefined) {
  315. item.delivery?.settle({ status: 'rejected' })
  316. discarded.push(item.item)
  317. }
  318. }
  319. this.rejectPendingAdmissions()
  320. // Clear before abort observers run: replacement work belongs to the next turn.
  321. this.queued.length = 0
  322. this.outbox.length = 0
  323. if (discarded.length > 0) emitAgentEvent(this.loopCtx, this, 'agent/inbox/discard', discarded)
  324. }
  325. const reason = Object.freeze({ kind: cause.kind })
  326. this.abort?.abort(reason)
  327. }
  328. /** Resolve at idle quiescence: no run driving and no waking prompt waiting. */
  329. async whenIdle(): Promise<void> {
  330. while (true) {
  331. // `done` is replaced per activity, so re-reading it follows chained turns.
  332. // Every driver failure today is contained before it can reject `done`,
  333. // but the waiter must not gamble quiescence on that: a future escape
  334. // still counts as settled activity.
  335. /* v8 ignore next 3 -- the catch arm backstops rejection paths that are all currently contained */
  336. while (this.busy || this.wakeScheduled || this.abort !== undefined || this.runnableWakingQueued) {
  337. await this.done.catch(() => undefined)
  338. }
  339. // A reservation is unfinished activity even with an empty queue, and a
  340. // prompt it withholds is not quiescent — but `done` never owns it, so
  341. // waiting on the queue alone would spin on an already-settled promise.
  342. const reservation = this.admission
  343. if (reservation === undefined) return
  344. await reservation.settled
  345. }
  346. }
  347. /** Whether a queued waking prompt may claim the driver now. */
  348. private get runnableWakingQueued(): boolean {
  349. return this.admission === undefined && this.queued.some(item => item.wakeup)
  350. }
  351. /** Defer idle admission while keeping {@link done} as its quiescence owner. */
  352. private scheduleKick(): void {
  353. // A held reservation keeps the item queued with no scheduled claim; its
  354. // release re-arms this path for whatever is queued by then.
  355. if (this.abort !== undefined || this.wakeScheduled || this.admission !== undefined) return
  356. this.wakeScheduled = true
  357. const pending = Promise.withResolvers<void>()
  358. const scheduled = pending.promise
  359. queueMicrotask(() => {
  360. this.wakeScheduled = false
  361. this.kick()
  362. const activity = this.done
  363. if (activity === scheduled) {
  364. pending.resolve()
  365. } else {
  366. void activity.then(
  367. () => { pending.resolve() },
  368. () => { pending.resolve() },
  369. )
  370. }
  371. })
  372. this.done = scheduled
  373. }
  374. /** Claim and admit the next queued prompt, then start its turn. */
  375. private kick(): void {
  376. if (this.abort !== undefined || !this.runnableWakingQueued) return
  377. // The some() guard above proves the queue is non-empty; the non-null
  378. // assertion expresses that invariant.
  379. // oxlint-disable-next-line typescript/no-non-null-assertion
  380. const pending = this.queued.shift()!
  381. const { item, delivery } = pending
  382. const { message } = item
  383. const inheritedOutboxLength = this.outbox.length
  384. const admission = new AbortController()
  385. this.abort = admission
  386. this.acceptsNextStep = true
  387. // Claimed admission is part of the running interval: it is cancellable
  388. // activity, so observers (and their cancel routing) must see it.
  389. if (!this.busy) {
  390. this.busy = true
  391. emitAgentEvent(this.loopCtx, this, 'agent/status', 'running')
  392. }
  393. // The admission body runs synchronously up to the prompt-submit
  394. // waterfall's first await, so the waterfall snapshots its listeners
  395. // before a disposal initiated by the running-status emit above can
  396. // unregister a vetoing plugin.
  397. this.done = this.loopCtx.agents.withInitiator(this, async () => {
  398. const signal = admission.signal
  399. const trigger: TurnTrigger = { kind: 'message', source: message.source }
  400. // Admitted input stays on the stack until its turn/start commits: the
  401. // turn owns it only once the turn exists in the log.
  402. let admitted: UserMessage[] | undefined
  403. try {
  404. signal.throwIfAborted()
  405. const decision = await this.loopCtx.waterfall(
  406. agentCarrier(this), 'agent/prompt-submit', this, message, signal,
  407. () => Promise.resolve<PromptDecision>({ kind: 'allow' }),
  408. )
  409. signal.throwIfAborted()
  410. if (decision.kind === 'allow') {
  411. admitted = [decision.content === undefined
  412. ? message
  413. : freezeMessage({ ...message, content: decision.content })]
  414. for (const context of decision.additionalContexts ?? []) {
  415. admitted.push(freezeMessage(context))
  416. }
  417. }
  418. } catch (error: unknown) {
  419. if (!signal.aborted) {
  420. this.loopCtx.logger.warn(`agent "${this.id}": prompt admission failed: ${errorChain(error)}`)
  421. }
  422. }
  423. // cancel() aborts but never clears the slot, and kick()/run()
  424. // all refuse to install a new owner while one exists, so the admission
  425. // still owns the slot here and releasing it unconditionally is exact.
  426. this.abort = undefined
  427. if (admitted === undefined) {
  428. delivery?.settle({ status: 'rejected' })
  429. this.acceptsNextStep = false
  430. try {
  431. this.flushRejectedAdmissionContexts()
  432. } catch (error: unknown) {
  433. // No turn exists for agent/error coordinates. Preserve the
  434. // uncommitted suffix for a later boundary and report locally.
  435. this.loopCtx.logger.warn(
  436. `agent "${this.id}": committing rejected-admission context failed: ${errorChain(error)}`,
  437. )
  438. }
  439. // A synchronously aborted admission would otherwise publish idle
  440. // inside send()'s own synchronous extent, before any post-send
  441. // subscriber could observe the transition.
  442. await Promise.resolve()
  443. this.continueOrIdle()
  444. return
  445. }
  446. await this.run(trigger, admitted, inheritedOutboxLength, Object.freeze([]), delivery)
  447. })
  448. // Published only after the abort owner and pending done are installed: a
  449. // dequeue listener that cancels or disposes must find live cancellation
  450. // and quiescence ownership, not the previous activity's settled state.
  451. emitAgentEvent(this.loopCtx, this, 'agent/inbox/dequeue', item)
  452. }
  453. /**
  454. * Run one turn and any request-error retry. `admitted` input enters the log
  455. * only after `turn/start` commits; until then it has no owner state to unwind.
  456. */
  457. private async run(
  458. trigger: TurnTrigger,
  459. admitted: UserMessage[] = [],
  460. inheritedOutboxLength = 0,
  461. priorFailures: readonly LlmFailure[] = Object.freeze([]),
  462. promptDelivery?: SteeringDelivery,
  463. ): Promise<void> {
  464. // Both entries hold the invariant: kick() clears the admission slot before
  465. // awaiting run(), and a retry is entered only after the prior run clears it.
  466. /* v8 ignore next -- unreachable guard: every caller clears or checks the abort slot first */
  467. if (this.abort !== undefined) throw new Error(`agent "${this.id}" is already running`)
  468. const controller = new AbortController()
  469. this.abort = controller
  470. this.preservePendingAdmissionsOnAbort = false
  471. this.acceptsNextStep = true
  472. const signal = controller.signal
  473. const turn = this.lastTurn + 1
  474. let step = 0
  475. let opened = false
  476. let reason: TurnEndReason = { kind: 'completed' }
  477. let settleReason: SettleReason = { kind: 'completed' }
  478. let requestFailureHistory = priorFailures
  479. let retryFailures: readonly LlmFailure[] | undefined
  480. const cancelRetry = (): void => { retryFailures = undefined }
  481. signal.addEventListener('abort', cancelRetry, { once: true })
  482. try {
  483. signal.throwIfAborted()
  484. this.session.append('turn/start', { turn, trigger })
  485. // Committed: publish the turn to the machine's own bookkeeping and let
  486. // the admitted input enter the log it now belongs to.
  487. this.turnOpen = true
  488. opened = true
  489. this.lastTurn = turn
  490. // Context or steering retained by an earlier rejected admission happened
  491. // before this prompt and must occupy the same order in durable history.
  492. this.drainOutbox(turn, inheritedOutboxLength)
  493. if (promptDelivery !== undefined) this.pendingAdmissions.push(promptDelivery)
  494. for (const input of admitted) {
  495. this.session.append('user/message', input, { surfaceOp: 'append' })
  496. }
  497. signal.throwIfAborted()
  498. steps: while (true) {
  499. step += 1
  500. const outcome = await this.step(turn, step, signal)
  501. switch (outcome.kind) {
  502. case 'completed':
  503. requestFailureHistory = Object.freeze([])
  504. if (outcome.maxTokens) reason = { kind: 'max-tokens' }
  505. // A concluding tool result is terminal: reject steering that did
  506. // not enter a request, while retaining same-boundary context in
  507. // durable history before the turn closes.
  508. if (outcome.concluded) {
  509. this.discardOutboxSteering()
  510. this.drainOutbox(turn)
  511. break steps
  512. }
  513. /* v8 ignore next -- step() folded the same steering predicate into continueTurn immediately before returning. */
  514. if (outcome.continueTurn || this.outbox.some(item => item.steering)) continue
  515. break
  516. case 'request-failed': {
  517. // step() reports request failures only after step/start commits
  518. // and before its own step/end, so the step is always open here.
  519. this.stepOpen = false
  520. this.session.append('step/end', { turn, step })
  521. if (!signal.aborted) {
  522. try {
  523. const action = await this.loopCtx.waterfall(
  524. agentCarrier(this), 'agent/request-error', this, turn, step, outcome.error,
  525. outcome.failure, requestFailureHistory, outcome.retryPolicy, signal,
  526. () => Promise.resolve<RequestErrorAction>(undefined),
  527. )
  528. // oxlint-disable-next-line typescript/no-unnecessary-condition -- signal can abort while recovery is awaited.
  529. if (action?.kind === 'retry' && !signal.aborted) {
  530. retryFailures = Object.freeze([...requestFailureHistory, outcome.failure])
  531. }
  532. } catch (recoveryError: unknown) {
  533. this.loopCtx.logger.warn(
  534. `agent "${this.id}": request recovery failed at turn ${turn}, step ${step}: ${errorChain(recoveryError)}`,
  535. )
  536. }
  537. }
  538. const settlement = this.settle(turn, step, outcome.error, signal, outcome.failure)
  539. reason = settlement.reason
  540. settleReason = settlement.settleReason
  541. break steps
  542. }
  543. /* v8 ignore next 2 -- closed-union exhaustiveness guard */
  544. default:
  545. assertNever(outcome)
  546. }
  547. await this.loopCtx.serial(agentCarrier(this), 'agent/turn-stopping', this, turn, signal)
  548. signal.throwIfAborted()
  549. this.drainOutboxContexts()
  550. if (!this.outbox.some(item => item.steering)) {
  551. break
  552. }
  553. }
  554. } catch (caught: unknown) {
  555. try {
  556. if (this.stepOpen) {
  557. this.stepOpen = false
  558. this.session.append('step/end', { turn, step })
  559. }
  560. } catch (closeError: unknown) {
  561. // Contained like the finally's turn close: a persistently rejecting
  562. // step boundary must not escape run(), or the post-finally tail would
  563. // never publish the terminal status and observers would see a
  564. // permanently running agent whose whenIdle() already resolved.
  565. this.loopCtx.logger.warn(`agent "${this.id}": closing step ${turn}/${step} failed: ${errorChain(closeError)}`)
  566. emitAgentEvent(this.loopCtx, this, 'agent/error', turn, step, closeError)
  567. }
  568. ({ reason, settleReason } = this.settle(turn, step, caught, signal))
  569. } finally {
  570. // Every step-close happens before this point on both success and
  571. // failure paths (step(), the request-failed branch, the catch), so the
  572. // finally owes only the turn boundary.
  573. this.acceptsNextStep = false
  574. try {
  575. if (this.turnOpen) {
  576. // Re-entrant turn/end listeners must route new input to a later turn.
  577. this.turnOpen = false
  578. this.session.append('turn/end', { turn, reason })
  579. }
  580. } catch (error: unknown) {
  581. retryFailures = undefined
  582. this.loopCtx.logger.warn(`agent "${this.id}": closing turn ${turn} failed: ${errorChain(error)}`)
  583. emitAgentEvent(this.loopCtx, this, 'agent/error', turn, step, error)
  584. }
  585. // cancel() aborts but never clears the slot, and no second run can
  586. // install a controller while this one is still unwinding, so the slot
  587. // is still this run's controller here.
  588. this.abort = undefined
  589. signal.removeEventListener('abort', cancelRetry)
  590. const preservePending = signal.aborted && this.preservePendingAdmissionsOnAbort
  591. this.preservePendingAdmissionsOnAbort = false
  592. // oxlint-disable-next-line typescript/no-unnecessary-condition -- keepInbox cancellation can set this while turn work is awaited.
  593. if (!preservePending) this.rejectPendingAdmissions()
  594. }
  595. if (opened) {
  596. try {
  597. await this.loopCtx.sessions.flush(this.session)
  598. } catch (error: unknown) {
  599. this.loopCtx.logger.warn(`agent "${this.id}": session/flush failed at turn ${turn}: ${errorChain(error)}`)
  600. emitAgentEvent(this.loopCtx, this, 'agent/error', turn, step, error)
  601. }
  602. }
  603. if (retryFailures !== undefined) {
  604. await this.run({ kind: 'retry' }, [], 0, retryFailures)
  605. } else {
  606. // agent/settled names only committed turns: a run aborted or rejected
  607. // before turn/start has no durable turn/end for consumers to settle
  608. // against, so it exits without the notification.
  609. if (opened) emitAgentEvent(this.loopCtx, this, 'agent/settled', turn, settleReason)
  610. this.continueOrIdle()
  611. }
  612. }
  613. /**
  614. * Run the `agent/step` extension point, commit pending input, derive one
  615. * request, and execute its tool calls inside one durable step boundary.
  616. */
  617. private async step(
  618. turn: number,
  619. step: number,
  620. signal: AbortSignal,
  621. ): Promise<StepOutcome> {
  622. const { session } = this
  623. // The single between-steps extension point: listeners inject, steer, or
  624. // edit the log here; the request derives from the log after this settles.
  625. await this.loopCtx.serial(agentCarrier(this), 'agent/step', this, turn, step, signal)
  626. signal.throwIfAborted()
  627. // Assemble request-owned prompt inputs fresh each step. Dynamic context is
  628. // committed at the tail before deriving history once, preserving the stable
  629. // system/history cache prefix while keeping every model-visible byte logged.
  630. const assembly = await this.loopCtx.systemPrompt.assemble(assembleContextFor(this, signal))
  631. signal.throwIfAborted()
  632. const system = renderPrompt(assembly)
  633. materializeRuntimeContext(session, renderContextSnapshot(assembly))
  634. // Commit the exact pending batch only after every asynchronous
  635. // pre-request contribution succeeded. Input accepted after this splice
  636. // remains pending for a later request.
  637. this.drainOutbox(turn)
  638. // Snapshot the exact log prefix: the reconstruction boundary. Appends
  639. // after this synchronous snapshot join the next request.
  640. const boundaryMessages = session.deriveMessages()
  641. session.append('step/start', { turn, step })
  642. this.stepOpen = true
  643. this.admitPendingAdmissions(turn, step)
  644. signal.throwIfAborted()
  645. const { request, preparedCall } = await this.buildRequest(
  646. turn, step, assembly.tools, system, boundaryMessages, signal,
  647. )
  648. const assembler = new BlockAssembler()
  649. const chunkSeqs: number[] = []
  650. const stream = preparedCall?.stream(request) ?? this.loopCtx.llm.stream(request)
  651. try {
  652. for await (const chunk of stream) {
  653. signal.throwIfAborted()
  654. const chunkEvent = session.append('assistant/chunk', { turn, step, chunk })
  655. chunkSeqs.push(chunkEvent.seq)
  656. assembler.push(chunk)
  657. }
  658. } catch (error: unknown) {
  659. const facts = llmFailureOf(stream, error)
  660. if (facts !== undefined && error instanceof Error) {
  661. return { kind: 'request-failed', error, failure: facts, retryPolicy: llmRetryPolicyOf(stream) }
  662. }
  663. throw error
  664. }
  665. signal.throwIfAborted()
  666. // Failure finish chunks take the same path as thrown stream errors.
  667. const finish = assembler.finish
  668. if (finish.kind === 'error' || finish.kind === 'aborted') {
  669. const error = new LlmError(finish.failure.message, finish.failure.code, finish.failure)
  670. return { kind: 'request-failed', error, failure: finish.failure, retryPolicy: llmRetryPolicyOf(stream) }
  671. }
  672. // Truncated (max-tokens) output cannot owe tool calls.
  673. const assembled = assembler.blocks()
  674. const content = finish.kind === 'max-tokens'
  675. ? assembled.filter(block => block.type !== 'tool-call')
  676. : assembled
  677. const message: AssistantMessage = createAssistantMessage({
  678. content,
  679. source: {
  680. provider: request.provider,
  681. model: request.model,
  682. ...assembler.replayState !== undefined ? { replayState: assembler.replayState } : {},
  683. },
  684. })
  685. session.append(
  686. 'assistant/message',
  687. {
  688. turn,
  689. step,
  690. message,
  691. ...assembler.usage === undefined ? {} : { usage: assembler.usage },
  692. },
  693. { surfaceOp: 'append', sourceEventSeqs: chunkSeqs },
  694. )
  695. const toolCalls = content.filter(block => block.type === 'tool-call')
  696. let concluded = false
  697. if (toolCalls.length > 0) {
  698. ({ concluded } = await executeToolCalls(
  699. this.loopCtx, turn, step, toolCalls, signal,
  700. context => this.outbox.push({ message: freezeMessage(context), steering: false }),
  701. ))
  702. }
  703. // Ordinary context keeps the base loop's result-adjacent commit point.
  704. // Steering remains provisional until the next request snapshot admits it.
  705. this.drainOutboxContexts()
  706. session.append('step/end', { turn, step })
  707. this.stepOpen = false
  708. return {
  709. kind: 'completed',
  710. continueTurn: (toolCalls.length > 0 && !concluded) || this.outbox.some(item => item.steering),
  711. concluded,
  712. maxTokens: finish.kind === 'max-tokens',
  713. }
  714. }
  715. /**
  716. * Compose one frozen request and bind it to the adapter registration that
  717. * resolved its exact-model defaults.
  718. */
  719. private async buildRequest(
  720. turn: number,
  721. step: number,
  722. tools: GenerateOptions['tools'] & object,
  723. system: string,
  724. boundaryMessages: Message[],
  725. signal: AbortSignal,
  726. ): Promise<{ request: GenerateOptions; preparedCall?: PreparedLlmCall }> {
  727. const { session } = this
  728. // A loop instance starts from its declared route, restoring only an explicit
  729. // effort owned by that exact model. Later steps re-resolve marked defaults.
  730. const persistedHeader = session.requestHeader()
  731. const persistedConfig = persistedHeader?.config
  732. const route = { provider: this.options.provider ?? '', model: this.options.model ?? '' }
  733. const reasoningEffort = persistedConfig?.provider === route.provider
  734. && persistedConfig.model === route.model
  735. && persistedHeader?.adapterDefaults?.reasoningEffort !== true
  736. ? persistedConfig.reasoningEffort
  737. : undefined
  738. const maxTokens = this.options.maxTokens
  739. const seedConfig = deepFreeze(structuredClone(
  740. this.requestHeaderLogged
  741. // oxlint-disable-next-line typescript/no-non-null-assertion -- the instance logged the header it now folds
  742. ? requestProposal(persistedHeader!)
  743. : {
  744. ...route,
  745. ...reasoningEffort === undefined ? {} : { reasoningEffort },
  746. ...maxTokens === undefined ? {} : { maxTokens },
  747. },
  748. ))
  749. const proposedConfig = await this.loopCtx.waterfall(
  750. agentCarrier(this), 'agent/request', this, turn, step, signal,
  751. () => Promise.resolve(seedConfig),
  752. )
  753. signal.throwIfAborted()
  754. if (!proposedConfig.provider || !proposedConfig.model) {
  755. throw new Error(`agent "${this.id}" has no provider/model: set AgentOptions.provider and AgentOptions.model or supply both via the agent/request waterfall`)
  756. }
  757. let config: LlmCallConfig
  758. let preparedCall: PreparedLlmCall | undefined
  759. try {
  760. preparedCall = await this.loopCtx.llm.prepareCall(proposedConfig, signal)
  761. config = preparedCall.config
  762. } catch (error: unknown) {
  763. // A llm/stream listener may own and short-circuit a route with no
  764. // adapter. Terminal dispatch still raises NO_ADAPTER when none does.
  765. if (!(error instanceof LlmError) || error.code !== 'NO_ADAPTER') throw error
  766. config = proposedConfig
  767. }
  768. signal.throwIfAborted()
  769. const header = canonicalHeader({
  770. config,
  771. ...preparedCall === undefined ? {} : { adapterDefaults: preparedCall.adapterDefaults },
  772. ...system ? { system } : {},
  773. ...tools.length > 0 ? { tools } : {},
  774. })
  775. const baseline = session.requestHeader()
  776. if (!this.requestHeaderLogged) {
  777. session.append('request/header', { header, reason: baseline === undefined ? 'initial' : 'resume' })
  778. this.requestHeaderLogged = true
  779. } else if (baseline === undefined || !headerEquals(baseline, header)) {
  780. session.append('request/header', { header, reason: 'change' })
  781. }
  782. // TODO: This looks like code smell.
  783. // Context metadata for the route this request resolved to, recorded from the same
  784. // registration-bound lookup that prepared the call (no second resolve).
  785. // A route with unknown capacity is still recorded so it clears any older
  786. // denominator; an unchanged route logs nothing.
  787. const contextWindow = preparedCall?.context?.contextWindow
  788. const requestContext: RequestContext = {
  789. provider: config.provider,
  790. model: config.model,
  791. ...contextWindow === undefined ? {} : { contextWindow },
  792. }
  793. const previous = session.requestContext()
  794. if (previous?.provider !== requestContext.provider
  795. || previous.model !== requestContext.model
  796. || previous.contextWindow !== requestContext.contextWindow) {
  797. session.append('request/context', requestContext)
  798. }
  799. const request = markAgentLoopRequest(deepFreeze({
  800. ...header.config,
  801. messages: boundaryMessages,
  802. ...header.system !== undefined ? { system: header.system } : {},
  803. ...header.tools !== undefined ? { tools: header.tools } : {},
  804. sessionId: session.id,
  805. signal,
  806. }))
  807. return { request, ...preparedCall === undefined ? {} : { preparedCall } }
  808. }
  809. /** Commit one stable outbox prefix and retain tracked delivery until snapshot admission. */
  810. private drainOutbox(turn: number, limit = this.outbox.length): void {
  811. const batch = this.outbox.splice(0, limit)
  812. for (let index = 0; index < batch.length; index += 1) {
  813. const item = batch[index]
  814. /* v8 ignore next -- the index walks the exact array length. */
  815. if (item === undefined) throw new Error(`agent "${this.id}" outbox item disappeared during drain`)
  816. try {
  817. if (item.steering) {
  818. /* v8 ignore next -- only inbox-backed steer entries carry steering:true. */
  819. if (item.item === undefined) throw new Error(`agent "${this.id}" steering outbox item has no inbox identity`)
  820. emitAgentEvent(this.loopCtx, this, 'agent/inbox/dequeue', item.item)
  821. this.session.append(
  822. 'steering/message',
  823. { turn, message: item.message },
  824. { surfaceOp: 'append' },
  825. )
  826. if (item.delivery !== undefined) this.pendingAdmissions.push(item.delivery)
  827. } else {
  828. this.session.append('user/message', item.message, { surfaceOp: 'append' })
  829. }
  830. } catch (error: unknown) {
  831. item.delivery?.settle({ status: 'rejected' })
  832. this.outbox.unshift(...batch.slice(item.steering ? index + 1 : index))
  833. throw error
  834. }
  835. }
  836. }
  837. /** Commit ordinary context while retaining provisional steering in order. */
  838. private drainOutboxContexts(): void {
  839. const pending = this.outbox
  840. this.outbox = []
  841. for (let index = 0; index < pending.length; index += 1) {
  842. const item = pending[index]
  843. /* v8 ignore next -- the index walks the exact array length. */
  844. if (item === undefined) throw new Error(`agent "${this.id}" outbox item disappeared during context drain`)
  845. if (item.steering) {
  846. this.outbox.push(item)
  847. continue
  848. }
  849. try {
  850. this.session.append('user/message', item.message, { surfaceOp: 'append' })
  851. } catch (error: unknown) {
  852. this.outbox.push(...pending.slice(index))
  853. throw error
  854. }
  855. }
  856. }
  857. /** Settle every committed steering item captured by this immutable request. */
  858. private admitPendingAdmissions(turn: number, step: number): void {
  859. const outcome: SteeringOutcome = { status: 'admitted', turn, step }
  860. for (const delivery of this.pendingAdmissions.splice(0)) delivery.settle(outcome)
  861. }
  862. /** Reject committed steering that left the inbox without reaching a request. */
  863. private rejectPendingAdmissions(): void {
  864. for (const delivery of this.pendingAdmissions.splice(0)) delivery.settle({ status: 'rejected' })
  865. }
  866. /** Discard uncommitted steering while retaining same-boundary injected context. */
  867. private discardOutboxSteering(): void {
  868. const contexts: typeof this.outbox = []
  869. const discarded: InboxItem[] = []
  870. for (const item of this.outbox) {
  871. if (!item.steering) {
  872. contexts.push(item)
  873. continue
  874. }
  875. item.delivery?.settle({ status: 'rejected' })
  876. /* v8 ignore next -- only inbox-backed steer entries carry steering:true. */
  877. if (item.item === undefined) throw new Error(`agent "${this.id}" steering outbox item has no inbox identity`)
  878. discarded.push(item.item)
  879. }
  880. this.outbox = contexts
  881. if (discarded.length > 0) emitAgentEvent(this.loopCtx, this, 'agent/inbox/discard', discarded)
  882. }
  883. /**
  884. * Give context-only input its ordinary idle placement when admission
  885. * produces no turn. Steering keeps the whole boundary staged so context
  886. * accepted beside it cannot split from the request it accompanies.
  887. */
  888. private flushRejectedAdmissionContexts(): void {
  889. if (this.outbox.some(item => item.steering)) return
  890. const contexts = this.outbox.splice(0)
  891. for (let index = 0; index < contexts.length; index += 1) {
  892. const item = contexts[index]
  893. /* v8 ignore next 2 -- the steering precheck proves this batch is context-only */
  894. if (item === undefined || item.steering) throw new Error('rejected-admission context batch changed')
  895. try {
  896. this.session.append('user/message', item.message, { surfaceOp: 'append' })
  897. } catch (error: unknown) {
  898. this.outbox.unshift(...contexts.slice(index))
  899. throw error
  900. }
  901. }
  902. }
  903. /**
  904. * The single settlement funnel: classify one turn failure (interruption
  905. * beats error) into the durable turn/end reason and live settlement report.
  906. */
  907. private settle(
  908. turn: number,
  909. step: number,
  910. error: unknown,
  911. signal: AbortSignal,
  912. failure?: LlmFailure,
  913. ): { reason: TurnEndReason; settleReason: SettleReason } {
  914. if (signal.aborted) {
  915. // Slot invariant, stated rather than re-validated: the turn controller
  916. // is machine-private and cancel() is its only aborter, always with one
  917. // frozen canonical cause as the reason.
  918. const interrupt = signal.reason as AgentInterruptReason
  919. return {
  920. reason: { kind: interrupt.kind === 'disposed' ? 'disposed' : 'aborted' },
  921. settleReason: { kind: 'aborted' },
  922. }
  923. }
  924. if (failure !== undefined) {
  925. emitAgentEvent(this.loopCtx, this, 'agent/error', turn, step, error)
  926. // The durable record renders the full cause chain: turn/end is the one
  927. // durable trace of the failure, so a wrapper message alone would lose
  928. // the transport detail the log exists to keep.
  929. const rendered = errorChain(error)
  930. return {
  931. reason: { kind: 'error', step, failure: { ...failure, ...rendered === '<unrenderable value>' ? {} : { message: rendered } } },
  932. settleReason: { kind: 'error', error, failure },
  933. }
  934. }
  935. emitAgentEvent(this.loopCtx, this, 'agent/error', turn, step, error)
  936. return {
  937. reason: { kind: 'error', step, message: errorChain(error), ...isHarnessError(error) ? { code: error.code } : {} },
  938. settleReason: { kind: 'error', error },
  939. }
  940. }
  941. /** Continue with a waking prompt, or publish the idle status. */
  942. private continueOrIdle(): void {
  943. if (this.runnableWakingQueued) {
  944. this.kick()
  945. } else {
  946. // Every caller sits inside an admission or run whose install marked the
  947. // interval busy, so the flag is still set here.
  948. this.busy = false
  949. emitAgentEvent(this.loopCtx, this, 'agent/status', 'idle')
  950. }
  951. }
  952. }