| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596 |
- /**
- * Concrete Agent loop over two pending-input lists: queued prompts each open a
- * turn that logs its admitted input after `turn/start` commits, while steering
- * and injected context enter through the outbox at step boundaries. Every
- * request is derived from the session log.
- *
- * @module dsh-agent-loop/agent
- */
- import { randomUUID } from 'node:crypto'
- import type { Context } from 'cordis'
- import { AgentMessageId, agentCarrier, agentInterruptReasonOf, assembleContextFor, emitAgentEvent } from '@deepseek-ai/dsh-agent'
- import { createScope } from '@deepseek-ai/dsh-scope'
- import type { Scope } from '@deepseek-ai/dsh-scope'
- import type {
- AgentMessage,
- Agent,
- CancelOptions,
- AgentInterruptReason,
- AgentOptions,
- AgentStatus,
- IdleReason,
- PromptDecision,
- RequestError,
- SendOptions,
- } from '@deepseek-ai/dsh-agent'
- import {
- BlockAssembler, LlmError, assertNever, deepFreeze, errorChain, isHarnessError, llmFailureOf, markAgentLoopRequest,
- } from '@deepseek-ai/dsh-llm'
- import type { GenerateOptions, LlmCallConfig, LlmFailure, Message } from '@deepseek-ai/dsh-llm'
- import { canonicalHeader, headerEquals } from '@deepseek-ai/dsh-session'
- import type { Session, SessionId, TurnEndReason, TurnTrigger, UserMessageData } from '@deepseek-ai/dsh-session'
- import { renderPrompt } from '@deepseek-ai/dsh-system-prompt'
- import type {} from '@deepseek-ai/dsh-tools'
- import { executeToolCalls } from './tool-calls.ts'
- /** One completed step or a final-adapter failure eligible for recovery. */
- type StepOutcome =
- | { kind: 'completed'; continueTurn: boolean; concluded: boolean; maxTokens: boolean }
- | { kind: 'request-failed'; error: RequestError; failure: LlmFailure }
- /**
- * The concrete {@link Agent}: each `run()` owns one turn and repeats model
- * steps while tools or steering require another request.
- */
- export class ReactLoopAgent implements Agent {
- /** Prompts awaiting individual turns. */
- private queued: { message: AgentMessage; wakeup: boolean }[] = []
- /** Input taken into the session log at step boundaries. */
- private outbox: (UserMessageData | AgentMessage)[] = []
- /** Whether observers see a running interval; consecutive turns share it. */
- private busy = false
- /** Abort owner for the current admission or turn. */
- private abort: AbortController | undefined
- /** Coalesced retry capability scoped to the active request-error waterfall. */
- private retryWindow: { requested: boolean } | undefined
- /** Resolves when the current admission and turn exit. */
- done: Promise<void> = Promise.resolve()
- /** The agent-scoped registration boundary; the lifecycle owner unwinds it after {@link done}. */
- readonly scope: Scope
- /** The agent's scoped composition context ({@link Agent.ctx}). */
- readonly ctx: Context
- /** Last turn number opened by this loop or present in its seeded log. */
- private lastTurn: number
- /** Whether the session log is owed a matching turn end event. */
- private turnOpen = false
- private stepOpen = false
- constructor(
- private loopCtx: Context,
- public readonly id: SessionId,
- public readonly options: AgentOptions,
- public readonly session: Session,
- ) {
- this.lastTurn = session.events.findLast(event => event.type === 'turn/start')?.data.turn ?? 0
- this.scope = createScope(loopCtx, this)
- this.ctx = this.scope.ctx.extend({ agent: this })
- }
- /** Last activity state published to observers. */
- get status(): AgentStatus {
- return this.busy ? 'running' : 'idle'
- }
- /** Accept and route one unified send item. */
- send(
- input: UserMessageData,
- options: SendOptions,
- ): AgentMessageId {
- const { content, source } = input
- const { target, wakeup } = options
- const id = AgentMessageId(randomUUID())
- if (target === 'next-step' && !wakeup) {
- if (this.turnOpen) {
- this.outbox.push({ content, source })
- return id
- }
- this.session.append('user/message', { content, source }, { surfaceOp: 'append' })
- return id
- }
- const steering = target === 'next-step' && this.turnOpen
- const message: AgentMessage = {
- id,
- content,
- source,
- }
- if (steering) {
- this.outbox.push(message)
- } else {
- this.queued.push({ message, wakeup })
- }
- emitAgentEvent(this.loopCtx, this, 'agent/inbox/enqueue', message)
- if (!steering && wakeup) this.kick()
- return id
- }
- /** Queue one ordinary prompt turn and wake the driver. */
- followup(input: UserMessageData): AgentMessageId {
- return this.send(input, {
- target: 'next-turn',
- wakeup: true,
- })
- }
- /** Steer the open turn, falling back to a waking prompt while idle. */
- steer(input: UserMessageData): AgentMessageId {
- return this.send(input, {
- target: 'next-step',
- wakeup: true,
- })
- }
- /** Append model-facing context without waking the driver. */
- inject(input: UserMessageData): AgentMessageId {
- return this.send(input, {
- target: 'next-step',
- wakeup: false,
- })
- }
- /**
- * Clear all pending work and abort the active turn; the first cause wins.
- * The cause is signal payload for observers and the durable turn/end
- * classification — it selects no machine behavior. Teardown is just
- * `cancel({kind:'disposed'})` + await {@link done} + {@link scope} dispose,
- * all owned by the factory.
- */
- cancel(cause: AgentInterruptReason, options: CancelOptions = {}): void {
- if (this.abort !== undefined || this.queued.length > 0 || this.outbox.length > 0) {
- // Observe-only: coordination consumers update their state before the
- // inboxes clear; listener failures are contained by the dispatcher.
- if (cause.kind !== 'disposed') emitAgentEvent(this.loopCtx, this, 'agent/cancel-requested', cause)
- }
- if (!options.keepInbox) {
- const discarded = this.queued.map(item => item.message)
- for (const message of this.outbox) {
- if ('id' in message) discarded.push(message)
- }
- // Clear before abort observers run: replacement work belongs to the next turn.
- this.queued.length = 0
- this.outbox.length = 0
- if (discarded.length > 0) emitAgentEvent(this.loopCtx, this, 'agent/inbox/discard', discarded)
- }
- if (this.retryWindow !== undefined) this.retryWindow.requested = false
- const reason = Object.freeze({ kind: cause.kind })
- this.abort?.abort(reason)
- }
- /**
- * Re-open a turn on the current session log without a new prompt — the
- * recovery verb. A request-error listener schedules the retry that follows
- * its failed turn; an idle caller starts one immediately.
- */
- retry(): void {
- if (this.abort !== undefined) {
- if (this.retryWindow === undefined) throw new Error(`agent "${this.id}" cannot retry while busy`)
- if (!this.abort.signal.aborted) this.retryWindow.requested = true
- return
- }
- this.done = this.loopCtx.agents.withInitiator(this, () => this.run({ kind: 'retry' }))
- }
- /** Resolve at idle quiescence: no run driving and no waking prompt waiting. */
- async whenIdle(): Promise<void> {
- // `done` is replaced per activity, so re-reading it follows chained turns;
- // a run failure still counts as quiescence for the waiter.
- while (this.abort !== undefined || this.queued.some(item => item.wakeup)) {
- await this.done.catch(() => undefined)
- }
- }
- /** Claim and admit the next queued prompt, then start its turn. */
- private kick(): void {
- if (this.abort !== undefined || !this.queued.some(item => item.wakeup)) return
- const item = this.queued.shift()
- /* v8 ignore next -- unreachable: the some() guard above proves the queue is non-empty */
- if (item === undefined) return
- const { message } = item
- emitAgentEvent(this.loopCtx, this, 'agent/inbox/dequeue', message)
- const admission = new AbortController()
- this.abort = admission
- this.done = this.loopCtx.agents.withInitiator(this, async () => {
- const signal = admission.signal
- const trigger: TurnTrigger = { kind: 'message', source: message.source }
- // Admitted input stays on the stack until its turn/start commits: the
- // turn owns it only once the turn exists in the log.
- let admitted: UserMessageData[] | undefined
- try {
- signal.throwIfAborted()
- const decision = await this.loopCtx.waterfall(
- agentCarrier(this), 'agent/prompt-submit', this, message.content, message.source, signal,
- () => Promise.resolve<PromptDecision>({ kind: 'allow' }),
- )
- signal.throwIfAborted()
- if (decision.kind === 'allow') {
- admitted = [{ content: decision.content ?? message.content, source: message.source }]
- for (const context of decision.additionalContexts ?? []) {
- admitted.push({ content: context.content, source: context.source })
- }
- }
- } catch (error: unknown) {
- if (agentInterruptReasonOf(signal) === undefined) {
- this.loopCtx.logger.warn(`agent "${this.id}": prompt admission failed: ${errorChain(error)}`)
- }
- }
- // cancel() aborts but never clears the slot, and kick()/run()/retry()
- // all refuse to install a new owner while one exists, so the admission
- // still owns the slot here.
- /* v8 ignore next -- unreachable false arm: no writer replaces the abort owner mid-admission */
- if (this.abort === admission) this.abort = undefined
- if (admitted === undefined) {
- this.continueOrIdle()
- return
- }
- await this.run(trigger, admitted)
- })
- }
- /**
- * Run one turn and any request-error retry. `admitted` input enters the log
- * only after `turn/start` commits; until then it has no owner state to unwind.
- */
- private async run(trigger: TurnTrigger, admitted: UserMessageData[] = []): Promise<void> {
- // Both entries hold the invariant: kick() clears the admission slot before
- // awaiting run(), and retry() returns early whenever a slot owner exists.
- /* v8 ignore next -- unreachable guard: every caller clears or checks the abort slot first */
- if (this.abort !== undefined) throw new Error(`agent "${this.id}" is already running`)
- const controller = new AbortController()
- this.abort = controller
- if (!this.busy) {
- this.busy = true
- emitAgentEvent(this.loopCtx, this, 'agent/status', 'running')
- }
- const signal = controller.signal
- const turn = this.lastTurn + 1
- let step = 0
- let reason: TurnEndReason = { kind: 'completed' }
- let idle: IdleReason = { kind: 'completed' }
- let retry = false
- const cancelRetry = (): void => { retry = false }
- signal.addEventListener('abort', cancelRetry, { once: true })
- try {
- signal.throwIfAborted()
- this.session.append('turn/start', { turn, trigger })
- // Committed: publish the turn to the machine's own bookkeeping and let
- // the admitted input enter the log it now belongs to.
- this.turnOpen = true
- this.lastTurn = turn
- for (const input of admitted) {
- this.session.append('user/message', input, { surfaceOp: 'append' })
- }
- signal.throwIfAborted()
- this.drainOutbox(turn)
- steps: while (true) {
- step += 1
- const outcome = await this.step(turn, step, signal)
- switch (outcome.kind) {
- case 'completed':
- if (outcome.maxTokens) reason = { kind: 'max-tokens' }
- // A concluding tool result is terminal: steering already in the
- // log waits for the next turn's request instead of reopening this
- // one, and the agent/stopping drain below is skipped for the same
- // reason.
- if (outcome.concluded) break steps
- if (outcome.continueTurn || this.outbox.some(item => 'id' in item)) continue
- break
- case 'request-failed': {
- if (this.stepOpen) {
- this.stepOpen = false
- this.session.append('step/end', { turn, step })
- }
- if (agentInterruptReasonOf(signal) === undefined) {
- const retryWindow = { requested: false }
- this.retryWindow = retryWindow
- let recoveryCompleted = false
- try {
- await this.loopCtx.waterfall(
- agentCarrier(this), 'agent/request-error', this, turn, step, outcome.error,
- outcome.failure, signal,
- () => Promise.resolve(),
- )
- recoveryCompleted = true
- } catch (recoveryError: unknown) {
- this.loopCtx.logger.warn(
- `agent "${this.id}": request recovery failed at turn ${turn}, step ${step}: ${errorChain(recoveryError)}`,
- )
- } finally {
- if (this.retryWindow === retryWindow) this.retryWindow = undefined
- }
- retry = recoveryCompleted
- && agentInterruptReasonOf(signal) === undefined
- && retryWindow.requested
- }
- const settlement = this.settle(turn, step, outcome.error, signal, outcome.failure)
- reason = settlement.reason
- idle = settlement.idle
- break steps
- }
- /* v8 ignore next 2 -- closed-union exhaustiveness guard */
- default:
- assertNever(outcome)
- }
- await this.loopCtx.serial(agentCarrier(this), 'agent/stopping', this, turn, signal)
- signal.throwIfAborted()
- if (!this.drainOutbox(turn)) break
- }
- } catch (caught: unknown) {
- if (this.stepOpen) {
- this.stepOpen = false
- this.session.append('step/end', { turn, step })
- }
- ({ reason, idle } = this.settle(turn, step, caught, signal))
- } finally {
- try {
- if (this.stepOpen) {
- this.stepOpen = false
- this.session.append('step/end', { turn, step })
- }
- if (this.turnOpen) {
- // Re-entrant turn/end listeners must route new input to a later turn.
- this.turnOpen = false
- this.session.append('turn/end', { turn, reason })
- }
- } catch (error: unknown) {
- retry = false
- this.loopCtx.logger.warn(`agent "${this.id}": closing turn ${turn} failed: ${errorChain(error)}`)
- emitAgentEvent(this.loopCtx, this, 'agent/error', turn, step, error)
- }
- this.retryWindow = undefined
- if (this.abort === controller) this.abort = undefined
- signal.removeEventListener('abort', cancelRetry)
- }
- if (retry) {
- await this.run({ kind: 'retry' })
- } else {
- emitAgentEvent(this.loopCtx, this, 'agent/idle', turn, idle)
- this.continueOrIdle()
- }
- }
- /**
- * Run the `agent/step` seam, commit pending input, derive one request, and
- * execute its tool calls inside one durable step boundary.
- */
- private async step(
- turn: number,
- step: number,
- signal: AbortSignal,
- ): Promise<StepOutcome> {
- const { session } = this
- // The single between-steps seam: listeners inject, steer, or edit the log
- // here; the request derives from the log after this settles.
- await this.loopCtx.serial(agentCarrier(this), 'agent/step', this, turn, step, signal)
- signal.throwIfAborted()
- // Take the outbox whole — same-boundary steering and context leave in
- // this request together.
- this.drainOutbox(turn)
- // Assemble the system prompt fresh each step (it may depend on log state).
- const assembly = await this.loopCtx.systemPrompt.assemble(assembleContextFor(this, signal))
- signal.throwIfAborted()
- const system = renderPrompt(assembly)
- // Snapshot the exact log prefix: the reconstruction boundary. Appends
- // after this synchronous snapshot join the next request.
- const boundaryMessages = session.deriveMessages()
- session.append('step/start', { turn, step })
- this.stepOpen = true
- signal.throwIfAborted()
- const request = await this.buildRequest(turn, step, assembly.tools, system, boundaryMessages, signal)
- const assembler = new BlockAssembler()
- const chunkSeqs: number[] = []
- const stream = this.loopCtx.llm.stream(request)
- try {
- for await (const chunk of stream) {
- signal.throwIfAborted()
- const chunkEvent = session.append('assistant/chunk', { turn, step, chunk })
- chunkSeqs.push(chunkEvent.seq)
- assembler.push(chunk)
- }
- } catch (error: unknown) {
- const facts = llmFailureOf(stream, error)
- if (facts !== undefined && error instanceof Error) {
- return { kind: 'request-failed', error, failure: facts }
- }
- throw error
- }
- signal.throwIfAborted()
- // Failure finish chunks take the same path as thrown stream errors.
- const finish = assembler.finish
- if (finish.kind === 'error' || finish.kind === 'aborted') {
- const error = new LlmError(finish.failure.message, finish.failure.code, finish.failure)
- return { kind: 'request-failed', error, failure: finish.failure }
- }
- // Truncated (max-tokens) output cannot owe tool calls.
- const assembled = assembler.message()
- const content = finish.kind === 'max-tokens'
- ? assembled.content.filter(block => block.type !== 'tool-call')
- : assembled.content
- session.append(
- 'assistant/message',
- {
- turn,
- step,
- content,
- provenance: {
- provider: request.provider,
- model: request.model,
- ...assembler.replayState !== undefined ? { replayState: assembler.replayState } : {},
- },
- ...assembler.usage === undefined ? {} : { usage: assembler.usage },
- },
- { surfaceOp: 'append', sourceEventSeqs: chunkSeqs },
- )
- const toolCalls = content.filter(block => block.type === 'tool-call')
- let concluded = false
- if (toolCalls.length > 0) {
- ({ concluded } = await executeToolCalls(
- this.loopCtx, turn, step, toolCalls, signal,
- context => this.outbox.push({ content: context.content, source: context.source }),
- ))
- }
- // Tool results stay adjacent to their calls; input accepted during the
- // request enters the log only after the complete result batch.
- const steered = this.drainOutbox(turn)
- session.append('step/end', { turn, step })
- this.stepOpen = false
- return {
- kind: 'completed',
- continueTurn: (toolCalls.length > 0 && !concluded) || steered,
- concluded,
- maxTokens: finish.kind === 'max-tokens',
- }
- }
- /**
- * Compose one frozen request: the `agent/request` config waterfall, the
- * canonical logged header, then the header plus the boundary snapshot,
- * byte-for-byte.
- */
- private async buildRequest(
- turn: number,
- step: number,
- tools: GenerateOptions['tools'] & object,
- system: string,
- boundaryMessages: Message[],
- signal: AbortSignal,
- ): Promise<GenerateOptions> {
- const { session } = this
- // Seed from the logged header when the log has one (the log is the
- // truth, across resumes too), else from agent options; freeze so
- // listeners must return a replacement.
- const seedConfig: LlmCallConfig = deepFreeze(structuredClone(
- session.requestHeader()?.config
- ?? { provider: this.options.provider ?? '', model: this.options.model ?? '' }))
- const config = await this.loopCtx.waterfall(
- agentCarrier(this), 'agent/request', this, turn, step, signal,
- () => Promise.resolve(seedConfig),
- )
- signal.throwIfAborted()
- if (!config.provider || !config.model) {
- throw new Error(`agent "${this.id}" has no provider/model: set AgentOptions.provider and AgentOptions.model or supply both via the agent/request waterfall`)
- }
- const header = canonicalHeader({
- config,
- ...system ? { system } : {},
- ...tools.length > 0 ? { tools } : {},
- })
- // Log the header the request will use only when it differs
- // from the folded baseline — reconstruction folds the log, so an
- // unchanged header needs no new snapshot.
- const baseline = session.requestHeader()
- if (baseline === undefined || !headerEquals(baseline, header)) {
- session.append('request/header', { header, reason: baseline === undefined ? 'initial' : 'change' })
- }
- return markAgentLoopRequest(deepFreeze({
- provider: header.config.provider,
- model: header.config.model,
- messages: boundaryMessages,
- ...header.system !== undefined ? { system: header.system } : {},
- ...header.tools !== undefined ? { tools: header.tools } : {},
- ...header.config.temperature !== undefined ? { temperature: header.config.temperature } : {},
- ...header.config.maxTokens !== undefined ? { maxTokens: header.config.maxTokens } : {},
- ...header.config.stop !== undefined ? { stop: header.config.stop } : {},
- sessionId: session.id,
- signal,
- }))
- }
- /** Commit the outbox and report whether it contained steering. */
- private drainOutbox(turn: number): boolean {
- let steered = false
- for (const message of this.outbox.splice(0)) {
- if ('id' in message) {
- steered = true
- emitAgentEvent(this.loopCtx, this, 'agent/inbox/dequeue', message)
- this.session.append(
- 'steering/message',
- { turn, content: message.content, source: message.source },
- { surfaceOp: 'append' },
- )
- } else {
- this.session.append('user/message', message, { surfaceOp: 'append' })
- }
- }
- return steered
- }
- /**
- * The single settlement funnel: classify one turn failure (interruption
- * beats error) into the durable turn/end reason and the live idle report.
- */
- private settle(
- turn: number,
- step: number,
- error: unknown,
- signal: AbortSignal,
- failure?: LlmFailure,
- ): { reason: TurnEndReason; idle: IdleReason } {
- const interrupt = agentInterruptReasonOf(signal)
- if (interrupt !== undefined) {
- return { reason: { kind: interrupt.kind === 'disposed' ? 'disposed' : 'aborted' }, idle: { kind: 'aborted' } }
- }
- if (failure !== undefined) {
- emitAgentEvent(this.loopCtx, this, 'agent/error', turn, step, error)
- // The durable record renders the full cause chain: turn/end is the one
- // durable trace of the failure, so a wrapper message alone would lose
- // the transport detail the log exists to keep.
- const rendered = errorChain(error)
- return {
- reason: { kind: 'error', step, failure: { ...failure, ...rendered === '<unrenderable value>' ? {} : { message: rendered } } },
- idle: { kind: 'error', error, failure },
- }
- }
- emitAgentEvent(this.loopCtx, this, 'agent/error', turn, step, error)
- return {
- reason: { kind: 'error', step, message: errorChain(error), ...isHarnessError(error) ? { code: error.code } : {} },
- idle: { kind: 'error', error },
- }
- }
- /** Continue with a waking prompt, or publish the idle status. */
- private continueOrIdle(): void {
- if (this.abort !== undefined) return
- if (this.queued.some(item => item.wakeup)) {
- this.kick()
- } else if (this.busy) {
- this.busy = false
- emitAgentEvent(this.loopCtx, this, 'agent/status', 'idle')
- }
- }
- }
|