1
/**2
* Schedules one assistant step's tool calls. Exclusive calls form barriers;3
* parallel calls use a bounded rolling pool and are reclassified before start.4
* Dispatch may overlap, while policy, results, and result context remain5
* model-ordered. Abort or an internal scheduler failure stops replenishment6
* and drains started calls.7
*8
* Abort records synthetic error results for skipped calls so replay stays9
* valid. A terminal scheduler failure rejects after draining; the owning step10
* records conservative recovery results before closing.11
* @module dsh-agent-loop/tool-calls12
*/14
import type { Context } from '@deepseek-ai/cordis'15
import { createToolResultMessage, type ToolCallBlock } from '@deepseek-ai/dsh-llm'16
import type { Session, SessionSeq, UserMessage } from '@deepseek-ai/dsh-session'17
import { TOOL_ABORTED_BEFORE_DISPATCH, TOOL_RUNTIME_SCHEDULER, type ToolExecutionInput, type ToolExecutionMode, type ToolExecutionResult, type ToolRunContext } from '@deepseek-ai/dsh-tools'18
import { assertNever } from '@deepseek-ai/dsh-util-values'20
/** One tool call after argument parsing, ready to schedule. */21
interface PlannedCall {22
block: ToolCallBlock23
exec: ToolExecutionInput24
}26
/** Settled dispatch awaiting model-order finalization. */27
interface Slot {28
exec: ToolRunContext29
result: ToolExecutionResult30
needsPost: boolean31
}33
/** One scheduler group outcome, including a drained cancellation. */34
interface GroupOutcome {35
consumed: number36
aborted: boolean37
/** Whether any committed result carried {@link ToolExecutionResult.concludesTurn}. */38
concluded: boolean39
}41
/**42
* Schedule one assistant step's tool calls by their live concurrency mode.43
* Ordinary completion and abort commit started-call results in order. Abort44
* drains them, records synthetic results for unstarted calls, and returns with45
* the signal still aborted after accepting started-call context through the46
* caller-supplied acceptor (the machine stages it in its next-step inbox for the47
* step boundary). An internal scheduler failure stops new dispatches, drains48
* already-started dispatches, and rejects with the first failure. The owning49
* step supplies error results for requests without a committed outcome.50
* The committed step's AgentLoop driver boundary supplies the initiating Agent51
* that becomes each explicit {@link ToolExecutionInput.agent}.52
*53
* @param ctx - loop context that owns the tool registry and carries the initiating Agent.54
* @param turn - current turn number.55
* @param step - current step number.56
* @param toolCalls - assistant calls in model order.57
* @param signal - abort signal shared by the step.58
* @param acceptContext - accepts committed result context for the next step boundary.59
*/60
export async function executeToolCalls(61
ctx: Context,62
turn: number,63
step: number,64
toolCalls: ToolCallBlock[],65
signal: AbortSignal,66
acceptContext: (context: UserMessage) => void,67
): Promise<{ concluded: boolean }> {68
const agent = ctx.agents.requireInitiator()69
const { session } = agent71
// Inputs are distinct because tools/execute wrappers may replace `exec.signal`.72
const planned: PlannedCall[] = toolCalls.map(block => ({73
block,74
exec: {75
callId: block.id,76
name: block.name,77
arguments: parseArguments(block.arguments),78
agent,79
signal,80
},81
}))83
let next = 084
let concluded = false85
while (next < planned.length) {86
// Commit before classifying again so registry changes affect unstarted calls.87
// oxlint-disable-next-line typescript/no-non-null-assertion -- bounded by the loop condition88
const first = planned[next]!89
const mode = ctx.tools.executionMode(first.exec).kind90
const group = mode === 'parallel' ? planned.slice(next) : [first]91
const outcome = await runGroup(92
ctx, turn, step, group, mode, signal, acceptContext,93
)94
next += outcome.consumed95
concluded ||= outcome.concluded96
if (outcome.aborted) {97
for (const call of planned.slice(next)) appendSkippedToolCall(session, turn, step, call.block)98
return { concluded }99
}100
}101
return { concluded }102
}104
/** Parse model arguments, preserving invalid JSON as text and mapping empty input to `{}`. */105
function parseArguments(raw: string): unknown {106
try {107
return raw ? JSON.parse(raw) : {}108
} catch {109
return raw110
}111
}113
/**114
* Run one exclusive barrier or parallel pool. Later calls are reclassified115
* before start; an exclusive reclassification waits for the current pool to116
* drain and remains for the caller's next barrier. Results and contexts commit117
* in model order. Abort stops starts, drains and commits started calls, accepts118
* their contexts into the owning batch, records results for skipped calls, and119
* returns an aborted outcome. Scheduler failure drains dispatches and rejects120
* for the owning step to record recovery results.121
*/122
async function runGroup(123
ctx: Context,124
turn: number,125
step: number,126
group: PlannedCall[],127
mode: ToolExecutionMode['kind'],128
signal: AbortSignal,129
acceptContext: (context: UserMessage) => void,130
): Promise<GroupOutcome> {131
const { session } = ctx.agents.requireInitiator()132
const maxParallelToolCalls = ctx.agentLoop.config.maxParallelToolCalls.get()133
const slots: (Slot | undefined)[] = group.map(() => undefined)134
// Started slots retain their `tool/call` seq so the result can cite it.135
const callSeqs: Array<SessionSeq | undefined> = group.map(() => undefined)136
let nextToStart = 0137
let committed = 0138
let started = 0139
let aborted: boolean = signal.aborted140
let concluded = false141
let schedulerFailure: { error: unknown } | undefined142
const throwSchedulerFailure = (): void => {143
if (schedulerFailure !== undefined) throw schedulerFailure.error144
}146
// `committed` advances only across contiguous model-order slots.147
const commitReady = async (): Promise<void> => {148
while (committed < group.length) {149
const slot = slots[committed]150
if (slot === undefined) break151
const call = group[committed]152
const result = slot.needsPost153
? await ctx.tools[TOOL_RUNTIME_SCHEDULER].finalize(slot.exec, slot.result)154
: ctx.tools[TOOL_RUNTIME_SCHEDULER].finish(slot.exec, slot.result)155
// oxlint-disable-next-line typescript/no-non-null-assertion -- bounded index156
appendToolResult(session, turn, step, call!.block, result, callSeqs[committed]!)157
for (const context of result.additionalContexts ?? []) acceptContext(context)158
concluded ||= result.concludesTurn === true159
committed++160
}161
}163
const inFlight = new Map<number, Promise<number>>()165
const startCall = async (index: number): Promise<void> => {166
// oxlint-disable-next-line typescript/no-non-null-assertion -- bounded index167
const call = group[index]!168
callSeqs[index] = appendToolCall(session, turn, step, call.block)169
started++170
const prepared = await ctx.tools[TOOL_RUNTIME_SCHEDULER].prepare(call.exec)171
throwSchedulerFailure()172
switch (prepared.kind) {173
case 'dispatch': {174
const promise = ctx.tools[TOOL_RUNTIME_SCHEDULER].dispatch(prepared.exec).then(175
(outcome) => {176
slots[index] = { exec: prepared.exec, result: outcome.result, needsPost: outcome.kind === 'post-result' }177
return index178
},179
(error: unknown) => {180
schedulerFailure ??= { error }181
return index182
},183
)184
inFlight.set(index, promise)185
break186
}187
case 'post-result':188
slots[index] = { exec: prepared.exec, result: prepared.result, needsPost: true }189
break190
case 'final-result':191
slots[index] = { exec: prepared.exec, result: prepared.result, needsPost: false }192
break193
/* v8 ignore next -- closed-union exhaustiveness guard */194
default:195
assertNever(prepared, 'tool-call scheduler prepare result')196
}197
}199
const fillPool = async (): Promise<void> => {200
while (!aborted && nextToStart < group.length && inFlight.size < maxParallelToolCalls) {201
// Re-read later modes after ordered commits so registry changes can create a barrier.202
// oxlint-disable-next-line typescript/no-non-null-assertion -- bounded by the loop condition203
const nextCall = group[nextToStart]!204
if (nextToStart > 0 && mode === 'parallel'205
&& ctx.tools.executionMode(nextCall.exec).kind !== 'parallel') break206
await startCall(nextToStart)207
nextToStart++208
throwSchedulerFailure()209
await commitReady()210
throwSchedulerFailure()211
// Abort may arrive while pre-execute awaits.212
if (signal.aborted) aborted = true213
}214
}216
// Ordered pre-execute may await; only dispatch/body overlaps. A scheduler217
// failure stops new dispatches and reaches the turn boundary after every218
// already-started dispatch settles.219
try {220
await fillPool()221
while (inFlight.size > 0) {222
const settledIndex = await Promise.race(inFlight.values())223
inFlight.delete(settledIndex)224
throwSchedulerFailure()225
await commitReady()226
throwSchedulerFailure()227
// Abort may arrive while a tool or ordered commit awaits.229
if (signal.aborted) aborted = true230
await fillPool()231
}232
} catch (error: unknown) {233
schedulerFailure ??= { error }234
await Promise.allSettled(inFlight.values())235
throw schedulerFailure.error236
}238
if (aborted) {239
// Started calls and accepted context settle first; every remaining model240
// call then receives an ordered synthetic result before the turn aborts.241
for (const call of group.slice(started)) appendSkippedToolCall(session, turn, step, call.block)242
return { consumed: group.length, aborted: true, concluded }243
}244
/* v8 ignore next -- unreachable: a non-aborted group commits every started call */245
if (committed !== started) throw new Error('tool-call scheduler: uncommitted settled calls')246
return { consumed: started, aborted: false, concluded }247
}249
/** Append the durable call/result pair for a model call skipped after cancellation. */250
function appendSkippedToolCall(session: Session, turn: number, step: number, block: ToolCallBlock): void {251
const callSeq = appendToolCall(session, turn, step, block)252
appendToolResult(session, turn, step, block, {253
content: [{ type: 'text', text: 'Error: tool call aborted before dispatch' }],254
isError: true,255
error: {256
message: 'tool call aborted before dispatch',257
info: { name: 'AbortError', code: TOOL_ABORTED_BEFORE_DISPATCH },258
},259
}, callSeq)260
}262
/** Append a started call and return the event seq that its result must cite. */263
function appendToolCall(session: Session, turn: number, step: number, block: ToolCallBlock): SessionSeq {264
const event = session.append('tool/call', { turn, step, callId: block.id, name: block.name, arguments: block.arguments })265
return event.seq266
}268
/** Append a model-ordered result linked to its call event. */269
function appendToolResult(270
session: Session,271
turn: number,272
step: number,273
block: ToolCallBlock,274
result: ToolExecutionResult,275
callSeq: SessionSeq,276
): void {277
const message = createToolResultMessage({278
callId: block.id,279
content: result.content,280
isError: result.isError,281
})282
session.append('tool/result', {283
turn, step,284
message,285
...result.error?.info ? { error: result.error.info } : {},286
// The tool's private presentation payload (e.g. a result-time diff),287
// persisted so a UI bridge reproduces the card on replay.288
...result.meta !== undefined ? { meta: result.meta } : {},289
}, { surfaceOp: 'append', sourceEventSeqs: [callSeq] })290
}