{"version":3,"file":"agent-loop.d.ts","sourceRoot":"","sources":["../src/agent-loop.ts"],"names":[],"mappings":"AAAA;;;GAGG;AAEH,OAAO,EAGN,WAAW,EAMX,MAAM,UAAU,CAAC;AAClB,OAAO,KAAK,EACX,YAAY,EACZ,UAAU,EACV,eAAe,EACf,YAAY,EAIZ,QAAQ,EAER,MAAM,YAAY,CAAC;AAEpB,MAAM,MAAM,cAAc,GAAG,CAAC,KAAK,EAAE,UAAU,KAAK,OAAO,CAAC,IAAI,CAAC,GAAG,IAAI,CAAC;AASzE;;;GAGG;AACH,wBAAgB,SAAS,CACxB,OAAO,EAAE,YAAY,EAAE,EACvB,OAAO,EAAE,YAAY,EACrB,MAAM,EAAE,eAAe,EACvB,MAAM,CAAC,EAAE,WAAW,EACpB,QAAQ,CAAC,EAAE,QAAQ,GACjB,WAAW,CAAC,UAAU,EAAE,YAAY,EAAE,CAAC,CAiBzC;AAED;;;;;;;GAOG;AACH,wBAAgB,iBAAiB,CAChC,OAAO,EAAE,YAAY,EACrB,MAAM,EAAE,eAAe,EACvB,MAAM,CAAC,EAAE,WAAW,EACpB,QAAQ,CAAC,EAAE,QAAQ,GACjB,WAAW,CAAC,UAAU,EAAE,YAAY,EAAE,CAAC,CAwBzC;AAED,wBAAsB,YAAY,CACjC,OAAO,EAAE,YAAY,EAAE,EACvB,OAAO,EAAE,YAAY,EACrB,MAAM,EAAE,eAAe,EACvB,IAAI,EAAE,cAAc,EACpB,MAAM,CAAC,EAAE,WAAW,EACpB,QAAQ,CAAC,EAAE,QAAQ,GACjB,OAAO,CAAC,YAAY,EAAE,CAAC,CAoBzB;AAED,wBAAsB,oBAAoB,CACzC,OAAO,EAAE,YAAY,EACrB,MAAM,EAAE,eAAe,EACvB,IAAI,EAAE,cAAc,EACpB,MAAM,CAAC,EAAE,WAAW,EACpB,QAAQ,CAAC,EAAE,QAAQ,GACjB,OAAO,CAAC,YAAY,EAAE,CAAC,CAqBzB","sourcesContent":["/**\n * Agent loop that works with AgentMessage throughout.\n * Transforms to Message[] only at the LLM call boundary.\n */\n\nimport {\n\ttype AssistantMessage,\n\ttype Context,\n\tEventStream,\n\tstreamSimple,\n\tsupportsMax,\n\tsupportsXhigh,\n\ttype ToolResultMessage,\n\tvalidateToolArguments,\n} from \"@dreb/ai\";\nimport type {\n\tAgentContext,\n\tAgentEvent,\n\tAgentLoopConfig,\n\tAgentMessage,\n\tAgentTool,\n\tAgentToolCall,\n\tAgentToolResult,\n\tStreamFn,\n\tThinkingLevel,\n} from \"./types.js\";\n\nexport type AgentEventSink = (event: AgentEvent) => Promise<void> | void;\n\nfunction getEffectiveThinkingLevel(config: AgentLoopConfig): ThinkingLevel {\n\tconst requested = config.reasoning ?? \"off\";\n\tif (!config.model.reasoning) return \"off\";\n\tif (requested === \"max\" && !supportsMax(config.model)) return supportsXhigh(config.model) ? \"xhigh\" : \"high\";\n\treturn requested === \"xhigh\" && !supportsXhigh(config.model) ? \"high\" : requested;\n}\n\n/**\n * Start an agent loop with a new prompt message.\n * The prompt is added to the context and events are emitted for it.\n */\nexport function agentLoop(\n\tprompts: AgentMessage[],\n\tcontext: AgentContext,\n\tconfig: AgentLoopConfig,\n\tsignal?: AbortSignal,\n\tstreamFn?: StreamFn,\n): EventStream<AgentEvent, AgentMessage[]> {\n\tconst stream = createAgentStream();\n\n\tvoid runAgentLoop(\n\t\tprompts,\n\t\tcontext,\n\t\tconfig,\n\t\tasync (event) => {\n\t\t\tstream.push(event);\n\t\t},\n\t\tsignal,\n\t\tstreamFn,\n\t).then((messages) => {\n\t\tstream.end(messages);\n\t});\n\n\treturn stream;\n}\n\n/**\n * Continue an agent loop from the current context without adding a new message.\n * Used for retries - context already has user message or tool results.\n *\n * **Important:** The last message in context must convert to a `user` or `toolResult` message\n * via `convertToLlm`. If it doesn't, the LLM provider will reject the request.\n * This cannot be validated here since `convertToLlm` is only called once per turn.\n */\nexport function agentLoopContinue(\n\tcontext: AgentContext,\n\tconfig: AgentLoopConfig,\n\tsignal?: AbortSignal,\n\tstreamFn?: StreamFn,\n): EventStream<AgentEvent, AgentMessage[]> {\n\tif (context.messages.length === 0) {\n\t\tthrow new Error(\"Cannot continue: no messages in context\");\n\t}\n\n\tif (context.messages[context.messages.length - 1].role === \"assistant\") {\n\t\tthrow new Error(\"Cannot continue from message role: assistant\");\n\t}\n\n\tconst stream = createAgentStream();\n\n\tvoid runAgentLoopContinue(\n\t\tcontext,\n\t\tconfig,\n\t\tasync (event) => {\n\t\t\tstream.push(event);\n\t\t},\n\t\tsignal,\n\t\tstreamFn,\n\t).then((messages) => {\n\t\tstream.end(messages);\n\t});\n\n\treturn stream;\n}\n\nexport async function runAgentLoop(\n\tprompts: AgentMessage[],\n\tcontext: AgentContext,\n\tconfig: AgentLoopConfig,\n\temit: AgentEventSink,\n\tsignal?: AbortSignal,\n\tstreamFn?: StreamFn,\n): Promise<AgentMessage[]> {\n\tconst newMessages: AgentMessage[] = [...prompts];\n\tconst currentContext: AgentContext = {\n\t\t...context,\n\t\tmessages: [...context.messages, ...prompts],\n\t};\n\n\tawait emit({\n\t\ttype: \"agent_start\",\n\t\tmodel: { provider: config.model.provider, id: config.model.id },\n\t\tthinkingLevel: getEffectiveThinkingLevel(config),\n\t});\n\tawait emit({ type: \"turn_start\" });\n\tfor (const prompt of prompts) {\n\t\tawait emit({ type: \"message_start\", message: prompt });\n\t\tawait emit({ type: \"message_end\", message: prompt });\n\t}\n\n\tawait runLoop(currentContext, newMessages, config, signal, emit, streamFn);\n\treturn newMessages;\n}\n\nexport async function runAgentLoopContinue(\n\tcontext: AgentContext,\n\tconfig: AgentLoopConfig,\n\temit: AgentEventSink,\n\tsignal?: AbortSignal,\n\tstreamFn?: StreamFn,\n): Promise<AgentMessage[]> {\n\tif (context.messages.length === 0) {\n\t\tthrow new Error(\"Cannot continue: no messages in context\");\n\t}\n\n\tif (context.messages[context.messages.length - 1].role === \"assistant\") {\n\t\tthrow new Error(\"Cannot continue from message role: assistant\");\n\t}\n\n\tconst newMessages: AgentMessage[] = [];\n\tconst currentContext: AgentContext = { ...context };\n\n\tawait emit({\n\t\ttype: \"agent_start\",\n\t\tmodel: { provider: config.model.provider, id: config.model.id },\n\t\tthinkingLevel: getEffectiveThinkingLevel(config),\n\t});\n\tawait emit({ type: \"turn_start\" });\n\n\tawait runLoop(currentContext, newMessages, config, signal, emit, streamFn);\n\treturn newMessages;\n}\n\nfunction createAgentStream(): EventStream<AgentEvent, AgentMessage[]> {\n\treturn new EventStream<AgentEvent, AgentMessage[]>(\n\t\t(event: AgentEvent) => event.type === \"agent_end\",\n\t\t(event: AgentEvent) => (event.type === \"agent_end\" ? event.messages : []),\n\t);\n}\n\n/**\n * Main loop logic shared by agentLoop and agentLoopContinue.\n */\nasync function runLoop(\n\tcurrentContext: AgentContext,\n\tnewMessages: AgentMessage[],\n\tconfig: AgentLoopConfig,\n\tsignal: AbortSignal | undefined,\n\temit: AgentEventSink,\n\tstreamFn?: StreamFn,\n): Promise<void> {\n\tlet firstTurn = true;\n\tlet endTurnRequested = false;\n\tlet llmCallCount = 0;\n\t// Check for steering messages at start (user may have typed while waiting)\n\tlet pendingMessages: AgentMessage[] = (await config.getSteeringMessages?.()) || [];\n\n\t// Outer loop: continues when queued follow-up messages arrive after agent would stop\n\twhile (true) {\n\t\tlet hasMoreToolCalls = true;\n\t\tendTurnRequested = false;\n\n\t\t// Inner loop: process tool calls and steering messages\n\t\twhile (hasMoreToolCalls || pendingMessages.length > 0) {\n\t\t\t// Process pending messages first (inject into context so they're never lost).\n\t\t\t// This happens before turn_start and shouldContinue so messages are always\n\t\t\t// preserved regardless of whether the turn proceeds.\n\t\t\tif (pendingMessages.length > 0) {\n\t\t\t\tfor (const message of pendingMessages) {\n\t\t\t\t\tawait emit({ type: \"message_start\", message });\n\t\t\t\t\tawait emit({ type: \"message_end\", message });\n\t\t\t\t\tcurrentContext.messages.push(message);\n\t\t\t\t\tnewMessages.push(message);\n\t\t\t\t}\n\t\t\t\tpendingMessages = [];\n\t\t\t}\n\n\t\t\t// Check shouldContinue before starting a new turn (not the very first).\n\t\t\t// This runs before turn_start so we never emit an orphaned turn_start.\n\t\t\tif (llmCallCount > 0 && config.shouldContinue && !config.shouldContinue()) {\n\t\t\t\tbreak;\n\t\t\t}\n\n\t\t\tif (config.beforeLlmCall) {\n\t\t\t\tconst prepared = await config.beforeLlmCall(\n\t\t\t\t\t{\n\t\t\t\t\t\t...currentContext,\n\t\t\t\t\t\tmessages: currentContext.messages.slice(),\n\t\t\t\t\t\ttools: currentContext.tools?.slice(),\n\t\t\t\t\t},\n\t\t\t\t\tsignal,\n\t\t\t\t);\n\t\t\t\tif (prepared?.messages) {\n\t\t\t\t\tcurrentContext.messages = prepared.messages.slice();\n\t\t\t\t}\n\t\t\t\tif (prepared?.model) {\n\t\t\t\t\tconfig.model = prepared.model;\n\t\t\t\t}\n\t\t\t}\n\n\t\t\tif (!firstTurn) {\n\t\t\t\tawait emit({ type: \"turn_start\" });\n\t\t\t} else {\n\t\t\t\tfirstTurn = false;\n\t\t\t}\n\n\t\t\t// Stream assistant response\n\t\t\tconst message = await streamAssistantResponse(currentContext, config, signal, emit, streamFn);\n\t\t\tnewMessages.push(message);\n\t\t\tllmCallCount++;\n\n\t\t\tif (message.stopReason === \"error\" || message.stopReason === \"aborted\") {\n\t\t\t\tawait emit({ type: \"turn_end\", message, toolResults: [] });\n\t\t\t\tawait emit({ type: \"agent_end\", messages: newMessages });\n\t\t\t\treturn;\n\t\t\t}\n\n\t\t\t// Check for tool calls\n\t\t\tconst toolCalls = message.content.filter((c) => c.type === \"toolCall\");\n\t\t\thasMoreToolCalls = toolCalls.length > 0;\n\n\t\t\tconst toolResults: ToolResultMessage[] = [];\n\t\t\tif (hasMoreToolCalls) {\n\t\t\t\tconst executionResult = await executeToolCalls(currentContext, message, config, signal, emit);\n\t\t\t\ttoolResults.push(...executionResult.results);\n\n\t\t\t\tfor (const result of toolResults) {\n\t\t\t\t\tcurrentContext.messages.push(result);\n\t\t\t\t\tnewMessages.push(result);\n\t\t\t\t}\n\n\t\t\t\tif (executionResult.endTurn) {\n\t\t\t\t\tendTurnRequested = true;\n\t\t\t\t}\n\t\t\t}\n\n\t\t\tawait emit({ type: \"turn_end\", message, toolResults });\n\n\t\t\t// If a tool requested endTurn, stop the loop (skip further LLM calls)\n\t\t\tif (endTurnRequested) {\n\t\t\t\thasMoreToolCalls = false;\n\t\t\t\tpendingMessages = [];\n\t\t\t\tbreak;\n\t\t\t}\n\n\t\t\tpendingMessages = (await config.getSteeringMessages?.()) || [];\n\t\t}\n\n\t\t// Agent would stop here. Drain any steering messages that arrived during\n\t\t// wind-down (e.g., background agent results delivered via steer() between\n\t\t// endTurn break and here). Then check for follow-up messages.\n\t\tconst leftoverSteering = (await config.getSteeringMessages?.()) || [];\n\t\tconst followUpMessages = (await config.getFollowUpMessages?.()) || [];\n\t\tconst allPending = [...leftoverSteering, ...followUpMessages];\n\t\tif (allPending.length > 0) {\n\t\t\t// Set as pending so inner loop processes them\n\t\t\tpendingMessages = allPending;\n\t\t\tcontinue;\n\t\t}\n\n\t\t// No more messages, exit\n\t\tbreak;\n\t}\n\n\tawait emit({ type: \"agent_end\", messages: newMessages });\n}\n\n/**\n * Check if an error is a stream-drop error (connection dropped before terminal event).\n * Only these errors trigger automatic retry — all other errors propagate as-is.\n */\nfunction isStreamDropError(error: unknown): boolean {\n\tconst msg = error instanceof Error ? error.message : typeof error === \"string\" ? error : \"\";\n\treturn (\n\t\tmsg.includes(\"Stream ended without\") ||\n\t\tmsg.includes(\"connection likely dropped\") ||\n\t\tmsg.includes(\"WebSocket stream closed before\")\n\t);\n}\n\n/**\n * Sleep for a duration, aborting early if the signal fires.\n */\nfunction sleep(ms: number, signal?: AbortSignal): Promise<void> {\n\treturn new Promise((resolve) => {\n\t\tif (signal?.aborted) {\n\t\t\tresolve();\n\t\t\treturn;\n\t\t}\n\t\tconst timer = setTimeout(resolve, ms);\n\t\tsignal?.addEventListener(\n\t\t\t\"abort\",\n\t\t\t() => {\n\t\t\t\tclearTimeout(timer);\n\t\t\t\tresolve();\n\t\t\t},\n\t\t\t{ once: true },\n\t\t);\n\t});\n}\n\n/**\n * Stream an assistant response from the LLM.\n * This is where AgentMessage[] gets transformed to Message[] for the LLM.\n */\nasync function streamAssistantResponse(\n\tcontext: AgentContext,\n\tconfig: AgentLoopConfig,\n\tsignal: AbortSignal | undefined,\n\temit: AgentEventSink,\n\tstreamFn?: StreamFn,\n): Promise<AssistantMessage> {\n\t// Apply context transform if configured (AgentMessage[] → AgentMessage[])\n\tlet messages = context.messages;\n\tif (config.transformContext) {\n\t\tmessages = await config.transformContext(messages, signal);\n\t}\n\n\t// Convert to LLM-compatible messages (AgentMessage[] → Message[])\n\tconst llmMessages = await config.convertToLlm(messages);\n\n\t// Build LLM context\n\tconst llmContext: Context = {\n\t\tsystemPrompt: context.systemPrompt,\n\t\tmessages: llmMessages,\n\t\ttools: context.tools,\n\t};\n\n\tconst streamFunction = streamFn || streamSimple;\n\n\t// Resolve API key (important for expiring tokens)\n\tconst resolvedApiKey =\n\t\t(config.getApiKey ? await config.getApiKey(config.model.provider) : undefined) || config.apiKey;\n\n\tconst streamStart = performance.now();\n\tlet partialMessage: AssistantMessage | null = null;\n\tlet addedPartial = false;\n\n\tconst finalizeMessage = async (finalMessage: AssistantMessage): Promise<AssistantMessage> => {\n\t\tfinalMessage.durationMs = Math.max(0, performance.now() - streamStart);\n\t\tif (addedPartial) {\n\t\t\tcontext.messages[context.messages.length - 1] = finalMessage;\n\t\t} else {\n\t\t\tcontext.messages.push(finalMessage);\n\t\t}\n\t\tif (!addedPartial) {\n\t\t\tawait emit({ type: \"message_start\", message: { ...finalMessage } });\n\t\t}\n\t\tawait emit({ type: \"message_end\", message: finalMessage });\n\t\treturn finalMessage;\n\t};\n\n\tconst createErrorMessage = (error: unknown): AssistantMessage => ({\n\t\trole: \"assistant\",\n\t\tcontent: partialMessage?.content ?? [{ type: \"text\", text: \"\" }],\n\t\tapi: config.model.api,\n\t\tprovider: config.model.provider,\n\t\tmodel: config.model.id,\n\t\tusage: partialMessage?.usage ?? {\n\t\t\tinput: 0,\n\t\t\toutput: 0,\n\t\t\tcacheRead: 0,\n\t\t\tcacheWrite: 0,\n\t\t\ttotalTokens: 0,\n\t\t\tcost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },\n\t\t},\n\t\tstopReason: signal?.aborted ? \"aborted\" : \"error\",\n\t\terrorMessage: error instanceof Error ? error.message : String(error),\n\t\ttimestamp: partialMessage?.timestamp ?? Date.now(),\n\t});\n\n\tconst maxRetries = config.streamRetries ?? 3;\n\tconst retryBaseDelay = config.streamRetryBaseDelayMs ?? 1000;\n\tconst lengthRetries = config.lengthRetries ?? 2;\n\t// Every attempt uses one fixed budget: an explicit call-site limit when\n\t// provided, otherwise the active model's configured output maximum.\n\tconst requestMaxTokens = config.maxTokens ?? config.model.maxTokens;\n\tlet lengthAttempts = 0;\n\tconst clonePartialForDebug = (): AssistantMessage | undefined => {\n\t\tif (!partialMessage) return undefined;\n\t\treturn {\n\t\t\t...partialMessage,\n\t\t\tcontent: partialMessage.content.map((block) => ({ ...block })),\n\t\t\tusage: { ...partialMessage.usage, cost: { ...partialMessage.usage.cost } },\n\t\t};\n\t};\n\n\t// Discard the truncated partial message (if any) before re-issuing a request.\n\tconst discardPartial = () => {\n\t\tif (addedPartial) {\n\t\t\tcontext.messages.pop();\n\t\t\taddedPartial = false;\n\t\t}\n\t\tpartialMessage = null;\n\t};\n\n\t// Build the truncation failure message after length retries are exhausted.\n\tconst createLengthExhaustedMessage = (attempts: number, detail?: string): AssistantMessage => {\n\t\tconst providerDetail = detail?.trim() || partialMessage?.errorMessage?.trim();\n\t\tconst message = `Response truncated at the configured output token limit after ${attempts} attempt${attempts === 1 ? \"\" : \"s\"}`;\n\t\tconst errorMessage = createErrorMessage(\n\t\t\tnew Error(providerDetail ? `${message}\\nProvider detail: ${providerDetail}` : message),\n\t\t);\n\t\t// Force \"error\" so the runLoop error guard terminates the turn loudly.\n\t\terrorMessage.stopReason = signal?.aborted ? \"aborted\" : \"error\";\n\t\treturn errorMessage;\n\t};\n\n\t// stream-drop retries are counted separately from length retries. A length retry\n\t// is a fresh request, so it resets the stream-drop attempt counter.\n\tlet streamDropAttempt = 0;\n\twhile (true) {\n\t\tlet response: Awaited<ReturnType<StreamFn>>;\n\t\ttry {\n\t\t\tresponse = await streamFunction(config.model, llmContext, {\n\t\t\t\t...config,\n\t\t\t\tmaxTokens: requestMaxTokens,\n\t\t\t\tapiKey: resolvedApiKey,\n\t\t\t\tsignal,\n\t\t\t});\n\t\t} catch (error) {\n\t\t\tif (isStreamDropError(error) && streamDropAttempt < maxRetries && !signal?.aborted) {\n\t\t\t\tawait emit({\n\t\t\t\t\ttype: \"stream_retry\",\n\t\t\t\t\tattempt: streamDropAttempt + 1,\n\t\t\t\t\tmaxAttempts: maxRetries,\n\t\t\t\t\terror: error instanceof Error ? error.message : String(error),\n\t\t\t\t\tdiscardedPartial: clonePartialForDebug(),\n\t\t\t\t});\n\t\t\t\tawait sleep(retryBaseDelay * 2 ** streamDropAttempt, signal);\n\t\t\t\tdiscardPartial();\n\t\t\t\tstreamDropAttempt++;\n\t\t\t\tcontinue;\n\t\t\t}\n\t\t\tconst surfacedError = isStreamDropError(error)\n\t\t\t\t? new Error(\"Stream dropped repeatedly — connection likely unstable\")\n\t\t\t\t: error instanceof Error\n\t\t\t\t\t? error\n\t\t\t\t\t: new Error(String(error));\n\t\t\treturn finalizeMessage(createErrorMessage(surfacedError));\n\t\t}\n\n\t\tlet shouldRetry = false;\n\t\tlet retryError = \"\";\n\t\t// Set when the stream completes with a \"length\" result that warrants another\n\t\t// attempt at the same configured output limit.\n\t\tlet lengthRetry = false;\n\t\ttry {\n\t\t\tfor await (const event of response) {\n\t\t\t\tswitch (event.type) {\n\t\t\t\t\tcase \"start\":\n\t\t\t\t\t\tpartialMessage = event.partial;\n\t\t\t\t\t\tcontext.messages.push(partialMessage);\n\t\t\t\t\t\taddedPartial = true;\n\t\t\t\t\t\tawait emit({ type: \"message_start\", message: { ...partialMessage } });\n\t\t\t\t\t\tbreak;\n\n\t\t\t\t\tcase \"text_start\":\n\t\t\t\t\tcase \"text_delta\":\n\t\t\t\t\tcase \"text_end\":\n\t\t\t\t\tcase \"thinking_start\":\n\t\t\t\t\tcase \"thinking_delta\":\n\t\t\t\t\tcase \"thinking_end\":\n\t\t\t\t\tcase \"toolcall_start\":\n\t\t\t\t\tcase \"toolcall_delta\":\n\t\t\t\t\tcase \"toolcall_end\":\n\t\t\t\t\t\tif (partialMessage) {\n\t\t\t\t\t\t\tpartialMessage = event.partial;\n\t\t\t\t\t\t\tcontext.messages[context.messages.length - 1] = partialMessage;\n\t\t\t\t\t\t\tawait emit({\n\t\t\t\t\t\t\t\ttype: \"message_update\",\n\t\t\t\t\t\t\t\tassistantMessageEvent: event,\n\t\t\t\t\t\t\t\tmessage: { ...partialMessage },\n\t\t\t\t\t\t\t});\n\t\t\t\t\t\t}\n\t\t\t\t\t\tbreak;\n\n\t\t\t\t\tcase \"done\":\n\t\t\t\t\tcase \"error\": {\n\t\t\t\t\t\tconst result = await response.result();\n\t\t\t\t\t\t// Check if this is a stream-drop error that should be retried\n\t\t\t\t\t\tconst isDrop = result.stopReason === \"error\" && isStreamDropError(result.errorMessage);\n\t\t\t\t\t\tif (isDrop && streamDropAttempt < maxRetries && !signal?.aborted) {\n\t\t\t\t\t\t\tshouldRetry = true;\n\t\t\t\t\t\t\tretryError =\n\t\t\t\t\t\t\t\tresult.errorMessage ?? \"Stream ended without terminal event — connection likely dropped\";\n\t\t\t\t\t\t} else if (isDrop) {\n\t\t\t\t\t\t\t// Final attempt exhausted — surface friendly message\n\t\t\t\t\t\t\treturn finalizeMessage(\n\t\t\t\t\t\t\t\tcreateErrorMessage(new Error(\"Stream dropped repeatedly — connection likely unstable\")),\n\t\t\t\t\t\t\t);\n\t\t\t\t\t\t} else if (result.stopReason === \"length\") {\n\t\t\t\t\t\t\t// Retry a bounded number of times at the same configured output limit.\n\t\t\t\t\t\t\t// A larger synthetic budget would override the active model setting and\n\t\t\t\t\t\t\t// discard paid work before eventually reaching the configured value.\n\t\t\t\t\t\t\tif (lengthAttempts < lengthRetries && !signal?.aborted) {\n\t\t\t\t\t\t\t\tlengthRetry = true;\n\t\t\t\t\t\t\t} else {\n\t\t\t\t\t\t\t\treturn finalizeMessage(createLengthExhaustedMessage(lengthAttempts + 1, result.errorMessage));\n\t\t\t\t\t\t\t}\n\t\t\t\t\t\t} else {\n\t\t\t\t\t\t\treturn finalizeMessage(result);\n\t\t\t\t\t\t}\n\t\t\t\t\t\tbreak;\n\t\t\t\t\t}\n\t\t\t\t}\n\t\t\t\tif (shouldRetry || lengthRetry) break;\n\t\t\t}\n\t\t} catch (error) {\n\t\t\tif (isStreamDropError(error) && streamDropAttempt < maxRetries && !signal?.aborted) {\n\t\t\t\tshouldRetry = true;\n\t\t\t\tretryError = error instanceof Error ? error.message : String(error);\n\t\t\t} else {\n\t\t\t\tconst surfacedError = isStreamDropError(error)\n\t\t\t\t\t? new Error(\"Stream dropped repeatedly — connection likely unstable\")\n\t\t\t\t\t: error instanceof Error\n\t\t\t\t\t\t? error\n\t\t\t\t\t\t: new Error(String(error));\n\t\t\t\treturn finalizeMessage(createErrorMessage(surfacedError));\n\t\t\t}\n\t\t}\n\n\t\tif (lengthRetry) {\n\t\t\tawait emit({\n\t\t\t\ttype: \"length_retry\",\n\t\t\t\tattempt: lengthAttempts + 1,\n\t\t\t\tmaxAttempts: lengthRetries,\n\t\t\t\tmaxTokens: requestMaxTokens,\n\t\t\t\tdiscardedPartial: clonePartialForDebug(),\n\t\t\t});\n\t\t\tdiscardPartial();\n\t\t\tlengthAttempts++;\n\t\t\t// A length retry is a fresh request; reset the stream-drop counter.\n\t\t\tstreamDropAttempt = 0;\n\t\t\tcontinue;\n\t\t}\n\n\t\tif (shouldRetry) {\n\t\t\tawait emit({\n\t\t\t\ttype: \"stream_retry\",\n\t\t\t\tattempt: streamDropAttempt + 1,\n\t\t\t\tmaxAttempts: maxRetries,\n\t\t\t\terror: retryError,\n\t\t\t\tdiscardedPartial: clonePartialForDebug(),\n\t\t\t});\n\t\t\tawait sleep(retryBaseDelay * 2 ** streamDropAttempt, signal);\n\t\t\tdiscardPartial();\n\t\t\tstreamDropAttempt++;\n\t\t\tcontinue;\n\t\t}\n\n\t\treturn finalizeMessage(await response.result());\n\t}\n}\n\ninterface ToolExecutionResult {\n\tresults: ToolResultMessage[];\n\tendTurn: boolean;\n}\n\n/**\n * Execute tool calls from an assistant message.\n */\nasync function executeToolCalls(\n\tcurrentContext: AgentContext,\n\tassistantMessage: AssistantMessage,\n\tconfig: AgentLoopConfig,\n\tsignal: AbortSignal | undefined,\n\temit: AgentEventSink,\n): Promise<ToolExecutionResult> {\n\tconst toolCalls = assistantMessage.content.filter((c) => c.type === \"toolCall\");\n\tif (config.toolExecution === \"sequential\") {\n\t\treturn executeToolCallsSequential(currentContext, assistantMessage, toolCalls, config, signal, emit);\n\t}\n\treturn executeToolCallsParallel(currentContext, assistantMessage, toolCalls, config, signal, emit);\n}\n\nasync function executeToolCallsSequential(\n\tcurrentContext: AgentContext,\n\tassistantMessage: AssistantMessage,\n\ttoolCalls: AgentToolCall[],\n\tconfig: AgentLoopConfig,\n\tsignal: AbortSignal | undefined,\n\temit: AgentEventSink,\n): Promise<ToolExecutionResult> {\n\tconst results: ToolResultMessage[] = [];\n\tlet endTurn = false;\n\n\tfor (const toolCall of toolCalls) {\n\t\tawait emit({\n\t\t\ttype: \"tool_execution_start\",\n\t\t\ttoolCallId: toolCall.id,\n\t\t\ttoolName: toolCall.name,\n\t\t\targs: toolCall.arguments,\n\t\t});\n\n\t\tconst preparation = await prepareToolCall(currentContext, assistantMessage, toolCall, config, signal);\n\t\tif (preparation.kind === \"immediate\") {\n\t\t\tif (preparation.result.endTurn) endTurn = true;\n\t\t\tresults.push(await emitToolCallOutcome(toolCall, preparation.result, preparation.isError, emit));\n\t\t} else {\n\t\t\tconst executed = await executePreparedToolCall(preparation, signal, emit);\n\t\t\tif (executed.result.endTurn) endTurn = true;\n\t\t\tresults.push(\n\t\t\t\tawait finalizeExecutedToolCall(\n\t\t\t\t\tcurrentContext,\n\t\t\t\t\tassistantMessage,\n\t\t\t\t\tpreparation,\n\t\t\t\t\texecuted,\n\t\t\t\t\tconfig,\n\t\t\t\t\tsignal,\n\t\t\t\t\temit,\n\t\t\t\t),\n\t\t\t);\n\t\t}\n\t}\n\n\treturn { results, endTurn };\n}\n\nasync function executeToolCallsParallel(\n\tcurrentContext: AgentContext,\n\tassistantMessage: AssistantMessage,\n\ttoolCalls: AgentToolCall[],\n\tconfig: AgentLoopConfig,\n\tsignal: AbortSignal | undefined,\n\temit: AgentEventSink,\n): Promise<ToolExecutionResult> {\n\tconst results: ToolResultMessage[] = [];\n\tconst runnableCalls: PreparedToolCall[] = [];\n\tlet endTurn = false;\n\n\tfor (const toolCall of toolCalls) {\n\t\tawait emit({\n\t\t\ttype: \"tool_execution_start\",\n\t\t\ttoolCallId: toolCall.id,\n\t\t\ttoolName: toolCall.name,\n\t\t\targs: toolCall.arguments,\n\t\t});\n\n\t\tconst preparation = await prepareToolCall(currentContext, assistantMessage, toolCall, config, signal);\n\t\tif (preparation.kind === \"immediate\") {\n\t\t\tif (preparation.result.endTurn) endTurn = true;\n\t\t\tresults.push(await emitToolCallOutcome(toolCall, preparation.result, preparation.isError, emit));\n\t\t} else {\n\t\t\trunnableCalls.push(preparation);\n\t\t}\n\t}\n\n\tconst runningCalls = runnableCalls.map((prepared) => ({\n\t\tprepared,\n\t\texecution: executePreparedToolCall(prepared, signal, emit),\n\t}));\n\n\tfor (const running of runningCalls) {\n\t\tconst executed = await running.execution;\n\t\tif (executed.result.endTurn) endTurn = true;\n\t\tresults.push(\n\t\t\tawait finalizeExecutedToolCall(\n\t\t\t\tcurrentContext,\n\t\t\t\tassistantMessage,\n\t\t\t\trunning.prepared,\n\t\t\t\texecuted,\n\t\t\t\tconfig,\n\t\t\t\tsignal,\n\t\t\t\temit,\n\t\t\t),\n\t\t);\n\t}\n\n\treturn { results, endTurn };\n}\n\ntype PreparedToolCall = {\n\tkind: \"prepared\";\n\ttoolCall: AgentToolCall;\n\ttool: AgentTool<any>;\n\targs: unknown;\n};\n\ntype ImmediateToolCallOutcome = {\n\tkind: \"immediate\";\n\tresult: AgentToolResult<any>;\n\tisError: boolean;\n};\n\ntype ExecutedToolCallOutcome = {\n\tresult: AgentToolResult<any>;\n\tisError: boolean;\n};\n\nasync function prepareToolCall(\n\tcurrentContext: AgentContext,\n\tassistantMessage: AssistantMessage,\n\ttoolCall: AgentToolCall,\n\tconfig: AgentLoopConfig,\n\tsignal: AbortSignal | undefined,\n): Promise<PreparedToolCall | ImmediateToolCallOutcome> {\n\tconst tool = currentContext.tools?.find((t) => t.name === toolCall.name);\n\tif (!tool) {\n\t\treturn {\n\t\t\tkind: \"immediate\",\n\t\t\tresult: createErrorToolResult(`Tool ${toolCall.name} not found`),\n\t\t\tisError: true,\n\t\t};\n\t}\n\n\ttry {\n\t\tconst validatedArgs = validateToolArguments(tool, toolCall);\n\t\tif (config.beforeToolCall) {\n\t\t\tconst beforeResult = await config.beforeToolCall(\n\t\t\t\t{\n\t\t\t\t\tassistantMessage,\n\t\t\t\t\ttoolCall,\n\t\t\t\t\targs: validatedArgs,\n\t\t\t\t\tcontext: currentContext,\n\t\t\t\t},\n\t\t\t\tsignal,\n\t\t\t);\n\t\t\tif (beforeResult?.block) {\n\t\t\t\treturn {\n\t\t\t\t\tkind: \"immediate\",\n\t\t\t\t\tresult: createErrorToolResult(beforeResult.reason || \"Tool execution was blocked\"),\n\t\t\t\t\tisError: true,\n\t\t\t\t};\n\t\t\t}\n\t\t}\n\t\treturn {\n\t\t\tkind: \"prepared\",\n\t\t\ttoolCall,\n\t\t\ttool,\n\t\t\targs: validatedArgs,\n\t\t};\n\t} catch (error) {\n\t\treturn {\n\t\t\tkind: \"immediate\",\n\t\t\tresult: createErrorToolResult(error instanceof Error ? error.message : String(error)),\n\t\t\tisError: true,\n\t\t};\n\t}\n}\n\nasync function executePreparedToolCall(\n\tprepared: PreparedToolCall,\n\tsignal: AbortSignal | undefined,\n\temit: AgentEventSink,\n): Promise<ExecutedToolCallOutcome> {\n\tconst updateEvents: Promise<void>[] = [];\n\n\ttry {\n\t\tconst result = await prepared.tool.execute(\n\t\t\tprepared.toolCall.id,\n\t\t\tprepared.args as never,\n\t\t\tsignal,\n\t\t\t(partialResult) => {\n\t\t\t\tupdateEvents.push(\n\t\t\t\t\tPromise.resolve(\n\t\t\t\t\t\temit({\n\t\t\t\t\t\t\ttype: \"tool_execution_update\",\n\t\t\t\t\t\t\ttoolCallId: prepared.toolCall.id,\n\t\t\t\t\t\t\ttoolName: prepared.toolCall.name,\n\t\t\t\t\t\t\targs: prepared.toolCall.arguments,\n\t\t\t\t\t\t\tpartialResult,\n\t\t\t\t\t\t}),\n\t\t\t\t\t),\n\t\t\t\t);\n\t\t\t},\n\t\t);\n\t\tawait Promise.all(updateEvents);\n\t\treturn { result, isError: false };\n\t} catch (error) {\n\t\tawait Promise.all(updateEvents);\n\t\treturn {\n\t\t\tresult: createErrorToolResult(error instanceof Error ? error.message : String(error)),\n\t\t\tisError: true,\n\t\t};\n\t}\n}\n\nasync function finalizeExecutedToolCall(\n\tcurrentContext: AgentContext,\n\tassistantMessage: AssistantMessage,\n\tprepared: PreparedToolCall,\n\texecuted: ExecutedToolCallOutcome,\n\tconfig: AgentLoopConfig,\n\tsignal: AbortSignal | undefined,\n\temit: AgentEventSink,\n): Promise<ToolResultMessage> {\n\tlet result = executed.result;\n\tlet isError = executed.isError;\n\n\tif (config.afterToolCall) {\n\t\tconst afterResult = await config.afterToolCall(\n\t\t\t{\n\t\t\t\tassistantMessage,\n\t\t\t\ttoolCall: prepared.toolCall,\n\t\t\t\targs: prepared.args,\n\t\t\t\tresult,\n\t\t\t\tisError,\n\t\t\t\tcontext: currentContext,\n\t\t\t},\n\t\t\tsignal,\n\t\t);\n\t\tif (afterResult) {\n\t\t\tresult = {\n\t\t\t\tcontent: afterResult.content ?? result.content,\n\t\t\t\tdetails: afterResult.details ?? result.details,\n\t\t\t};\n\t\t\tisError = afterResult.isError ?? isError;\n\t\t}\n\t}\n\n\treturn await emitToolCallOutcome(prepared.toolCall, result, isError, emit);\n}\n\nfunction createErrorToolResult(message: string): AgentToolResult<any> {\n\treturn {\n\t\tcontent: [{ type: \"text\", text: message }],\n\t\tdetails: {},\n\t};\n}\n\nasync function emitToolCallOutcome(\n\ttoolCall: AgentToolCall,\n\tresult: AgentToolResult<any>,\n\tisError: boolean,\n\temit: AgentEventSink,\n): Promise<ToolResultMessage> {\n\tawait emit({\n\t\ttype: \"tool_execution_end\",\n\t\ttoolCallId: toolCall.id,\n\t\ttoolName: toolCall.name,\n\t\tresult,\n\t\tisError,\n\t});\n\n\tconst toolResultMessage: ToolResultMessage = {\n\t\trole: \"toolResult\",\n\t\ttoolCallId: toolCall.id,\n\t\ttoolName: toolCall.name,\n\t\tcontent: result.content,\n\t\tdetails: result.details,\n\t\tisError,\n\t\ttimestamp: Date.now(),\n\t};\n\n\tawait emit({ type: \"message_start\", message: toolResultMessage });\n\tawait emit({ type: \"message_end\", message: toolResultMessage });\n\treturn toolResultMessage;\n}\n"]}