Files
EvoScientist-WebUI/src/app/hooks/useChat.ts
T

1229 lines
45 KiB
TypeScript

"use client";
import { useCallback, useEffect, useRef, useState } from "react";
import { type Message, type Assistant } from "@langchain/langgraph-sdk";
import { v4 as uuidv4 } from "uuid";
import type { TodoItem } from "@/app/types/types";
import { useQueryState } from "nuqs";
import { parseSummarizationEvent } from "@/lib/summarization";
import {
parseContextUsageEvent,
type ContextUsageEvent,
} from "@/lib/contextUsage";
import { findActiveTurnId } from "@/lib/usageTurn";
import { errorToast } from "@/lib/errorReporter";
import type { ModelRef, ThreadModelSelection } from "@/lib/modelRegistry";
import { setThreadModelSelection, setThreadRunning } from "@/app/hooks/useThreads";
import {
clearRunCursor,
getRunStreamCheckpoint,
isRunInProgress,
latestTurnId,
runRequestId,
runTurnId,
selectActiveRun,
setRunStreamCheckpoint,
type RecoverableRun,
} from "@/lib/runRecovery";
import { type ReviewMode } from "@/lib/reviewMode";
import {
StreamMessageAccumulator,
mergeStreamMessages,
type StreamMessageMetadata,
} from "@/lib/streamMessages";
import {
cancelConversationRun,
ConversationApiError,
createConversation,
createConversationRun,
getConversation,
getConversationRun,
joinConversationRunStream,
listConversationRuns,
putConversationFileState,
} from "@/lib/conversationApi";
export type StateType = {
messages: Message[];
todos: TodoItem[];
files: Record<string, string>;
email?: {
id?: string;
subject?: string;
page_content?: string;
};
// Background async sub-agents (writing-agent / data-analysis-agent) this
// conversation launched, keyed by task_id. Shape = deepagents' AsyncTask.
async_tasks?: Record<string, unknown>;
// Private state field set by the deepagents SummarizationMiddleware when the
// conversation is compacted. langgraph dev exposes it over the SDK; the UI
// surfaces it as a collapsible "Conversation compacted" block.
_summarization_event?: unknown;
ui?: any;
};
type InterruptLike = {
value?: unknown;
[key: string]: unknown;
};
/**
* Sanitize a raw interrupt returned by the conversation API before it is
* surfaced to the UI. The live SDK normalizes `stream.interrupt`, but the raw
* persisted task interrupt is unvalidated — if its `value.action_requests`
* (or `review_configs`) is present but NOT an array, ChatInterface's
* `actionRequests.map(...)` / `for (const rc of review_configs)` throws and
* blanks the entire page (the hard crash seen when deleting a file). Require an
* object with an object `value`, and coerce any malformed list field to `[]` so
* the worst case is "no card" instead of a render crash.
*/
function normalizePendingInterrupt(
pending: unknown
): { value: Record<string, unknown> } | undefined {
if (!pending || typeof pending !== "object") return undefined;
const value = (pending as { value?: unknown }).value;
if (!value || typeof value !== "object") return undefined;
const v = value as Record<string, unknown>;
const normalizedValue: Record<string, unknown> = { ...v };
if ("action_requests" in v && !Array.isArray(v.action_requests)) {
normalizedValue.action_requests = [];
}
if ("review_configs" in v && !Array.isArray(v.review_configs)) {
normalizedValue.review_configs = [];
}
// Preserve the interrupt's other fields (id, ns, …); only the value is fixed.
return { ...(pending as object), value: normalizedValue } as {
value: Record<string, unknown>;
};
}
/**
* Total visible text length across a message list. Used to detect when the live
* stream dropped tail CONTENT without dropping the message COUNT — e.g. the
* final assistant turn arrives as an empty/partial AI message (same count) while
* the persisted server snapshot has the full text. A pure length compare misses
* that; comparing total text catches it.
*/
function totalTextLength(msgs: Message[]): number {
let n = 0;
for (const m of msgs) {
const c = (m as { content?: unknown }).content;
if (typeof c === "string") {
n += c.length;
} else if (Array.isArray(c)) {
for (const part of c) {
const t = (part as { text?: unknown })?.text;
if (typeof t === "string") n += t.length;
}
}
}
return n;
}
/**
* A content key for an interrupt, used to tell "the stale interrupt the server
* already resolved" apart from "a genuinely new interrupt". We key on the
* `value` payload because both the live SDK interrupt and the getState-fetched
* one share it (and a fresh object identity each poll can't be compared).
*/
function interruptValueKey(i: unknown): string | null {
if (!i || typeof i !== "object") return null;
try {
return JSON.stringify((i as { value?: unknown }).value ?? null);
} catch {
return null;
}
}
function hasActionableInterrupt(i: unknown): boolean {
if (!i || typeof i !== "object") return false;
const value = (i as { value?: unknown }).value;
if (!value || typeof value !== "object") return false;
const v = value as { type?: unknown; action_requests?: unknown };
return (
v.type === "ask_user" ||
(Array.isArray(v.action_requests) && v.action_requests.length > 0)
);
}
function isModelRefShape(value: unknown): value is ModelRef {
if (!value || typeof value !== "object") return false;
const ref = value as { provider_id?: unknown; model_key?: unknown };
return (
typeof ref.provider_id === "string" && typeof ref.model_key === "string"
);
}
/** Client-side mirror of the thread metadata `model_selection` reader
* (design doc 7.2). Malformed payloads degrade to `inherit` for display; the
* BFF remains authoritative when the run snapshot is created. */
function readSelectionFromMetadata(metadata: Record<string, unknown> | undefined): {
selection: ThreadModelSelection;
revision: number;
} {
const source = metadata ?? {};
const rawRevision = source.model_selection_revision;
const revision =
typeof rawRevision === "number" &&
Number.isInteger(rawRevision) &&
rawRevision >= 0
? rawRevision
: 0;
const raw = source.model_selection;
if (raw === undefined || raw === "inherit") {
return { selection: "inherit", revision };
}
if (raw && typeof raw === "object") {
// Legacy `auxiliary` and flat `temperature`/`top_p` keys are tolerated
// and ignored.
const candidate = raw as {
primary?: unknown;
reasoning_effort?: unknown;
sampling_override?: unknown;
};
if (isModelRefShape(candidate.primary)) {
const effort = candidate.reasoning_effort;
const override = candidate.sampling_override as
| { kind?: unknown; value?: unknown }
| null
| undefined;
return {
selection: {
primary: candidate.primary,
reasoning_effort:
effort === "low" || effort === "medium" || effort === "high"
? effort
: null,
sampling_override:
override &&
(override.kind === "temperature" || override.kind === "top_p") &&
typeof override.value === "number" &&
Number.isFinite(override.value)
? { kind: override.kind, value: override.value }
: null,
},
revision,
};
}
}
return { selection: "inherit", revision };
}
function latestTaskInterrupt(
tasks: Array<{ interrupts?: unknown[] }> | undefined
): unknown {
if (!Array.isArray(tasks)) return undefined;
for (let i = tasks.length - 1; i >= 0; i--) {
const interrupts = tasks[i]?.interrupts;
if (Array.isArray(interrupts) && interrupts.length > 0) {
return interrupts[interrupts.length - 1];
}
}
return undefined;
}
// Build a human-readable summary from the SDK's `onError` payload, which can
// be a plain Error, a StreamError (structured `{ name, error, message }`),
// or a raw string. We try in order: structured `name: message`, plain
// `message`, JSON-of-`.error`, the raw string, finally a generic fallback.
// Capped at 300 chars so a giant stack trace doesn't blow up the toast; the
// full text is still available in the thread JSON via the export affordance.
function formatStreamError(error: unknown): string {
const cap = (s: string) => (s.length > 300 ? s.slice(0, 297) + "..." : s);
if (typeof error === "string" && error.trim()) return cap(error.trim());
if (error && typeof error === "object") {
const e = error as { name?: unknown; message?: unknown; error?: unknown };
const name = typeof e.name === "string" ? e.name.trim() : null;
const msg = typeof e.message === "string" ? e.message.trim() : null;
let inner: string | null = null;
if (typeof e.error === "string" && e.error.trim()) {
inner = e.error.trim();
} else if (e.error && typeof e.error === "object") {
try {
inner = JSON.stringify(e.error);
} catch {
inner = null;
}
}
const body = msg ?? inner;
const combined = name && body ? `${name}: ${body}` : name ?? body ?? "";
if (combined) return cap(combined);
}
return "Run failed.";
}
// Stable empty fallbacks. Returning `?? []` / `?? {}` inline produces a NEW
// array/object reference every render, which needlessly changes the context
// value identity and re-renders every consumer on each store notification.
const EMPTY_MESSAGES: Message[] = [];
const EMPTY_TODOS: TodoItem[] = [];
const EMPTY_FILES: Record<string, string> = {};
const EMPTY_ASYNC_TASKS: Record<string, unknown> = {};
const EMPTY_SUB_AGENT_ACTIVITY: Record<string, never[]> = {};
type TrackedRun = {
threadId: string;
runId: string;
turnId: string | null;
runRequestId: string | null;
};
type RunConnection = {
run: TrackedRun;
phase: "streaming" | "reconnecting";
};
export function useChat({
activeAssistant,
onHistoryRevalidate,
}: {
activeAssistant: Assistant | null;
onHistoryRevalidate?: () => void;
}) {
const [threadId, setThreadId] = useQueryState("threadId");
void activeAssistant;
const [runConnection, setRunConnection] = useState<RunConnection | null>(
null
);
const [isSubmitting, setIsSubmitting] = useState(false);
const [trackingRevision, setTrackingRevision] = useState(0);
const [recoveryRefreshVersion, setRecoveryRefreshVersion] = useState(0);
const lastErrorToastRef = useRef<{ key: string; at: number } | null>(null);
const activeRunRef = useRef<TrackedRun | null>(null);
const currentTurnIdRef = useRef<string | null>(null);
const currentTurnThreadIdRef = useRef<string | null>(null);
const resumeRequestIdsRef = useRef(new Map<string, string>());
const subscriptionControllerRef = useRef<AbortController | null>(null);
const stoppedRunIdsRef = useRef(new Set<string>());
const submittingRef = useRef(false);
const showError = useCallback((error: unknown, key: string) => {
const message = formatStreamError(error);
const now = Date.now();
if (
lastErrorToastRef.current?.key !== `${key}:${message}` ||
now - lastErrorToastRef.current.at > 10_000
) {
lastErrorToastRef.current = { key: `${key}:${message}`, at: now };
errorToast("chat", message);
}
}, []);
const isRunLoading = isSubmitting || runConnection !== null;
const isReconnecting = runConnection?.phase === "reconnecting";
// Zero-latency running indicator for the thread list: this client knows its
// own run lifecycle, no heartbeat wait. The 10s busy heartbeat reconciles
// any drift (e.g. SSE ending before the run actually pauses server-side).
useEffect(() => {
if (!threadId) return;
setThreadRunning(threadId, isRunLoading);
return () => setThreadRunning(threadId, false);
}, [threadId, isRunLoading]);
// Do not attach message modes to the SDK hook's synchronous external store.
// The background subscriber below batches messages-tuple chunks before
// publishing them, while this hook remains responsible for history/interrupts.
const liveInterrupt: InterruptLike | undefined = undefined;
const isThreadLoading = false;
const liveInterruptKey = interruptValueKey(liveInterrupt);
// --- Resilient pending-state fallback ------------------------------------
// The live SSE stream can end (isLoading flips false) BEFORE the run actually
// pauses on a tool-approval interrupt server-side — e.g. the backend's
// auxiliary tool-selector model emits into the stream and desyncs it. When
// that happens, `stream.interrupt` stays empty AND the live message values are
// stale (missing the final `execute` tool-call message), so the approval card
// never renders until a manual thread switch re-fetches history.
//
// Bridge it by reading thread state directly once the stream settles: while
// the run is still pending (`next` non-empty) but no live interrupt is shown,
// poll a BOUNDED number of times until the interrupt is persisted, then
// surface BOTH the interrupt and that snapshot's messages. Stops as soon as
// the interrupt is found or the run is truly done (`next` empty); a new run
// (isLoading→true) clears it. Not an unbounded poll — that would race the
// live stream and revive resolved interrupts.
const [fetchedInterrupt, setFetchedInterrupt] = useState<
InterruptLike | undefined
>(undefined);
// Content key of an interrupt the server confirmed RESOLVED. The getter
// suppresses a stale `stream.interrupt` (e.g. one re-surfaced from SDK history
// after approving) ONLY when it matches this key — so a genuinely new
// interrupt is never hidden (the old global-null sentinel hid everything).
const [resolvedInterruptKey, setResolvedInterruptKey] = useState<
string | null
>(null);
const [fetchedMessages, setFetchedMessages] = useState<Message[] | null>(
null
);
const [fetchedValues, setFetchedValues] = useState<Partial<StateType> | null>(
null
);
const [fetchedThreadId, setFetchedThreadId] = useState<string | null>(null);
// Latest context-occupancy event from the live run's custom stream, keyed
// by thread so a stale value never leaks across a thread switch.
const [contextUsageEntry, setContextUsageEntry] = useState<{
threadId: string;
usage: ContextUsageEvent;
} | null>(null);
const fetchedThreadIdRef = useRef<string | null>(null);
fetchedThreadIdRef.current = fetchedThreadId;
const recoveryRunRef = useRef(0);
// Per-thread model selection (design doc 7.2). The selection lives on the
// thread metadata (`model_selection` + `model_selection_revision`) and is
// written exclusively through the PATCH CAS endpoint; each new run freezes
// it into a runtime snapshot server-side.
//
// Fresh-chat wrinkle: the thread row doesn't exist server-side until the
// first message creates it, so a pre-thread pick is stashed in
// `pendingSelectionRef` and written through with
// `expected_selection_revision=0` once the thread id appears.
const [modelSelection, setModelSelectionState] =
useState<ThreadModelSelection>("inherit");
const [selectionRevision, setSelectionRevision] = useState(0);
const pendingSelectionRef = useRef<ThreadModelSelection | null>(null);
const createdThreadIdRef = useRef<string | null>(null);
const ensureThreadPromiseRef = useRef<Promise<string> | null>(null);
useEffect(() => {
if (!threadId) {
createdThreadIdRef.current = null;
setModelSelectionState(pendingSelectionRef.current ?? "inherit");
setSelectionRevision(0);
return;
}
// Thread just came into existence (or we switched onto an existing one).
// If we have a pending pre-thread pick, write it through to metadata with
// the initial revision; otherwise fetch the thread's persisted selection
// and seed local state from it.
if (
createdThreadIdRef.current === threadId &&
pendingSelectionRef.current
) {
const pending = pendingSelectionRef.current;
pendingSelectionRef.current = null;
createdThreadIdRef.current = null;
void (async () => {
try {
await setThreadModelSelection(threadId, pending, 0);
setModelSelectionState(pending);
setSelectionRevision(1);
} catch {
// The local state still reflects the pick; the next
// `setModelSelection` call (or thread reopen) resyncs from metadata.
}
})();
return;
}
// Opening an existing thread from the New Chat screen must not write a
// staged pick over that thread's own selection.
pendingSelectionRef.current = null;
createdThreadIdRef.current = null;
let cancelled = false;
void (async () => {
try {
const { thread: t } = await getConversation(threadId);
if (cancelled) return;
const stored = readSelectionFromMetadata(
(t.metadata ?? {}) as Record<string, unknown>
);
setModelSelectionState(stored.selection);
setSelectionRevision(stored.revision);
} catch (error) {
if (cancelled) return;
setModelSelectionState("inherit");
setSelectionRevision(0);
if (error instanceof ConversationApiError && error.status === 404) {
activeRunRef.current = null;
setRunConnection(null);
void setThreadId(null);
errorToast(
"chat.conversation",
"Conversation is no longer available. Started a new chat."
);
}
}
})();
return () => {
cancelled = true;
};
}, [setThreadId, threadId]);
// Persist + apply locally. Pre-thread (new chat with no threadId yet),
// stashes the pick in a ref so the thread-id effect can persist it as soon
// as the row is created server-side. Existing threads go through the CAS
// endpoint; a revision conflict refetches the current selection so the UI
// never diverges from the stored one.
const setModelSelection = useCallback(
async (next: ThreadModelSelection) => {
if (!threadId) {
pendingSelectionRef.current = next === "inherit" ? null : next;
setModelSelectionState(next);
setSelectionRevision(0);
return;
}
pendingSelectionRef.current = null;
try {
await setThreadModelSelection(threadId, next, selectionRevision);
setModelSelectionState(next);
setSelectionRevision((revision) => revision + 1);
} catch (error) {
if (
error instanceof ConversationApiError &&
error.code === "THREAD_MODEL_SELECTION_CONFLICT"
) {
try {
const { thread: t } = await getConversation(threadId);
const stored = readSelectionFromMetadata(
(t.metadata ?? {}) as Record<string, unknown>
);
setModelSelectionState(stored.selection);
setSelectionRevision(stored.revision);
} catch {
// The next thread load resyncs.
}
}
throw error;
}
},
[threadId, selectionRevision]
);
const liveMessages =
fetchedThreadId === threadId
? fetchedMessages ?? EMPTY_MESSAGES
: EMPTY_MESSAGES;
const liveMessageCount = liveMessages.length;
useEffect(() => {
if (!threadId) {
setFetchedInterrupt(undefined);
setFetchedMessages(null);
setFetchedThreadId(null);
setFetchedValues(null);
setResolvedInterruptKey(null);
return;
}
if (isRunLoading) {
recoveryRunRef.current += 1;
setFetchedInterrupt(undefined);
setResolvedInterruptKey(null);
return;
}
// The persisted record replaces the temporary chunk-assembled messages once
// the run settles, and remains the source of truth for history/interrupts.
const baseline = liveMessageCount;
const recoveryRunId = ++recoveryRunRef.current;
let cancelled = false;
let tries = 0;
let timer: ReturnType<typeof setTimeout> | undefined;
const MAX_TRIES = 15;
const attempt = async () => {
tries += 1;
try {
// `getState` returns the GRAPH CHECKPOINT state — which the backend
// windows/compacts for memory, so its `values.messages` is only the
// recent slice. `threads.get` returns the persisted THREAD RECORD with
// the full message history. We need both: state for run status
// (`next` / `tasks` / `interrupts`), record for the messages the UI
// displays. Done in parallel to keep the round trip tight.
const { state, thread: threadRecord } = await getConversation(threadId);
if (cancelled || recoveryRunRef.current !== recoveryRunId) return;
const values = threadRecord.values as Partial<StateType> | undefined;
const msgs = values?.messages;
const pending = latestTaskInterrupt(state.tasks);
const stillPending = Array.isArray(state.next) && state.next.length > 0;
const safePending = normalizePendingInterrupt(pending);
if (safePending && hasActionableInterrupt(safePending)) {
// Tool-approval interrupt reached — surface it and its matching message
// snapshot together. Mixing live messages with fetched interrupts is the
// race that hides approval cards for repeated execute calls.
setFetchedInterrupt(safePending as InterruptLike);
setResolvedInterruptKey(null);
if (Array.isArray(msgs)) {
setFetchedThreadId(threadId);
setFetchedMessages(msgs);
}
setFetchedValues(values ?? null);
return;
}
// Preserve the existing comparison so a running stream never overwrites a
// local optimistic message with an older persisted snapshot.
if (Array.isArray(msgs) && msgs.length > baseline) {
setFetchedThreadId(threadId);
setFetchedMessages(msgs);
setFetchedValues(values ?? null);
}
if (!stillPending) {
// The server has no pending task/interrupt anymore. Record the stale
// live interrupt's identity so the getter suppresses ONLY that one
// (composer unlocks after approving) — a new interrupt still shows.
setFetchedInterrupt(undefined);
setResolvedInterruptKey(liveInterruptKey);
if (Array.isArray(msgs)) {
setFetchedThreadId(threadId);
setFetchedMessages(msgs);
}
setFetchedValues(values ?? null);
return;
}
// Keep polling only while the run is still working server-side; a
// finished run (next empty) won't produce anything more.
if (stillPending && tries < MAX_TRIES && !cancelled) {
timer = setTimeout(attempt, 1000);
}
} catch {
if (
!cancelled &&
recoveryRunRef.current === recoveryRunId &&
tries < MAX_TRIES
) {
timer = setTimeout(attempt, 1000);
}
}
};
void attempt();
return () => {
cancelled = true;
if (timer) clearTimeout(timer);
};
// State is read from the BFF snapshot, so this effect never exposes a direct
// LangGraph client or deployment credential to the browser.
}, [
threadId,
liveInterruptKey,
isRunLoading,
liveMessageCount,
recoveryRefreshVersion,
]);
// Show the live interrupt unless it's the exact one the server told us was
// resolved (then fall through to the fetched one, usually undefined → composer
// unlocks). A new live interrupt has a different key, so it's never suppressed.
const interrupt =
!isRunLoading &&
liveInterrupt &&
(resolvedInterruptKey === null ||
interruptValueKey(liveInterrupt) !== resolvedInterruptKey)
? liveInterrupt
: fetchedThreadId === threadId
? fetchedInterrupt ?? undefined
: undefined;
// Prefer the backfilled snapshot when it is "ahead" of the live stream — i.e.
// the stream ended early and dropped the tail. "Ahead" means either MORE
// messages, or (once settled) the SAME number of messages but MORE total text:
// the final assistant turn often arrives as an empty/partial AI message with
// the right count but no content, so a pure length compare would keep showing
// the blank live version (the bug where the answer only appears after a manual
// refresh). The equal-count/more-text rule is gated on `!isLoading` so a
// mid-stream poll snapshot never flickers over the actively updating stream.
//
// Once the run has settled AND we have a snapshot, ALWAYS prefer the snapshot.
// Live values can carry transient stream-only messages and can lag the
// persisted main-thread state. Once settled, prefer the fetched thread record
// so the UI shows the stable server-side conversation.
const messages = (() => {
if (!fetchedMessages || fetchedThreadId !== threadId) return liveMessages;
if (fetchedInterrupt) return fetchedMessages;
if (!isRunLoading) return fetchedMessages;
if (fetchedMessages.length > liveMessages.length) return fetchedMessages;
if (
fetchedMessages.length === liveMessages.length &&
totalTextLength(fetchedMessages) > totalTextLength(liveMessages)
) {
return fetchedMessages;
}
return liveMessages;
})();
const messagesRef = useRef(messages);
messagesRef.current = messages;
const threadIdRef = useRef(threadId);
threadIdRef.current = threadId;
const onHistoryRevalidateRef = useRef(onHistoryRevalidate);
onHistoryRevalidateRef.current = onHistoryRevalidate;
const ensureThread = useCallback(async (): Promise<string> => {
const existingThreadId = threadIdRef.current;
if (existingThreadId) return existingThreadId;
if (ensureThreadPromiseRef.current) return ensureThreadPromiseRef.current;
const creation = (async () => {
const created = await createConversation();
const createdThreadId = created.threadId;
createdThreadIdRef.current = createdThreadId;
threadIdRef.current = createdThreadId;
setFetchedThreadId(createdThreadId);
// Flush a staged pre-thread pick before returning: the caller creates
// the first run immediately after, and the run snapshot freezes
// whatever selection the thread metadata holds at that moment.
const pending = pendingSelectionRef.current;
if (pending) {
pendingSelectionRef.current = null;
createdThreadIdRef.current = null;
try {
await setThreadModelSelection(createdThreadId, pending, 0);
setModelSelectionState(pending);
setSelectionRevision(1);
} catch {
// The thread-id effect resyncs from metadata.
}
}
await setThreadId(createdThreadId);
return createdThreadId;
})();
ensureThreadPromiseRef.current = creation;
try {
return await creation;
} finally {
if (ensureThreadPromiseRef.current === creation) {
ensureThreadPromiseRef.current = null;
}
}
}, [setThreadId]);
// Discover active runs from the backend, then subscribe independently. The
// run survives because it was created through `runs.create`; aborting this
// GET stream only drops the subscription. The checkpoint stores both the SSE
// cursor and assembled partial messages so a reload can resume without gaps.
useEffect(() => {
if (!threadId) {
subscriptionControllerRef.current?.abort();
subscriptionControllerRef.current = null;
activeRunRef.current = null;
setRunConnection(null);
return;
}
if (activeRunRef.current && activeRunRef.current.threadId !== threadId) {
activeRunRef.current = null;
setRunConnection(null);
}
if (currentTurnThreadIdRef.current !== threadId) {
currentTurnIdRef.current = null;
currentTurnThreadIdRef.current = threadId;
resumeRequestIdsRef.current.clear();
}
let disposed = false;
let wakeRetry: (() => void) | null = null;
let retryTimer: ReturnType<typeof setTimeout> | undefined;
let failedAttempts = 0;
let knownRun =
activeRunRef.current?.threadId === threadId ? activeRunRef.current : null;
const waitForRetry = (delay: number) =>
new Promise<void>((resolve) => {
let settled = false;
const finish = () => {
if (settled) return;
settled = true;
if (retryTimer) clearTimeout(retryTimer);
retryTimer = undefined;
wakeRetry = null;
resolve();
};
wakeRetry = finish;
retryTimer = setTimeout(finish, delay);
});
const rememberRun = (
run: RecoverableRun,
phase: RunConnection["phase"]
) => {
const tracked: TrackedRun = {
threadId,
runId: run.run_id,
turnId: runTurnId(run) ?? currentTurnIdRef.current,
runRequestId: runRequestId(run),
};
knownRun = tracked;
activeRunRef.current = tracked;
if (tracked.turnId) {
currentTurnIdRef.current = tracked.turnId;
currentTurnThreadIdRef.current = threadId;
}
setRunConnection({ run: tracked, phase });
return tracked;
};
const clearTrackedRun = (runId: string) => {
if (activeRunRef.current?.runId === runId) {
activeRunRef.current = null;
}
setRunConnection((current) =>
current?.run.runId === runId ? null : current
);
clearRunCursor(threadId, runId);
setRecoveryRefreshVersion((version) => version + 1);
onHistoryRevalidateRef.current?.();
};
const listRuns = async () =>
(await listConversationRuns(threadId)) as RecoverableRun[];
const publishStreamMessages = (streamed: Message[]) => {
if (streamed.length === 0 || disposed) return;
const sameThread = fetchedThreadIdRef.current === threadId;
fetchedThreadIdRef.current = threadId;
setFetchedThreadId(threadId);
setFetchedMessages((current) =>
mergeStreamMessages(
sameThread ? current ?? EMPTY_MESSAGES : EMPTY_MESSAGES,
streamed
)
);
};
const seedPersistedThread = async () => {
try {
const { thread: threadRecord } = await getConversation(threadId);
if (disposed) return;
const values = threadRecord.values as Partial<StateType> | undefined;
const persisted = values?.messages;
if (Array.isArray(persisted)) {
const sameThread = fetchedThreadIdRef.current === threadId;
fetchedThreadIdRef.current = threadId;
setFetchedThreadId(threadId);
setFetchedMessages((current) =>
mergeStreamMessages(
persisted,
sameThread ? current ?? EMPTY_MESSAGES : EMPTY_MESSAGES
)
);
}
setFetchedValues(values ?? null);
} catch {
// The resumable SSE stream can still reconstruct the active response.
}
};
const followRuns = async () => {
while (!disposed) {
let runs: RecoverableRun[];
try {
runs = await listRuns();
if (disposed) return;
} catch {
if (knownRun) {
setRunConnection({ run: knownRun, phase: "reconnecting" });
}
failedAttempts += 1;
await waitForRetry(
Math.min(1_000 * 2 ** Math.min(failedAttempts, 3), 10_000)
);
continue;
}
const discoveredTurnId = latestTurnId(runs);
if (discoveredTurnId) {
currentTurnIdRef.current = discoveredTurnId;
currentTurnThreadIdRef.current = threadId;
}
for (const run of runs) {
if (!isRunInProgress(run.status)) {
clearRunCursor(threadId, run.run_id);
}
}
let candidate = selectActiveRun(runs);
if (!candidate && knownRun) {
try {
const current = await getConversationRun(threadId, knownRun.runId);
if (isRunInProgress(current.status)) {
candidate = current as RecoverableRun;
}
} catch {
// The next list refresh is authoritative.
}
}
if (!candidate) {
if (knownRun) clearTrackedRun(knownRun.runId);
return;
}
if (stoppedRunIdsRef.current.has(candidate.run_id)) return;
const tracked = rememberRun(candidate, "reconnecting");
await seedPersistedThread();
if (disposed) return;
const controller = new AbortController();
subscriptionControllerRef.current?.abort();
subscriptionControllerRef.current = controller;
let streamFailure: unknown;
const checkpoint = getRunStreamCheckpoint(threadId, tracked.runId);
const messageAccumulator = new StreamMessageAccumulator(
checkpoint?.messages
);
if (checkpoint?.messages.length) {
publishStreamMessages(checkpoint.messages);
}
let latestEventId = checkpoint?.lastEventId ?? "-1";
let streamDirty = false;
let frameId: number | undefined;
let flushTimer: ReturnType<typeof setTimeout> | undefined;
let checkpointTimer: ReturnType<typeof setTimeout> | undefined;
let lastCheckpointAt = 0;
const persistStreamCheckpoint = () => {
if (latestEventId === "-1") return;
lastCheckpointAt = Date.now();
setRunStreamCheckpoint(threadId, tracked.runId, {
lastEventId: latestEventId,
messages: messageAccumulator.snapshot(),
});
};
const flushStreamMessages = (forceCheckpoint = false) => {
if (frameId !== undefined) cancelAnimationFrame(frameId);
if (flushTimer !== undefined) clearTimeout(flushTimer);
frameId = undefined;
flushTimer = undefined;
if (streamDirty) {
streamDirty = false;
publishStreamMessages(messageAccumulator.snapshot());
}
if (checkpointTimer !== undefined && forceCheckpoint) {
clearTimeout(checkpointTimer);
checkpointTimer = undefined;
}
const checkpointDelay = 500 - (Date.now() - lastCheckpointAt);
if (forceCheckpoint || checkpointDelay <= 0) {
persistStreamCheckpoint();
} else if (checkpointTimer === undefined) {
checkpointTimer = setTimeout(() => {
checkpointTimer = undefined;
persistStreamCheckpoint();
}, checkpointDelay);
}
};
const scheduleStreamFlush = () => {
streamDirty = true;
frameId ??= requestAnimationFrame(() => flushStreamMessages());
flushTimer ??= setTimeout(flushStreamMessages, 100);
};
try {
setRunConnection({ run: tracked, phase: "streaming" });
for await (const event of joinConversationRunStream(
threadId,
tracked.runId,
{
signal: controller.signal,
lastEventId: latestEventId,
}
)) {
if (event.id) latestEventId = event.id;
if (
event.event === "messages" &&
Array.isArray(event.data) &&
event.data.length >= 2
) {
const metadata =
event.data[1] && typeof event.data[1] === "object"
? (event.data[1] as StreamMessageMetadata)
: undefined;
messageAccumulator.add(event.data[0], metadata);
}
if (event.event === "custom") {
const usage = parseContextUsageEvent(event.data);
if (usage) setContextUsageEntry({ threadId: tracked.threadId, usage });
}
scheduleStreamFlush();
if (event.event === "error") streamFailure = event.data;
}
} catch (error) {
if (!(error instanceof Error && error.name === "AbortError")) {
streamFailure = error;
}
} finally {
flushStreamMessages(true);
if (subscriptionControllerRef.current === controller) {
subscriptionControllerRef.current = null;
}
}
if (disposed || stoppedRunIdsRef.current.has(tracked.runId)) return;
try {
const completed = await getConversationRun(threadId, tracked.runId);
if (isRunInProgress(completed.status)) {
knownRun = tracked;
setRunConnection({ run: tracked, phase: "reconnecting" });
failedAttempts += 1;
await waitForRetry(
Math.min(1_000 * 2 ** Math.min(failedAttempts, 3), 10_000)
);
continue;
}
clearTrackedRun(tracked.runId);
knownRun = null;
failedAttempts = 0;
if (completed.status === "error" || completed.status === "timeout") {
showError(
streamFailure ?? `Run ended with status ${completed.status}.`,
tracked.runId
);
}
// A queued run may already be waiting behind the one that completed.
continue;
} catch {
knownRun = tracked;
setRunConnection({ run: tracked, phase: "reconnecting" });
failedAttempts += 1;
await waitForRetry(
Math.min(1_000 * 2 ** Math.min(failedAttempts, 3), 10_000)
);
}
}
};
const recoverWhenOnline = () => {
failedAttempts = 0;
wakeRetry?.();
};
window.addEventListener("online", recoverWhenOnline);
void followRuns();
return () => {
disposed = true;
subscriptionControllerRef.current?.abort();
subscriptionControllerRef.current = null;
if (retryTimer) clearTimeout(retryTimer);
wakeRetry?.();
window.removeEventListener("online", recoverWhenOnline);
};
}, [showError, threadId, trackingRevision]);
// A closed tab recovers on mount via run discovery. An already-open tab also
// rechecks when it regains focus, covering runs started from another tab.
useEffect(() => {
if (!threadId) return;
const rediscover = () => setTrackingRevision((version) => version + 1);
const onVisibilityChange = () => {
if (document.visibilityState === "visible") rediscover();
};
window.addEventListener("focus", rediscover);
document.addEventListener("visibilitychange", onVisibilityChange);
return () => {
window.removeEventListener("focus", rediscover);
document.removeEventListener("visibilitychange", onVisibilityChange);
};
}, [threadId]);
const startBackgroundRun = useCallback(
async ({
input,
command,
turnId,
runRequestId: requestId,
interruptKey,
reviewMode,
optimisticMessageId,
}: {
input?: Record<string, unknown> | null;
command?: { resume: unknown };
turnId: string;
runRequestId: string;
interruptKey?: string | null;
reviewMode: ReviewMode;
optimisticMessageId?: string;
}) => {
if (submittingRef.current || activeRunRef.current) return;
submittingRef.current = true;
setIsSubmitting(true);
let usableThreadId = threadIdRef.current;
try {
if (!usableThreadId) {
usableThreadId = await ensureThread();
}
currentTurnIdRef.current = turnId;
currentTurnThreadIdRef.current = usableThreadId;
const createdRun = await createConversationRun(usableThreadId, {
turn_id: turnId,
run_request_id: requestId,
...(interruptKey ? { interrupt_key: interruptKey } : {}),
input,
command,
review_mode: reviewMode,
});
const tracked: TrackedRun = {
threadId: usableThreadId,
runId: createdRun.runId,
turnId,
runRequestId: createdRun.runRequestId,
};
activeRunRef.current = tracked;
setRunConnection({ run: tracked, phase: "reconnecting" });
setTrackingRevision((version) => version + 1);
onHistoryRevalidateRef.current?.();
} catch (error) {
// If the POST response was lost after the backend created the run, find
// it by this concrete request id. A logical turn can legitimately have
// several runs while approvals are being resolved.
let recovered = false;
if (usableThreadId) {
try {
const runs = (await listConversationRuns(
usableThreadId
)) as RecoverableRun[];
const matching = runs.find(
(run) => runRequestId(run) === requestId
);
if (matching) {
recovered = true;
if (isRunInProgress(matching.status)) {
const tracked: TrackedRun = {
threadId: usableThreadId,
runId: matching.run_id,
turnId,
runRequestId: requestId,
};
activeRunRef.current = tracked;
setRunConnection({ run: tracked, phase: "reconnecting" });
setTrackingRevision((version) => version + 1);
} else {
setRecoveryRefreshVersion((version) => version + 1);
}
}
} catch {
// Surface the original creation failure below.
}
}
if (!recovered) {
if (optimisticMessageId) {
setFetchedMessages(
(current) =>
current?.filter(
(message) => message.id !== optimisticMessageId
) ?? null
);
}
showError(error, `create:${requestId}`);
}
onHistoryRevalidateRef.current?.();
} finally {
submittingRef.current = false;
setIsSubmitting(false);
}
},
[ensureThread, showError]
);
const sendMessage = useCallback(
(content: string, reviewMode: ReviewMode) => {
// Keep the persisted snapshot visible and append the user message locally;
// the next settled thread-record refresh replaces this optimistic snapshot.
setFetchedInterrupt(undefined);
setFetchedThreadId(threadIdRef.current);
setResolvedInterruptKey(null);
recoveryRunRef.current += 1;
const newMessage: Message = { id: uuidv4(), type: "human", content };
setFetchedMessages((current) => [
...(current ?? EMPTY_MESSAGES),
newMessage,
]);
onHistoryRevalidateRef.current?.();
void startBackgroundRun({
input: { messages: [newMessage] },
turnId: newMessage.id!,
runRequestId: uuidv4(),
reviewMode,
optimisticMessageId: newMessage.id,
});
},
[startBackgroundRun]
);
const setFiles = useCallback(
async (files: Record<string, string>) => {
if (!threadId) return;
await putConversationFileState(threadId, files);
},
[threadId]
);
const resumeInterrupt = useCallback(
(value: any, reviewMode: ReviewMode, interruptKey?: string | null) => {
// A resume is a new backend run but remains part of the same logical turn.
setFetchedInterrupt(undefined);
setResolvedInterruptKey(null);
recoveryRunRef.current += 1;
const turnId =
(currentTurnThreadIdRef.current === threadIdRef.current
? currentTurnIdRef.current
: null) ??
findActiveTurnId(messagesRef.current) ??
uuidv4();
const requestId =
interruptKey && resumeRequestIdsRef.current.has(interruptKey)
? resumeRequestIdsRef.current.get(interruptKey)!
: uuidv4();
if (interruptKey) {
resumeRequestIdsRef.current.set(interruptKey, requestId);
}
onHistoryRevalidateRef.current?.();
void startBackgroundRun({
input: null,
command: { resume: value },
turnId,
runRequestId: requestId,
interruptKey,
reviewMode,
});
},
[startBackgroundRun]
);
const stopStream = useCallback(() => {
const active = activeRunRef.current;
if (!active) return;
stoppedRunIdsRef.current.add(active.runId);
subscriptionControllerRef.current?.abort();
subscriptionControllerRef.current = null;
activeRunRef.current = null;
setRunConnection(null);
clearRunCursor(active.threadId, active.runId);
setRecoveryRefreshVersion((version) => version + 1);
void cancelConversationRun(active.threadId, active.runId)
.then(() => stoppedRunIdsRef.current.delete(active.runId))
.catch((error) => {
stoppedRunIdsRef.current.delete(active.runId);
showError(error, `cancel:${active.runId}`);
setTrackingRevision((version) => version + 1);
})
.finally(() => onHistoryRevalidateRef.current?.());
}, [showError]);
return {
todos: fetchedValues?.todos ?? EMPTY_TODOS,
files: fetchedValues?.files ?? EMPTY_FILES,
email: fetchedValues?.email,
asyncTasks: fetchedValues?.async_tasks ?? EMPTY_ASYNC_TASKS,
summarizationEvent: parseSummarizationEvent(
fetchedValues?._summarization_event
),
ui: fetchedValues?.ui,
contextUsage:
contextUsageEntry && contextUsageEntry.threadId === threadId
? contextUsageEntry.usage
: null,
setFiles,
messages,
isLoading: isRunLoading,
isReconnecting,
isThreadLoading,
interrupt,
ensureThread,
sendMessage,
stopStream,
resumeInterrupt,
subAgentActivity: EMPTY_SUB_AGENT_ACTIVITY,
modelSelection,
setModelSelection,
};
}