返回源码地图

packages/core/agent-loop/src/tool-calls.ts

main snapshot · da00f7f5358f · 正文引用章节 07 / 07 / 08 / 08;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

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 remain
5 * model-ordered. Abort or an internal scheduler failure stops replenishment
6 * and drains started calls.
7 *
8 * Abort records synthetic error results for skipped calls so replay stays
9 * valid. A terminal scheduler failure rejects after draining; the owning step
10 * records conservative recovery results before closing.
11 * @module dsh-agent-loop/tool-calls
12 */
13
14import type { Context } from '@deepseek-ai/cordis'
15import { createToolResultMessage, type ToolCallBlock } from '@deepseek-ai/dsh-llm'
16import type { Session, SessionSeq, UserMessage } from '@deepseek-ai/dsh-session'
17import { TOOL_ABORTED_BEFORE_DISPATCH, TOOL_RUNTIME_SCHEDULER, type ToolExecutionInput, type ToolExecutionMode, type ToolExecutionResult, type ToolRunContext } from '@deepseek-ai/dsh-tools'
18import { assertNever } from '@deepseek-ai/dsh-util-values'
19
20/** One tool call after argument parsing, ready to schedule. */
21interface PlannedCall {
22 block: ToolCallBlock
23 exec: ToolExecutionInput
24}
25
26/** Settled dispatch awaiting model-order finalization. */
27interface Slot {
28 exec: ToolRunContext
29 result: ToolExecutionResult
30 needsPost: boolean
31}
32
33/** One scheduler group outcome, including a drained cancellation. */
34interface GroupOutcome {
35 consumed: number
36 aborted: boolean
37 /** Whether any committed result carried {@link ToolExecutionResult.concludesTurn}. */
38 concluded: boolean
39}
40
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. Abort
44 * drains them, records synthetic results for unstarted calls, and returns with
45 * the signal still aborted after accepting started-call context through the
46 * caller-supplied acceptor (the machine stages it in its next-step inbox for the
47 * step boundary). An internal scheduler failure stops new dispatches, drains
48 * already-started dispatches, and rejects with the first failure. The owning
49 * step supplies error results for requests without a committed outcome.
50 * The committed step's AgentLoop driver boundary supplies the initiating Agent
51 * 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 */
60export 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 } = agent
70
71 // 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 }))
82
83 let next = 0
84 let concluded = false
85 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 condition
88 const first = planned[next]!
89 const mode = ctx.tools.executionMode(first.exec).kind
90 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.consumed
95 concluded ||= outcome.concluded
96 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}
103
104/** Parse model arguments, preserving invalid JSON as text and mapping empty input to `{}`. */
105function parseArguments(raw: string): unknown {
106 try {
107 return raw ? JSON.parse(raw) : {}
108 } catch {
109 return raw
110 }
111}
112
113/**
114 * Run one exclusive barrier or parallel pool. Later calls are reclassified
115 * before start; an exclusive reclassification waits for the current pool to
116 * drain and remains for the caller's next barrier. Results and contexts commit
117 * in model order. Abort stops starts, drains and commits started calls, accepts
118 * their contexts into the owning batch, records results for skipped calls, and
119 * returns an aborted outcome. Scheduler failure drains dispatches and rejects
120 * for the owning step to record recovery results.
121 */
122async 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 = 0
137 let committed = 0
138 let started = 0
139 let aborted: boolean = signal.aborted
140 let concluded = false
141 let schedulerFailure: { error: unknown } | undefined
142 const throwSchedulerFailure = (): void => {
143 if (schedulerFailure !== undefined) throw schedulerFailure.error
144 }
145
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) break
151 const call = group[committed]
152 const result = slot.needsPost
153 ? 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 index
156 appendToolResult(session, turn, step, call!.block, result, callSeqs[committed]!)
157 for (const context of result.additionalContexts ?? []) acceptContext(context)
158 concluded ||= result.concludesTurn === true
159 committed++
160 }
161 }
162
163 const inFlight = new Map<number, Promise<number>>()
164
165 const startCall = async (index: number): Promise<void> => {
166 // oxlint-disable-next-line typescript/no-non-null-assertion -- bounded index
167 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 index
178 },
179 (error: unknown) => {
180 schedulerFailure ??= { error }
181 return index
182 },
183 )
184 inFlight.set(index, promise)
185 break
186 }
187 case 'post-result':
188 slots[index] = { exec: prepared.exec, result: prepared.result, needsPost: true }
189 break
190 case 'final-result':
191 slots[index] = { exec: prepared.exec, result: prepared.result, needsPost: false }
192 break
193 /* v8 ignore next -- closed-union exhaustiveness guard */
194 default:
195 assertNever(prepared, 'tool-call scheduler prepare result')
196 }
197 }
198
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 condition
203 const nextCall = group[nextToStart]!
204 if (nextToStart > 0 && mode === 'parallel'
205 && ctx.tools.executionMode(nextCall.exec).kind !== 'parallel') break
206 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 = true
213 }
214 }
215
216 // Ordered pre-execute may await; only dispatch/body overlaps. A scheduler
217 // failure stops new dispatches and reaches the turn boundary after every
218 // 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.
228
229 if (signal.aborted) aborted = true
230 await fillPool()
231 }
232 } catch (error: unknown) {
233 schedulerFailure ??= { error }
234 await Promise.allSettled(inFlight.values())
235 throw schedulerFailure.error
236 }
237
238 if (aborted) {
239 // Started calls and accepted context settle first; every remaining model
240 // 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}
248
249/** Append the durable call/result pair for a model call skipped after cancellation. */
250function 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}
261
262/** Append a started call and return the event seq that its result must cite. */
263function 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.seq
266}
267
268/** Append a model-ordered result linked to its call event. */
269function 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}