fix: stabilize chat stream updates
CI / Format, lint & build (push) Has been cancelled

This commit is contained in:
m4
2026-07-10 16:26:46 +08:00
parent 363f15d850
commit 14b246f4e7
11 changed files with 285 additions and 251 deletions
-90
View File
@@ -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": [
+57
View File
@@ -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<string, string> = {};
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);
}
}
+1 -1
View File
@@ -34,7 +34,7 @@ interface ActionGroupProps {
submittedActionRequestKeys: Set<string>;
onActionRequestSubmitted: (key: string) => void;
reviewConfigsMap: Map<string, ReviewConfig> | null;
stream: unknown;
stream?: unknown;
onResumeInterrupt: (value: unknown) => void;
graphId?: string;
onEditMessage: (content: string) => void;
+63 -19
View File
@@ -230,6 +230,41 @@ function getMessageToolCalls(message: Message): Array<{
});
}
function stableStringify(value: unknown): string {
const seen = new WeakSet<object>();
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<Record<string, unknown>>((out, key) => {
out[key] = (v as Record<string, unknown>)[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<ChatInterfaceProps>(
);
});
}, [pickerModels, modelSearch]);
const autoApprovedRef = useRef<unknown>(null);
const autoApprovedRef = useRef<string | null>(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<ChatInterfaceProps>(
// 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<ChatInterfaceProps>(
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<ChatInterfaceProps>(
setAutoApproveState(getThreadAutoApprove(threadId));
autoApprovedRef.current = null;
autoApproveResumePendingRef.current = false;
setAutoApproveDialogOpen(false);
setPendingFiles([]);
migrateAutoApproveForCreatedThreadRef.current = false;
@@ -1019,15 +1059,28 @@ export const ChatInterface = React.memo<ChatInterfaceProps>(
!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<AskUserQuestion[] | null>(() => {
@@ -1074,12 +1127,8 @@ export const ChatInterface = React.memo<ChatInterfaceProps>(
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:<id>|…") 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<string>();
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<ChatInterfaceProps>(
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<ChatInterfaceProps>(
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<ChatInterfaceProps>(
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<ChatInterfaceProps>(
isLastMessage ? reviewConfigsMap : undefined
}
ui={messageUi}
stream={stream}
onResumeInterrupt={resumeInterrupt}
graphId={assistant?.graph_id}
onEditMessage={handleEditMessage}
+4 -1
View File
@@ -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);
+9 -2
View File
@@ -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);
+18 -11
View File
@@ -45,13 +45,18 @@ function fetchRegistry(
const headers: Record<string, string> = {};
if (apiKey) headers["X-Api-Key"] = apiKey;
const p = fetch(`${key}/api/models`, { headers })
const fetchJson = async (url: string): Promise<RegistryResponse> => {
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
+122 -122
View File
@@ -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<string, string> = {};
const EMPTY_ASYNC_TASKS: Record<string, unknown> = {};
const EMPTY_SUB_AGENT_ACTIVITY: Record<string, never[]> = {};
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:<id>"). Ephemeral: it resets when the chat
// session remounts on thread switch, and is not persisted (lost on reload).
const [subAgentActivity, setSubAgentActivity] = useState<
Record<string, SubAgentStep[]>
>({});
const stream = useStream<StateType>({
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:<id>"]) 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<typeof stream.interrupt>(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<Message[] | null>(
null
);
const [fetchedValues, setFetchedValues] = useState<Partial<StateType> | null>(
null
);
const [fetchedThreadId, setFetchedThreadId] = useState<string | null>(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<StateType>;
}>,
]);
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<string, string>) => {
@@ -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,
};
+4 -1
View File
@@ -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);
+3 -2
View File
@@ -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));
}
+4 -2
View File
@@ -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 (