import { DynamicWorkerExecutor } from "@cloudflare/codemode"; import { createCodeTool } from "@cloudflare/codemode/ai"; import { tool, type Tool, type ToolExecutionOptions } from "ai"; import * as Effect from "effect/Effect"; import { z } from "zod"; import { WebDeadline } from "./browser-session"; import type { createWebTools } from "./web-tools"; const RESEARCH_TIMEOUT_MS = 45_000; const MAX_CALLS = 12; const CONCURRENCY = 3; type ResearchTools = ReturnType; export type ResearchObservation = { id: string; name: keyof ResearchTools; input: unknown; signal: AbortSignal; }; type Observe = ( call: ResearchObservation, run: () => Promise, ) => Promise; export function createResearchTool( tools: ResearchTools, loader: WorkerLoader, observe: Observe, ) { const executor = new DynamicWorkerExecutor({ loader, timeout: RESEARCH_TIMEOUT_MS, globalOutbound: null, }); const description = createCodeTool({ tools, executor }).description; return tool({ description: `Batch web searches, page reads and browser renders in one code execution. Only the three research tools below are available. Return source evidence and URLs, not just a conclusion.\n${description}`, inputSchema: z.object({ code: z.string().min(1).max(32_768) }).strict(), execute: async (input, options) => { const deadline = new WebDeadline( RESEARCH_TIMEOUT_MS, options.abortSignal, ); const slots: Promise[] = Array.from( { length: CONCURRENCY }, () => Promise.resolve(), ); let calls = 0; function wrap( name: keyof ResearchTools, original: Tool, ) { return { description: original.description, inputSchema: original.inputSchema, outputSchema: original.outputSchema, // Code Mode calls execute with only the input. Capture this invocation's // AI SDK context explicitly so cancellation and browser progress survive. execute: async (args: INPUT) => { deadline.remaining(); if (calls >= MAX_CALLS) throw new Error("Research batch exceeds 12 tool calls"); const sequence = calls++; const slot = sequence % CONCURRENCY; const id = `${options.toolCallId}:research:${sequence}`; const pending = slots[slot].then(async () => { deadline.remaining(); const context: ToolExecutionOptions = { ...options, toolCallId: id, abortSignal: deadline.signal, }; return observe( { id, name, input: args, signal: deadline.signal }, async () => (await original.execute!(args, context)) as OUTPUT, ); }); slots[slot] = pending.catch(() => {}); return pending; }, }; } try { const codeTool = createCodeTool({ tools: { web_search: wrap("web_search", tools.web_search), read_url: wrap("read_url", tools.read_url), browser_read: wrap("browser_read", tools.browser_read), }, executor, }); return await Effect.runPromise( deadline.limit( Effect.tryPromise(() => Promise.resolve(codeTool.execute!(input, options)), ), ), ); } finally { // Also stop unawaited calls or calls left running after generated code fails. deadline.dispose(); await Promise.allSettled(slots); } }, }); }