agent.ts 23 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596
  1. /**
  2. * Concrete Agent loop over two pending-input lists: queued prompts each open a
  3. * turn that logs its admitted input after `turn/start` commits, while steering
  4. * and injected context enter through the outbox at step boundaries. Every
  5. * request is derived from the session log.
  6. *
  7. * @module dsh-agent-loop/agent
  8. */
  9. import { randomUUID } from 'node:crypto'
  10. import type { Context } from 'cordis'
  11. import { AgentMessageId, agentCarrier, agentInterruptReasonOf, assembleContextFor, emitAgentEvent } from '@deepseek-ai/dsh-agent'
  12. import { createScope } from '@deepseek-ai/dsh-scope'
  13. import type { Scope } from '@deepseek-ai/dsh-scope'
  14. import type {
  15. AgentMessage,
  16. Agent,
  17. CancelOptions,
  18. AgentInterruptReason,
  19. AgentOptions,
  20. AgentStatus,
  21. IdleReason,
  22. PromptDecision,
  23. RequestError,
  24. SendOptions,
  25. } from '@deepseek-ai/dsh-agent'
  26. import {
  27. BlockAssembler, LlmError, assertNever, deepFreeze, errorChain, isHarnessError, llmFailureOf, markAgentLoopRequest,
  28. } from '@deepseek-ai/dsh-llm'
  29. import type { GenerateOptions, LlmCallConfig, LlmFailure, Message } from '@deepseek-ai/dsh-llm'
  30. import { canonicalHeader, headerEquals } from '@deepseek-ai/dsh-session'
  31. import type { Session, SessionId, TurnEndReason, TurnTrigger, UserMessageData } from '@deepseek-ai/dsh-session'
  32. import { renderPrompt } from '@deepseek-ai/dsh-system-prompt'
  33. import type {} from '@deepseek-ai/dsh-tools'
  34. import { executeToolCalls } from './tool-calls.ts'
  35. /** One completed step or a final-adapter failure eligible for recovery. */
  36. type StepOutcome =
  37. | { kind: 'completed'; continueTurn: boolean; concluded: boolean; maxTokens: boolean }
  38. | { kind: 'request-failed'; error: RequestError; failure: LlmFailure }
  39. /**
  40. * The concrete {@link Agent}: each `run()` owns one turn and repeats model
  41. * steps while tools or steering require another request.
  42. */
  43. export class ReactLoopAgent implements Agent {
  44. /** Prompts awaiting individual turns. */
  45. private queued: { message: AgentMessage; wakeup: boolean }[] = []
  46. /** Input taken into the session log at step boundaries. */
  47. private outbox: (UserMessageData | AgentMessage)[] = []
  48. /** Whether observers see a running interval; consecutive turns share it. */
  49. private busy = false
  50. /** Abort owner for the current admission or turn. */
  51. private abort: AbortController | undefined
  52. /** Coalesced retry capability scoped to the active request-error waterfall. */
  53. private retryWindow: { requested: boolean } | undefined
  54. /** Resolves when the current admission and turn exit. */
  55. done: Promise<void> = Promise.resolve()
  56. /** The agent-scoped registration boundary; the lifecycle owner unwinds it after {@link done}. */
  57. readonly scope: Scope
  58. /** The agent's scoped composition context ({@link Agent.ctx}). */
  59. readonly ctx: Context
  60. /** Last turn number opened by this loop or present in its seeded log. */
  61. private lastTurn: number
  62. /** Whether the session log is owed a matching turn end event. */
  63. private turnOpen = false
  64. private stepOpen = false
  65. constructor(
  66. private loopCtx: Context,
  67. public readonly id: SessionId,
  68. public readonly options: AgentOptions,
  69. public readonly session: Session,
  70. ) {
  71. this.lastTurn = session.events.findLast(event => event.type === 'turn/start')?.data.turn ?? 0
  72. this.scope = createScope(loopCtx, this)
  73. this.ctx = this.scope.ctx.extend({ agent: this })
  74. }
  75. /** Last activity state published to observers. */
  76. get status(): AgentStatus {
  77. return this.busy ? 'running' : 'idle'
  78. }
  79. /** Accept and route one unified send item. */
  80. send(
  81. input: UserMessageData,
  82. options: SendOptions,
  83. ): AgentMessageId {
  84. const { content, source } = input
  85. const { target, wakeup } = options
  86. const id = AgentMessageId(randomUUID())
  87. if (target === 'next-step' && !wakeup) {
  88. if (this.turnOpen) {
  89. this.outbox.push({ content, source })
  90. return id
  91. }
  92. this.session.append('user/message', { content, source }, { surfaceOp: 'append' })
  93. return id
  94. }
  95. const steering = target === 'next-step' && this.turnOpen
  96. const message: AgentMessage = {
  97. id,
  98. content,
  99. source,
  100. }
  101. if (steering) {
  102. this.outbox.push(message)
  103. } else {
  104. this.queued.push({ message, wakeup })
  105. }
  106. emitAgentEvent(this.loopCtx, this, 'agent/inbox/enqueue', message)
  107. if (!steering && wakeup) this.kick()
  108. return id
  109. }
  110. /** Queue one ordinary prompt turn and wake the driver. */
  111. followup(input: UserMessageData): AgentMessageId {
  112. return this.send(input, {
  113. target: 'next-turn',
  114. wakeup: true,
  115. })
  116. }
  117. /** Steer the open turn, falling back to a waking prompt while idle. */
  118. steer(input: UserMessageData): AgentMessageId {
  119. return this.send(input, {
  120. target: 'next-step',
  121. wakeup: true,
  122. })
  123. }
  124. /** Append model-facing context without waking the driver. */
  125. inject(input: UserMessageData): AgentMessageId {
  126. return this.send(input, {
  127. target: 'next-step',
  128. wakeup: false,
  129. })
  130. }
  131. /**
  132. * Clear all pending work and abort the active turn; the first cause wins.
  133. * The cause is signal payload for observers and the durable turn/end
  134. * classification — it selects no machine behavior. Teardown is just
  135. * `cancel({kind:'disposed'})` + await {@link done} + {@link scope} dispose,
  136. * all owned by the factory.
  137. */
  138. cancel(cause: AgentInterruptReason, options: CancelOptions = {}): void {
  139. if (this.abort !== undefined || this.queued.length > 0 || this.outbox.length > 0) {
  140. // Observe-only: coordination consumers update their state before the
  141. // inboxes clear; listener failures are contained by the dispatcher.
  142. if (cause.kind !== 'disposed') emitAgentEvent(this.loopCtx, this, 'agent/cancel-requested', cause)
  143. }
  144. if (!options.keepInbox) {
  145. const discarded = this.queued.map(item => item.message)
  146. for (const message of this.outbox) {
  147. if ('id' in message) discarded.push(message)
  148. }
  149. // Clear before abort observers run: replacement work belongs to the next turn.
  150. this.queued.length = 0
  151. this.outbox.length = 0
  152. if (discarded.length > 0) emitAgentEvent(this.loopCtx, this, 'agent/inbox/discard', discarded)
  153. }
  154. if (this.retryWindow !== undefined) this.retryWindow.requested = false
  155. const reason = Object.freeze({ kind: cause.kind })
  156. this.abort?.abort(reason)
  157. }
  158. /**
  159. * Re-open a turn on the current session log without a new prompt — the
  160. * recovery verb. A request-error listener schedules the retry that follows
  161. * its failed turn; an idle caller starts one immediately.
  162. */
  163. retry(): void {
  164. if (this.abort !== undefined) {
  165. if (this.retryWindow === undefined) throw new Error(`agent "${this.id}" cannot retry while busy`)
  166. if (!this.abort.signal.aborted) this.retryWindow.requested = true
  167. return
  168. }
  169. this.done = this.loopCtx.agents.withInitiator(this, () => this.run({ kind: 'retry' }))
  170. }
  171. /** Resolve at idle quiescence: no run driving and no waking prompt waiting. */
  172. async whenIdle(): Promise<void> {
  173. // `done` is replaced per activity, so re-reading it follows chained turns;
  174. // a run failure still counts as quiescence for the waiter.
  175. while (this.abort !== undefined || this.queued.some(item => item.wakeup)) {
  176. await this.done.catch(() => undefined)
  177. }
  178. }
  179. /** Claim and admit the next queued prompt, then start its turn. */
  180. private kick(): void {
  181. if (this.abort !== undefined || !this.queued.some(item => item.wakeup)) return
  182. const item = this.queued.shift()
  183. /* v8 ignore next -- unreachable: the some() guard above proves the queue is non-empty */
  184. if (item === undefined) return
  185. const { message } = item
  186. emitAgentEvent(this.loopCtx, this, 'agent/inbox/dequeue', message)
  187. const admission = new AbortController()
  188. this.abort = admission
  189. this.done = this.loopCtx.agents.withInitiator(this, async () => {
  190. const signal = admission.signal
  191. const trigger: TurnTrigger = { kind: 'message', source: message.source }
  192. // Admitted input stays on the stack until its turn/start commits: the
  193. // turn owns it only once the turn exists in the log.
  194. let admitted: UserMessageData[] | undefined
  195. try {
  196. signal.throwIfAborted()
  197. const decision = await this.loopCtx.waterfall(
  198. agentCarrier(this), 'agent/prompt-submit', this, message.content, message.source, signal,
  199. () => Promise.resolve<PromptDecision>({ kind: 'allow' }),
  200. )
  201. signal.throwIfAborted()
  202. if (decision.kind === 'allow') {
  203. admitted = [{ content: decision.content ?? message.content, source: message.source }]
  204. for (const context of decision.additionalContexts ?? []) {
  205. admitted.push({ content: context.content, source: context.source })
  206. }
  207. }
  208. } catch (error: unknown) {
  209. if (agentInterruptReasonOf(signal) === undefined) {
  210. this.loopCtx.logger.warn(`agent "${this.id}": prompt admission failed: ${errorChain(error)}`)
  211. }
  212. }
  213. // cancel() aborts but never clears the slot, and kick()/run()/retry()
  214. // all refuse to install a new owner while one exists, so the admission
  215. // still owns the slot here.
  216. /* v8 ignore next -- unreachable false arm: no writer replaces the abort owner mid-admission */
  217. if (this.abort === admission) this.abort = undefined
  218. if (admitted === undefined) {
  219. this.continueOrIdle()
  220. return
  221. }
  222. await this.run(trigger, admitted)
  223. })
  224. }
  225. /**
  226. * Run one turn and any request-error retry. `admitted` input enters the log
  227. * only after `turn/start` commits; until then it has no owner state to unwind.
  228. */
  229. private async run(trigger: TurnTrigger, admitted: UserMessageData[] = []): Promise<void> {
  230. // Both entries hold the invariant: kick() clears the admission slot before
  231. // awaiting run(), and retry() returns early whenever a slot owner exists.
  232. /* v8 ignore next -- unreachable guard: every caller clears or checks the abort slot first */
  233. if (this.abort !== undefined) throw new Error(`agent "${this.id}" is already running`)
  234. const controller = new AbortController()
  235. this.abort = controller
  236. if (!this.busy) {
  237. this.busy = true
  238. emitAgentEvent(this.loopCtx, this, 'agent/status', 'running')
  239. }
  240. const signal = controller.signal
  241. const turn = this.lastTurn + 1
  242. let step = 0
  243. let reason: TurnEndReason = { kind: 'completed' }
  244. let idle: IdleReason = { kind: 'completed' }
  245. let retry = false
  246. const cancelRetry = (): void => { retry = false }
  247. signal.addEventListener('abort', cancelRetry, { once: true })
  248. try {
  249. signal.throwIfAborted()
  250. this.session.append('turn/start', { turn, trigger })
  251. // Committed: publish the turn to the machine's own bookkeeping and let
  252. // the admitted input enter the log it now belongs to.
  253. this.turnOpen = true
  254. this.lastTurn = turn
  255. for (const input of admitted) {
  256. this.session.append('user/message', input, { surfaceOp: 'append' })
  257. }
  258. signal.throwIfAborted()
  259. this.drainOutbox(turn)
  260. steps: while (true) {
  261. step += 1
  262. const outcome = await this.step(turn, step, signal)
  263. switch (outcome.kind) {
  264. case 'completed':
  265. if (outcome.maxTokens) reason = { kind: 'max-tokens' }
  266. // A concluding tool result is terminal: steering already in the
  267. // log waits for the next turn's request instead of reopening this
  268. // one, and the agent/stopping drain below is skipped for the same
  269. // reason.
  270. if (outcome.concluded) break steps
  271. if (outcome.continueTurn || this.outbox.some(item => 'id' in item)) continue
  272. break
  273. case 'request-failed': {
  274. if (this.stepOpen) {
  275. this.stepOpen = false
  276. this.session.append('step/end', { turn, step })
  277. }
  278. if (agentInterruptReasonOf(signal) === undefined) {
  279. const retryWindow = { requested: false }
  280. this.retryWindow = retryWindow
  281. let recoveryCompleted = false
  282. try {
  283. await this.loopCtx.waterfall(
  284. agentCarrier(this), 'agent/request-error', this, turn, step, outcome.error,
  285. outcome.failure, signal,
  286. () => Promise.resolve(),
  287. )
  288. recoveryCompleted = true
  289. } catch (recoveryError: unknown) {
  290. this.loopCtx.logger.warn(
  291. `agent "${this.id}": request recovery failed at turn ${turn}, step ${step}: ${errorChain(recoveryError)}`,
  292. )
  293. } finally {
  294. if (this.retryWindow === retryWindow) this.retryWindow = undefined
  295. }
  296. retry = recoveryCompleted
  297. && agentInterruptReasonOf(signal) === undefined
  298. && retryWindow.requested
  299. }
  300. const settlement = this.settle(turn, step, outcome.error, signal, outcome.failure)
  301. reason = settlement.reason
  302. idle = settlement.idle
  303. break steps
  304. }
  305. /* v8 ignore next 2 -- closed-union exhaustiveness guard */
  306. default:
  307. assertNever(outcome)
  308. }
  309. await this.loopCtx.serial(agentCarrier(this), 'agent/stopping', this, turn, signal)
  310. signal.throwIfAborted()
  311. if (!this.drainOutbox(turn)) break
  312. }
  313. } catch (caught: unknown) {
  314. if (this.stepOpen) {
  315. this.stepOpen = false
  316. this.session.append('step/end', { turn, step })
  317. }
  318. ({ reason, idle } = this.settle(turn, step, caught, signal))
  319. } finally {
  320. try {
  321. if (this.stepOpen) {
  322. this.stepOpen = false
  323. this.session.append('step/end', { turn, step })
  324. }
  325. if (this.turnOpen) {
  326. // Re-entrant turn/end listeners must route new input to a later turn.
  327. this.turnOpen = false
  328. this.session.append('turn/end', { turn, reason })
  329. }
  330. } catch (error: unknown) {
  331. retry = false
  332. this.loopCtx.logger.warn(`agent "${this.id}": closing turn ${turn} failed: ${errorChain(error)}`)
  333. emitAgentEvent(this.loopCtx, this, 'agent/error', turn, step, error)
  334. }
  335. this.retryWindow = undefined
  336. if (this.abort === controller) this.abort = undefined
  337. signal.removeEventListener('abort', cancelRetry)
  338. }
  339. if (retry) {
  340. await this.run({ kind: 'retry' })
  341. } else {
  342. emitAgentEvent(this.loopCtx, this, 'agent/idle', turn, idle)
  343. this.continueOrIdle()
  344. }
  345. }
  346. /**
  347. * Run the `agent/step` seam, commit pending input, derive one request, and
  348. * execute its tool calls inside one durable step boundary.
  349. */
  350. private async step(
  351. turn: number,
  352. step: number,
  353. signal: AbortSignal,
  354. ): Promise<StepOutcome> {
  355. const { session } = this
  356. // The single between-steps seam: listeners inject, steer, or edit the log
  357. // here; the request derives from the log after this settles.
  358. await this.loopCtx.serial(agentCarrier(this), 'agent/step', this, turn, step, signal)
  359. signal.throwIfAborted()
  360. // Take the outbox whole — same-boundary steering and context leave in
  361. // this request together.
  362. this.drainOutbox(turn)
  363. // Assemble the system prompt fresh each step (it may depend on log state).
  364. const assembly = await this.loopCtx.systemPrompt.assemble(assembleContextFor(this, signal))
  365. signal.throwIfAborted()
  366. const system = renderPrompt(assembly)
  367. // Snapshot the exact log prefix: the reconstruction boundary. Appends
  368. // after this synchronous snapshot join the next request.
  369. const boundaryMessages = session.deriveMessages()
  370. session.append('step/start', { turn, step })
  371. this.stepOpen = true
  372. signal.throwIfAborted()
  373. const request = await this.buildRequest(turn, step, assembly.tools, system, boundaryMessages, signal)
  374. const assembler = new BlockAssembler()
  375. const chunkSeqs: number[] = []
  376. const stream = this.loopCtx.llm.stream(request)
  377. try {
  378. for await (const chunk of stream) {
  379. signal.throwIfAborted()
  380. const chunkEvent = session.append('assistant/chunk', { turn, step, chunk })
  381. chunkSeqs.push(chunkEvent.seq)
  382. assembler.push(chunk)
  383. }
  384. } catch (error: unknown) {
  385. const facts = llmFailureOf(stream, error)
  386. if (facts !== undefined && error instanceof Error) {
  387. return { kind: 'request-failed', error, failure: facts }
  388. }
  389. throw error
  390. }
  391. signal.throwIfAborted()
  392. // Failure finish chunks take the same path as thrown stream errors.
  393. const finish = assembler.finish
  394. if (finish.kind === 'error' || finish.kind === 'aborted') {
  395. const error = new LlmError(finish.failure.message, finish.failure.code, finish.failure)
  396. return { kind: 'request-failed', error, failure: finish.failure }
  397. }
  398. // Truncated (max-tokens) output cannot owe tool calls.
  399. const assembled = assembler.message()
  400. const content = finish.kind === 'max-tokens'
  401. ? assembled.content.filter(block => block.type !== 'tool-call')
  402. : assembled.content
  403. session.append(
  404. 'assistant/message',
  405. {
  406. turn,
  407. step,
  408. content,
  409. provenance: {
  410. provider: request.provider,
  411. model: request.model,
  412. ...assembler.replayState !== undefined ? { replayState: assembler.replayState } : {},
  413. },
  414. ...assembler.usage === undefined ? {} : { usage: assembler.usage },
  415. },
  416. { surfaceOp: 'append', sourceEventSeqs: chunkSeqs },
  417. )
  418. const toolCalls = content.filter(block => block.type === 'tool-call')
  419. let concluded = false
  420. if (toolCalls.length > 0) {
  421. ({ concluded } = await executeToolCalls(
  422. this.loopCtx, turn, step, toolCalls, signal,
  423. context => this.outbox.push({ content: context.content, source: context.source }),
  424. ))
  425. }
  426. // Tool results stay adjacent to their calls; input accepted during the
  427. // request enters the log only after the complete result batch.
  428. const steered = this.drainOutbox(turn)
  429. session.append('step/end', { turn, step })
  430. this.stepOpen = false
  431. return {
  432. kind: 'completed',
  433. continueTurn: (toolCalls.length > 0 && !concluded) || steered,
  434. concluded,
  435. maxTokens: finish.kind === 'max-tokens',
  436. }
  437. }
  438. /**
  439. * Compose one frozen request: the `agent/request` config waterfall, the
  440. * canonical logged header, then the header plus the boundary snapshot,
  441. * byte-for-byte.
  442. */
  443. private async buildRequest(
  444. turn: number,
  445. step: number,
  446. tools: GenerateOptions['tools'] & object,
  447. system: string,
  448. boundaryMessages: Message[],
  449. signal: AbortSignal,
  450. ): Promise<GenerateOptions> {
  451. const { session } = this
  452. // Seed from the logged header when the log has one (the log is the
  453. // truth, across resumes too), else from agent options; freeze so
  454. // listeners must return a replacement.
  455. const seedConfig: LlmCallConfig = deepFreeze(structuredClone(
  456. session.requestHeader()?.config
  457. ?? { provider: this.options.provider ?? '', model: this.options.model ?? '' }))
  458. const config = await this.loopCtx.waterfall(
  459. agentCarrier(this), 'agent/request', this, turn, step, signal,
  460. () => Promise.resolve(seedConfig),
  461. )
  462. signal.throwIfAborted()
  463. if (!config.provider || !config.model) {
  464. throw new Error(`agent "${this.id}" has no provider/model: set AgentOptions.provider and AgentOptions.model or supply both via the agent/request waterfall`)
  465. }
  466. const header = canonicalHeader({
  467. config,
  468. ...system ? { system } : {},
  469. ...tools.length > 0 ? { tools } : {},
  470. })
  471. // Log the header the request will use only when it differs
  472. // from the folded baseline — reconstruction folds the log, so an
  473. // unchanged header needs no new snapshot.
  474. const baseline = session.requestHeader()
  475. if (baseline === undefined || !headerEquals(baseline, header)) {
  476. session.append('request/header', { header, reason: baseline === undefined ? 'initial' : 'change' })
  477. }
  478. return markAgentLoopRequest(deepFreeze({
  479. provider: header.config.provider,
  480. model: header.config.model,
  481. messages: boundaryMessages,
  482. ...header.system !== undefined ? { system: header.system } : {},
  483. ...header.tools !== undefined ? { tools: header.tools } : {},
  484. ...header.config.temperature !== undefined ? { temperature: header.config.temperature } : {},
  485. ...header.config.maxTokens !== undefined ? { maxTokens: header.config.maxTokens } : {},
  486. ...header.config.stop !== undefined ? { stop: header.config.stop } : {},
  487. sessionId: session.id,
  488. signal,
  489. }))
  490. }
  491. /** Commit the outbox and report whether it contained steering. */
  492. private drainOutbox(turn: number): boolean {
  493. let steered = false
  494. for (const message of this.outbox.splice(0)) {
  495. if ('id' in message) {
  496. steered = true
  497. emitAgentEvent(this.loopCtx, this, 'agent/inbox/dequeue', message)
  498. this.session.append(
  499. 'steering/message',
  500. { turn, content: message.content, source: message.source },
  501. { surfaceOp: 'append' },
  502. )
  503. } else {
  504. this.session.append('user/message', message, { surfaceOp: 'append' })
  505. }
  506. }
  507. return steered
  508. }
  509. /**
  510. * The single settlement funnel: classify one turn failure (interruption
  511. * beats error) into the durable turn/end reason and the live idle report.
  512. */
  513. private settle(
  514. turn: number,
  515. step: number,
  516. error: unknown,
  517. signal: AbortSignal,
  518. failure?: LlmFailure,
  519. ): { reason: TurnEndReason; idle: IdleReason } {
  520. const interrupt = agentInterruptReasonOf(signal)
  521. if (interrupt !== undefined) {
  522. return { reason: { kind: interrupt.kind === 'disposed' ? 'disposed' : 'aborted' }, idle: { kind: 'aborted' } }
  523. }
  524. if (failure !== undefined) {
  525. emitAgentEvent(this.loopCtx, this, 'agent/error', turn, step, error)
  526. // The durable record renders the full cause chain: turn/end is the one
  527. // durable trace of the failure, so a wrapper message alone would lose
  528. // the transport detail the log exists to keep.
  529. const rendered = errorChain(error)
  530. return {
  531. reason: { kind: 'error', step, failure: { ...failure, ...rendered === '<unrenderable value>' ? {} : { message: rendered } } },
  532. idle: { kind: 'error', error, failure },
  533. }
  534. }
  535. emitAgentEvent(this.loopCtx, this, 'agent/error', turn, step, error)
  536. return {
  537. reason: { kind: 'error', step, message: errorChain(error), ...isHarnessError(error) ? { code: error.code } : {} },
  538. idle: { kind: 'error', error },
  539. }
  540. }
  541. /** Continue with a waking prompt, or publish the idle status. */
  542. private continueOrIdle(): void {
  543. if (this.abort !== undefined) return
  544. if (this.queued.some(item => item.wakeup)) {
  545. this.kick()
  546. } else if (this.busy) {
  547. this.busy = false
  548. emitAgentEvent(this.loopCtx, this, 'agent/status', 'idle')
  549. }
  550. }
  551. }