// Shared plumbing for the native function-calling gather transports. Each // provider file (anthropic/openai/gemini) implements a non-streaming tool loop: // call the model with a `search_transcripts` tool + a `finish` tool, run the // searches it requests, feed results back, and stop when it finishes (or the // search budget is reached). Only the request/response shapes differ per // provider; this module holds the common types, tool text, and POST helper. import type { ChatMessage, DebugCall, Provider } from "../askProvider"; import { ContextTooLargeError, isContextLengthError, PausableError, } from "../askProvider"; import { acquire, noteRateLimited, parseRetryAfterMs, withRateLimitRetry, } from "../rateLimit"; // Re-exported for existing importers: PausableError now lives in askProvider (so // ensureOk can throw it without a cycle), but the native gather code + the sweep // driver still import it from here. export { PausableError } from "../askProvider"; export type NativeGatherContext = { apiKey: string; model: string; // Full gather system prompt (built by the caller via gatherSystemPrompt). system: string; // Prior turns (already grounded), replayed for context. history: ChatMessage[]; question: string; // Max searches this turn. budget: number; // Runs one search and returns the results text to feed back to the model. runSearch: (query: string) => Promise; // Reads more of a pinned video's transcript around a moment and returns the // window text. Present only in expand mode with a pinned set (native only); // when undefined the fetch_context tool is not offered. runFetchContext?: (video: string, aroundSeconds?: number) => Promise; // Upserts a section of the running report document (report mode, native only). // When undefined the update_report tool is not offered. runUpdateReport?: (section: string, content: string) => Promise; // Optional capture of each gather round (request + raw response) for the debug // export. Never receives the API key. onDebug?: (rec: DebugCall) => void; signal?: AbortSignal; }; export const SEARCH_TOOL_NAME = "search_transcripts"; export const FETCH_TOOL_NAME = "fetch_context"; export const REPORT_TOOL_NAME = "update_report"; export const FINISH_TOOL_NAME = "finish"; // Bounds transcript reads per turn (independent of the search budget) so the // user's key isn't spent on unbounded fetching. export const FETCH_BUDGET = 4; // Bounds report writes per turn (independent of the search budget). Report mode // only needs a handful of section upserts per turn to keep the document current. export const REPORT_BUDGET = 6; export const SEARCH_TOOL_DESCRIPTION = "Search the video-transcript archive for excerpts relevant to a query. " + "Returns matching videos with timestamped snippet lines. Use focused keyword " + "or name queries; call it again to refine based on what you find."; export const FETCH_TOOL_DESCRIPTION = "Read more of a specific pinned video's transcript around a moment, to get the " + "context surrounding a snippet. Pass the video's ref (shown with each pinned " + "result) and optionally a timestamp in seconds to centre on. The returned " + "lines are added to your citable excerpts for that video."; export const REPORT_TOOL_DESCRIPTION = "Create or update a section of the running report document that persists " + "across turns. Pass the section heading and its full new markdown content; " + "the section is replaced (or appended if new)."; export const FINISH_TOOL_DESCRIPTION = "Call this when you have gathered enough excerpts (or none are needed) and " + "are ready to answer."; // JSON-Schema for the search tool's input (provider files wrap this in their // own tool envelope). export const SEARCH_PARAMS = { type: "object", properties: { query: { type: "string", description: "Keywords or a name/phrase to search for.", }, }, required: ["query"], } as const; // JSON-Schema for the fetch_context tool's input. export const FETCH_PARAMS = { type: "object", properties: { video: { type: "string", description: "The video ref to read, exactly as shown with the pinned results.", }, aroundSeconds: { type: "number", description: "Centre the excerpt on this timestamp (seconds). Optional; defaults to " + "the video's first matched moment.", }, }, required: ["video"], } as const; // JSON-Schema for the update_report tool's input. export const REPORT_PARAMS = { type: "object", properties: { section: { type: "string", description: "The section heading (without the leading ##).", }, content: { type: "string", description: "The full new markdown content for the section.", }, }, required: ["section", "content"], } as const; export const EMPTY_PARAMS = { type: "object", properties: {} } as const; // Shape returned by every provider's tool-call parser: `query` for search, // `video`/`aroundSeconds` for fetch_context, `section`/`content` for // update_report (all optional, filled per tool — omitted keys stay absent so // deep-equality against a search/finish call still matches). export type ParsedToolCall = { id: string; name: string; query: string; video: string; aroundSeconds?: number; section?: string; content?: string; }; // Thrown when a native tool-calling request fails in a way that suggests the // model/endpoint doesn't support tools (HTTP 400/404). The agent catches this // and retries the turn with the scripted transport. export class ToolsUnavailableError extends Error {} export function abortError(): DOMException { return new DOMException("Aborted", "AbortError"); } // POST JSON and parse the response. 400/404 → ToolsUnavailableError (so the agent // can fall back to scripted); other non-OK statuses throw a normal Error the chat // surfaces (e.g. 401 bad key). Paced by the shared rate limiter: `acquire` spaces // this call to the model's RPM, and the fetch is wrapped in `withRateLimitRetry` // so a 429 retries the provider's hint before falling through to a PausableError // (which the sweep driver still turns into a resumable pause once retries run out). export async function postJson( url: string, headers: Record, body: unknown, provider: string, signal?: AbortSignal, debug?: { onDebug?: (rec: DebugCall) => void; model?: string }, ): Promise { const providerId = provider.toLowerCase() as Provider; // Model comes from the request body (anthropic/openai) or, for gemini (model in // the URL), the debug hint — either way it keys the per-model pacing gate. const bodyModel = (body as { model?: unknown })?.model; const model = (typeof bodyModel === "string" ? bodyModel : undefined) ?? debug?.model ?? ""; if (signal?.aborted) throw abortError(); await acquire(providerId, model, signal); return withRateLimitRetry( () => postJsonOnce(url, headers, body, provider, providerId, model, signal, debug), { provider: providerId, model, signal }, ); } // One request attempt, re-issuable by withRateLimitRetry. async function postJsonOnce( url: string, headers: Record, body: unknown, provider: string, providerId: Provider, model: string, signal?: AbortSignal, debug?: { onDebug?: (rec: DebugCall) => void; model?: string }, ): Promise { if (signal?.aborted) throw abortError(); const t0 = Date.now(); const res = await fetch(url, { method: "POST", signal, headers: { "content-type": "application/json", ...headers }, body: JSON.stringify(body), }); if (!res.ok) { let detail = ""; try { detail = await res.text(); } catch { /* ignore */ } const msg = `${provider} request failed (${res.status}). ${detail.slice(0, 300)}`; // A 400 caused by an over-large context must NOT be mistaken for "this model // doesn't support tools" (which would pointlessly retry the same oversized // request via the scripted transport). Disambiguate before falling back. if (res.status === 400 && isContextLengthError(detail)) { throw new ContextTooLargeError(msg); } if (res.status === 400 || res.status === 404) { throw new ToolsUnavailableError(msg); } // Rate-limit / quota → resumable. Parse the provider's retry hint + step the // session RPM down; the retry wrapper re-issues this request, and only after // its attempts are exhausted does this bubble to the sweep driver's pause. if (res.status === 429) { const retryAfterMs = parseRetryAfterMs(res.headers, detail) ?? undefined; noteRateLimited(providerId, model); throw new PausableError(msg, "usage limit", retryAfterMs); } throw new Error(msg); } const json = await res.json(); debug?.onDebug?.({ phase: "gather", provider, model: debug.model, request: body, response: json, ms: Date.now() - t0, }); return json; }