agent.ts 43 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036
  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. case 'steer': {
  240. if (!this.acceptsNextStep) return 'steer-unavailable'
  241. this.queued.splice(queuedIndex, 1)
  242. const item: InboxItem = Object.freeze({
  243. id: InboxItemId(randomUUID()),
  244. message: pending.item.message,
  245. placement: 'steering',
  246. })
  247. this.outbox.push({
  248. message: item.message,
  249. steering: true,
  250. item,
  251. ...pending.delivery === undefined ? {} : { delivery: pending.delivery },
  252. })
  253. // Publish the replacement only after it is owned by the outbox. Its
  254. // enqueue precedes the old occurrence's discard so reentrant
  255. // cancellation can terminally account for both occurrences.
  256. emitAgentEvent(this.loopCtx, this, 'agent/inbox/enqueue', item)
  257. emitAgentEvent(this.loopCtx, this, 'agent/inbox/discard', [pending.item])
  258. return 'applied'
  259. }
  260. default:
  261. /* v8 ignore next -- InboxAction is a closed discriminated union. */
  262. return assertNever(action)
  263. }
  264. }
  265. /** Queue one ordinary prompt turn and wake the driver. */
  266. followup(input: UserMessage): void {
  267. this.send(input, {
  268. target: 'next-turn',
  269. wakeup: true,
  270. })
  271. }
  272. /** Steer the open turn, falling back to a tracked waking prompt while idle. */
  273. steer(input: UserMessage): SteeringReceipt {
  274. const delivery = createSteeringDelivery()
  275. this.route(input, {
  276. target: 'next-step',
  277. wakeup: true,
  278. }, delivery)
  279. return delivery.receipt
  280. }
  281. /** Append model-facing context without waking the driver. */
  282. inject(input: UserMessage): void {
  283. this.send(input, {
  284. target: 'next-step',
  285. wakeup: false,
  286. })
  287. }
  288. /**
  289. * Hold the idle admission boundary so no queued prompt can open a turn until
  290. * the returned release runs. Later sends keep their ordinary placement and
  291. * `wakeup` facts; only the driver's claim waits.
  292. * @returns the idempotent release, or `undefined` when the driver is active or already committed to waking work.
  293. */
  294. reserveTurnAdmission(): (() => void) | undefined {
  295. // `busy` covers every abort owner: kick() and run() mark the interval
  296. // running before they install one. `wakeScheduled` is the same-tick state
  297. // of an accepted waking prompt whose claim is still a pending microtask.
  298. if (this.busy || this.wakeScheduled || this.admission !== undefined
  299. || this.queued.some(item => item.wakeup)) return undefined
  300. const pending = Promise.withResolvers<void>()
  301. const reservation = { settled: pending.promise, settle: pending.resolve }
  302. this.admission = reservation
  303. return () => {
  304. // Idempotent, and inert once a later reservation owns the boundary.
  305. if (this.admission !== reservation) return
  306. this.admission = undefined
  307. // Re-arm the ordinary path first, so an idle waiter released below
  308. // re-reads live admission activity instead of settled state.
  309. if (this.queued.some(item => item.wakeup)) this.scheduleKick()
  310. reservation.settle()
  311. }
  312. }
  313. /**
  314. * Clear all pending work and abort the active turn; the first cause wins.
  315. * The cause is signal payload for observers and the durable turn/end
  316. * classification — it selects no machine behavior. Teardown is just
  317. * `cancel({kind:'disposed'})` + await {@link done} + {@link scope} dispose,
  318. * all owned by the factory.
  319. */
  320. cancel(cause: AgentInterruptReason, options: CancelOptions = {}): void {
  321. // Effective only when it aborts the active turn or actually discards
  322. // pending work: a keepInbox call with no active turn is a documented
  323. // no-op, so it must not emit cancel-requested for consumers to misread.
  324. const discards = !options.keepInbox && (this.queued.length > 0 || this.outbox.length > 0)
  325. if (this.abort !== undefined || discards) {
  326. // Observe-only: coordination consumers update their state before the
  327. // inboxes clear; listener failures are contained by the dispatcher.
  328. if (cause.kind !== 'disposed') emitAgentEvent(this.loopCtx, this, 'agent/cancel-requested', cause)
  329. }
  330. if (options.keepInbox && this.abort !== undefined) this.preservePendingAdmissionsOnAbort = true
  331. if (!options.keepInbox) {
  332. const discarded = this.queued.map(item => item.item)
  333. for (const item of this.queued) item.delivery?.settle({ status: 'rejected' })
  334. for (const item of this.outbox) {
  335. if (item.steering && item.item !== undefined) {
  336. item.delivery?.settle({ status: 'rejected' })
  337. discarded.push(item.item)
  338. }
  339. }
  340. this.rejectPendingAdmissions()
  341. // Clear before abort observers run: replacement work belongs to the next turn.
  342. this.queued.length = 0
  343. this.outbox.length = 0
  344. if (discarded.length > 0) emitAgentEvent(this.loopCtx, this, 'agent/inbox/discard', discarded)
  345. }
  346. const reason = Object.freeze({ kind: cause.kind })
  347. this.abort?.abort(reason)
  348. }
  349. /** Resolve at idle quiescence: no run driving and no waking prompt waiting. */
  350. async whenIdle(): Promise<void> {
  351. while (true) {
  352. // `done` is replaced per activity, so re-reading it follows chained turns.
  353. // Every driver failure today is contained before it can reject `done`,
  354. // but the waiter must not gamble quiescence on that: a future escape
  355. // still counts as settled activity.
  356. /* v8 ignore next 3 -- the catch arm backstops rejection paths that are all currently contained */
  357. while (this.busy || this.wakeScheduled || this.abort !== undefined || this.runnableWakingQueued) {
  358. await this.done.catch(() => undefined)
  359. }
  360. // A reservation is unfinished activity even with an empty queue, and a
  361. // prompt it withholds is not quiescent — but `done` never owns it, so
  362. // waiting on the queue alone would spin on an already-settled promise.
  363. const reservation = this.admission
  364. if (reservation === undefined) return
  365. await reservation.settled
  366. }
  367. }
  368. /** Whether a queued waking prompt may claim the driver now. */
  369. private get runnableWakingQueued(): boolean {
  370. return this.admission === undefined && this.queued.some(item => item.wakeup)
  371. }
  372. /** Defer idle admission while keeping {@link done} as its quiescence owner. */
  373. private scheduleKick(): void {
  374. // A held reservation keeps the item queued with no scheduled claim; its
  375. // release re-arms this path for whatever is queued by then.
  376. if (this.abort !== undefined || this.wakeScheduled || this.admission !== undefined) return
  377. this.wakeScheduled = true
  378. const pending = Promise.withResolvers<void>()
  379. const scheduled = pending.promise
  380. queueMicrotask(() => {
  381. this.wakeScheduled = false
  382. this.kick()
  383. const activity = this.done
  384. if (activity === scheduled) {
  385. pending.resolve()
  386. } else {
  387. void activity.then(
  388. () => { pending.resolve() },
  389. () => { pending.resolve() },
  390. )
  391. }
  392. })
  393. this.done = scheduled
  394. }
  395. /** Claim and admit the next queued prompt, then start its turn. */
  396. private kick(): void {
  397. if (this.abort !== undefined || !this.runnableWakingQueued) return
  398. // The some() guard above proves the queue is non-empty; the non-null
  399. // assertion expresses that invariant.
  400. // oxlint-disable-next-line typescript/no-non-null-assertion
  401. const pending = this.queued.shift()!
  402. const { item, delivery } = pending
  403. const { message } = item
  404. const inheritedOutboxLength = this.outbox.length
  405. const admission = new AbortController()
  406. this.abort = admission
  407. this.acceptsNextStep = true
  408. // Claimed admission is part of the running interval: it is cancellable
  409. // activity, so observers (and their cancel routing) must see it.
  410. if (!this.busy) {
  411. this.busy = true
  412. emitAgentEvent(this.loopCtx, this, 'agent/status', 'running')
  413. }
  414. // The admission body runs synchronously up to the prompt-submit
  415. // waterfall's first await, so the waterfall snapshots its listeners
  416. // before a disposal initiated by the running-status emit above can
  417. // unregister a vetoing plugin.
  418. this.done = this.loopCtx.agents.withInitiator(this, async () => {
  419. const signal = admission.signal
  420. const trigger: TurnTrigger = { kind: 'message', source: message.source }
  421. // Admitted input stays on the stack until its turn/start commits: the
  422. // turn owns it only once the turn exists in the log.
  423. let admitted: UserMessage[] | undefined
  424. try {
  425. signal.throwIfAborted()
  426. const decision = await this.loopCtx.waterfall(
  427. agentCarrier(this), 'agent/prompt-submit', this, message, signal,
  428. () => Promise.resolve<PromptDecision>({ kind: 'allow' }),
  429. )
  430. signal.throwIfAborted()
  431. if (decision.kind === 'allow') {
  432. admitted = [decision.content === undefined
  433. ? message
  434. : freezeMessage({ ...message, content: decision.content })]
  435. for (const context of decision.additionalContexts ?? []) {
  436. admitted.push(freezeMessage(context))
  437. }
  438. }
  439. } catch (error: unknown) {
  440. if (!signal.aborted) {
  441. this.loopCtx.logger.warn(`agent "${this.id}": prompt admission failed: ${errorChain(error)}`)
  442. }
  443. }
  444. // cancel() aborts but never clears the slot, and kick()/run()
  445. // all refuse to install a new owner while one exists, so the admission
  446. // still owns the slot here and releasing it unconditionally is exact.
  447. this.abort = undefined
  448. if (admitted === undefined) {
  449. delivery?.settle({ status: 'rejected' })
  450. this.acceptsNextStep = false
  451. try {
  452. this.flushRejectedAdmissionContexts()
  453. } catch (error: unknown) {
  454. // No turn exists for agent/error coordinates. Preserve the
  455. // uncommitted suffix for a later boundary and report locally.
  456. this.loopCtx.logger.warn(
  457. `agent "${this.id}": committing rejected-admission context failed: ${errorChain(error)}`,
  458. )
  459. }
  460. // A synchronously aborted admission would otherwise publish idle
  461. // inside send()'s own synchronous extent, before any post-send
  462. // subscriber could observe the transition.
  463. await Promise.resolve()
  464. this.continueOrIdle()
  465. return
  466. }
  467. await this.run(trigger, admitted, inheritedOutboxLength, Object.freeze([]), delivery)
  468. })
  469. // Published only after the abort owner and pending done are installed: a
  470. // dequeue listener that cancels or disposes must find live cancellation
  471. // and quiescence ownership, not the previous activity's settled state.
  472. emitAgentEvent(this.loopCtx, this, 'agent/inbox/dequeue', item)
  473. }
  474. /**
  475. * Run one turn and any request-error retry. `admitted` input enters the log
  476. * only after `turn/start` commits; until then it has no owner state to unwind.
  477. */
  478. private async run(
  479. trigger: TurnTrigger,
  480. admitted: UserMessage[] = [],
  481. inheritedOutboxLength = 0,
  482. priorFailures: readonly LlmFailure[] = Object.freeze([]),
  483. promptDelivery?: SteeringDelivery,
  484. ): Promise<void> {
  485. // Both entries hold the invariant: kick() clears the admission slot before
  486. // awaiting run(), and a retry is entered only after the prior run clears it.
  487. /* v8 ignore next -- unreachable guard: every caller clears or checks the abort slot first */
  488. if (this.abort !== undefined) throw new Error(`agent "${this.id}" is already running`)
  489. const controller = new AbortController()
  490. this.abort = controller
  491. this.preservePendingAdmissionsOnAbort = false
  492. this.acceptsNextStep = true
  493. const signal = controller.signal
  494. const turn = this.lastTurn + 1
  495. let step = 0
  496. let opened = false
  497. let reason: TurnEndReason = { kind: 'completed' }
  498. let settleReason: SettleReason = { kind: 'completed' }
  499. let requestFailureHistory = priorFailures
  500. let retryFailures: readonly LlmFailure[] | undefined
  501. const cancelRetry = (): void => { retryFailures = undefined }
  502. signal.addEventListener('abort', cancelRetry, { once: true })
  503. try {
  504. signal.throwIfAborted()
  505. this.session.append('turn/start', { turn, trigger })
  506. // Committed: publish the turn to the machine's own bookkeeping and let
  507. // the admitted input enter the log it now belongs to.
  508. this.turnOpen = true
  509. opened = true
  510. this.lastTurn = turn
  511. // Context or steering retained by an earlier rejected admission happened
  512. // before this prompt and must occupy the same order in durable history.
  513. this.drainOutbox(turn, inheritedOutboxLength)
  514. if (promptDelivery !== undefined) this.pendingAdmissions.push(promptDelivery)
  515. for (const input of admitted) {
  516. this.session.append('user/message', input, { surfaceOp: 'append' })
  517. }
  518. signal.throwIfAborted()
  519. steps: while (true) {
  520. step += 1
  521. const outcome = await this.step(turn, step, signal)
  522. switch (outcome.kind) {
  523. case 'completed':
  524. requestFailureHistory = Object.freeze([])
  525. if (outcome.maxTokens) reason = { kind: 'max-tokens' }
  526. // A concluding tool result is terminal: reject steering that did
  527. // not enter a request, while retaining same-boundary context in
  528. // durable history before the turn closes.
  529. if (outcome.concluded) {
  530. this.discardOutboxSteering()
  531. this.drainOutbox(turn)
  532. break steps
  533. }
  534. /* v8 ignore next -- step() folded the same steering predicate into continueTurn immediately before returning. */
  535. if (outcome.continueTurn || this.outbox.some(item => item.steering)) continue
  536. break
  537. case 'request-failed': {
  538. // step() reports request failures only after step/start commits
  539. // and before its own step/end, so the step is always open here.
  540. this.stepOpen = false
  541. this.session.append('step/end', { turn, step })
  542. if (!signal.aborted) {
  543. try {
  544. const action = await this.loopCtx.waterfall(
  545. agentCarrier(this), 'agent/request-error', this, turn, step, outcome.error,
  546. outcome.failure, requestFailureHistory, outcome.retryPolicy, signal,
  547. () => Promise.resolve<RequestErrorAction>(undefined),
  548. )
  549. // oxlint-disable-next-line typescript/no-unnecessary-condition -- signal can abort while recovery is awaited.
  550. if (action?.kind === 'retry' && !signal.aborted) {
  551. retryFailures = Object.freeze([...requestFailureHistory, outcome.failure])
  552. }
  553. } catch (recoveryError: unknown) {
  554. this.loopCtx.logger.warn(
  555. `agent "${this.id}": request recovery failed at turn ${turn}, step ${step}: ${errorChain(recoveryError)}`,
  556. )
  557. }
  558. }
  559. const settlement = this.settle(turn, step, outcome.error, signal, outcome.failure)
  560. reason = settlement.reason
  561. settleReason = settlement.settleReason
  562. break steps
  563. }
  564. /* v8 ignore next 2 -- closed-union exhaustiveness guard */
  565. default:
  566. assertNever(outcome)
  567. }
  568. await this.loopCtx.serial(agentCarrier(this), 'agent/turn-stopping', this, turn, signal)
  569. signal.throwIfAborted()
  570. this.drainOutboxContexts()
  571. if (!this.outbox.some(item => item.steering)) {
  572. break
  573. }
  574. }
  575. } catch (caught: unknown) {
  576. try {
  577. if (this.stepOpen) {
  578. this.stepOpen = false
  579. this.session.append('step/end', { turn, step })
  580. }
  581. } catch (closeError: unknown) {
  582. // Contained like the finally's turn close: a persistently rejecting
  583. // step boundary must not escape run(), or the post-finally tail would
  584. // never publish the terminal status and observers would see a
  585. // permanently running agent whose whenIdle() already resolved.
  586. this.loopCtx.logger.warn(`agent "${this.id}": closing step ${turn}/${step} failed: ${errorChain(closeError)}`)
  587. emitAgentEvent(this.loopCtx, this, 'agent/error', turn, step, closeError)
  588. }
  589. ({ reason, settleReason } = this.settle(turn, step, caught, signal))
  590. } finally {
  591. // Every step-close happens before this point on both success and
  592. // failure paths (step(), the request-failed branch, the catch), so the
  593. // finally owes only the turn boundary.
  594. this.acceptsNextStep = false
  595. try {
  596. if (this.turnOpen) {
  597. // Re-entrant turn/end listeners must route new input to a later turn.
  598. this.turnOpen = false
  599. this.session.append('turn/end', { turn, reason })
  600. }
  601. } catch (error: unknown) {
  602. retryFailures = undefined
  603. this.loopCtx.logger.warn(`agent "${this.id}": closing turn ${turn} failed: ${errorChain(error)}`)
  604. emitAgentEvent(this.loopCtx, this, 'agent/error', turn, step, error)
  605. }
  606. // cancel() aborts but never clears the slot, and no second run can
  607. // install a controller while this one is still unwinding, so the slot
  608. // is still this run's controller here.
  609. this.abort = undefined
  610. signal.removeEventListener('abort', cancelRetry)
  611. const preservePending = signal.aborted && this.preservePendingAdmissionsOnAbort
  612. this.preservePendingAdmissionsOnAbort = false
  613. // oxlint-disable-next-line typescript/no-unnecessary-condition -- keepInbox cancellation can set this while turn work is awaited.
  614. if (!preservePending) this.rejectPendingAdmissions()
  615. }
  616. if (opened) {
  617. try {
  618. await this.loopCtx.sessions.flush(this.session)
  619. } catch (error: unknown) {
  620. this.loopCtx.logger.warn(`agent "${this.id}": session/flush failed at turn ${turn}: ${errorChain(error)}`)
  621. emitAgentEvent(this.loopCtx, this, 'agent/error', turn, step, error)
  622. }
  623. }
  624. if (retryFailures !== undefined) {
  625. await this.run({ kind: 'retry' }, [], 0, retryFailures)
  626. } else {
  627. // agent/settled names only committed turns: a run aborted or rejected
  628. // before turn/start has no durable turn/end for consumers to settle
  629. // against, so it exits without the notification.
  630. if (opened) emitAgentEvent(this.loopCtx, this, 'agent/settled', turn, settleReason)
  631. this.continueOrIdle()
  632. }
  633. }
  634. /**
  635. * Run the `agent/step` extension point, commit pending input, derive one
  636. * request, and execute its tool calls inside one durable step boundary.
  637. */
  638. private async step(
  639. turn: number,
  640. step: number,
  641. signal: AbortSignal,
  642. ): Promise<StepOutcome> {
  643. const { session } = this
  644. // The single between-steps extension point: listeners inject, steer, or
  645. // edit the log here; the request derives from the log after this settles.
  646. await this.loopCtx.serial(agentCarrier(this), 'agent/step', this, turn, step, signal)
  647. signal.throwIfAborted()
  648. // Assemble request-owned prompt inputs fresh each step. Dynamic context is
  649. // committed at the tail before deriving history once, preserving the stable
  650. // system/history cache prefix while keeping every model-visible byte logged.
  651. const assembly = await this.loopCtx.systemPrompt.assemble(assembleContextFor(this, signal))
  652. signal.throwIfAborted()
  653. const system = renderPrompt(assembly)
  654. materializeRuntimeContext(session, renderContextSnapshot(assembly))
  655. // Commit the exact pending batch only after every asynchronous
  656. // pre-request contribution succeeded. Input accepted after this splice
  657. // remains pending for a later request.
  658. this.drainOutbox(turn)
  659. // Snapshot the exact log prefix: the reconstruction boundary. Appends
  660. // after this synchronous snapshot join the next request.
  661. const boundaryMessages = session.deriveMessages()
  662. session.append('step/start', { turn, step })
  663. this.stepOpen = true
  664. this.admitPendingAdmissions(turn, step)
  665. signal.throwIfAborted()
  666. const { request, preparedCall } = await this.buildRequest(
  667. turn, step, assembly.tools, system, boundaryMessages, signal,
  668. )
  669. const assembler = new BlockAssembler()
  670. const chunkSeqs: number[] = []
  671. const stream = preparedCall?.stream(request) ?? this.loopCtx.llm.stream(request)
  672. try {
  673. for await (const chunk of stream) {
  674. signal.throwIfAborted()
  675. const chunkEvent = session.append('assistant/chunk', { turn, step, chunk })
  676. chunkSeqs.push(chunkEvent.seq)
  677. assembler.push(chunk)
  678. }
  679. } catch (error: unknown) {
  680. const facts = llmFailureOf(stream, error)
  681. if (facts !== undefined && error instanceof Error) {
  682. return { kind: 'request-failed', error, failure: facts, retryPolicy: llmRetryPolicyOf(stream) }
  683. }
  684. throw error
  685. }
  686. signal.throwIfAborted()
  687. // Failure finish chunks take the same path as thrown stream errors.
  688. const finish = assembler.finish
  689. if (finish.kind === 'error' || finish.kind === 'aborted') {
  690. const error = new LlmError(finish.failure.message, finish.failure.code, finish.failure)
  691. return { kind: 'request-failed', error, failure: finish.failure, retryPolicy: llmRetryPolicyOf(stream) }
  692. }
  693. // Truncated (max-tokens) output cannot owe tool calls.
  694. const assembled = assembler.blocks()
  695. const content = finish.kind === 'max-tokens'
  696. ? assembled.filter(block => block.type !== 'tool-call')
  697. : assembled
  698. const message: AssistantMessage = createAssistantMessage({
  699. content,
  700. source: {
  701. provider: request.provider,
  702. model: request.model,
  703. ...assembler.replayState !== undefined ? { replayState: assembler.replayState } : {},
  704. },
  705. })
  706. session.append(
  707. 'assistant/message',
  708. {
  709. turn,
  710. step,
  711. message,
  712. ...assembler.usage === undefined ? {} : { usage: assembler.usage },
  713. },
  714. { surfaceOp: 'append', sourceEventSeqs: chunkSeqs },
  715. )
  716. const toolCalls = content.filter(block => block.type === 'tool-call')
  717. let concluded = false
  718. if (toolCalls.length > 0) {
  719. ({ concluded } = await executeToolCalls(
  720. this.loopCtx, turn, step, toolCalls, signal,
  721. context => this.outbox.push({ message: freezeMessage(context), steering: false }),
  722. ))
  723. }
  724. // Ordinary context keeps the base loop's result-adjacent commit point.
  725. // Steering remains provisional until the next request snapshot admits it.
  726. this.drainOutboxContexts()
  727. session.append('step/end', { turn, step })
  728. this.stepOpen = false
  729. return {
  730. kind: 'completed',
  731. continueTurn: (toolCalls.length > 0 && !concluded) || this.outbox.some(item => item.steering),
  732. concluded,
  733. maxTokens: finish.kind === 'max-tokens',
  734. }
  735. }
  736. /**
  737. * Compose one frozen request and bind it to the adapter registration that
  738. * resolved its exact-model defaults.
  739. */
  740. private async buildRequest(
  741. turn: number,
  742. step: number,
  743. tools: GenerateOptions['tools'] & object,
  744. system: string,
  745. boundaryMessages: Message[],
  746. signal: AbortSignal,
  747. ): Promise<{ request: GenerateOptions; preparedCall?: PreparedLlmCall }> {
  748. const { session } = this
  749. // A loop instance starts from its declared route, restoring only an explicit
  750. // effort owned by that exact model. Later steps re-resolve marked defaults.
  751. const persistedHeader = session.requestHeader()
  752. const persistedConfig = persistedHeader?.config
  753. const route = { provider: this.options.provider ?? '', model: this.options.model ?? '' }
  754. const reasoningEffort = persistedConfig?.provider === route.provider
  755. && persistedConfig.model === route.model
  756. && persistedHeader?.adapterDefaults?.reasoningEffort !== true
  757. ? persistedConfig.reasoningEffort
  758. : undefined
  759. const maxTokens = this.options.maxTokens
  760. const seedConfig = deepFreeze(structuredClone(
  761. this.requestHeaderLogged
  762. // oxlint-disable-next-line typescript/no-non-null-assertion -- the instance logged the header it now folds
  763. ? requestProposal(persistedHeader!)
  764. : {
  765. ...route,
  766. ...reasoningEffort === undefined ? {} : { reasoningEffort },
  767. ...maxTokens === undefined ? {} : { maxTokens },
  768. },
  769. ))
  770. const proposedConfig = await this.loopCtx.waterfall(
  771. agentCarrier(this), 'agent/request', this, turn, step, signal,
  772. () => Promise.resolve(seedConfig),
  773. )
  774. signal.throwIfAborted()
  775. if (!proposedConfig.provider || !proposedConfig.model) {
  776. throw new Error(`agent "${this.id}" has no provider/model: set AgentOptions.provider and AgentOptions.model or supply both via the agent/request waterfall`)
  777. }
  778. let config: LlmCallConfig
  779. let preparedCall: PreparedLlmCall | undefined
  780. try {
  781. preparedCall = await this.loopCtx.llm.prepareCall(proposedConfig, signal)
  782. config = preparedCall.config
  783. } catch (error: unknown) {
  784. // A llm/stream listener may own and short-circuit a route with no
  785. // adapter. Terminal dispatch still raises NO_ADAPTER when none does.
  786. if (!(error instanceof LlmError) || error.code !== 'NO_ADAPTER') throw error
  787. config = proposedConfig
  788. }
  789. signal.throwIfAborted()
  790. const header = canonicalHeader({
  791. config,
  792. ...preparedCall === undefined ? {} : { adapterDefaults: preparedCall.adapterDefaults },
  793. ...system ? { system } : {},
  794. ...tools.length > 0 ? { tools } : {},
  795. })
  796. const baseline = session.requestHeader()
  797. if (!this.requestHeaderLogged) {
  798. session.append('request/header', { header, reason: baseline === undefined ? 'initial' : 'resume' })
  799. this.requestHeaderLogged = true
  800. } else if (baseline === undefined || !headerEquals(baseline, header)) {
  801. session.append('request/header', { header, reason: 'change' })
  802. }
  803. // TODO: This looks like code smell.
  804. // Context metadata for the route this request resolved to, recorded from the same
  805. // registration-bound lookup that prepared the call (no second resolve).
  806. // A route with unknown capacity is still recorded so it clears any older
  807. // denominator; an unchanged route logs nothing.
  808. const contextWindow = preparedCall?.context?.contextWindow
  809. const requestContext: RequestContext = {
  810. provider: config.provider,
  811. model: config.model,
  812. ...contextWindow === undefined ? {} : { contextWindow },
  813. }
  814. const previous = session.requestContext()
  815. if (previous?.provider !== requestContext.provider
  816. || previous.model !== requestContext.model
  817. || previous.contextWindow !== requestContext.contextWindow) {
  818. session.append('request/context', requestContext)
  819. }
  820. const request = markAgentLoopRequest(deepFreeze({
  821. ...header.config,
  822. messages: boundaryMessages,
  823. ...header.system !== undefined ? { system: header.system } : {},
  824. ...header.tools !== undefined ? { tools: header.tools } : {},
  825. sessionId: session.id,
  826. signal,
  827. }))
  828. return { request, ...preparedCall === undefined ? {} : { preparedCall } }
  829. }
  830. /** Commit one stable outbox prefix and retain tracked delivery until snapshot admission. */
  831. private drainOutbox(turn: number, limit = this.outbox.length): void {
  832. const batch = this.outbox.splice(0, limit)
  833. for (let index = 0; index < batch.length; index += 1) {
  834. const item = batch[index]
  835. /* v8 ignore next -- the index walks the exact array length. */
  836. if (item === undefined) throw new Error(`agent "${this.id}" outbox item disappeared during drain`)
  837. try {
  838. if (item.steering) {
  839. /* v8 ignore next -- only inbox-backed steer entries carry steering:true. */
  840. if (item.item === undefined) throw new Error(`agent "${this.id}" steering outbox item has no inbox identity`)
  841. emitAgentEvent(this.loopCtx, this, 'agent/inbox/dequeue', item.item)
  842. this.session.append(
  843. 'steering/message',
  844. { turn, message: item.message },
  845. { surfaceOp: 'append' },
  846. )
  847. if (item.delivery !== undefined) this.pendingAdmissions.push(item.delivery)
  848. } else {
  849. this.session.append('user/message', item.message, { surfaceOp: 'append' })
  850. }
  851. } catch (error: unknown) {
  852. item.delivery?.settle({ status: 'rejected' })
  853. this.outbox.unshift(...batch.slice(item.steering ? index + 1 : index))
  854. throw error
  855. }
  856. }
  857. }
  858. /** Commit ordinary context while retaining provisional steering in order. */
  859. private drainOutboxContexts(): void {
  860. const pending = this.outbox
  861. this.outbox = []
  862. for (let index = 0; index < pending.length; index += 1) {
  863. const item = pending[index]
  864. /* v8 ignore next -- the index walks the exact array length. */
  865. if (item === undefined) throw new Error(`agent "${this.id}" outbox item disappeared during context drain`)
  866. if (item.steering) {
  867. this.outbox.push(item)
  868. continue
  869. }
  870. try {
  871. this.session.append('user/message', item.message, { surfaceOp: 'append' })
  872. } catch (error: unknown) {
  873. this.outbox.push(...pending.slice(index))
  874. throw error
  875. }
  876. }
  877. }
  878. /** Settle every committed steering item captured by this immutable request. */
  879. private admitPendingAdmissions(turn: number, step: number): void {
  880. const outcome: SteeringOutcome = { status: 'admitted', turn, step }
  881. for (const delivery of this.pendingAdmissions.splice(0)) delivery.settle(outcome)
  882. }
  883. /** Reject committed steering that left the inbox without reaching a request. */
  884. private rejectPendingAdmissions(): void {
  885. for (const delivery of this.pendingAdmissions.splice(0)) delivery.settle({ status: 'rejected' })
  886. }
  887. /** Discard uncommitted steering while retaining same-boundary injected context. */
  888. private discardOutboxSteering(): void {
  889. const contexts: typeof this.outbox = []
  890. const discarded: InboxItem[] = []
  891. for (const item of this.outbox) {
  892. if (!item.steering) {
  893. contexts.push(item)
  894. continue
  895. }
  896. item.delivery?.settle({ status: 'rejected' })
  897. /* v8 ignore next -- only inbox-backed steer entries carry steering:true. */
  898. if (item.item === undefined) throw new Error(`agent "${this.id}" steering outbox item has no inbox identity`)
  899. discarded.push(item.item)
  900. }
  901. this.outbox = contexts
  902. if (discarded.length > 0) emitAgentEvent(this.loopCtx, this, 'agent/inbox/discard', discarded)
  903. }
  904. /**
  905. * Give context-only input its ordinary idle placement when admission
  906. * produces no turn. Steering keeps the whole boundary staged so context
  907. * accepted beside it cannot split from the request it accompanies.
  908. */
  909. private flushRejectedAdmissionContexts(): void {
  910. if (this.outbox.some(item => item.steering)) return
  911. const contexts = this.outbox.splice(0)
  912. for (let index = 0; index < contexts.length; index += 1) {
  913. const item = contexts[index]
  914. /* v8 ignore next 2 -- the steering precheck proves this batch is context-only */
  915. if (item === undefined || item.steering) throw new Error('rejected-admission context batch changed')
  916. try {
  917. this.session.append('user/message', item.message, { surfaceOp: 'append' })
  918. } catch (error: unknown) {
  919. this.outbox.unshift(...contexts.slice(index))
  920. throw error
  921. }
  922. }
  923. }
  924. /**
  925. * The single settlement funnel: classify one turn failure (interruption
  926. * beats error) into the durable turn/end reason and live settlement report.
  927. */
  928. private settle(
  929. turn: number,
  930. step: number,
  931. error: unknown,
  932. signal: AbortSignal,
  933. failure?: LlmFailure,
  934. ): { reason: TurnEndReason; settleReason: SettleReason } {
  935. if (signal.aborted) {
  936. // Slot invariant, stated rather than re-validated: the turn controller
  937. // is machine-private and cancel() is its only aborter, always with one
  938. // frozen canonical cause as the reason.
  939. const interrupt = signal.reason as AgentInterruptReason
  940. return {
  941. reason: { kind: interrupt.kind === 'disposed' ? 'disposed' : 'aborted' },
  942. settleReason: { kind: 'aborted' },
  943. }
  944. }
  945. if (failure !== undefined) {
  946. emitAgentEvent(this.loopCtx, this, 'agent/error', turn, step, error)
  947. // The durable record renders the full cause chain: turn/end is the one
  948. // durable trace of the failure, so a wrapper message alone would lose
  949. // the transport detail the log exists to keep.
  950. const rendered = errorChain(error)
  951. return {
  952. reason: { kind: 'error', step, failure: { ...failure, ...rendered === '<unrenderable value>' ? {} : { message: rendered } } },
  953. settleReason: { kind: 'error', error, failure },
  954. }
  955. }
  956. emitAgentEvent(this.loopCtx, this, 'agent/error', turn, step, error)
  957. return {
  958. reason: { kind: 'error', step, message: errorChain(error), ...isHarnessError(error) ? { code: error.code } : {} },
  959. settleReason: { kind: 'error', error },
  960. }
  961. }
  962. /** Continue with a waking prompt, or publish the idle status. */
  963. private continueOrIdle(): void {
  964. if (this.runnableWakingQueued) {
  965. this.kick()
  966. } else {
  967. // Every caller sits inside an admission or run whose install marked the
  968. // interval busy, so the flag is still set here.
  969. this.busy = false
  970. emitAgentEvent(this.loopCtx, this, 'agent/status', 'idle')
  971. }
  972. }
  973. }