tool-calls.ts 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290
  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 preserves already-recorded `tool/call`
  10. * events without fabricating results.
  11. * @module dsh-agent-loop/tool-calls
  12. */
  13. import type { Context } from '@deepseek-ai/cordis'
  14. import { createToolResultMessage, type ToolCallBlock } from '@deepseek-ai/dsh-llm'
  15. import type { Session, SessionSeq, UserMessage } from '@deepseek-ai/dsh-session'
  16. import { TOOL_ABORTED_BEFORE_DISPATCH, TOOL_RUNTIME_SCHEDULER, type ToolExecutionInput, type ToolExecutionMode, type ToolExecutionResult, type ToolRunContext } from '@deepseek-ai/dsh-tools'
  17. import { assertNever } from '@deepseek-ai/dsh-util-values'
  18. /** One tool call after argument parsing, ready to schedule. */
  19. interface PlannedCall {
  20. block: ToolCallBlock
  21. exec: ToolExecutionInput
  22. }
  23. /** Settled dispatch awaiting model-order finalization. */
  24. interface Slot {
  25. exec: ToolRunContext
  26. result: ToolExecutionResult
  27. needsPost: boolean
  28. }
  29. /** One scheduler group outcome, including a drained cancellation. */
  30. interface GroupOutcome {
  31. consumed: number
  32. aborted: boolean
  33. /** Whether any committed result carried {@link ToolExecutionResult.concludesTurn}. */
  34. concluded: boolean
  35. }
  36. /**
  37. * Schedule one assistant step's tool calls by their live concurrency mode.
  38. * Ordinary completion and abort commit started-call results in order. Abort
  39. * drains them, records synthetic results for unstarted calls, and returns with
  40. * the signal still aborted after accepting started-call context through the
  41. * caller-supplied acceptor (the machine stages it in its next-step inbox for the
  42. * step boundary). An internal scheduler failure stops new dispatches, drains
  43. * already-started dispatches, and rejects with the first failure without
  44. * fabricating tool results.
  45. * The committed step's AgentLoop driver boundary supplies the initiating Agent
  46. * that becomes each explicit {@link ToolExecutionInput.agent}.
  47. *
  48. * @param ctx - loop context that owns the tool registry and carries the initiating Agent.
  49. * @param turn - current turn number.
  50. * @param step - current step number.
  51. * @param toolCalls - assistant calls in model order.
  52. * @param signal - abort signal shared by the step.
  53. * @param acceptContext - accepts committed result context for the next step boundary.
  54. */
  55. export async function executeToolCalls(
  56. ctx: Context,
  57. turn: number,
  58. step: number,
  59. toolCalls: ToolCallBlock[],
  60. signal: AbortSignal,
  61. acceptContext: (context: UserMessage) => void,
  62. ): Promise<{ concluded: boolean }> {
  63. const agent = ctx.agents.requireInitiator()
  64. const { session } = agent
  65. // Inputs are distinct because tools/execute wrappers may replace `exec.signal`.
  66. const planned: PlannedCall[] = toolCalls.map(block => ({
  67. block,
  68. exec: {
  69. callId: block.id,
  70. name: block.name,
  71. arguments: parseArguments(block.arguments),
  72. agent,
  73. signal,
  74. },
  75. }))
  76. let next = 0
  77. let concluded = false
  78. while (next < planned.length) {
  79. // Commit before classifying again so registry changes affect unstarted calls.
  80. // oxlint-disable-next-line typescript/no-non-null-assertion -- bounded by the loop condition
  81. const first = planned[next]!
  82. const mode = ctx.tools.executionMode(first.exec).kind
  83. const group = mode === 'parallel' ? planned.slice(next) : [first]
  84. const outcome = await runGroup(
  85. ctx, turn, step, group, mode, signal, acceptContext,
  86. )
  87. next += outcome.consumed
  88. concluded ||= outcome.concluded
  89. if (outcome.aborted) {
  90. for (const call of planned.slice(next)) appendSkippedToolCall(session, turn, step, call.block)
  91. return { concluded }
  92. }
  93. }
  94. return { concluded }
  95. }
  96. /** Parse model arguments, preserving invalid JSON as text and mapping empty input to `{}`. */
  97. function parseArguments(raw: string): unknown {
  98. try {
  99. return raw ? JSON.parse(raw) : {}
  100. } catch {
  101. return raw
  102. }
  103. }
  104. /**
  105. * Run one exclusive barrier or parallel pool. Later calls are reclassified
  106. * before start; an exclusive reclassification waits for the current pool to
  107. * drain and remains for the caller's next barrier. Results and contexts commit
  108. * in model order. Abort stops starts, drains and commits started calls, accepts
  109. * their contexts into the owning batch, records results for skipped calls, and
  110. * returns an aborted outcome. Scheduler failure drains dispatches without
  111. * committing synthetic recovery results.
  112. */
  113. async function runGroup(
  114. ctx: Context,
  115. turn: number,
  116. step: number,
  117. group: PlannedCall[],
  118. mode: ToolExecutionMode['kind'],
  119. signal: AbortSignal,
  120. acceptContext: (context: UserMessage) => void,
  121. ): Promise<GroupOutcome> {
  122. const { session } = ctx.agents.requireInitiator()
  123. const { maxParallelToolCalls } = ctx.agentLoop.config
  124. const slots: (Slot | undefined)[] = group.map(() => undefined)
  125. // Started slots retain their `tool/call` seq so the result can cite it.
  126. const callSeqs: Array<SessionSeq | undefined> = group.map(() => undefined)
  127. let nextToStart = 0
  128. let committed = 0
  129. let started = 0
  130. let aborted: boolean = signal.aborted
  131. let concluded = false
  132. let schedulerFailure: { error: unknown } | undefined
  133. const throwSchedulerFailure = (): void => {
  134. if (schedulerFailure !== undefined) throw schedulerFailure.error
  135. }
  136. // `committed` advances only across contiguous model-order slots.
  137. const commitReady = async (): Promise<void> => {
  138. while (committed < group.length) {
  139. const slot = slots[committed]
  140. if (slot === undefined) break
  141. const call = group[committed]
  142. const result = slot.needsPost
  143. ? await ctx.tools[TOOL_RUNTIME_SCHEDULER].finalize(slot.exec, slot.result)
  144. : ctx.tools[TOOL_RUNTIME_SCHEDULER].finish(slot.exec, slot.result)
  145. // oxlint-disable-next-line typescript/no-non-null-assertion -- bounded index
  146. appendToolResult(session, turn, step, call!.block, result, callSeqs[committed]!)
  147. for (const context of result.additionalContexts ?? []) acceptContext(context)
  148. concluded ||= result.concludesTurn === true
  149. committed++
  150. }
  151. }
  152. const inFlight = new Map<number, Promise<number>>()
  153. const startCall = async (index: number): Promise<void> => {
  154. // oxlint-disable-next-line typescript/no-non-null-assertion -- bounded index
  155. const call = group[index]!
  156. callSeqs[index] = appendToolCall(session, turn, step, call.block)
  157. started++
  158. const prepared = await ctx.tools[TOOL_RUNTIME_SCHEDULER].prepare(call.exec)
  159. throwSchedulerFailure()
  160. switch (prepared.kind) {
  161. case 'dispatch': {
  162. const promise = ctx.tools[TOOL_RUNTIME_SCHEDULER].dispatch(prepared.exec).then(
  163. (outcome) => {
  164. slots[index] = { exec: prepared.exec, result: outcome.result, needsPost: outcome.kind === 'post-result' }
  165. return index
  166. },
  167. (error: unknown) => {
  168. schedulerFailure ??= { error }
  169. return index
  170. },
  171. )
  172. inFlight.set(index, promise)
  173. break
  174. }
  175. case 'post-result':
  176. slots[index] = { exec: prepared.exec, result: prepared.result, needsPost: true }
  177. break
  178. case 'final-result':
  179. slots[index] = { exec: prepared.exec, result: prepared.result, needsPost: false }
  180. break
  181. /* v8 ignore next -- closed-union exhaustiveness guard */
  182. default:
  183. assertNever(prepared, 'tool-call scheduler prepare result')
  184. }
  185. }
  186. const fillPool = async (): Promise<void> => {
  187. while (!aborted && nextToStart < group.length && inFlight.size < maxParallelToolCalls) {
  188. // Re-read later modes after ordered commits so registry changes can create a barrier.
  189. // oxlint-disable-next-line typescript/no-non-null-assertion -- bounded by the loop condition
  190. const nextCall = group[nextToStart]!
  191. if (nextToStart > 0 && mode === 'parallel'
  192. && ctx.tools.executionMode(nextCall.exec).kind !== 'parallel') break
  193. await startCall(nextToStart)
  194. nextToStart++
  195. throwSchedulerFailure()
  196. await commitReady()
  197. throwSchedulerFailure()
  198. // Abort may arrive while pre-execute awaits.
  199. if (signal.aborted) aborted = true
  200. }
  201. }
  202. // Ordered pre-execute may await; only dispatch/body overlaps. A scheduler
  203. // failure stops new dispatches and reaches the turn boundary after every
  204. // already-started dispatch settles.
  205. try {
  206. await fillPool()
  207. while (inFlight.size > 0) {
  208. const settledIndex = await Promise.race(inFlight.values())
  209. inFlight.delete(settledIndex)
  210. throwSchedulerFailure()
  211. await commitReady()
  212. throwSchedulerFailure()
  213. // Abort may arrive while a tool or ordered commit awaits.
  214. if (signal.aborted) aborted = true
  215. await fillPool()
  216. }
  217. } catch (error: unknown) {
  218. schedulerFailure ??= { error }
  219. await Promise.allSettled(inFlight.values())
  220. throw schedulerFailure.error
  221. }
  222. if (aborted) {
  223. // Started calls and accepted context settle first; every remaining model
  224. // call then receives an ordered synthetic result before the turn aborts.
  225. for (const call of group.slice(started)) appendSkippedToolCall(session, turn, step, call.block)
  226. return { consumed: group.length, aborted: true, concluded }
  227. }
  228. /* v8 ignore next -- unreachable: a non-aborted group commits every started call */
  229. if (committed !== started) throw new Error('tool-call scheduler: uncommitted settled calls')
  230. return { consumed: started, aborted: false, concluded }
  231. }
  232. /** Append the durable call/result pair for a model call skipped after cancellation. */
  233. function appendSkippedToolCall(session: Session, turn: number, step: number, block: ToolCallBlock): void {
  234. const callSeq = appendToolCall(session, turn, step, block)
  235. appendToolResult(session, turn, step, block, {
  236. content: [{ type: 'text', text: 'Error: tool call aborted before dispatch' }],
  237. isError: true,
  238. error: {
  239. message: 'tool call aborted before dispatch',
  240. info: { name: 'AbortError', code: TOOL_ABORTED_BEFORE_DISPATCH },
  241. },
  242. }, callSeq)
  243. }
  244. /** Append a started call and return the event seq that its result must cite. */
  245. function appendToolCall(session: Session, turn: number, step: number, block: ToolCallBlock): SessionSeq {
  246. const event = session.append('tool/call', { turn, step, callId: block.id, name: block.name, arguments: block.arguments })
  247. return event.seq
  248. }
  249. /** Append a model-ordered result linked to its call event. */
  250. function appendToolResult(
  251. session: Session,
  252. turn: number,
  253. step: number,
  254. block: ToolCallBlock,
  255. result: ToolExecutionResult,
  256. callSeq: SessionSeq,
  257. ): void {
  258. const message = createToolResultMessage({
  259. callId: block.id,
  260. content: result.content,
  261. isError: result.isError,
  262. })
  263. session.append('tool/result', {
  264. turn, step,
  265. message,
  266. ...result.error?.info ? { error: result.error.info } : {},
  267. // The tool's private presentation payload (e.g. a result-time diff),
  268. // persisted so a UI bridge reproduces the card on replay.
  269. ...result.meta !== undefined ? { meta: result.meta } : {},
  270. }, { surfaceOp: 'append', sourceEventSeqs: [callSeq] })
  271. }