| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262 |
- /**
- * Schedules one assistant step's tool calls. Exclusive calls form barriers;
- * parallel calls use a bounded rolling pool and are reclassified before start.
- * Dispatch may overlap, while policy, results, and result context remain
- * model-ordered. Abort stops replenishment and drains started calls.
- *
- * Each advertised call records a balanced `tool/call`/`tool/result` pair. Calls
- * skipped after abort receive synthetic error results so replay stays valid.
- * @module dsh-agent-loop/tool-calls
- */
- import type { Context } from 'cordis'
- import { assertNever, createToolResultMessage, type ToolCallBlock } from '@deepseek-ai/dsh-llm'
- import type { Session, UserMessage } from '@deepseek-ai/dsh-session'
- import { TOOL_ABORTED_BEFORE_DISPATCH, TOOL_REGISTRY_SCHEDULER, type ToolExecutionInput, type ToolExecutionMode, type ToolExecutionResult, type ToolRunContext } from '@deepseek-ai/dsh-tools'
- /** One tool call after argument parsing, ready to schedule. */
- interface PlannedCall {
- block: ToolCallBlock
- exec: ToolExecutionInput
- }
- /** Settled dispatch awaiting model-order finalization. */
- interface Slot {
- exec: ToolRunContext
- result: ToolExecutionResult
- needsPost: boolean
- }
- /** One scheduler group outcome, including a drained cancellation. */
- interface GroupOutcome {
- consumed: number
- aborted: boolean
- /** Whether any committed result carried {@link ToolExecutionResult.concludesTurn}. */
- concluded: boolean
- }
- /**
- * Schedule one assistant step's tool calls by their live concurrency mode.
- * Started calls receive ordered results. Abort drains them, records synthetic
- * results for unstarted calls, and returns with the signal still aborted after
- * accepting started-call context through the caller-supplied acceptor (the
- * machine stages it on its outbox for the next step boundary).
- * The committed step's AgentLoop driver boundary supplies the initiating Agent
- * that becomes each explicit {@link ToolExecutionInput.agent}.
- *
- * @param ctx - loop context that owns the tool registry and carries the initiating Agent.
- * @param turn - current turn number.
- * @param step - current step number.
- * @param toolCalls - assistant calls in model order.
- * @param signal - abort signal shared by the step.
- * @param acceptContext - accepts committed result context for the next step boundary.
- */
- export async function executeToolCalls(
- ctx: Context,
- turn: number,
- step: number,
- toolCalls: ToolCallBlock[],
- signal: AbortSignal,
- acceptContext: (context: UserMessage) => void,
- ): Promise<{ concluded: boolean }> {
- const agent = ctx.agents.requireInitiator()
- const { session } = agent
- // Inputs are distinct because tools/execute wrappers may replace `exec.signal`.
- const planned: PlannedCall[] = toolCalls.map(block => ({
- block,
- exec: {
- callId: block.id,
- name: block.name,
- arguments: parseArguments(block.arguments),
- agent,
- signal,
- },
- }))
- let next = 0
- let concluded = false
- while (next < planned.length) {
- // Commit before classifying again so registry changes affect unstarted calls.
- // eslint-disable-next-line @typescript-eslint/no-non-null-assertion -- bounded by the loop condition
- const first = planned[next]!
- const mode = ctx.tools.executionMode(first.exec).kind
- const group = mode === 'parallel' ? planned.slice(next) : [first]
- const outcome = await runGroup(
- ctx, turn, step, group, mode, signal, acceptContext,
- )
- next += outcome.consumed
- concluded ||= outcome.concluded
- if (outcome.aborted) {
- for (const call of planned.slice(next)) appendSkippedToolCall(session, turn, step, call.block)
- return { concluded }
- }
- }
- return { concluded }
- }
- /** Parse model arguments, preserving invalid JSON as text and mapping empty input to `{}`. */
- function parseArguments(raw: string): unknown {
- try {
- return raw ? JSON.parse(raw) : {}
- } catch {
- return raw
- }
- }
- /**
- * Run one exclusive barrier or parallel pool. Later calls are reclassified
- * before start; an exclusive reclassification waits for the current pool to
- * drain and remains for the caller's next barrier. Results and contexts commit
- * in model order. Abort stops starts, drains and commits started calls, accepts
- * their contexts into the owning batch, records results for skipped calls, and
- * returns an aborted outcome.
- */
- async function runGroup(
- ctx: Context,
- turn: number,
- step: number,
- group: PlannedCall[],
- mode: ToolExecutionMode['kind'],
- signal: AbortSignal,
- acceptContext: (context: UserMessage) => void,
- ): Promise<GroupOutcome> {
- const { session } = ctx.agents.requireInitiator()
- const { maxParallelToolCalls } = ctx.agentLoop.config
- const slots: (Slot | undefined)[] = group.map(() => undefined)
- // Started slots retain their tool/call seq for result provenance.
- const callSeqs: number[] = group.map(() => -1)
- let nextToStart = 0
- let committed = 0
- let started = 0
- let aborted: boolean = signal.aborted
- let concluded = false
- // `committed` advances only across contiguous model-order slots.
- const commitReady = async (): Promise<void> => {
- while (committed < group.length) {
- const slot = slots[committed]
- if (slot === undefined) break
- const call = group[committed]
- const result = slot.needsPost
- ? await ctx.tools[TOOL_REGISTRY_SCHEDULER].finalize(slot.exec, slot.result)
- : ctx.tools[TOOL_REGISTRY_SCHEDULER].finish(slot.exec, slot.result)
- // eslint-disable-next-line @typescript-eslint/no-non-null-assertion -- bounded index
- appendToolResult(session, turn, step, call!.block, result, callSeqs[committed]!)
- for (const context of result.additionalContexts ?? []) acceptContext(context)
- concluded ||= result.concludesTurn === true
- committed++
- }
- }
- const inFlight = new Map<number, Promise<number>>()
- const startCall = async (index: number): Promise<void> => {
- // eslint-disable-next-line @typescript-eslint/no-non-null-assertion -- bounded index
- const call = group[index]!
- callSeqs[index] = appendToolCall(session, turn, step, call.block)
- started++
- const prepared = await ctx.tools[TOOL_REGISTRY_SCHEDULER].prepare(call.exec)
- switch (prepared.kind) {
- case 'dispatch': {
- const promise = ctx.tools[TOOL_REGISTRY_SCHEDULER].dispatch(prepared.exec).then((outcome) => {
- slots[index] = { exec: prepared.exec, result: outcome.result, needsPost: outcome.kind === 'post-result' }
- return index
- })
- inFlight.set(index, promise)
- break
- }
- case 'post-result':
- slots[index] = { exec: prepared.exec, result: prepared.result, needsPost: true }
- break
- case 'final-result':
- slots[index] = { exec: prepared.exec, result: prepared.result, needsPost: false }
- break
- /* v8 ignore next -- closed-union exhaustiveness guard */
- default:
- assertNever(prepared, 'tool-call scheduler prepare result')
- }
- }
- const fillPool = async (): Promise<void> => {
- while (!aborted && nextToStart < group.length && inFlight.size < maxParallelToolCalls) {
- // Re-read later modes after ordered commits so registry changes can create a barrier.
- // eslint-disable-next-line @typescript-eslint/no-non-null-assertion -- bounded by the loop condition
- const nextCall = group[nextToStart]!
- if (nextToStart > 0 && mode === 'parallel'
- && ctx.tools.executionMode(nextCall.exec).kind !== 'parallel') break
- await startCall(nextToStart)
- nextToStart++
- await commitReady()
- // Abort may arrive while pre-execute awaits.
- if (signal.aborted) aborted = true
- }
- }
- // Ordered pre-execute may await; only dispatch/body overlaps.
- // TODO: Drain every started call before rethrowing a scheduler error; tool
- // bodies must not outlive the failed turn.
- await fillPool()
- while (inFlight.size > 0) {
- const settledIndex = await Promise.race(inFlight.values())
- inFlight.delete(settledIndex)
- await commitReady()
- // Abort may arrive while a tool or ordered commit awaits.
- if (signal.aborted) aborted = true
- await fillPool()
- }
- if (aborted) {
- // Started calls and accepted context settle first; every remaining model
- // call then receives an ordered synthetic result before the turn aborts.
- for (const call of group.slice(started)) appendSkippedToolCall(session, turn, step, call.block)
- return { consumed: group.length, aborted: true, concluded }
- }
- /* v8 ignore next -- unreachable: a non-aborted group commits every started call */
- if (committed !== started) throw new Error('tool-call scheduler: uncommitted settled calls')
- return { consumed: started, aborted: false, concluded }
- }
- /** Append the durable call/result pair for a model call skipped after cancellation. */
- function appendSkippedToolCall(session: Session, turn: number, step: number, block: ToolCallBlock): void {
- const callSeq = appendToolCall(session, turn, step, block)
- appendToolResult(session, turn, step, block, {
- content: [{ type: 'text', text: 'Error: tool call aborted before dispatch' }],
- isError: true,
- error: {
- message: 'tool call aborted before dispatch',
- info: { name: 'AbortError', code: TOOL_ABORTED_BEFORE_DISPATCH },
- },
- }, callSeq)
- }
- /** Append a started call and return its provenance sequence. */
- function appendToolCall(session: Session, turn: number, step: number, block: ToolCallBlock): number {
- const event = session.append('tool/call', { turn, step, callId: block.id, name: block.name, arguments: block.arguments })
- return event.seq
- }
- /** Append a model-ordered result linked to its call event. */
- function appendToolResult(
- session: Session,
- turn: number,
- step: number,
- block: ToolCallBlock,
- result: ToolExecutionResult,
- callSeq: number,
- ): void {
- const message = createToolResultMessage({
- callId: block.id,
- content: result.content,
- isError: result.isError,
- })
- session.append('tool/result', {
- turn, step,
- message,
- ...result.error?.info ? { error: result.error.info } : {},
- // The tool's private presentation payload (e.g. a result-time diff),
- // persisted so a UI bridge reproduces the card on replay.
- ...result.meta !== undefined ? { meta: result.meta } : {},
- }, { surfaceOp: 'append', sourceEventSeqs: [callSeq] })
- }
|