返回源码地图

packages/agent/src/agent-loop.ts

v1.0.0 · a13d35a742c6 · 05:完整主循环和工具执行管线

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

1/**
2 * Agent loop that works with AgentMessage throughout.
3 * Transforms to Message[] only at the LLM call boundary.
4 */
5
6import {
7 type AssistantMessage,
8 EventStream,
9 getCurrentTools,
10 getToolStateChanges,
11 normalizeContext,
12 type SystemMessage,
13 type ToolResultMessage,
14 type ToolStateChanges,
15 toToolDeclaration,
16 validateToolArguments,
17} from "@earendil-works/pi-ai";
18import { getDefaultStreamFn } from "./stream-fn.ts";
19import type {
20 AgentContext,
21 AgentEvent,
22 AgentLoopConfig,
23 AgentMessage,
24 AgentTool,
25 AgentToolCall,
26 AgentToolCallOutcome,
27 AgentToolResult,
28 PrepareNextTurnContext,
29 StreamFn,
30} from "./types.ts";
31
32export type AgentEventSink = (event: AgentEvent) => Promise<void> | void;
33
34/**
35 * Start an agent loop with a new prompt message.
36 * The prompt is added to the context and events are emitted for it.
37 */
38export function agentLoop(
39 prompts: AgentMessage[],
40 context: AgentContext,
41 config: AgentLoopConfig,
42 signal: AbortSignal | undefined,
43 streamFn: StreamFn,
44): EventStream<AgentEvent, AgentMessage[]> {
45 const stream = createAgentStream();
46
47 void runAgentLoop(
48 prompts,
49 context,
50 config,
51 async (event) => {
52 stream.push(event);
53 },
54 signal,
55 streamFn,
56 ).then((messages) => {
57 stream.end(messages);
58 });
59
60 return stream;
61}
62
63/**
64 * Continue an agent loop from the current context without adding a new message.
65 * Used for retries - context already has user message or tool results.
66 *
67 * **Important:** The last message in context must convert to a `user` or `toolResult` message
68 * via `convertToLlm`. If it doesn't, the LLM provider will reject the request.
69 * This cannot be validated here since `convertToLlm` is only called once per turn.
70 */
71export function agentLoopContinue(
72 context: AgentContext,
73 config: AgentLoopConfig,
74 signal: AbortSignal | undefined,
75 streamFn: StreamFn,
76): EventStream<AgentEvent, AgentMessage[]> {
77 if (context.messages.length === 0) {
78 throw new Error("Cannot continue: no messages in context");
79 }
80
81 if (context.messages[context.messages.length - 1].role === "assistant") {
82 throw new Error("Cannot continue from message role: assistant");
83 }
84
85 const stream = createAgentStream();
86
87 void runAgentLoopContinue(
88 context,
89 config,
90 async (event) => {
91 stream.push(event);
92 },
93 signal,
94 streamFn,
95 ).then((messages) => {
96 stream.end(messages);
97 });
98
99 return stream;
100}
101
102export async function runAgentLoop(
103 prompts: AgentMessage[],
104 context: AgentContext,
105 config: AgentLoopConfig,
106 emit: AgentEventSink,
107 signal: AbortSignal | undefined,
108 streamFn: StreamFn,
109): Promise<AgentMessage[]> {
110 const initialMessages = declareToolChanges(context, prompts);
111 const newMessages: AgentMessage[] = [...initialMessages];
112 const currentContext: AgentContext = {
113 ...context,
114 messages: [...context.messages, ...initialMessages],
115 };
116
117 await emit({ type: "agent_start" });
118 await emit({ type: "turn_start" });
119 for (const message of initialMessages) {
120 await emit({ type: "message_start", message });
121 await emit({ type: "message_end", message });
122 }
123
124 await runLoop(currentContext, newMessages, config, signal, emit, streamFn ?? getDefaultStreamFn());
125 return newMessages;
126}
127
128export async function runAgentLoopContinue(
129 context: AgentContext,
130 config: AgentLoopConfig,
131 emit: AgentEventSink,
132 signal: AbortSignal | undefined,
133 streamFn: StreamFn,
134): Promise<AgentMessage[]> {
135 if (context.messages.length === 0) {
136 throw new Error("Cannot continue: no messages in context");
137 }
138
139 if (context.messages[context.messages.length - 1].role === "assistant") {
140 throw new Error("Cannot continue from message role: assistant");
141 }
142
143 const newMessages: AgentMessage[] = [];
144 const currentContext: AgentContext = { ...context };
145
146 await emit({ type: "agent_start" });
147 await emit({ type: "turn_start" });
148
149 await runLoop(currentContext, newMessages, config, signal, emit, streamFn ?? getDefaultStreamFn());
150 return newMessages;
151}
152
153function createAgentStream(): EventStream<AgentEvent, AgentMessage[]> {
154 return new EventStream<AgentEvent, AgentMessage[]>(
155 (event: AgentEvent) => event.type === "agent_end",
156 (event: AgentEvent) => (event.type === "agent_end" ? event.messages : []),
157 );
158}
159
160/**
161 * Main loop logic shared by agentLoop and agentLoopContinue.
162 */
163async function runLoop(
164 initialContext: AgentContext,
165 newMessages: AgentMessage[],
166 initialConfig: AgentLoopConfig,
167 signal: AbortSignal | undefined,
168 emit: AgentEventSink,
169 streamFunction: StreamFn,
170): Promise<void> {
171 let currentContext = initialContext;
172 let config = initialConfig;
173 let lastCompletedTurn: PrepareNextTurnContext | undefined;
174 let explicitContinuation = false;
175 // Check for steering messages at start (user may have typed while waiting)
176 let pendingMessages: AgentMessage[] = (await config.getSteeringMessages?.()) || [];
177
178 // Outer loop: continues when queued follow-up messages arrive after agent would stop
179 while (true) {
180 let hasMoreToolCalls = true;
181
182 // Inner loop: process tool calls and steering messages
183 while (hasMoreToolCalls || pendingMessages.length > 0) {
184 let preparedMessages: AgentMessage[] = [];
185 if (lastCompletedTurn) {
186 const nextTurnSnapshot = await config.prepareNextTurn?.(lastCompletedTurn);
187 if (nextTurnSnapshot) {
188 currentContext = nextTurnSnapshot.context ?? currentContext;
189 preparedMessages = nextTurnSnapshot.messages ?? [];
190 config = {
191 ...config,
192 model: nextTurnSnapshot.model ?? config.model,
193 reasoning:
194 nextTurnSnapshot.thinkingLevel === undefined
195 ? config.reasoning
196 : nextTurnSnapshot.thinkingLevel === "off"
197 ? undefined
198 : nextTurnSnapshot.thinkingLevel,
199 };
200 }
201 // Preparation can be long-running (for example, compaction). Pick up steering
202 // queued while it ran. Only poll again if the earlier poll returned nothing;
203 // otherwise one-at-a-time mode would deliver two messages in this turn.
204 if (pendingMessages.length === 0) {
205 pendingMessages = (await config.getSteeringMessages?.()) || [];
206 }
207 await emit({ type: "turn_start" });
208 }
209
210 // Process prepared and queued messages before the next assistant response.
211 for (const message of declareToolChanges(currentContext, [...preparedMessages, ...pendingMessages])) {
212 await emit({ type: "message_start", message });
213 await emit({ type: "message_end", message });
214 currentContext.messages.push(message);
215 newMessages.push(message);
216 }
217 pendingMessages = [];
218
219 const requestUpdate = await config.prepareRequest?.(
220 {
221 context: currentContext,
222 model: config.model,
223 thinkingLevel: config.reasoning ?? "off",
224 },
225 signal,
226 );
227 if (requestUpdate) {
228 currentContext = requestUpdate.context ?? currentContext;
229 config = {
230 ...config,
231 model: requestUpdate.model ?? config.model,
232 reasoning:
233 requestUpdate.thinkingLevel === undefined
234 ? config.reasoning
235 : requestUpdate.thinkingLevel === "off"
236 ? undefined
237 : requestUpdate.thinkingLevel,
238 };
239 }
240
241 // Stream assistant response
242 const message = await streamAssistantResponse(currentContext, config, signal, emit, streamFunction);
243 newMessages.push(message);
244
245 if (message.stopReason === "error" || message.stopReason === "aborted") {
246 lastCompletedTurn = {
247 message,
248 toolResults: [],
249 context: currentContext,
250 newMessages,
251 };
252 await config.finishTurn?.(lastCompletedTurn, signal);
253 await emit({ type: "turn_end", message, toolResults: [] });
254 await emit({ type: "agent_end", messages: newMessages });
255 return;
256 }
257
258 // Check for tool calls
259 const toolCalls = message.content.filter((c) => c.type === "toolCall");
260
261 const toolResults: ToolResultMessage[] = [];
262 hasMoreToolCalls = false;
263 if (toolCalls.length > 0) {
264 // A "length" stop means the output was cut off by the token limit, so
265 // every tool call in the message may carry truncated arguments. Fail
266 // them all instead of executing potentially borked calls.
267 const executedToolBatch =
268 message.stopReason === "length"
269 ? await failToolCallsFromTruncatedMessage(toolCalls, emit)
270 : await executeToolCalls(currentContext, message, config, signal, emit);
271 toolResults.push(...executedToolBatch.messages);
272 hasMoreToolCalls = !executedToolBatch.terminate;
273
274 for (const result of toolResults) {
275 currentContext.messages.push(result);
276 newMessages.push(result);
277 }
278 }
279
280 lastCompletedTurn = {
281 message,
282 toolResults,
283 context: currentContext,
284 newMessages,
285 };
286 const decision = await config.finishTurn?.(lastCompletedTurn, signal);
287 await emit({ type: "turn_end", message, toolResults });
288
289 if (decision?.action === "end") {
290 await emit({ type: "agent_end", messages: newMessages });
291 return;
292 }
293
294 explicitContinuation = decision?.action === "continue";
295 pendingMessages = (await config.getSteeringMessages?.()) || [];
296 if (hasMoreToolCalls || pendingMessages.length > 0) {
297 explicitContinuation = false;
298 }
299 }
300
301 // Agent would stop here. Check for follow-up messages.
302 const followUpMessages = (await config.getFollowUpMessages?.()) || [];
303 if (followUpMessages.length > 0) {
304 // Set as pending so inner loop processes them
305 explicitContinuation = false;
306 pendingMessages = followUpMessages;
307 continue;
308 }
309
310 // No natural request was selected, so fulfill the continuation decision with one context-only turn.
311 if (explicitContinuation) {
312 explicitContinuation = false;
313 continue;
314 }
315
316 // No more messages, exit
317 break;
318 }
319
320 await emit({ type: "agent_end", messages: newMessages });
321}
322
323/**
324 * Declare tool loadout changes to the model.
325 *
326 * `context.tools` is what the runtime can execute; the transcript's system messages declare
327 * what the model may call. Before each request the difference becomes `toolsAdded` and
328 * `toolsRemoved` on a system message. When a pending system message exists, its tool fields
329 * are treated as intent and replaced with the delta between the committed transcript and
330 * the executable set, so replay always yields exactly `context.tools`. Otherwise a new
331 * system message is inserted before the first non-system pending message.
332 */
333function declareToolChanges(context: AgentContext, pendingMessages: AgentMessage[]): AgentMessage[] {
334 let systemIndex = -1;
335 for (let i = pendingMessages.length - 1; i >= 0; i--) {
336 if (pendingMessages[i].role === "system") {
337 systemIndex = i;
338 break;
339 }
340 }
341 const pending = pendingMessages[systemIndex] as SystemMessage | undefined;
342 const baseline = pending
343 ? pendingMessages.map((message, index) =>
344 index === systemIndex ? withToolChanges(pending, NO_CHANGES) : message,
345 )
346 : pendingMessages;
347 const changes = getToolStateChanges(
348 getCurrentTools([...context.messages, ...baseline]),
349 (context.tools ?? []).map(toToolDeclaration),
350 );
351 const unchanged = changes.toolsAdded.length === 0 && changes.toolsRemoved.length === 0;
352
353 if (pending) {
354 // Keep the caller's message object when it already declares no tool changes.
355 if (unchanged && !pending.toolsAdded?.length && !pending.toolsRemoved?.length) return pendingMessages;
356 return baseline.map((message, index) => (index === systemIndex ? withToolChanges(pending, changes) : message));
357 }
358 if (unchanged) return pendingMessages;
359 const update = withToolChanges({ role: "system", content: "", timestamp: Date.now() }, changes);
360 const insertIndex = pendingMessages.findIndex((message) => message.role !== "system");
361 const index = insertIndex === -1 ? pendingMessages.length : insertIndex;
362 return [...pendingMessages.slice(0, index), update, ...pendingMessages.slice(index)];
363}
364
365const NO_CHANGES: ToolStateChanges = { toolsAdded: [], toolsRemoved: [] };
366
367/** Copy a system message with its tool fields replaced by `changes`; empty lists omit the field. */
368function withToolChanges(message: SystemMessage, { toolsAdded, toolsRemoved }: ToolStateChanges): SystemMessage {
369 const { toolsAdded: _added, toolsRemoved: _removed, ...rest } = message;
370 return {
371 ...rest,
372 ...(toolsAdded.length > 0 ? { toolsAdded } : {}),
373 ...(toolsRemoved.length > 0 ? { toolsRemoved } : {}),
374 };
375}
376
377/**
378 * Stream an assistant response from the LLM.
379 * This is where AgentMessage[] gets transformed to Message[] for the LLM.
380 */
381async function streamAssistantResponse(
382 context: AgentContext,
383 config: AgentLoopConfig,
384 signal: AbortSignal | undefined,
385 emit: AgentEventSink,
386 streamFunction: StreamFn,
387): Promise<AssistantMessage> {
388 // Apply context transform if configured (AgentMessage[] → AgentMessage[])
389 let messages = context.messages;
390 if (config.transformContext) {
391 messages = await config.transformContext(messages, signal);
392 }
393
394 // Convert to LLM-compatible messages (AgentMessage[] → Message[])
395 const llmMessages = await config.convertToLlm(messages);
396
397 const llmContext = normalizeContext({ messages: llmMessages });
398
399 // Resolve API key (important for expiring tokens)
400 const resolvedApiKey =
401 (config.getApiKey ? await config.getApiKey(config.model.provider) : undefined) || config.apiKey;
402
403 const response = await streamFunction(config.model, llmContext, {
404 ...config,
405 apiKey: resolvedApiKey,
406 signal,
407 });
408 // Record the requested level, whichever stream function answered.
409 const result = async () => Object.assign(await response.result(), { thinkingLevel: config.reasoning ?? "off" });
410
411 let partialMessage: AssistantMessage | null = null;
412 let addedPartial = false;
413
414 for await (const event of response) {
415 switch (event.type) {
416 case "start":
417 partialMessage = event.partial;
418 context.messages.push(partialMessage);
419 addedPartial = true;
420 await emit({ type: "message_start", message: { ...partialMessage } });
421 break;
422
423 case "text_start":
424 case "text_delta":
425 case "text_end":
426 case "thinking_start":
427 case "thinking_delta":
428 case "thinking_end":
429 case "toolcall_start":
430 case "toolcall_delta":
431 case "toolcall_end":
432 if (partialMessage) {
433 partialMessage = event.partial;
434 context.messages[context.messages.length - 1] = partialMessage;
435 await emit({
436 type: "message_update",
437 assistantMessageEvent: event,
438 message: { ...partialMessage },
439 });
440 }
441 break;
442
443 case "done":
444 case "error": {
445 const finalMessage = await result();
446 if (addedPartial) {
447 context.messages[context.messages.length - 1] = finalMessage;
448 } else {
449 context.messages.push(finalMessage);
450 }
451 if (!addedPartial) {
452 await emit({ type: "message_start", message: { ...finalMessage } });
453 }
454 await emit({ type: "message_end", message: finalMessage });
455 return finalMessage;
456 }
457 }
458 }
459
460 const finalMessage = await result();
461 if (addedPartial) {
462 context.messages[context.messages.length - 1] = finalMessage;
463 } else {
464 context.messages.push(finalMessage);
465 await emit({ type: "message_start", message: { ...finalMessage } });
466 }
467 await emit({ type: "message_end", message: finalMessage });
468 return finalMessage;
469}
470
471/**
472 * Fail all tool calls from an assistant message that was truncated by the
473 * output token limit. Streamed tool-call arguments are finalized with a
474 * best-effort JSON salvage parser, so a truncated message can yield tool calls
475 * whose arguments parse and validate but are silently incomplete. None of them
476 * are safe to execute; report each as an error so the model can re-issue them.
477 */
478async function failToolCallsFromTruncatedMessage(
479 toolCalls: AgentToolCall[],
480 emit: AgentEventSink,
481): Promise<ExecutedToolCallBatch> {
482 const messages: ToolResultMessage[] = [];
483 for (const toolCall of toolCalls) {
484 await emit({
485 type: "tool_execution_start",
486 toolCallId: toolCall.id,
487 toolName: toolCall.name,
488 args: toolCall.arguments,
489 });
490 const finalized: FinalizedToolCallOutcome = {
491 toolCall,
492 result: createErrorToolResult(
493 `Tool call "${toolCall.name}" was not executed: the response hit the output token limit, so its arguments may be truncated. Re-issue the tool call with complete arguments.`,
494 ),
495 isError: true,
496 };
497 await emitToolExecutionEnd(finalized, emit);
498 const toolResultMessage = createToolResultMessage(finalized);
499 await emitToolResultMessage(toolResultMessage, emit);
500 messages.push(toolResultMessage);
501 }
502 return { messages, terminate: false };
503}
504
505/**
506 * Execute tool calls from an assistant message.
507 */
508async function executeToolCalls(
509 currentContext: AgentContext,
510 assistantMessage: AssistantMessage,
511 config: AgentLoopConfig,
512 signal: AbortSignal | undefined,
513 emit: AgentEventSink,
514): Promise<ExecutedToolCallBatch> {
515 const toolCalls = assistantMessage.content.filter((c) => c.type === "toolCall");
516 const hasSequentialToolCall = toolCalls.some(
517 (tc) => currentContext.tools?.find((t) => t.name === tc.name)?.executionMode === "sequential",
518 );
519 if (config.toolExecution === "sequential" || hasSequentialToolCall) {
520 return executeToolCallsSequential(currentContext, assistantMessage, toolCalls, config, signal, emit);
521 }
522 return executeToolCallsParallel(currentContext, assistantMessage, toolCalls, config, signal, emit);
523}
524
525type ExecutedToolCallBatch = {
526 messages: ToolResultMessage[];
527 terminate: boolean;
528};
529
530async function executeToolCallsSequential(
531 currentContext: AgentContext,
532 assistantMessage: AssistantMessage,
533 toolCalls: AgentToolCall[],
534 config: AgentLoopConfig,
535 signal: AbortSignal | undefined,
536 emit: AgentEventSink,
537): Promise<ExecutedToolCallBatch> {
538 const finalizedCalls: FinalizedToolCallOutcome[] = [];
539 const messages: ToolResultMessage[] = [];
540
541 for (const toolCall of toolCalls) {
542 await emit({
543 type: "tool_execution_start",
544 toolCallId: toolCall.id,
545 toolName: toolCall.name,
546 args: toolCall.arguments,
547 });
548
549 const preparation = await prepareToolCall(currentContext, assistantMessage, toolCall, config, signal);
550 let finalized: FinalizedToolCallOutcome;
551 if (preparation.kind === "immediate") {
552 finalized = {
553 toolCall,
554 result: preparation.result,
555 isError: preparation.isError,
556 };
557 } else {
558 const executed = await executePreparedToolCall(preparation, signal, emitToolExecutionUpdate(toolCall, emit));
559 finalized = await finalizeExecutedToolCall(
560 currentContext,
561 assistantMessage,
562 preparation,
563 executed,
564 config,
565 signal,
566 );
567 }
568
569 await emitToolExecutionEnd(finalized, emit);
570 const toolResultMessage = createToolResultMessage(finalized);
571 await emitToolResultMessage(toolResultMessage, emit);
572 finalizedCalls.push(finalized);
573 messages.push(toolResultMessage);
574
575 if (signal?.aborted) {
576 break;
577 }
578 }
579
580 return {
581 messages,
582 terminate: shouldTerminateToolBatch(finalizedCalls),
583 };
584}
585
586async function executeToolCallsParallel(
587 currentContext: AgentContext,
588 assistantMessage: AssistantMessage,
589 toolCalls: AgentToolCall[],
590 config: AgentLoopConfig,
591 signal: AbortSignal | undefined,
592 emit: AgentEventSink,
593): Promise<ExecutedToolCallBatch> {
594 const finalizedCalls: FinalizedToolCallEntry[] = [];
595
596 for (const toolCall of toolCalls) {
597 await emit({
598 type: "tool_execution_start",
599 toolCallId: toolCall.id,
600 toolName: toolCall.name,
601 args: toolCall.arguments,
602 });
603
604 const preparation = await prepareToolCall(currentContext, assistantMessage, toolCall, config, signal);
605 if (preparation.kind === "immediate") {
606 const finalized = {
607 toolCall,
608 result: preparation.result,
609 isError: preparation.isError,
610 } satisfies FinalizedToolCallOutcome;
611 await emitToolExecutionEnd(finalized, emit);
612 finalizedCalls.push(finalized);
613 if (signal?.aborted) {
614 break;
615 }
616 continue;
617 }
618
619 finalizedCalls.push(async () => {
620 if (signal?.aborted) {
621 const finalized = {
622 toolCall,
623 result: createErrorToolResult("Operation aborted"),
624 isError: true,
625 } satisfies FinalizedToolCallOutcome;
626 await emitToolExecutionEnd(finalized, emit);
627 return finalized;
628 }
629 const executed = await executePreparedToolCall(preparation, signal, emitToolExecutionUpdate(toolCall, emit));
630 const finalized = await finalizeExecutedToolCall(
631 currentContext,
632 assistantMessage,
633 preparation,
634 executed,
635 config,
636 signal,
637 );
638 await emitToolExecutionEnd(finalized, emit);
639 return finalized;
640 });
641 if (signal?.aborted) {
642 break;
643 }
644 }
645
646 const orderedFinalizedCalls = await Promise.all(
647 finalizedCalls.map((entry) => (typeof entry === "function" ? entry() : Promise.resolve(entry))),
648 );
649 const messages: ToolResultMessage[] = [];
650 for (const finalized of orderedFinalizedCalls) {
651 const toolResultMessage = createToolResultMessage(finalized);
652 await emitToolResultMessage(toolResultMessage, emit);
653 messages.push(toolResultMessage);
654 }
655
656 return {
657 messages,
658 terminate: shouldTerminateToolBatch(orderedFinalizedCalls),
659 };
660}
661
662type PreparedToolCall = {
663 kind: "prepared";
664 toolCall: AgentToolCall;
665 tool: AgentTool<any>;
666 args: unknown;
667};
668
669type ImmediateToolCallOutcome = {
670 kind: "immediate";
671 result: AgentToolResult<any>;
672 isError: boolean;
673};
674
675type ExecutedToolCallOutcome = {
676 result: AgentToolResult<any>;
677 isError: boolean;
678};
679
680type FinalizedToolCallOutcome = AgentToolCallOutcome;
681
682/** The `beforeToolCall` and `afterToolCall` hooks of {@link AgentLoopConfig}. */
683export type ToolCallHooks = Pick<AgentLoopConfig, "beforeToolCall" | "afterToolCall">;
684
685type ToolUpdateSink = (partialResult: AgentToolResult<any>) => Promise<void> | void;
686
687type FinalizedToolCallEntry = FinalizedToolCallOutcome | (() => Promise<FinalizedToolCallOutcome>);
688
689function shouldTerminateToolBatch(finalizedCalls: FinalizedToolCallOutcome[]): boolean {
690 return finalizedCalls.length > 0 && finalizedCalls.every((finalized) => finalized.result.terminate === true);
691}
692
693function prepareToolCallArguments(tool: AgentTool<any>, toolCall: AgentToolCall): AgentToolCall {
694 if (!tool.prepareArguments) {
695 return toolCall;
696 }
697 const preparedArguments = tool.prepareArguments(toolCall.arguments);
698 if (preparedArguments === toolCall.arguments) {
699 return toolCall;
700 }
701 return {
702 ...toolCall,
703 arguments: preparedArguments as Record<string, any>,
704 };
705}
706
707async function prepareToolCall(
708 currentContext: AgentContext,
709 assistantMessage: AssistantMessage,
710 toolCall: AgentToolCall,
711 config: ToolCallHooks,
712 signal: AbortSignal | undefined,
713 tools: readonly AgentTool<any>[] = currentContext.tools ?? [],
714): Promise<PreparedToolCall | ImmediateToolCallOutcome> {
715 const tool = tools.find((t) => t.name === toolCall.name);
716 if (!tool) {
717 return {
718 kind: "immediate",
719 result: createErrorToolResult(`Tool ${toolCall.name} not found`),
720 isError: true,
721 };
722 }
723
724 try {
725 const preparedToolCall = prepareToolCallArguments(tool, toolCall);
726 const validatedArgs = validateToolArguments(tool, preparedToolCall);
727 if (config.beforeToolCall) {
728 const beforeResult = await config.beforeToolCall(
729 {
730 assistantMessage,
731 toolCall,
732 args: validatedArgs,
733 context: currentContext,
734 },
735 signal,
736 );
737 if (signal?.aborted) {
738 return {
739 kind: "immediate",
740 result: createErrorToolResult("Operation aborted"),
741 isError: true,
742 };
743 }
744 if (beforeResult?.block) {
745 const result = createErrorToolResult(beforeResult.reason || "Tool execution was blocked");
746 if (beforeResult.terminate === true) {
747 result.terminate = true;
748 }
749 return {
750 kind: "immediate",
751 result,
752 isError: true,
753 };
754 }
755 }
756 if (signal?.aborted) {
757 return {
758 kind: "immediate",
759 result: createErrorToolResult("Operation aborted"),
760 isError: true,
761 };
762 }
763 return {
764 kind: "prepared",
765 toolCall,
766 tool,
767 args: validatedArgs,
768 };
769 } catch (error) {
770 return {
771 kind: "immediate",
772 result: createErrorToolResult(error instanceof Error ? error.message : String(error)),
773 isError: true,
774 };
775 }
776}
777
778function emitToolExecutionUpdate(toolCall: AgentToolCall, emit: AgentEventSink): ToolUpdateSink {
779 return (partialResult) =>
780 emit({
781 type: "tool_execution_update",
782 toolCallId: toolCall.id,
783 toolName: toolCall.name,
784 args: toolCall.arguments,
785 partialResult,
786 });
787}
788
789/** Options for {@link runToolCall}. */
790export interface RunToolCallOptions extends ToolCallHooks {
791 /** Tools the call resolves against. */
792 tools: readonly AgentTool<any>[];
793 /** Passed to the hooks as the message that issued the call. */
794 assistantMessage: AssistantMessage;
795 /** Passed to the hooks as the current agent context. */
796 context: AgentContext;
797 signal?: AbortSignal;
798 onUpdate?: ToolUpdateSink;
799}
800
801/**
802 * Run one tool call through the same steps as a model-issued call: argument preparation, schema
803 * validation, `beforeToolCall`, execution, and `afterToolCall`. Emits no events and adds no
804 * messages. Tools that call other tools use this so the hooks (for example permission checks)
805 * apply to those calls too.
806 *
807 * Never rejects for tool failures: unknown tools, validation errors, blocked calls, and thrown
808 * errors come back as `isError: true`.
809 */
810export async function runToolCall(toolCall: AgentToolCall, options: RunToolCallOptions): Promise<AgentToolCallOutcome> {
811 const { assistantMessage, context, signal } = options;
812 const preparation = await prepareToolCall(context, assistantMessage, toolCall, options, signal, options.tools);
813 if (preparation.kind === "immediate") {
814 return { toolCall, result: preparation.result, isError: preparation.isError };
815 }
816 const executed = await executePreparedToolCall(preparation, signal, options.onUpdate ?? (() => {}));
817 return finalizeExecutedToolCall(context, assistantMessage, preparation, executed, options, signal);
818}
819
820async function executePreparedToolCall(
821 prepared: PreparedToolCall,
822 signal: AbortSignal | undefined,
823 onUpdate: ToolUpdateSink,
824): Promise<ExecutedToolCallOutcome> {
825 const updateEvents: Promise<void>[] = [];
826 let acceptingUpdates = true;
827
828 try {
829 const result = await prepared.tool.execute(
830 prepared.toolCall.id,
831 prepared.args as never,
832 signal,
833 (partialResult) => {
834 if (!acceptingUpdates) return;
835 updateEvents.push(Promise.resolve(onUpdate(partialResult)));
836 },
837 );
838 acceptingUpdates = false;
839 await Promise.all(updateEvents);
840 return { result, isError: result.isError === true };
841 } catch (error) {
842 acceptingUpdates = false;
843 await Promise.all(updateEvents);
844 return {
845 result: createErrorToolResult(error instanceof Error ? error.message : String(error)),
846 isError: true,
847 };
848 } finally {
849 acceptingUpdates = false;
850 }
851}
852
853async function finalizeExecutedToolCall(
854 currentContext: AgentContext,
855 assistantMessage: AssistantMessage,
856 prepared: PreparedToolCall,
857 executed: ExecutedToolCallOutcome,
858 config: ToolCallHooks,
859 signal: AbortSignal | undefined,
860): Promise<FinalizedToolCallOutcome> {
861 let result = executed.result;
862 let isError = executed.isError;
863
864 if (config.afterToolCall) {
865 try {
866 const afterResult = await config.afterToolCall(
867 {
868 assistantMessage,
869 toolCall: prepared.toolCall,
870 args: prepared.args,
871 result,
872 isError,
873 context: currentContext,
874 },
875 signal,
876 );
877 if (afterResult) {
878 // Structured content not replaced along with the content may no longer match it.
879 const structuredContent =
880 afterResult.structuredContent ?? (afterResult.content ? undefined : result.structuredContent);
881 result = {
882 ...result,
883 content: afterResult.content ?? result.content,
884 details: afterResult.details ?? result.details,
885 usage: afterResult.usage ?? result.usage,
886 terminate: afterResult.terminate ?? result.terminate,
887 };
888 if (structuredContent === undefined) delete result.structuredContent;
889 else result.structuredContent = structuredContent;
890 isError = afterResult.isError ?? isError;
891 }
892 } catch (error) {
893 result = createErrorToolResult(error instanceof Error ? error.message : String(error));
894 isError = true;
895 }
896 }
897
898 return {
899 toolCall: prepared.toolCall,
900 result,
901 isError,
902 };
903}
904
905function createErrorToolResult(message: string): AgentToolResult<any> {
906 return {
907 content: [{ type: "text", text: message }],
908 details: {},
909 };
910}
911
912async function emitToolExecutionEnd(finalized: FinalizedToolCallOutcome, emit: AgentEventSink): Promise<void> {
913 await emit({
914 type: "tool_execution_end",
915 toolCallId: finalized.toolCall.id,
916 toolName: finalized.toolCall.name,
917 result: finalized.result,
918 isError: finalized.isError,
919 });
920}
921
922function createToolResultMessage(finalized: FinalizedToolCallOutcome): ToolResultMessage {
923 return {
924 role: "toolResult",
925 toolCallId: finalized.toolCall.id,
926 toolName: finalized.toolCall.name,
927 // Untyped tools (JS extensions) can return results without content; normalize
928 // so the null never enters session history or provider payloads.
929 content: finalized.result.content ?? [],
930 details: finalized.result.details,
931 usage: finalized.result.usage,
932 isError: finalized.isError,
933 timestamp: Date.now(),
934 };
935}
936
937async function emitToolResultMessage(toolResultMessage: ToolResultMessage, emit: AgentEventSink): Promise<void> {
938 await emit({ type: "message_start", message: toolResultMessage });
939 await emit({ type: "message_end", message: toolResultMessage });
940}