"use client"; import { useCallback, useEffect, useMemo, useRef, useState } from "react"; import { useSearchData } from "yt-dlp-transcript-common/components/SearchDataContext"; import { useSearchSession } from "yt-dlp-transcript-common/components/SearchSessionContext"; import type { SearchHandoff } from "yt-dlp-transcript-common/lib/aiHandoff"; import { fetchTranscript } from "yt-dlp-transcript-common/components/transcriptCache"; import { useChannelPostsManifests } from "yt-dlp-transcript-common/components/postsCache"; import { makeId, splitId } from "yt-dlp-transcript-common/components/originId"; import { cuesToSnippets, mergeSnippets, windowCues, } from "yt-dlp-transcript-common/lib/transcriptWindow"; import { PROVIDERS, type DebugCall, type Provider } from "../lib/askProvider"; import { DEFAULT_ANSWER_TOKENS, runAskTurn, runReportChunk, supportsNativeTools, type AgentEvent, type AgentMode, } from "../lib/searchAgent"; import type { RetrievedVideo } from "../lib/askRetrieval"; import { PausableError } from "../lib/nativeTools/shared"; import { defaultRpmFor, effectiveRpm as computeEffectiveRpm, getRpmOverride, onRateLimitEvent, setRpmOverride, } from "../lib/rateLimit"; import { buildDebugExport } from "../lib/askDebug"; import { buildApiMessages, buildGroundedContent, chunk, mergeReportSources, parseContext, reportNumbering, serializeContext, type SearchStep, type UiMessage, } from "../lib/askConversation"; import { loadSavedChats, saveSavedChats, savedChatsEqual, type SavedChat, type SavedChatsState, type SweepCheckpoint, } from "./askChatStorage"; import { ensureSweepNotifyPermission, notifySweep } from "./askNotify"; const K_PROVIDER = "ytdlp-tb:ai:provider"; const K_REMEMBER = "ytdlp-tb:ai:remember"; const K_SEARCHMODE = "ytdlp-tb:ai:searchmode"; const K_MARKDOWN = "ytdlp-tb:ai:md"; const K_CONVO = "ytdlp-tb:ai:conversation"; const K_MAXANSWER = "ytdlp-tb:ai:maxanswer"; // Optional test/power-user override for the sweep batch size (see sweepChunkSize). const K_SWEEPCHUNK = "ytdlp-tb:ai:sweepchunk"; const keyFor = (p: Provider) => `ytdlp-tb:ai:key:${p}`; const modelFor = (p: Provider) => `ytdlp-tb:ai:model:${p}`; // Bounds for the user-set "Max answer length" (output tokens). const MAX_ANSWER_MIN = 256; const MAX_ANSWER_MAX = 32768; // Whole-corpus sweep: how many result videos each batch folds into the report. // Small enough to stay within a native gather's token budget per pass, large // enough that a big result set is a handful of batches, not hundreds. const SWEEP_CHUNK_SIZE = 10; // Directive fallback when the composer is empty at sweep time. export const DEFAULT_SWEEP_DIRECTIVE = "key claims & contradictions"; // How many recent provider calls to retain for the debug export (ring buffer). const DEBUG_TRACE_CAP = 8; type Snip = { clock: string; seconds: number; text: string }; type PersistedConvo = { messages: UiMessage[]; contextOverride: string | null; // When grounded in the live search: answer only from it (skip the AI's own // search). strictGrounding?: boolean; // The user detached from the live search to ask a free-form question. detached?: boolean; // Per-video extra transcript pulled in via "Load context", overlaid onto the // live search grounding (keyed by video slug). enrichments?: Record; // Report mode: the running report document + whether the mode is on. Persisted // so a long research session's report survives reloads. report?: string; reportMode?: boolean; // The report's SOURCE REGISTRY — every video the report's `[n]` citations index // into, ordered + deduped, trimmed (no snippet text). Persisted alongside the // report so its clickable citations survive a reload. reportSources?: RetrievedVideo[]; // A paused sweep's checkpoint (remaining batches + directive), so a paused sweep // survives a reload and can be resumed — even under a different model. checkpoint?: SweepCheckpoint | null; }; const NON_TERMINAL_PHASES = new Set(["gathering", "answering", "streaming"]); // A turn that was mid-flight when the tab closed/reloaded can't resume — mark it // interrupted (shows a Retry) instead of restoring a phase that spins forever. function normalizeRestoredMessage(m: UiMessage): UiMessage { if (m.role === "assistant" && m.phase && NON_TERMINAL_PHASES.has(m.phase)) { return { ...m, phase: "error", error: true, content: m.content || "This answer was interrupted. Retry to continue.", }; } return m; } // Write the conversation, catching a storage-quota overflow. Returns true if it // persisted, false if the quota was exceeded (so the caller can trim + retry). function writeConvo(payload: PersistedConvo): boolean { try { localStorage.setItem(K_CONVO, JSON.stringify(payload)); return true; } catch (e) { const quota = e instanceof DOMException && (e.name === "QuotaExceededError" || e.name === "NS_ERROR_DOM_QUOTA_REACHED"); return !quota; // non-quota errors: treat as "done" (storage unavailable) } } // State + turn orchestration for the /ask chat. Owns the conversation (persisted // across reloads), the retrieval agent turns (gather → answer), a Retry path that // reuses already-gathered excerpts, and a human-editable context override. export function useAskChat() { const [provider, setProviderState] = useState("anthropic"); const [apiKey, setApiKey] = useState(""); const [model, setModel] = useState(PROVIDERS.anthropic.defaultModel); const [remember, setRemember] = useState(true); const [showKey, setShowKey] = useState(false); const [searchMode, setSearchModeState] = useState("auto"); const [markdownOn, setMarkdownOnState] = useState(true); const { summariesState, channels, aliases, postsManifest } = useSearchData(); const { summaries, summariesReady } = summariesState; const corpusError = summariesState.error?.message ?? null; // The social-post corpus, so the chat's retrieval is grounded in BOTH // datasets. Flattened from the per-channel posts manifests into the slug set // the search engine scopes posts leaves to. Empty on a video-only site. const postsRefs = useMemo( () => postsManifest ? postsManifest.channels.map((c) => { const { origin, slug } = splitId(c.slug); return { origin, channelSlug: slug }; }) : [], [postsManifest], ); const channelPostsQueries = useChannelPostsManifests(postsRefs); const postScopeSlugs = useMemo>(() => { const set = new Set(); for (let i = 0; i < channelPostsQueries.length; i++) { const data = channelPostsQueries[i]?.data; const ref = postsRefs[i]; if (!data || !ref) continue; for (const id of Object.keys(data.slugToPage)) { set.add(makeId(ref.origin, `${ref.channelSlug}/${id}`)); } } return set; // eslint-disable-next-line react-hooks/exhaustive-deps }, [postsRefs, channelPostsQueries.map((q) => (q.data ? 1 : 0)).join("")]); // The live "active search" this workspace shares — the chat auto-grounds in it. // `activeGrounding` is the grounding for the ACTIVE mode (whole search ⇄ // selection); `getGrounding(target)` lazily builds a FULL-excerpt handoff over // a target for the corpus sweep. The selection API + mode drive the palette. const { activeGrounding, getGrounding, capped, groundingMode, setGroundingMode, effectiveGroundingMode, selectedSlugs, selectedCount, toggleSelected, clearSelection, } = useSearchSession(); const [messages, setMessages] = useState([]); // A human-pruned context that seeds the conversation: prepended to the replayed // history on every turn. Set by applying an edit in the context panel. const [contextOverride, setContextOverride] = useState(null); // The user detached from the live search to ask a free-form question (the chat // then searches on its own, ignoring the current search). const [detached, setDetached] = useState(false); // Per-video extra transcript pulled in via "Load context", overlaid onto the // live grounding (keyed by video slug) so an expansion survives search updates. const [enrichments, setEnrichments] = useState>({}); // Grounded mode: answer strictly from the live search (skip gather) vs. use it // as a starting point the AI may expand with its own searches. const [strictGrounding, setStrictGroundingState] = useState(true); // Report mode: an opt-in mode (native tool-capable providers only) where the // model maintains a persistent markdown report via update_report, and the // conversation is compacted into that report instead of replaying every turn. const [report, setReport] = useState(""); // The report's source registry: every video cited across all report writes // (report-mode turns + sweep batches), ordered + deduped by key, trimmed. The // report's `[n]`/`[n @ mm:ss]` citations resolve against this. Grows as writes // land; cleared on New chat / a fresh sweep. const [reportSources, setReportSources] = useState([]); const [reportMode, setReportModeState] = useState(false); // True while a turn is writing to the report (drives the panel shimmer); // cleared when the turn ends. const [reportUpdating, setReportUpdating] = useState(false); // Whole-corpus sweep: true while a sweep runs, plus its per-batch progress // ({done, total} batches + running section-write count) driving the progress // strip in the pinned panel and the ReportPanel shimmer. const [sweeping, setSweeping] = useState(false); const [sweepProgress, setSweepProgress] = useState({ done: 0, total: 0, sections: 0, }); // User-settable sweep batch size (how many result videos each batch folds into // the report). Restored from the K_SWEEPCHUNK key on mount (which e2e/power-users // also set directly). Bigger → fewer, heavier batches; smaller → more, lighter. const [sweepBatch, setSweepBatchState] = useState(SWEEP_CHUNK_SIZE); const sweepBatchRef = useRef(sweepBatch); sweepBatchRef.current = sweepBatch; // Set while a manual Pause has been requested; the sweep loop checkpoints and // stops cleanly at the next batch boundary (vs. Stop, which aborts + discards). const sweepPauseRef = useRef(false); // A paused sweep's checkpoint (remaining batches + directive + reason). Non-null // ⇒ the sweep paused and can be resumed. Persisted with the working convo AND // any saved chat, so a pause survives a reload or a model switch. const [sweepCheckpoint, setSweepCheckpoint] = useState(null); // Saved chats library: named conversation copies + which one the working chat // derives from (for the dirty dot). The live "working" chat stays in K_CONVO. const [savedChats, setSavedChats] = useState>({}); const [activeChatName, setActiveChatName] = useState(null); // Max output tokens for the answer phase (user-tunable; default 8192). Raising // it lets long reports finish instead of truncating at the model's output cap. const [maxAnswerTokens, setMaxAnswerTokensState] = useState(DEFAULT_ANSWER_TOKENS); // Ring buffer of the most recent provider calls (request/response/finishReason), // for the debug export. Runtime-only (not persisted); trimmed to DEBUG_TRACE_CAP. const debugTraceRef = useRef([]); // A short human notice while the rate limiter is pacing/backing off (e.g. // "Rate limited — retrying in 12s"), driven by onRateLimitEvent. Cleared when // the wait ends. Surfaced in the sweep progress strip so a paced run never // looks frozen. Runtime-only. const [rateLimitNotice, setRateLimitNotice] = useState(null); // Bumped after an RPM override is written so the derived effective/override/ // default values recompute (the store itself lives in localStorage). const [rpmTick, setRpmTick] = useState(0); // False until after hydration. The RPM values read localStorage, which is absent // during SSR — gate on this so the first client render matches the server (empty) // and only then reflect any stored override, avoiding a hydration mismatch. const [hydrated, setHydrated] = useState(false); // The effective pinned grounding = the ACTIVE-mode grounding (whole search or // the hand-picked selection), overlaid with any per-video "Load context" // enrichments, unless the user detached. const pinned = useMemo(() => { if (detached || !activeGrounding) return null; if (Object.keys(enrichments).length === 0) return activeGrounding; return { ...activeGrounding, videos: activeGrounding.videos.map((v) => enrichments[v.key] ? { ...v, snippets: mergeSnippets(v.snippets, enrichments[v.key], 30) } : v, ), }; }, [detached, activeGrounding, enrichments]); // Per-video key → true while its "Load context" fetch is in flight. const [expanding, setExpanding] = useState>({}); // Per-video key → true when its last "Load context" fetch failed (so the panel // can say so, instead of the spinner just quietly stopping). const [expandError, setExpandError] = useState>({}); // Set when the conversation outgrew the storage quota and older turns had to be // dropped from what's saved — surfaced so a reload's missing history isn't a // silent surprise. const [storageWarning, setStorageWarning] = useState(false); const [input, setInput] = useState(""); const [busy, setBusy] = useState(false); const abortRef = useRef(null); // Guards against a rapid double-submit: `busy` flips asynchronously inside // runTurn, so two fast Enter presses could both pass the `!busy` check before // it settled. This ref is set synchronously in send() and cleared in runTurn. const sendingRef = useRef(false); // Per-provider models the user typed this session (survives provider switches // even when "remember" is off). const sessionModels = useRef>>({}); // Always-current mirrors so send()/runTurn() never act on a stale snapshot // (e.g. sending immediately after applying an edited context). const messagesRef = useRef(messages); messagesRef.current = messages; const contextOverrideRef = useRef(contextOverride); contextOverrideRef.current = contextOverride; const pinnedRef = useRef(pinned); pinnedRef.current = pinned; const strictRef = useRef(strictGrounding); strictRef.current = strictGrounding; const reportRef = useRef(report); reportRef.current = report; const reportSourcesRef = useRef(reportSources); reportSourcesRef.current = reportSources; const reportModeRef = useRef(reportMode); reportModeRef.current = reportMode; const maxAnswerTokensRef = useRef(maxAnswerTokens); maxAnswerTokensRef.current = maxAnswerTokens; // Append a captured provider call to the debug ring buffer (keep the last N). const pushDebug = useCallback((rec: DebugCall) => { const next = [...debugTraceRef.current, rec]; debugTraceRef.current = next.length > DEBUG_TRACE_CAP ? next.slice(-DEBUG_TRACE_CAP) : next; }, []); // Surface the rate limiter's pacing/backoff waits as a short notice (the // progress strip shows it during a sweep). "resumed" clears it. useEffect(() => { return onRateLimitEvent((e) => { if (e.type === "resumed") { setRateLimitNotice(null); return; } const secs = Math.max(1, Math.round(e.ms / 1000)); setRateLimitNotice( e.reason === "retry" ? `Rate limited — retrying in ${secs}s` : `Pacing requests — next in ${secs}s`, ); }); }, []); // Restore saved preferences + (if remembered) the provider's key/model, plus a // persisted conversation. useEffect(() => { /* eslint-disable react-hooks/set-state-in-effect -- one-time restore from localStorage after hydration; it must run in an effect, not a lazy useState initializer, because there is no localStorage during SSR. */ try { const savedProvider = localStorage.getItem(K_PROVIDER) as Provider | null; const rememberSaved = localStorage.getItem(K_REMEMBER) !== "0"; const p = savedProvider && PROVIDERS[savedProvider] ? savedProvider : "anthropic"; setProviderState(p); setRemember(rememberSaved); const savedMode = localStorage.getItem(K_SEARCHMODE) as AgentMode | null; if (savedMode === "auto" || savedMode === "native" || savedMode === "scripted") { setSearchModeState(savedMode); } setMarkdownOnState(localStorage.getItem(K_MARKDOWN) !== "0"); const savedMax = Number(localStorage.getItem(K_MAXANSWER)); if (Number.isFinite(savedMax) && savedMax >= MAX_ANSWER_MIN) { setMaxAnswerTokensState(Math.min(MAX_ANSWER_MAX, savedMax)); } const savedBatch = Number(localStorage.getItem(K_SWEEPCHUNK)); if (Number.isFinite(savedBatch) && savedBatch >= 1) { setSweepBatchState(Math.min(100, Math.floor(savedBatch))); } loadProviderCreds(p, rememberSaved); const rawConvo = localStorage.getItem(K_CONVO); if (rawConvo) { const parsed = JSON.parse(rawConvo) as PersistedConvo; if (Array.isArray(parsed.messages)) { setMessages(parsed.messages.map(normalizeRestoredMessage)); } if (typeof parsed.contextOverride === "string") { setContextOverride(parsed.contextOverride); } if (typeof parsed.strictGrounding === "boolean") { setStrictGroundingState(parsed.strictGrounding); } if (typeof parsed.detached === "boolean") setDetached(parsed.detached); if (parsed.enrichments && typeof parsed.enrichments === "object") { setEnrichments(parsed.enrichments); } if (typeof parsed.report === "string") setReport(parsed.report); if (Array.isArray(parsed.reportSources)) { setReportSources(parsed.reportSources); } if (typeof parsed.reportMode === "boolean") { setReportModeState(parsed.reportMode); } if (parsed.checkpoint && typeof parsed.checkpoint === "object") { setSweepCheckpoint(parsed.checkpoint); } } // Saved chats library (named copies + active pointer). const chats = loadSavedChats(); setSavedChats(chats.items); setActiveChatName(chats.activeName); } catch { /* storage unavailable / corrupt */ } // Now safe to read localStorage-backed RPM overrides in render. setHydrated(true); /* eslint-enable react-hooks/set-state-in-effect */ }, []); // Persist the conversation (messages + override), debounced so streaming deltas // don't thrash storage. Survives reloads and hard errors. const persistTimer = useRef | null>(null); useEffect(() => { if (persistTimer.current) clearTimeout(persistTimer.current); persistTimer.current = setTimeout(() => { try { const hasEnrichments = Object.keys(enrichments).length > 0; if ( messages.length === 0 && !contextOverride && !detached && !hasEnrichments && !report && !reportMode && !sweepCheckpoint ) { localStorage.removeItem(K_CONVO); setStorageWarning(false); return; } // Persist, trimming the oldest turn-pairs if the blob exceeds the quota, // so a long conversation keeps its most recent history rather than // silently failing to save anything. let msgs = messages; for (;;) { if ( writeConvo({ messages: msgs, contextOverride, strictGrounding, detached, enrichments, report, reportSources, reportMode, checkpoint: sweepCheckpoint, }) ) { setStorageWarning(msgs.length < messages.length); break; } if (msgs.length <= 2) { try { localStorage.removeItem(K_CONVO); } catch { /* ignore */ } setStorageWarning(true); break; } msgs = msgs.slice(2); } } catch { /* ignore */ } }, 400); return () => { if (persistTimer.current) clearTimeout(persistTimer.current); }; }, [ messages, contextOverride, strictGrounding, detached, enrichments, report, reportSources, reportMode, sweepCheckpoint, ]); function loadProviderCreds(p: Provider, rememberOn: boolean) { // Prefer a model the user typed for this provider earlier this session, so a // hand-typed model isn't discarded by a round-trip through the dropdown. const session = sessionModels.current[p]; if (rememberOn) { setApiKey(localStorage.getItem(keyFor(p)) ?? ""); setModel( session ?? localStorage.getItem(modelFor(p)) ?? PROVIDERS[p].defaultModel, ); } else { setApiKey(""); setModel(session ?? PROVIDERS[p].defaultModel); } } function setProvider(p: Provider) { // Stash the current (possibly hand-typed) model before switching away. sessionModels.current[provider] = model; setProviderState(p); try { localStorage.setItem(K_PROVIDER, p); } catch { /* ignore */ } loadProviderCreds(p, remember); } function persistKey(p: Provider, key: string, mdl: string, rememberOn: boolean) { try { if (rememberOn) { localStorage.setItem(keyFor(p), key); localStorage.setItem(modelFor(p), mdl); } else { localStorage.removeItem(keyFor(p)); localStorage.removeItem(modelFor(p)); } localStorage.setItem(K_REMEMBER, rememberOn ? "1" : "0"); } catch { /* ignore */ } } function setSearchMode(m: AgentMode) { setSearchModeState(m); try { localStorage.setItem(K_SEARCHMODE, m); } catch { /* ignore */ } } function setMarkdownOn(on: boolean) { setMarkdownOnState(on); try { localStorage.setItem(K_MARKDOWN, on ? "1" : "0"); } catch { /* ignore */ } } function setMaxAnswerTokens(n: number) { const clamped = Math.max( MAX_ANSWER_MIN, Math.min(MAX_ANSWER_MAX, Math.round(n) || DEFAULT_ANSWER_TOKENS), ); setMaxAnswerTokensState(clamped); try { localStorage.setItem(K_MAXANSWER, String(clamped)); } catch { /* ignore */ } } function setSweepBatch(n: number) { const clamped = Math.max(1, Math.min(100, Math.floor(n) || SWEEP_CHUNK_SIZE)); setSweepBatchState(clamped); try { localStorage.setItem(K_SWEEPCHUNK, String(clamped)); } catch { /* ignore */ } } // Requests-per-minute pacing for the CURRENT provider+model. `rpmOverride` is // the user's pinned value (null ⇒ using the default/self-tuned rate); `rpmDefault` // is the free-tier-safe default (shown as the field placeholder); `effectiveRpm` // drives the "≈Ns between requests" hint. All recompute on provider/model change // and after a write (rpmTick). const rpmModel = model.trim() || PROVIDERS[provider].defaultModel; // rpmOverride / effectiveRpm read imperative stores (localStorage + the session // step-down map) React can't track, so rpmTick — bumped on every write — and // rateLimitNotice — bumped on every pacing event, i.e. after a step-down — are // listed as recompute triggers even though the bodies don't reference them. const rpmOverride = useMemo( () => (hydrated ? getRpmOverride(provider, rpmModel) : null), // eslint-disable-next-line react-hooks/exhaustive-deps [provider, rpmModel, rpmTick, hydrated], ); const rpmDefault = useMemo( () => defaultRpmFor(provider, rpmModel), [provider, rpmModel], ); const effectiveRpm = useMemo( () => hydrated ? computeEffectiveRpm(provider, rpmModel) : defaultRpmFor(provider, rpmModel), // eslint-disable-next-line react-hooks/exhaustive-deps [provider, rpmModel, rpmTick, rateLimitNotice, hydrated], ); // Set (n) or clear (null) the RPM override for the current provider+model. const setRpm = useCallback( (n: number | null) => { const clean = n == null || !Number.isFinite(n) || n <= 0 ? null : Math.min(100000, Math.floor(n)); setRpmOverride(provider, rpmModel, clean); setRpmTick((t) => t + 1); }, [provider, rpmModel], ); // Assemble a redacted debug snapshot (chat state + recent raw provider calls) // for offline debugging, and hand it back as a pretty JSON string. The API key // is never included. const buildDebugJson = useCallback((): string => { return buildDebugExport({ provider, model: model.trim() || PROVIDERS[provider].defaultModel, searchMode, markdownOn, reportMode, strictGrounding, detached, maxAnswerTokens, messages: messagesRef.current, report, grounding: pinned ? { label: pinned.label, videoCount: pinned.videos.length, truncated: pinned.truncated } : null, contextOverride, trace: debugTraceRef.current, }); }, [ provider, model, searchMode, markdownOn, reportMode, strictGrounding, detached, maxAnswerTokens, report, pinned, contextOverride, ]); // Report mode is only meaningful on a tool-capable transport: Scripted has no // tool to maintain the report, so the toggle is disabled there. (Persistence is // handled by the conversation-persist effect, which lists reportMode as a dep.) const reportModeAvailable = searchMode !== "scripted" && (searchMode === "native" || supportsNativeTools(provider, model)); function setReportMode(on: boolean) { setReportModeState(on); } // Patch one message by index, immutably. const patchAt = (index: number, fn: (m: UiMessage) => UiMessage) => setMessages((prev) => { if (index < 0 || index >= prev.length) return prev; const next = prev.slice(); next[index] = fn(next[index]); return next; }); // Core turn runner shared by send() and retry(). Patches the assistant message // at `assistantIndex` (and its user turn at `userIndex`) as the agent streams. const runTurn = useCallback( async (args: { question: string; prior: UiMessage[]; userIndex: number; assistantIndex: number; precomputedGrounding?: { videos: UiMessage["sources"]; groundedContent: string; truncated?: boolean; }; seedVideos?: RetrievedVideo[]; }) => { const { question, prior, assistantIndex } = args; setBusy(true); const ac = new AbortController(); abortRef.current = ac; const override = contextOverrideRef.current; const historyOverride = override ? [...parseContext(override), ...buildApiMessages(prior)] : undefined; const onEvent = (e: AgentEvent) => { switch (e.type) { case "search_start": patchAt(assistantIndex, (m) => ({ ...m, phase: "gathering", searchSteps: [...(m.searchSteps ?? []), { query: e.query }], })); break; case "search_done": patchAt(assistantIndex, (m) => { const steps = (m.searchSteps ?? []).slice(); for (let i = steps.length - 1; i >= 0; i--) { if (steps[i].query === e.query && steps[i].count === undefined) { steps[i] = { ...steps[i], count: e.count }; break; } } return { ...m, searchSteps: steps }; }); break; case "fetch_start": patchAt(assistantIndex, (m) => ({ ...m, phase: "gathering", searchSteps: [ ...(m.searchSteps ?? []), { query: e.label, kind: "fetch" }, ], })); break; case "fetch_done": patchAt(assistantIndex, (m) => { const steps = (m.searchSteps ?? []).slice(); for (let i = steps.length - 1; i >= 0; i--) { if (steps[i].kind === "fetch" && steps[i].count === undefined) { steps[i] = { ...steps[i], count: e.count }; break; } } return { ...m, searchSteps: steps }; }); break; case "answer_start": // Store the grounding videos NOW (before streaming) so a failed // stream can be retried without re-searching, and follow-ups can // reuse the excerpts via the cumulative pool. The full grounded text // is rebuilt from `sources` on demand — not persisted. patchAt(assistantIndex, (m) => ({ ...m, phase: "answering", sources: e.videos, })); break; case "report_update": // The model wrote a section — flag the live "updating report" state // (cleared when the turn ends). The updated report text arrives with // the turn's result below. setReportUpdating(true); break; case "delta": patchAt(assistantIndex, (m) => ({ ...m, phase: "streaming", content: m.content + e.text, })); break; } }; // Report mode only takes effect on a tool-capable transport; on scripted it // must stay off (there's no update_report tool, and dropping the replayed // history would lose context with nothing to compact into). const effectiveReportMode = reportModeRef.current && reportModeAvailable; try { const result = await runAskTurn({ provider, apiKey: apiKey.trim(), model: model.trim() || PROVIDERS[provider].defaultModel, mode: searchMode, question, prior, summaries, postScopeSlugs, aliases, signal: ac.signal, onEvent, historyOverride, seedVideos: args.seedVideos, report: reportRef.current, reportMode: effectiveReportMode, reportSources: reportSourcesRef.current, maxAnswerTokens: maxAnswerTokensRef.current, onDebug: pushDebug, precomputedGrounding: args.precomputedGrounding?.groundedContent ? { videos: args.precomputedGrounding.videos ?? [], groundedContent: args.precomputedGrounding.groundedContent, truncated: args.precomputedGrounding.truncated, } : undefined, }); // Persist the (possibly updated) report from this turn. if (effectiveReportMode && result.report !== reportRef.current) { setReport(result.report); } // Grow the report source registry with this turn's gathered videos (in the // global-index order the report writes cited them by), trimmed for storage. if (effectiveReportMode && result.reportSources) { setReportSources(mergeReportSources([], result.reportSources)); } patchAt(assistantIndex, (m) => ({ ...m, content: m.content || result.answer, sources: result.videos, truncated: result.truncated, finishReason: result.finishReason, blockReason: result.blockReason, error: false, phase: "done", model: model.trim() || PROVIDERS[provider].defaultModel, })); } catch (e) { const err = e as Error; if (err.name === "AbortError") { patchAt(assistantIndex, (m) => ({ ...m, content: m.content + "\n\n_(stopped)_", phase: "stopped", })); } else { // Surface the failure even mid-stream: if tokens already arrived, keep // them AND append the error (previously the error was dropped whenever // any content had streamed, so the user never learned what went wrong). // Grounding captured at answer_start is preserved so Retry can skip the // search. patchAt(assistantIndex, (m) => ({ ...m, content: m.content ? `${m.content}\n\n_${err.message}_` : err.message, error: true, phase: "error", })); } } finally { setBusy(false); setReportUpdating(false); sendingRef.current = false; abortRef.current = null; } }, [ provider, apiKey, model, searchMode, summaries, postScopeSlugs, aliases, reportModeAvailable, pushDebug, ], ); const send = useCallback(async () => { const question = input.trim(); if (!question || busy || sendingRef.current || !apiKey.trim() || !summariesReady) { return; } sendingRef.current = true; persistKey(provider, apiKey, model, remember); const prior = messagesRef.current; const userIndex = prior.length; const assistantIndex = prior.length + 1; // Pinned search results ground the turn. Strict → every question answers // only from them (skip gather entirely). Expand → seed the AI's search with // them on the first turn, then let it search freely (the pin stays in history). const pin = pinnedRef.current; let precomputedGrounding: | { videos: RetrievedVideo[]; groundedContent: string; truncated?: boolean } | undefined; let seedVideos: RetrievedVideo[] | undefined; if (pin && pin.videos.length) { const videos = pin.videos as RetrievedVideo[]; if (strictRef.current) { precomputedGrounding = { videos, groundedContent: buildGroundedContent(question, videos), truncated: pin.truncated, }; } else if (prior.length === 0) { seedVideos = videos; } } setInput(""); setMessages((prev) => [ ...prev, { role: "user", content: question }, { role: "assistant", content: "", phase: "gathering", searchSteps: [] }, ]); await runTurn({ question, prior, userIndex, assistantIndex, precomputedGrounding, seedVideos, }); }, [input, busy, apiKey, summariesReady, provider, model, remember, runTurn]); // The shared batch loop behind both a fresh/extend sweep AND a resume. It folds // `remaining` videos into `startReport` batch by batch, updating the sweep // message + progress strip. Three exits: // • completed → the "Built a report…" summary, checkpoint cleared; // • STOP → abort (the ✕): keep the partial report, DISCARD the // checkpoint (a deliberate give-up); // • PAUSE → a manual Pause OR a rate-limit/model warning (PausableError): // checkpoint the remaining batches + directive so the run can // resume later — even after a reload or a model switch. const driveSweep = useCallback( async (args: { remaining: RetrievedVideo[]; startReport: string; directive: string; // Batches already folded before this run (0 for a fresh sweep; the paused // count for a resume) — keeps "Batch i of n" continuous across a pause. doneOffset: number; assistantIndex: number; resultLabel: string; }) => { const { remaining, startReport, directive, doneOffset, assistantIndex, resultLabel } = args; const batches = chunk(remaining, sweepBatchRef.current); const total = doneOffset + batches.length; const mdl = model.trim() || PROVIDERS[provider].defaultModel; setBusy(true); setSweeping(true); setReportModeState(true); sweepPauseRef.current = false; const ac = new AbortController(); abortRef.current = ac; let report = startReport; let sections = 0; let done = doneOffset; setSweepProgress({ done, total, sections }); // Checkpoint the batches from `i` onward (flattened to slugs) so a resume can // rebuild the grounding for exactly what's left, then surface a paused line. const pauseAt = (i: number, reason: string) => { const remainingSlugs = batches.slice(i).flat().map((v) => v.key); setSweepCheckpoint({ directive, remainingSlugs, done, total, reason }); notifySweep( "Sweep paused", `${reason} — ${remainingSlugs.length} result${remainingSlugs.length === 1 ? "" : "s"} left. Resume when it clears.`, ); patchAt(assistantIndex, (m) => ({ ...m, phase: "stopped", content: `Paused (${reason}) at ${done} of ${total} batch${total === 1 ? "" : "es"} — ` + `resume when ready. The partial report is saved.`, })); }; try { for (let i = 0; i < batches.length; i++) { if (ac.signal.aborted) throw new DOMException("Aborted", "AbortError"); // A manual Pause requested between batches → checkpoint + stop cleanly. if (sweepPauseRef.current) { pauseAt(i, "paused"); return; } setReportUpdating(true); // Merge this batch into the report source registry BEFORE building its // context, so its excerpts are numbered by their GLOBAL registry index — // making the report's `[n]` stable across every batch and resolving // against the persisted source list. Keep the ref current synchronously // so the next batch numbers on top of this one. const mergedSources = mergeReportSources( reportSourcesRef.current, batches[i], ); reportSourcesRef.current = mergedSources; setReportSources(mergedSources); const numberOf = reportNumbering(mergedSources); let r; try { r = await runReportChunk({ provider, apiKey: apiKey.trim(), model: mdl, directive, videos: batches[i], index: done + 1, count: total, report, aliases, numberOf, signal: ac.signal, onEvent: (e) => { if (e.type === "report_update") { sections += 1; setSweepProgress((p) => ({ ...p, sections })); } }, onDebug: pushDebug, }); } catch (e) { // A rate-limit / quota (429, surfaced as PausableError) is resumable — // pause + checkpoint rather than lose the whole run. if (e instanceof PausableError) { pauseAt(i, e.reason); return; } throw e; } report = r.report; setReport(report); done += 1; setSweepProgress({ done, total, sections }); } // Completed cleanly → no checkpoint to resume. setSweepCheckpoint(null); notifySweep( "Sweep finished", `${resultLabel} folded into one report.`, ); patchAt(assistantIndex, (m) => ({ ...m, phase: "done", content: `Built a report from ${resultLabel} ` + `across ${total} batch${total === 1 ? "" : "es"}. See the Report panel.`, })); } catch (e) { const err = e as Error; if (err.name === "AbortError") { // Stop is a deliberate give-up → discard the checkpoint, keep the report. setSweepCheckpoint(null); patchAt(assistantIndex, (m) => ({ ...m, phase: "stopped", content: `Stopped after ${done} of ${total} batches — the partial report is saved.`, })); } else { patchAt(assistantIndex, (m) => ({ ...m, phase: "error", error: true, content: `Couldn't finish the sweep (${err.message}). The report has what completed so far.`, })); } } finally { setBusy(false); setSweeping(false); setReportUpdating(false); setRateLimitNotice(null); sweepPauseRef.current = false; sendingRef.current = false; abortRef.current = null; } }, [provider, apiKey, model, aliases, pushDebug], ); // Whole-corpus "sweep": read the active target (whole search or the hand-picked // selection) in batches, folding each batch's findings into the running report // and discarding its raw excerpts. `mode` "fresh" starts a new report; "extend" // folds into the existing one. Native-only (report mode). const runAccumulationReport = useCallback( async ({ directive, mode, }: { directive: string; mode: "fresh" | "extend"; }) => { if ( busy || sendingRef.current || !apiKey.trim() || !summariesReady || !reportModeAvailable ) { return; } const full = getGrounding(effectiveGroundingMode); if (!full || full.videos.length === 0) return; sendingRef.current = true; persistKey(provider, apiKey, model, remember); // Best-effort: on a PWA site, offer to notify when this (possibly long) run // finishes or pauses. No-op on non-PWA sites / if already decided. void ensureSweepNotifyPermission(); const videos = full.videos as RetrievedVideo[]; const cleanDirective = directive.trim() || DEFAULT_SWEEP_DIRECTIVE; const resultLabel = `${videos.length} result${videos.length === 1 ? "" : "s"}`; // A sweep "assistant" message tracks the run in the transcript. const assistantIndex = messagesRef.current.length; setMessages((prev) => [ ...prev, { role: "assistant", content: `Sweeping ${videos.length} result${videos.length === 1 ? "" : "s"}…`, phase: "gathering", searchSteps: [], }, ]); const startReport = mode === "fresh" ? "" : reportRef.current; if (mode === "fresh") { setReport(""); // A fresh report starts from an empty source registry too (a new document's // `[n]` must not resolve against the previous report's sources). setReportSources([]); reportSourcesRef.current = []; } await driveSweep({ remaining: videos, startReport, directive: cleanDirective, doneOffset: 0, assistantIndex, resultLabel, }); }, [ busy, apiKey, summariesReady, reportModeAvailable, getGrounding, effectiveGroundingMode, provider, model, remember, driveSweep, ], ); // Resume a paused sweep from its checkpoint. Rebuilds the grounding for the // remaining slugs from the CURRENT active target (so it works after a reload, // and under any provider/model — the checkpoint carries no key), then continues // folding into the saved report from where it paused. const resumeSweep = useCallback(async () => { const cp = sweepCheckpoint; if ( !cp || busy || sendingRef.current || !apiKey.trim() || !summariesReady || !reportModeAvailable ) { return; } const full = getGrounding(effectiveGroundingMode) ?? getGrounding("search"); const bySlug = new Map((full?.videos ?? []).map((v) => [v.key, v])); const remaining = cp.remainingSlugs .map((s) => bySlug.get(s)) .filter((v): v is RetrievedVideo => Boolean(v)); if (remaining.length === 0) { // The grounding changed out from under the checkpoint — nothing to resume. setSweepCheckpoint(null); return; } sendingRef.current = true; persistKey(provider, apiKey, model, remember); const resultLabel = `${remaining.length} remaining result${remaining.length === 1 ? "" : "s"}`; const assistantIndex = messagesRef.current.length; setMessages((prev) => [ ...prev, { role: "assistant", content: `Resuming the sweep — ${remaining.length} result${remaining.length === 1 ? "" : "s"} left…`, phase: "gathering", searchSteps: [], }, ]); await driveSweep({ remaining, startReport: reportRef.current, directive: cp.directive, doneOffset: cp.done, assistantIndex, resultLabel, }); }, [ sweepCheckpoint, busy, apiKey, summariesReady, reportModeAvailable, getGrounding, effectiveGroundingMode, provider, model, remember, driveSweep, ]); // Request a courtesy pause: the sweep loop checkpoints + stops at the next batch // boundary (the in-flight batch finishes first, so nothing is wasted). const pauseSweep = useCallback(() => { if (!sweeping) return; sweepPauseRef.current = true; }, [sweeping]); // Discard a paused sweep's checkpoint (give up on resuming; the report stays). const discardCheckpoint = useCallback(() => { setSweepCheckpoint(null); }, []); // Re-run a failed assistant turn. Reuses the grounding captured before the // answer stream failed (no re-search); falls back to a full turn otherwise. const retry = useCallback( async (assistantIndex: number) => { if (busy || !apiKey.trim() || !summariesReady) return; const current = messagesRef.current; const userIndex = assistantIndex - 1; const userMsg = current[userIndex]; const failed = current[assistantIndex]; if (!userMsg || userMsg.role !== "user" || !failed) return; const prior = current.slice(0, userIndex); // Reuse the grounding captured before the answer stream failed — rebuilt // from the stored `sources` (no re-search). const precomputedGrounding = failed.sources && failed.sources.length ? { videos: failed.sources, groundedContent: buildGroundedContent(userMsg.content, failed.sources), } : undefined; // Reset the failed message to a pending state (keep captured grounding). patchAt(assistantIndex, (m) => ({ ...m, content: "", error: false, phase: "gathering", searchSteps: precomputedGrounding ? m.searchSteps : [], })); await runTurn({ question: userMsg.content, prior, userIndex, assistantIndex, precomputedGrounding, }); }, [busy, apiKey, summariesReady, runTurn], ); // Edit a prior user turn: pull its text back into the composer and drop it (and // every turn after it) so the next send re-asks from that point. const editUserMessage = useCallback( (userIndex: number) => { if (busy) return; const current = messagesRef.current; const msg = current[userIndex]; if (!msg || msg.role !== "user") return; setInput(msg.content); setMessages(current.slice(0, userIndex)); }, [busy], ); const stop = () => abortRef.current?.abort(); const reset = () => { if (busy) return; setMessages([]); setContextOverride(null); setDetached(false); setEnrichments({}); setStrictGroundingState(true); setReport(""); setReportSources([]); setReportModeState(false); setSweepCheckpoint(null); // A new chat is a blank working conversation not derived from any saved chat — // clear the active pointer in state AND in the persisted store, so a reload // doesn't reattach to the previously-active saved chat. setActiveChatName(null); try { const cur = loadSavedChats(); if (cur.activeName !== null) { cur.activeName = null; saveSavedChats(cur); } } catch { /* ignore */ } }; // Detach from the live search → ask a free-form question (the AI searches on // its own). Re-attach binds the current search back as grounding. const detach = () => { if (busy) return; setDetached(true); }; const reattach = () => { if (busy) return; setDetached(false); }; // Whether there's a grounding available to ground in (regardless of detach) — // the whole search or a hand-picked selection. const hasLiveSearch = !!activeGrounding; const setStrictGrounding = (on: boolean) => setStrictGroundingState(on); // User-driven "expand this hit": read the surrounding transcript for a grounded // video (client-side, IDB-cached) and stash it as an enrichment overlaid onto // the live grounding, so the assistant reads the fuller context on the next // question. Works in strict AND expand mode, on every provider — no LLM call. // Persisted with the convo; survives the live search updating around it. const expandPinnedVideo = useCallback( async (key: string, aroundSeconds?: number) => { const v = pinnedRef.current?.videos.find((x) => x.key === key); if (!v) return; const center = typeof aroundSeconds === "number" ? aroundSeconds : v.snippets[0]?.seconds ?? 0; setExpanding((e) => ({ ...e, [key]: true })); setExpandError((e) => { if (!e[key]) return e; const next = { ...e }; delete next[key]; return next; }); try { const detail = await fetchTranscript(key); const snips = cuesToSnippets( windowCues(detail.cues ?? [], center, { before: 45, after: 45, maxCues: 60, }), ); setEnrichments((e) => ({ ...e, [key]: mergeSnippets(e[key] ?? [], snips, 30), })); } catch { // Transcript unavailable — flag it so the panel can surface "couldn't // load" instead of the spinner silently stopping. setExpandError((e) => ({ ...e, [key]: true })); } finally { setExpanding((e) => { const next = { ...e }; delete next[key]; return next; }); } }, [], ); // The effective context that will be sent next turn, serialized for the panel: // the override (if any) prepended to the replayed conversation. Held stable // while a turn streams — recomputing it per token thrashed the panel's // (uncontrolled) textarea and blew away any in-progress edit. const contextTextRef = useRef(""); const contextText = useMemo(() => { if (busy) return contextTextRef.current; const base = contextOverride ? parseContext(contextOverride) : []; const text = serializeContext([...base, ...buildApiMessages(messages)]); contextTextRef.current = text; return text; }, [contextOverride, messages, busy]); // Apply an edited context as the new starting point: it becomes the override // (prepended to future turns) and the visible conversation resets to it. const applyContext = (text: string) => { if (busy) return; const trimmed = text.trim(); setContextOverride(trimmed ? trimmed : null); setMessages([]); }; const clearContextOverride = () => { if (busy) return; setContextOverride(null); }; // ── Saved chats library ── // The current working conversation as a SavedChat snapshot (what a save writes). const workingChat = useMemo( () => ({ messages, contextOverride, strictGrounding, detached, enrichments, report, reportSources, reportMode, checkpoint: sweepCheckpoint, provider, model: model.trim() || PROVIDERS[provider].defaultModel, }), [ messages, contextOverride, strictGrounding, detached, enrichments, report, reportSources, reportMode, sweepCheckpoint, provider, model, ], ); // Dirty when the working chat diverges from the saved copy it derives from // (drives the warning dot next to the Save action). const chatDirty = activeChatName != null && savedChats[activeChatName] != null && !savedChatsEqual(savedChats[activeChatName], workingChat); // One funnel for every mutation of the saved-chats store (mirrors the search // profiles' writeStorage): load current → mutate → persist → reflect in state. const writeChats = useCallback( (mut: (s: SavedChatsState) => void) => { const current: SavedChatsState = { v: 1, items: { ...savedChats }, activeName: activeChatName, }; mut(current); saveSavedChats(current); setSavedChats(current.items); setActiveChatName(current.activeName); }, [savedChats, activeChatName], ); const saveChat = useCallback(() => { if (!activeChatName) return; const snap: SavedChat = { ...workingChat, savedAt: Date.now() }; writeChats((s) => { s.items[activeChatName] = snap; }); }, [activeChatName, workingChat, writeChats]); const saveChatAs = useCallback(() => { if (typeof window === "undefined") return; const raw = window.prompt("Save this chat as:"); if (raw == null) return; const name = raw.trim(); if (!name) return; if (Object.prototype.hasOwnProperty.call(savedChats, name)) { window.alert(`A chat named "${name}" already exists.`); return; } const snap: SavedChat = { ...workingChat, savedAt: Date.now() }; writeChats((s) => { s.items[name] = snap; s.activeName = name; }); }, [savedChats, workingChat, writeChats]); const renameChat = useCallback(() => { if (typeof window === "undefined" || !activeChatName) return; const raw = window.prompt("Rename chat:", activeChatName); if (raw == null) return; const next = raw.trim(); if (!next || next === activeChatName) return; if (Object.prototype.hasOwnProperty.call(savedChats, next)) { window.alert(`A chat named "${next}" already exists.`); return; } writeChats((s) => { const copy: Record = {}; for (const [k, v] of Object.entries(s.items)) { copy[k === activeChatName ? next : k] = v; } s.items = copy; s.activeName = next; }); }, [activeChatName, savedChats, writeChats]); const deleteChat = useCallback(() => { if (typeof window === "undefined" || !activeChatName) return; if (!window.confirm(`Delete chat "${activeChatName}"?`)) return; writeChats((s) => { delete s.items[activeChatName]; s.activeName = null; }); }, [activeChatName, writeChats]); // Load a saved chat as the working conversation (replaces it). Resumes under the // CURRENT provider/model — the saved model is info only — so a paused sweep can // be resumed after a model switch. const loadChat = useCallback( (name: string) => { if (busy) return; const chat = savedChats[name]; if (!chat) return; setMessages((chat.messages ?? []).map(normalizeRestoredMessage)); setContextOverride(chat.contextOverride ?? null); setStrictGroundingState(chat.strictGrounding ?? true); setDetached(chat.detached ?? false); setEnrichments(chat.enrichments ?? {}); setReport(chat.report ?? ""); setReportSources(chat.reportSources ?? []); setReportModeState(chat.reportMode ?? false); setSweepCheckpoint(chat.checkpoint ?? null); writeChats((s) => { s.activeName = name; }); }, [busy, savedChats, writeChats], ); return { // provider settings provider, setProvider, apiKey, setApiKey, model, setModel, remember, setRemember, showKey, setShowKey, persistKey, searchMode, setSearchMode, markdownOn, setMarkdownOn, // report mode report, reportSources, reportMode, setReportMode, reportModeAvailable, reportUpdating, // whole-corpus sweep sweep: runAccumulationReport, sweeping, sweepProgress, // configurable batch size sweepBatch, setSweepBatch, // requests-per-minute pacing (current provider+model) rpmOverride, rpmDefault, effectiveRpm, setRpm, // rate-limiter pacing/backoff notice for the progress strip rateLimitNotice, // pausable sweep: a checkpoint (non-null ⇒ paused/resumable) + controls. sweepCheckpoint, sweepPaused: sweepCheckpoint != null && !sweeping, resumeSweep, pauseSweep, discardCheckpoint, // Native-tool-capable → the sweep can run (report mode is available). canSweep: reportModeAvailable, // The active target's size (before any cap) + whether the sweep can't cover // everything. Selection is a curated subset (never "capped"); only a whole- // search sweep is bounded by the engine's true hit cap (`capped`), NOT by // liveGrounding.truncated (which trips merely at >100 videos even though the // sweep covers every discovered video). totalGroundingVideos: activeGrounding?.totalVideos ?? 0, sweepCapped: effectiveGroundingMode === "search" ? capped : false, // ── Grounding palette (mode toggle + selection) ── groundingMode, setGroundingMode, effectiveGroundingMode, selectedCount, selectedSlugs, toggleSelected, clearSelection, // answer length + debug export maxAnswerTokens, setMaxAnswerTokens, buildDebugJson, // data summariesReady, corpusError, channels, // conversation messages, input, setInput, busy, send, stop, reset, retry, editUserMessage, storageWarning, // saved chats library savedChats, activeChatName, chatDirty, saveChat, saveChatAs, renameChat, deleteChat, loadChat, // context editing contextText, contextOverride, applyContext, clearContextOverride, // pinned search results pinned, strictGrounding, setStrictGrounding, detached, detach, reattach, hasLiveSearch, expandPinnedVideo, expanding, expandError, }; } export type AskChatState = ReturnType; export type { SearchStep };