From 14b246f4e7afb7740d391cdef4a40e88ee14b1e2 Mon Sep 17 00:00:00 2001 From: m4 Date: Fri, 10 Jul 2026 16:26:46 +0800 Subject: [PATCH] fix: stabilize chat stream updates --- package-lock.json | 90 --------- src/app/api/models/route.ts | 57 ++++++ src/app/components/ActionGroup.tsx | 2 +- src/app/components/ChatInterface.tsx | 82 ++++++-- src/app/components/IdentityTab.tsx | 5 +- src/app/components/ObservationGraph.tsx | 11 +- src/app/hooks/useAvailableModels.ts | 29 +-- src/app/hooks/useChat.ts | 244 ++++++++++++------------ src/app/page.tsx | 5 +- src/lib/uiSettings.ts | 5 +- src/providers/ThemeProvider.tsx | 6 +- 11 files changed, 285 insertions(+), 251 deletions(-) create mode 100644 src/app/api/models/route.ts diff --git a/package-lock.json b/package-lock.json index 1e60834..84b8a56 100644 --- a/package-lock.json +++ b/package-lock.json @@ -765,9 +765,6 @@ "cpu": [ "arm" ], - "libc": [ - "glibc" - ], "license": "LGPL-3.0-or-later", "optional": true, "os": [ @@ -784,9 +781,6 @@ "cpu": [ "arm64" ], - "libc": [ - "glibc" - ], "license": "LGPL-3.0-or-later", "optional": true, "os": [ @@ -803,9 +797,6 @@ "cpu": [ "ppc64" ], - "libc": [ - "glibc" - ], "license": "LGPL-3.0-or-later", "optional": true, "os": [ @@ -822,9 +813,6 @@ "cpu": [ "riscv64" ], - "libc": [ - "glibc" - ], "license": "LGPL-3.0-or-later", "optional": true, "os": [ @@ -841,9 +829,6 @@ "cpu": [ "s390x" ], - "libc": [ - "glibc" - ], "license": "LGPL-3.0-or-later", "optional": true, "os": [ @@ -860,9 +845,6 @@ "cpu": [ "x64" ], - "libc": [ - "glibc" - ], "license": "LGPL-3.0-or-later", "optional": true, "os": [ @@ -879,9 +861,6 @@ "cpu": [ "arm64" ], - "libc": [ - "musl" - ], "license": "LGPL-3.0-or-later", "optional": true, "os": [ @@ -898,9 +877,6 @@ "cpu": [ "x64" ], - "libc": [ - "musl" - ], "license": "LGPL-3.0-or-later", "optional": true, "os": [ @@ -917,9 +893,6 @@ "cpu": [ "arm" ], - "libc": [ - "glibc" - ], "license": "Apache-2.0", "optional": true, "os": [ @@ -942,9 +915,6 @@ "cpu": [ "arm64" ], - "libc": [ - "glibc" - ], "license": "Apache-2.0", "optional": true, "os": [ @@ -967,9 +937,6 @@ "cpu": [ "ppc64" ], - "libc": [ - "glibc" - ], "license": "Apache-2.0", "optional": true, "os": [ @@ -992,9 +959,6 @@ "cpu": [ "riscv64" ], - "libc": [ - "glibc" - ], "license": "Apache-2.0", "optional": true, "os": [ @@ -1017,9 +981,6 @@ "cpu": [ "s390x" ], - "libc": [ - "glibc" - ], "license": "Apache-2.0", "optional": true, "os": [ @@ -1042,9 +1003,6 @@ "cpu": [ "x64" ], - "libc": [ - "glibc" - ], "license": "Apache-2.0", "optional": true, "os": [ @@ -1067,9 +1025,6 @@ "cpu": [ "arm64" ], - "libc": [ - "musl" - ], "license": "Apache-2.0", "optional": true, "os": [ @@ -1092,9 +1047,6 @@ "cpu": [ "x64" ], - "libc": [ - "musl" - ], "license": "Apache-2.0", "optional": true, "os": [ @@ -1393,9 +1345,6 @@ "cpu": [ "arm64" ], - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -1412,9 +1361,6 @@ "cpu": [ "arm64" ], - "libc": [ - "musl" - ], "license": "MIT", "optional": true, "os": [ @@ -1431,9 +1377,6 @@ "cpu": [ "x64" ], - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -1450,9 +1393,6 @@ "cpu": [ "x64" ], - "libc": [ - "musl" - ], "license": "MIT", "optional": true, "os": [ @@ -3097,9 +3037,6 @@ "arm64" ], "dev": true, - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -3114,9 +3051,6 @@ "arm64" ], "dev": true, - "libc": [ - "musl" - ], "license": "MIT", "optional": true, "os": [ @@ -3131,9 +3065,6 @@ "loong64" ], "dev": true, - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -3148,9 +3079,6 @@ "loong64" ], "dev": true, - "libc": [ - "musl" - ], "license": "MIT", "optional": true, "os": [ @@ -3165,9 +3093,6 @@ "ppc64" ], "dev": true, - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -3182,9 +3107,6 @@ "riscv64" ], "dev": true, - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -3199,9 +3121,6 @@ "riscv64" ], "dev": true, - "libc": [ - "musl" - ], "license": "MIT", "optional": true, "os": [ @@ -3216,9 +3135,6 @@ "s390x" ], "dev": true, - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -3233,9 +3149,6 @@ "x64" ], "dev": true, - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -3250,9 +3163,6 @@ "x64" ], "dev": true, - "libc": [ - "musl" - ], "license": "MIT", "optional": true, "os": [ diff --git a/src/app/api/models/route.ts b/src/app/api/models/route.ts new file mode 100644 index 0000000..0822784 --- /dev/null +++ b/src/app/api/models/route.ts @@ -0,0 +1,57 @@ +import { NextRequest, NextResponse } from "next/server"; +import { isCrossOrigin } from "@/lib/server/workspace"; + +export const runtime = "nodejs"; + +function fail(error: unknown, status = 400) { + return NextResponse.json( + { + error: + error instanceof Error + ? error.message + : "Model registry request failed.", + }, + { status } + ); +} + +function resolveDeploymentUrl(request: NextRequest): URL { + const raw = request.nextUrl.searchParams.get("deploymentUrl"); + if (!raw?.trim()) throw new Error("A deployment URL is required."); + const url = new URL(raw); + if (url.protocol !== "http:" && url.protocol !== "https:") { + throw new Error("Deployment URL must use http or https."); + } + url.search = ""; + url.hash = ""; + return url; +} + +export async function GET(request: NextRequest) { + try { + if (isCrossOrigin(request)) { + return fail("Cross-origin model registry access is not allowed.", 403); + } + const deploymentUrl = resolveDeploymentUrl(request) + .toString() + .replace(/\/$/, ""); + const headers: Record = {}; + const apiKey = request.headers.get("x-api-key"); + if (apiKey) headers["X-Api-Key"] = apiKey; + + const upstream = await fetch(`${deploymentUrl}/api/models`, { + headers, + cache: "no-store", + }); + const body = await upstream.text(); + return new NextResponse(body, { + status: upstream.status, + headers: { + "content-type": + upstream.headers.get("content-type") ?? "application/json", + }, + }); + } catch (error) { + return fail(error); + } +} diff --git a/src/app/components/ActionGroup.tsx b/src/app/components/ActionGroup.tsx index e6fa2ce..5890353 100644 --- a/src/app/components/ActionGroup.tsx +++ b/src/app/components/ActionGroup.tsx @@ -34,7 +34,7 @@ interface ActionGroupProps { submittedActionRequestKeys: Set; onActionRequestSubmitted: (key: string) => void; reviewConfigsMap: Map | null; - stream: unknown; + stream?: unknown; onResumeInterrupt: (value: unknown) => void; graphId?: string; onEditMessage: (content: string) => void; diff --git a/src/app/components/ChatInterface.tsx b/src/app/components/ChatInterface.tsx index 845fa44..0ca7628 100644 --- a/src/app/components/ChatInterface.tsx +++ b/src/app/components/ChatInterface.tsx @@ -230,6 +230,41 @@ function getMessageToolCalls(message: Message): Array<{ }); } +function stableStringify(value: unknown): string { + const seen = new WeakSet(); + try { + return JSON.stringify(value, (_key, v) => { + if (!v || typeof v !== "object") return v; + if (seen.has(v)) return "[Circular]"; + seen.add(v); + if (Array.isArray(v)) return v; + return Object.keys(v) + .sort() + .reduce>((out, key) => { + out[key] = (v as Record)[key]; + return out; + }, {}); + }); + } catch { + return String(value); + } +} + +function autoApproveInterruptKey( + interrupt: unknown, + actionRequests: unknown[] +): string { + const ir = + interrupt && typeof interrupt === "object" + ? (interrupt as { ns?: unknown; scope?: unknown; value?: unknown }) + : {}; + return stableStringify({ + ns: ir.ns, + scope: ir.scope, + action_requests: actionRequests, + }); +} + const getStatusIcon = (status: TodoItem["status"], className?: string) => { switch (status) { case "completed": @@ -384,14 +419,14 @@ export const ChatInterface = React.memo( ); }); }, [pickerModels, modelSearch]); - const autoApprovedRef = useRef(null); + const autoApprovedRef = useRef(null); + const autoApproveResumePendingRef = useRef(false); const previousThreadIdRef = useRef(threadId); const migrateAutoApproveForCreatedThreadRef = useRef(false); const { scrollRef, contentRef, scrollToBottom, isAtBottom } = useStickToBottom(); const { - stream, messages, todos, files, @@ -511,8 +546,11 @@ export const ChatInterface = React.memo( // resuming an interrupt → isLoading flips true). Without this, if the user had // drifted even slightly off the bottom after the previous answer, a short new // reply would render below the fold and look like nothing happened. + const wasLoadingRef = useRef(isLoading); useEffect(() => { - if (isLoading) void scrollToBottom(); + const startedLoading = isLoading && !wasLoadingRef.current; + wasLoadingRef.current = isLoading; + if (startedLoading) void scrollToBottom(); }, [isLoading, scrollToBottom]); // Register a "notify the main agent" hook up to page (Agents board → "Notify @@ -722,6 +760,7 @@ export const ChatInterface = React.memo( setThreadAutoApprove(threadId, false); setAutoApproveDialogOpen(false); autoApprovedRef.current = null; + autoApproveResumePendingRef.current = false; }, [threadId]); // Follow the thread: when the active thread changes, load THAT thread's saved @@ -745,6 +784,7 @@ export const ChatInterface = React.memo( setAutoApproveState(getThreadAutoApprove(threadId)); autoApprovedRef.current = null; + autoApproveResumePendingRef.current = false; setAutoApproveDialogOpen(false); setPendingFiles([]); migrateAutoApproveForCreatedThreadRef.current = false; @@ -1019,15 +1059,28 @@ export const ChatInterface = React.memo( !Array.isArray(actionRequests) || actionRequests.length === 0 ) { - autoApprovedRef.current = null; + if (!isLoading && !autoApproveResumePendingRef.current) { + autoApprovedRef.current = null; + } return; } - if (autoApprovedRef.current === ir) return; - autoApprovedRef.current = ir; + if (isLoading) { + return; + } + const interruptKey = autoApproveInterruptKey(ir, actionRequests); + if (autoApprovedRef.current === interruptKey) return; + autoApprovedRef.current = interruptKey; + autoApproveResumePendingRef.current = true; resumeInterrupt({ decisions: actionRequests.map(() => ({ type: "approve" })), }); - }, [autoApprove, interrupt, resumeInterrupt]); + }, [autoApprove, interrupt, isLoading, resumeInterrupt]); + + useEffect(() => { + if (isLoading) { + autoApproveResumePendingRef.current = false; + } + }, [isLoading]); // ask_user: the agent is asking the user structured questions. const askUserQuestions = useMemo(() => { @@ -1074,12 +1127,8 @@ export const ChatInterface = React.memo( string, { message: Message; toolCalls: ToolCall[] } >(); - // Sub-agent (subgraph) messages stream in alongside the main conversation - // when streamSubgraphs is on. They carry a NESTED langgraph_checkpoint_ns - // ("tools:|…") while the main agent's own messages are single-segment. - // Keep them OUT of the main flow — they render under each sub-agent block's - // "Steps" instead. (streamMetadata is live-only; once complete these messages - // aren't in thread state anyway.) + // The chat uses updates-only streaming. Persisted thread records contain + // main-thread messages only, so there is no live subgraph metadata to filter. const seenAsyncUpdates = new Set(); const visibleMessages = messages.filter((message: Message) => { // Humans are always user-typed (or our injected async-update pills) — @@ -1095,9 +1144,6 @@ export const ChatInterface = React.memo( seenAsyncUpdates.add(key); return true; } - const meta = stream.getMessagesMetadata(message)?.streamMetadata; - const ns = meta?.["langgraph_checkpoint_ns"]; - if (typeof ns === "string" && ns.includes("|")) return false; // The conversation-compaction summary is generated by a SEPARATE LLM // call (its own "Context Extraction Assistant" system prompt, like the // tool-selector). Its output transiently leaks into the raw stream as an @@ -1194,7 +1240,7 @@ export const ChatInterface = React.memo( showAvatar: data.message.type !== prevMessage?.type, }; }); - }, [messages, actionRequests, interrupt, isLoading, stream]); + }, [messages, actionRequests, interrupt, isLoading]); // UI preference: auto-collapse completed agent-action groups. The user can // turn this off in ConfigDialog; default is on. @@ -1571,7 +1617,6 @@ export const ChatInterface = React.memo( submittedActionRequestKeys={submittedActionRequestKeys} onActionRequestSubmitted={markActionRequestSubmitted} reviewConfigsMap={reviewConfigsMap} - stream={stream} onResumeInterrupt={resumeInterrupt} graphId={assistant?.graph_id} onEditMessage={handleEditMessage} @@ -1613,7 +1658,6 @@ export const ChatInterface = React.memo( isLastMessage ? reviewConfigsMap : undefined } ui={messageUi} - stream={stream} onResumeInterrupt={resumeInterrupt} graphId={assistant?.graph_id} onEditMessage={handleEditMessage} diff --git a/src/app/components/IdentityTab.tsx b/src/app/components/IdentityTab.tsx index a8088df..7c4a18f 100644 --- a/src/app/components/IdentityTab.tsx +++ b/src/app/components/IdentityTab.tsx @@ -270,7 +270,10 @@ export function IdentityTab({ listing, listingLoading }: IdentityTabProps) { useEffect(() => { const mq = window.matchMedia("(min-width: 768px)"); - const update = () => setIsDesktop(mq.matches); + const update = () => + setIsDesktop((current) => + current === mq.matches ? current : mq.matches + ); update(); mq.addEventListener("change", update); return () => mq.removeEventListener("change", update); diff --git a/src/app/components/ObservationGraph.tsx b/src/app/components/ObservationGraph.tsx index 1dd819a..7052f1f 100644 --- a/src/app/components/ObservationGraph.tsx +++ b/src/app/components/ObservationGraph.tsx @@ -811,7 +811,10 @@ export function ObservationGraph({ useEffect(() => { const mq = window.matchMedia("(min-width: 768px)"); - const update = () => setIsDesktop(mq.matches); + const update = () => + setIsDesktop((current) => + current === mq.matches ? current : mq.matches + ); update(); mq.addEventListener("change", update); return () => mq.removeEventListener("change", update); @@ -844,9 +847,13 @@ export function ObservationGraph({ const el = containerRef.current; if (!el) return; const ro = new ResizeObserver(([entry]) => { - setDims({ + const next = { w: entry.contentRect.width, h: entry.contentRect.height, + }; + setDims((current) => { + if (current.w === next.w && current.h === next.h) return current; + return next; }); }); ro.observe(el); diff --git a/src/app/hooks/useAvailableModels.ts b/src/app/hooks/useAvailableModels.ts index 03a0be8..0a0a595 100644 --- a/src/app/hooks/useAvailableModels.ts +++ b/src/app/hooks/useAvailableModels.ts @@ -45,13 +45,18 @@ function fetchRegistry( const headers: Record = {}; if (apiKey) headers["X-Api-Key"] = apiKey; - const p = fetch(`${key}/api/models`, { headers }) + const fetchJson = async (url: string): Promise => { + const r = await fetch(url, { headers }); + if (!r.ok) throw new Error(`HTTP ${r.status}`); + return (await r.json()) as RegistryResponse; + }; + + const p = fetchJson(`/api/models?deploymentUrl=${encodeURIComponent(key)}`) + .catch(() => fetchJson(`${key}/api/models`)) .then(async (r) => { - if (!r.ok) throw new Error(`HTTP ${r.status}`); - const body = (await r.json()) as RegistryResponse; const entries: ModelRegistryEntry[] = []; - if (Array.isArray(body.entries)) { - for (const raw of body.entries) { + if (Array.isArray(r.entries)) { + for (const raw of r.entries) { if (!raw || typeof raw !== "object") continue; const e = raw as { name?: unknown; @@ -73,8 +78,8 @@ function fetchRegistry( } } let defaultEntry: ModelRegistry["defaultEntry"] = null; - if (body.default && typeof body.default === "object") { - const d = body.default as { name?: unknown; provider?: unknown }; + if (r.default && typeof r.default === "object") { + const d = r.default as { name?: unknown; provider?: unknown }; if (typeof d.name === "string" && d.name) { defaultEntry = { name: d.name, @@ -95,10 +100,12 @@ function fetchRegistry( } /** - * Fetch the backend's authoritative model registry from - * `GET ${deploymentUrl}/api/models`. Results are cached at module level — - * the registry is static between deployment restarts, so remounting - * ChatInterface never triggers a redundant network request. + * Fetch the backend's authoritative model registry through the WebUI's + * same-origin `/api/models` proxy first, falling back to + * `GET ${deploymentUrl}/api/models` for older/alternate deployments. Results + * are cached at module level — the registry is static between deployment + * restarts, so remounting ChatInterface never triggers a redundant network + * request. * * Failures are non-fatal: the picker falls back to its curated * `COMMON_MODELS` list when `entries` is empty. Failed fetches are evicted diff --git a/src/app/hooks/useChat.ts b/src/app/hooks/useChat.ts index 4437c8f..b452f90 100644 --- a/src/app/hooks/useChat.ts +++ b/src/app/hooks/useChat.ts @@ -8,10 +8,6 @@ import type { UseStreamThread } from "@langchain/langgraph-sdk/react"; import type { TodoItem } from "@/app/types/types"; import { useClient } from "@/providers/ClientProvider"; import { useQueryState } from "nuqs"; -import { - extractSubAgentSteps, - type SubAgentStep, -} from "@/lib/subAgentActivity"; import { parseSummarizationEvent } from "@/lib/summarization"; import { toast } from "sonner"; import { @@ -39,6 +35,11 @@ export type StateType = { ui?: any; }; +type InterruptLike = { + value?: unknown; + [key: string]: unknown; +}; + /** * Sanitize a raw interrupt pulled from `client.threads.getState` before it is * surfaced to the UI. The live SDK normalizes `stream.interrupt`, but the raw @@ -161,6 +162,15 @@ function formatStreamError(error: unknown): string { 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 = {}; +const EMPTY_ASYNC_TASKS: Record = {}; +const EMPTY_SUB_AGENT_ACTIVITY: Record = {}; + export function useChat({ activeAssistant, onHistoryRevalidate, @@ -173,17 +183,10 @@ export function useChat({ const [threadId, setThreadId] = useQueryState("threadId"); const client = useClient(); - // Live sub-agent activity captured from subgraph stream events, keyed by the - // subgraph namespace (e.g. "tools:"). Ephemeral: it resets when the chat - // session remounts on thread switch, and is not persisted (lost on reload). - const [subAgentActivity, setSubAgentActivity] = useState< - Record - >({}); - const stream = useStream({ assistantId: activeAssistant?.assistant_id || "", client: client ?? undefined, - reconnectOnMount: true, + reconnectOnMount: false, threadId: threadId ?? null, onThreadId: setThreadId, defaultHeaders: { "x-auth-scheme": "langsmith" }, @@ -200,30 +203,23 @@ export function useChat({ toast.error(formatStreamError(error)); }, onCreated: onHistoryRevalidate, - // Capture sub-agent (subgraph) node outputs as they stream. `namespace` is - // non-empty (e.g. ["tools:"]) for subgraphs and empty for the main graph, - // which we skip. - onUpdateEvent: (data, options) => { - const ns = options?.namespace; - if (!ns || ns.length === 0) return; - const steps = extractSubAgentSteps(data); - if (steps.length === 0) return; - const key = ns.join("|"); - setSubAgentActivity((prev) => ({ - ...prev, - [key]: [...(prev[key] ?? []), ...steps], - })); - }, experimental_thread: thread, }); + // Do not read `stream.values` or `stream.messages`. Both getters add a + // high-frequency stream mode to the SDK request; each chunk then synchronously + // re-notifies React's external store. The UI refreshes from the persisted + // thread record below instead, while `isLoading` still drives run controls. + const liveInterrupt = stream.interrupt as InterruptLike | undefined; + 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 `stream.messages` is stale - // (missing the final `execute` tool-call message), so the approval card never - // renders until a manual thread switch re-fetches history. + // 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, @@ -232,8 +228,9 @@ export function useChat({ // 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(undefined); + 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 @@ -244,6 +241,9 @@ export function useChat({ const [fetchedMessages, setFetchedMessages] = useState( null ); + const [fetchedValues, setFetchedValues] = useState | null>( + null + ); const [fetchedThreadId, setFetchedThreadId] = useState(null); const recoveryRunRef = useRef(0); @@ -339,11 +339,16 @@ export function useChat({ }, [threadId] ); + + const liveMessages = EMPTY_MESSAGES; + const liveMessageCount = liveMessages.length; + useEffect(() => { if (!threadId) { setFetchedInterrupt(undefined); setFetchedMessages(null); setFetchedThreadId(null); + setFetchedValues(null); setResolvedInterruptKey(null); return; } @@ -353,12 +358,10 @@ export function useChat({ setResolvedInterruptKey(null); return; } - // The live stream count at the moment it settled. If the server's persisted - // state has MORE messages than this, the stream ended early and dropped the - // tail — either the final assistant text, or the `execute` tool-call message - // plus its approval interrupt. Either way we backfill from thread state - // (the same data a thread-switch re-fetch would pull in). - const baseline = stream.messages.length; + // Updates-only mode intentionally leaves the live message array empty. The + // persisted thread record is therefore the source of truth for the rendered + // conversation once the SDK settles. + const baseline = liveMessageCount; const recoveryRunId = ++recoveryRunRef.current; let cancelled = false; let tries = 0; @@ -380,11 +383,12 @@ export function useChat({ values?: { messages?: Message[] }; }>, client.threads.get(threadId) as Promise<{ - values?: { messages?: Message[] }; + values?: Partial; }>, ]); if (cancelled || recoveryRunRef.current !== recoveryRunId) return; - const msgs = threadRecord.values?.messages; + const values = threadRecord.values; + const msgs = values?.messages; const pending = latestTaskInterrupt(state.tasks); const stillPending = Array.isArray(state.next) && state.next.length > 0; const safePending = normalizePendingInterrupt(pending); @@ -393,32 +397,34 @@ export function useChat({ // snapshot together. Mixing live messages with fetched interrupts is the // race that hides approval cards for repeated execute calls. setFetchedInterrupt( - safePending as unknown as typeof stream.interrupt + safePending as InterruptLike ); setResolvedInterruptKey(null); if (Array.isArray(msgs)) { setFetchedThreadId(threadId); setFetchedMessages(msgs); } + setFetchedValues(values ?? null); return; } - // Backfill only after the live stream is idle. During active streaming the - // live message list owns rendering; this recovery loop is for dropped tail - // state after the stream has settled. + // 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(interruptValueKey(stream.interrupt)); + 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 @@ -443,13 +449,11 @@ export function useChat({ }; // Precise deps on purpose: re-running on the whole `stream` object (new each // render) would loop the getState fetch. - // eslint-disable-next-line react-hooks/exhaustive-deps - }, [threadId, stream.interrupt, stream.isLoading, client]); + }, [threadId, liveInterruptKey, stream.isLoading, client, liveMessageCount]); // 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 liveInterrupt = stream.interrupt; const interrupt = liveInterrupt && (resolvedInterruptKey === null || @@ -468,25 +472,21 @@ export function useChat({ // 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. - // `stream.messages` can carry subgraph noise (streamSubgraphs: true) plus - // stale per-message metadata from earlier runs, inflating its length above the - // persisted main-thread state. A pure `>` compare against that bloated count - // would keep us on the stream — which makes the downstream subgraph-namespace - // filter (ChatInterface.processedMessages) drop legitimate main-thread - // history that's only tagged subgraph in stale stream metadata. + // 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 stream.messages; + if (!fetchedMessages || fetchedThreadId !== threadId) return liveMessages; if (fetchedInterrupt) return fetchedMessages; if (!stream.isLoading) return fetchedMessages; - if (fetchedMessages.length > stream.messages.length) return fetchedMessages; + if (fetchedMessages.length > liveMessages.length) return fetchedMessages; if ( - fetchedMessages.length === stream.messages.length && - totalTextLength(fetchedMessages) > totalTextLength(stream.messages) + fetchedMessages.length === liveMessages.length && + totalTextLength(fetchedMessages) > totalTextLength(liveMessages) ) { return fetchedMessages; } - return stream.messages; + return liveMessages; })(); // Fold the per-thread model override into the assistant's base config. The @@ -507,36 +507,43 @@ export function useChat({ return { ...base, configurable, recursion_limit: 100 }; }, [activeAssistant?.config, modelOverride]); - const sendMessage = useCallback( - (content: string) => { - // Drop any settled-run snapshot up front. Otherwise, until `isLoading` - // flips true (and the effect above clears it), a previous run's - // `fetchedMessages` can still out-count `stream.messages` and shadow the - // just-added optimistic user message — making it flicker/vanish. - setFetchedInterrupt(undefined); - setFetchedMessages(null); - setFetchedThreadId(null); - setResolvedInterruptKey(null); - recoveryRunRef.current += 1; - const newMessage: Message = { id: uuidv4(), type: "human", content }; - stream.submit( - { messages: [newMessage] }, - { - optimisticValues: (prev) => ({ - messages: [...(prev.messages ?? []), newMessage], - }), - config: buildRunConfig(), - streamSubgraphs: true, - streamMode: ["updates"], - streamResumable: true, - onDisconnect: "continue", - } - ); - // Update thread list immediately when sending a message - onHistoryRevalidate?.(); - }, - [stream, buildRunConfig, onHistoryRevalidate] - ); + const streamRef = useRef(stream); + streamRef.current = stream; + const threadIdRef = useRef(threadId); + threadIdRef.current = threadId; + const buildRunConfigRef = useRef(buildRunConfig); + buildRunConfigRef.current = buildRunConfig; + const onHistoryRevalidateRef = useRef(onHistoryRevalidate); + onHistoryRevalidateRef.current = onHistoryRevalidate; + + const sendMessage = useCallback((content: string) => { + // 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, + ]); + streamRef.current.submit( + { messages: [newMessage] }, + { + optimisticValues: (prev) => ({ + messages: [...(prev.messages ?? []), newMessage], + }), + config: buildRunConfigRef.current(), + streamSubgraphs: false, + streamMode: ["updates"], + streamResumable: false, + onDisconnect: "cancel", + } + ); + // Update thread list immediately when sending a message + onHistoryRevalidateRef.current?.(); + }, []); const setFiles = useCallback( async (files: Record) => { @@ -548,44 +555,37 @@ export function useChat({ [client, threadId] ); - const resumeInterrupt = useCallback( - (value: any) => { - // Same as sendMessage: clear the prior snapshot before resuming so a stale - // fetchedInterrupt/fetchedMessages can't briefly re-surface a resolved - // approval card or shadow the resumed run's messages. - setFetchedInterrupt(undefined); - setFetchedMessages(null); - setFetchedThreadId(null); - setResolvedInterruptKey(null); - recoveryRunRef.current += 1; - stream.submit(null, { - command: { resume: value }, - config: buildRunConfig(), - streamSubgraphs: true, - streamMode: ["updates"], - streamResumable: true, - onDisconnect: "continue", - }); - // Update thread list when resuming from interrupt - onHistoryRevalidate?.(); - }, - [stream, buildRunConfig, onHistoryRevalidate] - ); + const resumeInterrupt = useCallback((value: any) => { + // Keep the transcript snapshot while the approval resumes; the next thread + // record refresh replaces it with the completed turn. + setFetchedInterrupt(undefined); + setResolvedInterruptKey(null); + recoveryRunRef.current += 1; + streamRef.current.submit(null, { + command: { resume: value }, + config: buildRunConfigRef.current(), + streamSubgraphs: false, + streamMode: ["updates"], + streamResumable: false, + onDisconnect: "cancel", + }); + // Update thread list when resuming from interrupt + onHistoryRevalidateRef.current?.(); + }, []); const stopStream = useCallback(() => { - stream.stop(); - }, [stream]); + streamRef.current.stop(); + }, []); return { - stream, - todos: stream.values.todos ?? [], - files: stream.values.files ?? {}, - email: stream.values.email, - asyncTasks: stream.values.async_tasks ?? {}, + todos: fetchedValues?.todos ?? EMPTY_TODOS, + files: fetchedValues?.files ?? EMPTY_FILES, + email: fetchedValues?.email, + asyncTasks: fetchedValues?.async_tasks ?? EMPTY_ASYNC_TASKS, summarizationEvent: parseSummarizationEvent( - stream.values._summarization_event + fetchedValues?._summarization_event ), - ui: stream.values.ui, + ui: fetchedValues?.ui, setFiles, messages, isLoading: stream.isLoading, @@ -594,7 +594,7 @@ export function useChat({ sendMessage, stopStream, resumeInterrupt, - subAgentActivity, + subAgentActivity: EMPTY_SUB_AGENT_ACTIVITY, modelOverride, setModelOverride, }; diff --git a/src/app/page.tsx b/src/app/page.tsx index d7da906..d5b333c 100644 --- a/src/app/page.tsx +++ b/src/app/page.tsx @@ -131,7 +131,10 @@ function HomePageInner({ useEffect(() => { const mediaQuery = window.matchMedia("(min-width: 768px)"); - const updateLayout = () => setIsDesktopLayout(mediaQuery.matches); + const updateLayout = () => + setIsDesktopLayout((current) => + current === mediaQuery.matches ? current : mediaQuery.matches + ); updateLayout(); mediaQuery.addEventListener("change", updateLayout); diff --git a/src/lib/uiSettings.ts b/src/lib/uiSettings.ts index 3086951..1fe5ffe 100644 --- a/src/lib/uiSettings.ts +++ b/src/lib/uiSettings.ts @@ -27,14 +27,15 @@ export function useCollapseAgentActions(): { useEffect(() => { const onStorage = (e: StorageEvent) => { if (e.key !== COLLAPSE_AGENT_ACTIONS_KEY) return; - setValueState(readCollapseAgentActions()); + const next = readCollapseAgentActions(); + setValueState((current) => (current === next ? current : next)); }; window.addEventListener("storage", onStorage); return () => window.removeEventListener("storage", onStorage); }, []); const setValue = useCallback((next: boolean) => { - setValueState(next); + setValueState((current) => (current === next ? current : next)); if (typeof window !== "undefined") { window.localStorage.setItem(COLLAPSE_AGENT_ACTIONS_KEY, String(next)); } diff --git a/src/providers/ThemeProvider.tsx b/src/providers/ThemeProvider.tsx index 98c0123..a1bfca3 100644 --- a/src/providers/ThemeProvider.tsx +++ b/src/providers/ThemeProvider.tsx @@ -46,7 +46,9 @@ export function ThemeProvider({ children }: { children: React.ReactNode }) { const apply = () => { const resolved: ResolvedTheme = theme === "system" ? (mql.matches ? "dark" : "light") : theme; - setResolvedTheme(resolved); + setResolvedTheme((current) => + current === resolved ? current : resolved + ); applyTheme(resolved); }; apply(); @@ -62,7 +64,7 @@ export function ThemeProvider({ children }: { children: React.ReactNode }) { } catch { // Private mode / storage disabled — fall back to in-memory only. } - setThemeState(next); + setThemeState((current) => (current === next ? current : next)); }, []); return (