1
/**2
* Agent loop that works with AgentMessage throughout.3
* Transforms to Message[] only at the LLM call boundary.4
*/6
import {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";18
import { getDefaultStreamFn } from "./stream-fn.ts";19
import 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";32
export type AgentEventSink = (event: AgentEvent) => Promise<void> | void;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
*/38
export 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();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
});60
return stream;61
}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` message68
* 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
*/71
export 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
}81
if (context.messages[context.messages.length - 1].role === "assistant") {82
throw new Error("Cannot continue from message role: assistant");83
}85
const stream = createAgentStream();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
});99
return stream;100
}102
export 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
};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
}124
await runLoop(currentContext, newMessages, config, signal, emit, streamFn ?? getDefaultStreamFn());125
return newMessages;126
}128
export 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
}139
if (context.messages[context.messages.length - 1].role === "assistant") {140
throw new Error("Cannot continue from message role: assistant");141
}143
const newMessages: AgentMessage[] = [];144
const currentContext: AgentContext = { ...context };146
await emit({ type: "agent_start" });147
await emit({ type: "turn_start" });149
await runLoop(currentContext, newMessages, config, signal, emit, streamFn ?? getDefaultStreamFn());150
return newMessages;151
}153
function 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
}160
/**161
* Main loop logic shared by agentLoop and agentLoopContinue.162
*/163
async 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?.()) || [];178
// Outer loop: continues when queued follow-up messages arrive after agent would stop179
while (true) {180
let hasMoreToolCalls = true;182
// Inner loop: process tool calls and steering messages183
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 === undefined195
? config.reasoning196
: nextTurnSnapshot.thinkingLevel === "off"197
? undefined198
: nextTurnSnapshot.thinkingLevel,199
};200
}201
// Preparation can be long-running (for example, compaction). Pick up steering202
// 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
}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 = [];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 === undefined234
? config.reasoning235
: requestUpdate.thinkingLevel === "off"236
? undefined237
: requestUpdate.thinkingLevel,238
};239
}241
// Stream assistant response242
const message = await streamAssistantResponse(currentContext, config, signal, emit, streamFunction);243
newMessages.push(message);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
}258
// Check for tool calls259
const toolCalls = message.content.filter((c) => c.type === "toolCall");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, so265
// every tool call in the message may carry truncated arguments. Fail266
// 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;274
for (const result of toolResults) {275
currentContext.messages.push(result);276
newMessages.push(result);277
}278
}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 });289
if (decision?.action === "end") {290
await emit({ type: "agent_end", messages: newMessages });291
return;292
}294
explicitContinuation = decision?.action === "continue";295
pendingMessages = (await config.getSteeringMessages?.()) || [];296
if (hasMoreToolCalls || pendingMessages.length > 0) {297
explicitContinuation = false;298
}299
}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 them305
explicitContinuation = false;306
pendingMessages = followUpMessages;307
continue;308
}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
}316
// No more messages, exit317
break;318
}320
await emit({ type: "agent_end", messages: newMessages });321
}323
/**324
* Declare tool loadout changes to the model.325
*326
* `context.tools` is what the runtime can execute; the transcript's system messages declare327
* what the model may call. Before each request the difference becomes `toolsAdded` and328
* `toolsRemoved` on a system message. When a pending system message exists, its tool fields329
* are treated as intent and replaced with the delta between the committed transcript and330
* the executable set, so replay always yields exactly `context.tools`. Otherwise a new331
* system message is inserted before the first non-system pending message.332
*/333
function 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 = pending343
? 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;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
}365
const NO_CHANGES: ToolStateChanges = { toolsAdded: [], toolsRemoved: [] };367
/** Copy a system message with its tool fields replaced by `changes`; empty lists omit the field. */368
function 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
}377
/**378
* Stream an assistant response from the LLM.379
* This is where AgentMessage[] gets transformed to Message[] for the LLM.380
*/381
async 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
}394
// Convert to LLM-compatible messages (AgentMessage[] → Message[])395
const llmMessages = await config.convertToLlm(messages);397
const llmContext = normalizeContext({ messages: llmMessages });399
// Resolve API key (important for expiring tokens)400
const resolvedApiKey =401
(config.getApiKey ? await config.getApiKey(config.model.provider) : undefined) || config.apiKey;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" });411
let partialMessage: AssistantMessage | null = null;412
let addedPartial = false;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;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;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
}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
}471
/**472
* Fail all tool calls from an assistant message that was truncated by the473
* output token limit. Streamed tool-call arguments are finalized with a474
* best-effort JSON salvage parser, so a truncated message can yield tool calls475
* whose arguments parse and validate but are silently incomplete. None of them476
* are safe to execute; report each as an error so the model can re-issue them.477
*/478
async 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
}505
/**506
* Execute tool calls from an assistant message.507
*/508
async 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
}525
type ExecutedToolCallBatch = {526
messages: ToolResultMessage[];527
terminate: boolean;528
};530
async 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[] = [];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
});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
}569
await emitToolExecutionEnd(finalized, emit);570
const toolResultMessage = createToolResultMessage(finalized);571
await emitToolResultMessage(toolResultMessage, emit);572
finalizedCalls.push(finalized);573
messages.push(toolResultMessage);575
if (signal?.aborted) {576
break;577
}578
}580
return {581
messages,582
terminate: shouldTerminateToolBatch(finalizedCalls),583
};584
}586
async 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[] = [];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
});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
}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
}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
}656
return {657
messages,658
terminate: shouldTerminateToolBatch(orderedFinalizedCalls),659
};660
}662
type PreparedToolCall = {663
kind: "prepared";664
toolCall: AgentToolCall;665
tool: AgentTool<any>;666
args: unknown;667
};669
type ImmediateToolCallOutcome = {670
kind: "immediate";671
result: AgentToolResult<any>;672
isError: boolean;673
};675
type ExecutedToolCallOutcome = {676
result: AgentToolResult<any>;677
isError: boolean;678
};680
type FinalizedToolCallOutcome = AgentToolCallOutcome;682
/** The `beforeToolCall` and `afterToolCall` hooks of {@link AgentLoopConfig}. */683
export type ToolCallHooks = Pick<AgentLoopConfig, "beforeToolCall" | "afterToolCall">;685
type ToolUpdateSink = (partialResult: AgentToolResult<any>) => Promise<void> | void;687
type FinalizedToolCallEntry = FinalizedToolCallOutcome | (() => Promise<FinalizedToolCallOutcome>);689
function shouldTerminateToolBatch(finalizedCalls: FinalizedToolCallOutcome[]): boolean {690
return finalizedCalls.length > 0 && finalizedCalls.every((finalized) => finalized.result.terminate === true);691
}693
function 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
}707
async 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
}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
}778
function 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
}789
/** Options for {@link runToolCall}. */790
export 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
}801
/**802
* Run one tool call through the same steps as a model-issued call: argument preparation, schema803
* validation, `beforeToolCall`, execution, and `afterToolCall`. Emits no events and adds no804
* 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 thrown808
* errors come back as `isError: true`.809
*/810
export 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
}820
async function executePreparedToolCall(821
prepared: PreparedToolCall,822
signal: AbortSignal | undefined,823
onUpdate: ToolUpdateSink,824
): Promise<ExecutedToolCallOutcome> {825
const updateEvents: Promise<void>[] = [];826
let acceptingUpdates = true;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
}853
async 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;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
}898
return {899
toolCall: prepared.toolCall,900
result,901
isError,902
};903
}905
function createErrorToolResult(message: string): AgentToolResult<any> {906
return {907
content: [{ type: "text", text: message }],908
details: {},909
};910
}912
async 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
}922
function 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; normalize928
// 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
}937
async 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
}