Files
deer-flow/frontend/src/core/threads/hooks.ts
T
Ryker_FengandGitHub 6a4e5a3bb2 feat(frontend): add side conversations for quoted follow-ups (#3934)
* feat(frontend): add side conversations for quoted follow-ups

* style(frontend): apply prettier formatting to sidecar-chat files

* fix(frontend): surface sidecar cascade cleanup failures via console.warn

Previously deleteSidecarThreadsForParent silently swallowed both
lookup errors and per-thread deletion failures, so parent thread
deletions could succeed while orphaning sidecar threads with no
signal to the caller. Log a warning that includes the parent id
and the failed thread ids/reasons so the leak is discoverable in
telemetry, matching the existing console.warn/error pattern in
this file.

* fix(frontend): address all sidecar review feedback

Resolve every reviewer comment on PR #3934:

- input-box/hooks/sidecar-panel: clear quoted references only via an
  `onSent` callback that fires after the in-flight guard, so a dropped
  send no longer silently discards quotes (willem-bd #3550).
- message-list: flip the selection toolbar below the selection when it
  would clip above the viewport (willem-bd #3551).
- reference-metadata/thread/input-box: keep referenced ids, roles, and
  count arrays 1:1 parallel instead of deduping ids (willem-bd #3552).
- message-list: widen selection containment to the shared assistant-turn
  container and hint when a selection crosses messages (willem-bd #3553).
- sidecar/api: coalesce concurrent sidecar creates for one parent behind
  a single in-flight promise to prevent duplicates (willem-bd #3554).
- sidecar-trigger/context: force-restore on trigger click so a sidecar
  deleted elsewhere self-heals instead of opening a dead thread
  (willem-bd #3555).
- threads/hooks: surface sidecar cascade cleanup failures via
  console.warn for both lookup and per-thread deletes (Copilot).

Add unit + e2e coverage for parallel metadata, atomic create, and
trigger self-healing.
2026-07-05 00:12:16 +08:00

2232 lines
67 KiB
TypeScript

import type { AIMessage, Message, Run } from "@langchain/langgraph-sdk";
import type { ThreadsClient } from "@langchain/langgraph-sdk/client";
import { useStream } from "@langchain/langgraph-sdk/react";
import {
type QueryClient,
type InfiniteData,
useInfiniteQuery,
useMutation,
useQuery,
useQueryClient,
} from "@tanstack/react-query";
import { useCallback, useEffect, useMemo, useRef, useState } from "react";
import { toast } from "sonner";
import type { PromptInputMessage } from "@/components/ai-elements/prompt-input";
import { getAPIClient } from "../api";
import { fetch } from "../api/fetcher";
import { getBackendBaseURL } from "../config";
import { useI18n } from "../i18n/hooks";
import { isHiddenFromUIMessage } from "../messages/utils";
import type { FileInMessage } from "../messages/utils";
import type { LocalSettings } from "../settings";
import { isSidecarThread, SIDECAR_METADATA_KEY } from "../sidecar/thread";
import { useUpdateSubtask } from "../tasks/context";
import { messageToStep } from "../tasks/steps";
import type { UploadedFileInfo } from "../uploads";
import { promptInputFilePartToFile, uploadFiles } from "../uploads";
import { fetchThreadTokenUsage } from "./api";
import {
buildThreadsSearchQueryOptions,
DEFAULT_THREAD_SEARCH_PARAMS,
filterThreadSearchResults,
type ThreadSearchParams,
} from "./thread-search-query";
import { threadTokenUsageQueryKey } from "./token-usage";
import type {
AgentThread,
AgentThreadState,
RunMessage,
ThreadTokenUsageResponse,
} from "./types";
export type ToolEndEvent = {
name: string;
data: unknown;
};
export type ThreadStreamOptions = {
threadId?: string | null | undefined;
displayThreadId?: string | null | undefined;
context: LocalSettings["context"];
isMock?: boolean;
onSend?: (threadId: string) => void;
onStart?: (threadId: string, runId: string) => void;
onFinish?: (state: AgentThreadState) => void;
onToolEnd?: (event: ToolEndEvent) => void;
};
type SendMessageOptions = {
additionalKwargs?: Record<string, unknown>;
additionalInputMessages?: Message[];
/**
* Invoked exactly once when the send passes the in-flight guard and is
* genuinely dispatched. It never fires on the early-return path, so callers
* can safely perform one-time cleanup (e.g. clearing quoted references)
* without losing state when a concurrent send is dropped.
*/
onSent?: () => void;
};
type ThreadDeleteClient = {
threads: {
delete: (threadId: string) => Promise<unknown>;
search: (query: Record<string, unknown>) => Promise<AgentThread[]>;
};
};
type ThreadSidecarSearchClient = {
threads: {
search: (query: Record<string, unknown>) => Promise<AgentThread[]>;
};
};
type RegeneratePrepareResponse = {
input: Partial<AgentThreadState>;
checkpoint: {
checkpoint_ns: string;
checkpoint_id: string;
checkpoint_map: Record<string, unknown> | null;
};
metadata: Record<string, unknown>;
target_run_id: string;
};
export function buildThreadSubmitMessages({
text,
additionalKwargs,
additionalInputMessages = [],
filesForSubmit = [],
}: {
text: string;
additionalKwargs?: Record<string, unknown>;
additionalInputMessages?: Message[];
filesForSubmit?: FileInMessage[];
}): Message[] {
return [
...additionalInputMessages,
{
type: "human",
content: [
{
type: "text",
text,
},
],
additional_kwargs: {
...additionalKwargs,
...(filesForSubmit.length > 0 ? { files: filesForSubmit } : {}),
},
} as Message,
];
}
const EMPTY_THREAD_VALUES: AgentThreadState = {
title: "",
messages: [],
artifacts: [],
todos: [],
};
function isNonEmptyString(value: string | undefined): value is string {
return typeof value === "string" && value.length > 0;
}
const SUMMARIZATION_MIDDLEWARE_UPDATE_KEYS = new Set([
"SummarizationMiddleware.before_model",
"DeerFlowSummarizationMiddleware.before_model",
]);
function messageIdentity(message: Message): string | undefined {
if (
"tool_call_id" in message &&
typeof message.tool_call_id === "string" &&
message.tool_call_id.length > 0
) {
return `tool:${message.tool_call_id}`;
}
if (typeof message.id === "string" && message.id.length > 0) {
return `message:${message.id}`;
}
return undefined;
}
function dedupeMessagesByIdentity(messages: Message[]): Message[] {
const lastIndexByIdentity = new Map<string, number>();
const lastVisibleIndexByIdentity = new Map<string, number>();
// This is a UI-display dedupe rule, not a general LangChain message-stream
// contract. Hidden messages that share an identity with a visible message are
// treated as control messages for this merged view; hidden messages carrying
// independent tracing/task semantics should use a distinct id or a custom
// stream/state channel instead of relying on message dedupe preservation.
const preservedTurnDurations = new Map<string, number>();
messages.forEach((message, index) => {
const identity = messageIdentity(message);
if (identity) {
lastIndexByIdentity.set(identity, index);
if (!isHiddenFromUIMessage(message)) {
lastVisibleIndexByIdentity.set(identity, index);
}
if (message.additional_kwargs?.turn_duration !== undefined) {
preservedTurnDurations.set(
identity,
message.additional_kwargs.turn_duration as number,
);
}
}
});
return messages
.filter((message, index) => {
const identity = messageIdentity(message);
if (!identity) {
return true;
}
const visibleIndex = lastVisibleIndexByIdentity.get(identity);
if (visibleIndex !== undefined) {
return visibleIndex === index;
}
return lastIndexByIdentity.get(identity) === index;
})
.map((message) => {
const identity = messageIdentity(message);
if (
identity &&
preservedTurnDurations.has(identity) &&
message.additional_kwargs?.turn_duration === undefined
) {
return {
...message,
additional_kwargs: {
...message.additional_kwargs,
turn_duration: preservedTurnDurations.get(identity),
},
} as Message;
}
return message;
});
}
function dedupeRunMessagesByIdentity(messages: RunMessage[]): RunMessage[] {
const lastIndexByIdentity = new Map<string, number>();
messages.forEach((message, index) => {
const identity = messageIdentity(message.content);
if (identity) {
lastIndexByIdentity.set(`${message.run_id}:${identity}`, index);
}
});
return messages.filter((message, index) => {
const identity = messageIdentity(message.content);
if (!identity) {
return true;
}
return lastIndexByIdentity.get(`${message.run_id}:${identity}`) === index;
});
}
export function getSupersededRunIds(
runs: Run[] | undefined,
pendingSupersededRunIds?: ReadonlySet<string>,
) {
const ids = new Set(pendingSupersededRunIds ?? []);
for (const run of runs ?? []) {
if (run.status !== "success") {
continue;
}
const metadata = run.metadata;
if (metadata && typeof metadata === "object") {
const fromRunId = Reflect.get(metadata, "regenerate_from_run_id");
if (typeof fromRunId === "string" && fromRunId) {
ids.add(fromRunId);
}
}
}
return ids;
}
export function removeSetItems<T>(
values: ReadonlySet<T>,
itemsToRemove: Iterable<T>,
) {
const next = new Set(values);
for (const item of itemsToRemove) {
next.delete(item);
}
return next;
}
export function buildVisibleHistoryMessages(
messageRows: RunMessage[],
supersededRunIds: ReadonlySet<string>,
appendedMessages: Message[],
) {
const visibleRows = messageRows.filter(
(message) => !supersededRunIds.has(message.run_id),
);
return dedupeMessagesByIdentity([
// Carry the owning run_id onto the content message so historical subtask
// cards can fetch their persisted step history on expand (#3779). run_id
// lives on the RunMessage wrapper and would otherwise be dropped here.
...visibleRows.map((message) => ({
...message.content,
run_id: message.run_id,
})),
...appendedMessages,
]);
}
export function findLatestUnloadedRunIndex(
runs: Run[],
loadedRunIds: ReadonlySet<string>,
): number {
for (let i = 0; i < runs.length; i++) {
const run = runs[i];
if (run && !loadedRunIds.has(run.run_id)) {
return i;
}
}
return -1;
}
export const MAX_CONSECUTIVE_EMPTY_RUN_LOADS = 5;
export function shouldAutoContinueOnEmptyRun(
fetchedMessageCount: number,
consecutiveEmptyLoads: number,
maxConsecutiveEmptyLoads: number = MAX_CONSECUTIVE_EMPTY_RUN_LOADS,
): boolean {
return (
fetchedMessageCount === 0 &&
consecutiveEmptyLoads < maxConsecutiveEmptyLoads
);
}
type RunMessagesPageResponse = {
data: RunMessage[];
has_more?: boolean;
hasMore?: boolean;
};
export function runMessagesPageHasMore(result: RunMessagesPageResponse) {
return result.has_more ?? result.hasMore ?? false;
}
export function getOldestRunMessageSeq(messages: RunMessage[]) {
let oldestSeq: number | null = null;
for (const message of messages) {
if (typeof message.seq !== "number") {
continue;
}
oldestSeq =
oldestSeq === null ? message.seq : Math.min(oldestSeq, message.seq);
}
return oldestSeq;
}
export function getNextRunMessagesBeforeSeq(
result: RunMessagesPageResponse,
): number | null | undefined {
if (!runMessagesPageHasMore(result)) {
return null;
}
return getOldestRunMessageSeq(result.data) ?? undefined;
}
export function buildRunMessagesUrl(
baseUrl: string,
threadId: string,
runId: string,
beforeSeq?: number,
) {
const normalizedBaseUrl = baseUrl.replace(/\/$/, "");
const path = `/api/threads/${encodeURIComponent(threadId)}/runs/${encodeURIComponent(runId)}/messages`;
const url = new URL(
`${normalizedBaseUrl}${path}`,
typeof window !== "undefined" ? window.location.origin : "http://localhost",
);
if (beforeSeq !== undefined) {
url.searchParams.set("before_seq", String(beforeSeq));
}
return normalizedBaseUrl ? url.toString() : `${url.pathname}${url.search}`;
}
export function mergeMessages(
historyMessages: Message[],
threadMessages: Message[],
optimisticMessages: Message[],
): Message[] {
// Only visible live messages should trim overlapping history. Hidden messages
// are UI control messages in this path, not observability records; any hidden
// message that must survive as task/tracing data should use custom events or a
// separate state channel instead of participating in this overlap heuristic.
const savedTurnDurations = new Map<string, number>();
for (const msg of historyMessages) {
const identity = messageIdentity(msg);
if (identity && msg.additional_kwargs?.turn_duration !== undefined) {
savedTurnDurations.set(
identity,
msg.additional_kwargs.turn_duration as number,
);
}
}
const threadMessageIds = new Set(
threadMessages
.filter((message) => !isHiddenFromUIMessage(message))
.map(messageIdentity)
.filter(isNonEmptyString),
);
// The overlap is a contiguous suffix of historyMessages (newest history == oldest thread).
// Scan from the end: shrink cutoff while messages are already in thread, stop as soon as
// we hit one that isn't — everything before that point is non-overlapping.
let cutoff = historyMessages.length;
for (let i = historyMessages.length - 1; i >= 0; i--) {
const msg = historyMessages[i];
if (!msg) {
continue;
}
const identity = messageIdentity(msg);
if (identity && threadMessageIds.has(identity)) {
cutoff = i;
} else {
break;
}
}
const merged = dedupeMessagesByIdentity([
...historyMessages.slice(0, cutoff),
...threadMessages,
...optimisticMessages,
]);
return merged.map((message) => {
const identity = messageIdentity(message);
if (
identity &&
savedTurnDurations.has(identity) &&
message.additional_kwargs?.turn_duration === undefined
) {
return {
...message,
additional_kwargs: {
...message.additional_kwargs,
turn_duration: savedTurnDurations.get(identity),
},
} as Message;
}
return message;
});
}
/**
* Derive the live turns that context summarization is about to drop and that
* therefore must be re-archived into history.
*
* Summarization emits `RemoveMessage(ALL)` + a hidden summary + the retained
* tail. Everything in the current live thread before the first retained visible
* message is being removed; we keep those (minus the summary control messages
* already tracked) so the UI can still show the full conversation (#3825).
*/
export function computeSummarizationMovedMessages(
currentMessages: Message[],
summarizationMessages: Message[],
summarizedMessageIds: ReadonlySet<string>,
): Message[] {
const firstRetainedVisibleIdentity = summarizationMessages
.filter((message) => message.type !== "remove")
.filter((message) => !isHiddenFromUIMessage(message))
.map(messageIdentity)
.find(isNonEmptyString);
const moved: Message[] = [];
for (const message of currentMessages) {
if (
firstRetainedVisibleIdentity &&
messageIdentity(message) === firstRetainedVisibleIdentity
) {
break;
}
if (!summarizedMessageIds.has(message.id ?? "")) {
moved.push(message);
}
}
return moved;
}
/**
* Overlay the messages rescued from context summarization on top of the
* (possibly stale) visible history so the merged view never drops them.
*
* Background (#3825): after summarization the backend removes every live
* message (`RemoveMessage(ALL)`) and `onUpdateEvent` re-archives the removed
* messages into history through an async `setState`. The live thread messages
* are owned by the LangGraph SDK external store while the archived history is
* React state, so a render can observe the post-summary (shrunk) thread before
* the archive `setState` commits — leaving the rescued messages in neither
* merge input. Reading them from a synchronous buffer here keeps the merge
* correct at every render regardless of how the two state channels interleave.
*
* The rescued messages are the oldest live turns, so they follow whatever the
* already-loaded history holds. Only messages still missing from history are
* appended: once history absorbs a rescued message, its live copy stays
* authoritative (the buffered copy is an older snapshot and must never overwrite
* it), and ordering is preserved.
*/
export function resolvePreservedHistory(
visibleHistory: Message[],
pendingArchivedMessages: Message[],
): Message[] {
if (pendingArchivedMessages.length === 0) {
return visibleHistory;
}
const presentIdentities = new Set(
visibleHistory.map(messageIdentity).filter(isNonEmptyString),
);
const missing = pendingArchivedMessages.filter((message) => {
const identity = messageIdentity(message);
// Identity-less messages are intentionally skipped: without a stable
// identity they cannot be matched against history to drain or dedupe, so
// overlaying them would risk a permanent duplicate. They are still archived
// through appendMessages and surface via the normal history path instead.
return identity !== undefined && !presentIdentities.has(identity);
});
if (missing.length === 0) {
return visibleHistory;
}
return [...visibleHistory, ...missing];
}
/**
* Drop the archive-buffer entries that the canonical history state has already
* absorbed. This keeps the buffer a transient bridge across the async gap
* rather than a second long-lived source of truth — otherwise a stale copy
* could resurrect a message that history later filtered out (e.g. a superseded
* or regenerated run).
*/
export function pruneConfirmedArchivedMessages(
pendingArchivedMessages: Message[],
visibleHistory: Message[],
): Message[] {
if (pendingArchivedMessages.length === 0) {
return pendingArchivedMessages;
}
const confirmedIdentities = new Set(
visibleHistory.map(messageIdentity).filter(isNonEmptyString),
);
return pendingArchivedMessages.filter((message) => {
const identity = messageIdentity(message);
return !identity || !confirmedIdentities.has(identity);
});
}
function getMessagesAfterBaseline(
messages: Message[],
baselineMessageIds: ReadonlySet<string>,
): Message[] {
return messages.filter((message) => {
const id = messageIdentity(message);
return !id || !baselineMessageIds.has(id);
});
}
export function getVisibleOptimisticMessages(
optimisticMessages: Message[],
previousHumanMessageCount: number,
currentHumanMessageCount: number,
): Message[] {
if (
optimisticMessages.some((message) => message.type === "human") &&
currentHumanMessageCount > previousHumanMessageCount
) {
return [];
}
return optimisticMessages;
}
export function getSummarizationMiddlewareMessages(
data: unknown,
): Message[] | undefined {
if (typeof data !== "object" || data === null) {
return undefined;
}
for (const [key, update] of Object.entries(data)) {
if (!SUMMARIZATION_MIDDLEWARE_UPDATE_KEYS.has(key)) {
continue;
}
if (typeof update !== "object" || update === null) {
continue;
}
const messages = Reflect.get(update, "messages");
if (Array.isArray(messages)) {
return [...messages] as Message[];
}
}
return undefined;
}
export function upsertThreadInSearchCache(
queryClient: QueryClient,
thread: AgentThread,
) {
queryClient.setQueriesData(
{
queryKey: ["threads", "search"],
exact: false,
},
(oldData: Array<AgentThread> | undefined) => {
if (!oldData) {
return [thread];
}
const existingIndex = oldData.findIndex(
(t) => t.thread_id === thread.thread_id,
);
if (existingIndex === -1) {
return [thread, ...oldData];
}
return oldData.map((t, index) => {
if (index !== existingIndex) {
return t;
}
return {
...thread,
...t,
metadata: {
...(thread.metadata ?? {}),
...(t.metadata ?? {}),
},
values: {
...thread.values,
...t.values,
},
};
});
},
);
}
export function upsertThreadInInfiniteCache(
queryClient: QueryClient,
thread: AgentThread,
) {
queryClient.setQueriesData(
{
queryKey: INFINITE_THREADS_QUERY_KEY_PREFIX,
exact: false,
},
(oldData: InfiniteData<AgentThread[]> | undefined) => {
if (!oldData) {
return oldData;
}
const merged = oldData.pages.map((page) =>
page.map((t) =>
t.thread_id === thread.thread_id
? {
...thread,
...t,
metadata: {
...(thread.metadata ?? {}),
...(t.metadata ?? {}),
},
values: {
...thread.values,
...t.values,
},
}
: t,
),
);
const exists = merged.some((page) =>
page.some((t) => t.thread_id === thread.thread_id),
);
if (exists) {
return { ...oldData, pages: merged };
}
const firstPage = merged[0] ?? [];
const restPages = merged.slice(1);
return {
...oldData,
pages: [[thread, ...firstPage], ...restPages],
};
},
);
}
export function invalidateStoppedThreadCaches(
queryClient: QueryClient,
threadId: string | null | undefined,
isMock = false,
) {
void queryClient.invalidateQueries({ queryKey: ["threads", "search"] });
void queryClient.invalidateQueries({
queryKey: INFINITE_THREADS_QUERY_KEY_PREFIX,
});
if (!threadId || isMock) {
return;
}
void queryClient.invalidateQueries({ queryKey: ["thread", threadId] });
void queryClient.invalidateQueries({
queryKey: ["thread", "metadata", threadId, isMock],
});
void queryClient.invalidateQueries({
queryKey: threadTokenUsageQueryKey(threadId),
});
}
export const STOP_THREAD_FINALIZATION_REFETCH_DELAY_MS = 1500;
function scheduleStoppedThreadFinalizationRefetch(
queryClient: QueryClient,
threadId: string | null | undefined,
isMock = false,
) {
if (isMock) {
return;
}
globalThis.setTimeout(() => {
invalidateStoppedThreadCaches(queryClient, threadId, isMock);
}, STOP_THREAD_FINALIZATION_REFETCH_DELAY_MS);
}
export async function stopThreadAndInvalidateCaches(
queryClient: QueryClient,
stop: () => Promise<void> | void,
threadId: string | null | undefined,
isMock = false,
) {
try {
await stop();
} finally {
invalidateStoppedThreadCaches(queryClient, threadId, isMock);
scheduleStoppedThreadFinalizationRefetch(queryClient, threadId, isMock);
}
}
function getStreamErrorMessage(error: unknown): string {
if (typeof error === "string" && error.trim()) {
return error;
}
if (error instanceof Error && error.message.trim()) {
return error.message;
}
if (typeof error === "object" && error !== null) {
const message = Reflect.get(error, "message");
if (typeof message === "string" && message.trim()) {
return message;
}
const nestedError = Reflect.get(error, "error");
if (nestedError instanceof Error && nestedError.message.trim()) {
return nestedError.message;
}
if (typeof nestedError === "string" && nestedError.trim()) {
return nestedError;
}
}
return "Request failed.";
}
async function readResponseErrorMessage(
response: Response,
fallback = "Request failed.",
) {
try {
const data = await response.json();
if (typeof data?.detail === "string" && data.detail.trim()) {
return data.detail;
}
} catch {
// Use the fallback below when the response body is not JSON.
}
return response.statusText || fallback;
}
function getHttpStatus(error: unknown): number | undefined {
if (typeof error !== "object" || error === null) {
return undefined;
}
const status = Reflect.get(error, "status");
if (typeof status === "number") {
return status;
}
const response = Reflect.get(error, "response");
if (typeof response === "object" && response !== null) {
const responseStatus = Reflect.get(response, "status");
if (typeof responseStatus === "number") {
return responseStatus;
}
}
return undefined;
}
function isThreadMissingError(error: unknown): boolean {
const status = getHttpStatus(error);
// Treat 403 like 404 here to avoid disclosing whether an inaccessible thread
// exists; callers redirect stale/inaccessible URLs back to a blank chat.
return status === 403 || status === 404;
}
export function useThreadStream({
threadId,
displayThreadId,
context,
isMock,
onSend,
onStart,
onFinish,
onToolEnd,
}: ThreadStreamOptions) {
const { t } = useI18n();
const currentViewThreadId = displayThreadId ?? threadId ?? null;
const currentViewThreadIdRef = useRef(currentViewThreadId);
currentViewThreadIdRef.current = currentViewThreadId;
// Optimistic messages shown before the server stream responds.
const [optimisticMessages, setOptimisticMessages] = useState<Message[]>([]);
const [optimisticThreadId, setOptimisticThreadId] = useState<string | null>(
null,
);
const [liveMessagesThreadId, setLiveMessagesThreadId] = useState<
string | null
>(null);
const [pendingSupersededRunIds, setPendingSupersededRunIds] = useState<
ReadonlySet<string>
>(() => new Set());
const [pendingSupersededMessageIds, setPendingSupersededMessageIds] =
useState<ReadonlySet<string>>(() => new Set());
const [isUploading, setIsUploading] = useState(false);
// Track the thread ID that is currently streaming to handle thread changes during streaming
const [onStreamThreadId, setOnStreamThreadId] = useState(() => threadId);
// Ref to track current thread ID across async callbacks without causing re-renders,
// and to allow access to the current thread id in onUpdateEvent
const threadIdRef = useRef<string | null>(threadId ?? null);
const startedRef = useRef(false);
const pendingUsageBaselineMessageIdsRef = useRef<Set<string>>(new Set());
const listeners = useRef({
onSend,
onStart,
onFinish,
onToolEnd,
});
const {
messages: history,
hasMore: hasMoreHistory,
loadMore: loadMoreHistory,
loading: isHistoryLoading,
appendMessages,
} = useThreadHistory(onStreamThreadId ?? "", {
enabled: !isMock,
pendingSupersededRunIds,
});
// Keep listeners ref updated with latest callbacks
useEffect(() => {
listeners.current = { onSend, onStart, onFinish, onToolEnd };
}, [onSend, onStart, onFinish, onToolEnd]);
useEffect(() => {
const normalizedThreadId = threadId ?? null;
if (!normalizedThreadId) {
// Reset when the UI moves back to a brand new unsaved thread.
startedRef.current = false;
setOnStreamThreadId(normalizedThreadId);
} else {
setOnStreamThreadId(normalizedThreadId);
}
threadIdRef.current = normalizedThreadId;
}, [threadId]);
const handleStreamStart = useCallback((_threadId: string, _runId: string) => {
threadIdRef.current = _threadId;
setOptimisticThreadId((currentOptimisticThreadId) => {
const currentView = currentViewThreadIdRef.current;
if (
currentOptimisticThreadId &&
(currentOptimisticThreadId === currentView ||
currentOptimisticThreadId === _threadId)
) {
return _threadId;
}
return currentOptimisticThreadId;
});
setLiveMessagesThreadId((currentLiveMessagesThreadId) => {
const currentView = currentViewThreadIdRef.current;
if (
currentLiveMessagesThreadId &&
(currentLiveMessagesThreadId === currentView ||
currentLiveMessagesThreadId === _threadId)
) {
return _threadId;
}
return currentLiveMessagesThreadId;
});
if (!startedRef.current) {
listeners.current.onStart?.(_threadId, _runId);
startedRef.current = true;
}
setOnStreamThreadId(_threadId);
}, []);
const queryClient = useQueryClient();
const updateSubtask = useUpdateSubtask();
const thread = useStream<AgentThreadState>({
client: getAPIClient(isMock),
assistantId: "lead_agent",
threadId: onStreamThreadId,
reconnectOnMount: true,
fetchStateHistory: { limit: 1 },
onCreated(meta) {
handleStreamStart(meta.thread_id, meta.run_id);
const now = new Date().toISOString();
upsertThreadInSearchCache(queryClient, {
thread_id: meta.thread_id,
created_at: now,
updated_at: now,
metadata: context.agent_name ? { agent_name: context.agent_name } : {},
status: "busy",
values: {
title: t.pages.newChat,
messages: [],
artifacts: [],
},
interrupts: {},
});
upsertThreadInInfiniteCache(queryClient, {
thread_id: meta.thread_id,
created_at: now,
updated_at: now,
metadata: context.agent_name ? { agent_name: context.agent_name } : {},
status: "busy",
values: {
title: t.pages.newChat,
messages: [],
artifacts: [],
},
interrupts: {},
});
if (context.agent_name && !isMock) {
void getAPIClient()
.threads.update(meta.thread_id, {
metadata: { agent_name: context.agent_name },
})
.catch(() => ({}));
}
},
onLangChainEvent(event) {
if (event.event === "on_tool_end") {
listeners.current.onToolEnd?.({
name: event.name,
data: event.data,
});
}
},
onUpdateEvent(data) {
const _messages = getSummarizationMiddlewareMessages(data);
if (_messages && _messages.length >= 2) {
for (const m of _messages) {
// Backward-compat shim: pre-PR2 threads may still carry a synthetic
// HumanMessage(name="summary") from the old summarization path. New
// threads keep the summary in ThreadState.summary_text instead.
if (m.name === "summary" && m.type === "human") {
summarizedRef.current?.add(m.id ?? "");
}
}
const _movedMessages = computeSummarizationMovedMessages(
messagesRef.current,
_messages,
summarizedRef.current ?? new Set<string>(),
);
// Buffer the rescued messages synchronously so the merge can keep
// displaying them immediately, even though appendMessages below only
// updates the archived-history state asynchronously (#3825).
pendingArchivedMessagesRef.current = dedupeMessagesByIdentity([
...pendingArchivedMessagesRef.current,
..._movedMessages,
]);
pendingArchiveThreadIdRef.current = threadIdRef.current;
appendMessages(_movedMessages);
messagesRef.current = [];
}
const updates: Array<Partial<AgentThreadState> | null> = Object.values(
data || {},
);
for (const update of updates) {
if (update && "title" in update && update.title) {
void queryClient.setQueriesData(
{
queryKey: ["threads", "search"],
exact: false,
},
(oldData: Array<AgentThread> | undefined) => {
return oldData?.map((t) => {
if (t.thread_id === threadIdRef.current) {
return {
...t,
values: {
...t.values,
title: update.title,
},
};
}
return t;
});
},
);
const nextTitle: string = update.title;
void queryClient.setQueriesData(
{
queryKey: INFINITE_THREADS_QUERY_KEY_PREFIX,
exact: false,
},
(oldData: InfiniteData<AgentThread[]> | undefined) =>
mapInfiniteThreadsCache(
oldData,
(t): AgentThread =>
t.thread_id === threadIdRef.current
? {
...t,
values: {
...t.values,
title: nextTitle,
},
}
: t,
),
);
}
}
},
onCustomEvent(event: unknown) {
if (
typeof event === "object" &&
event !== null &&
"type" in event &&
event.type === "task_running"
) {
const e = event as {
type: "task_running";
task_id: string;
message: AIMessage;
message_index?: number;
};
// Accumulate the full step history instead of overwriting (#3779): keep
// latestMessage for the collapsed-header tool-call hint, and append the
// normalized step (assistant turn or tool output) to the timeline.
updateSubtask({
id: e.task_id,
latestMessage: e.message,
steps: [messageToStep(e.message, e.message_index ?? 0)],
});
return;
}
if (
typeof event === "object" &&
event !== null &&
"type" in event &&
event.type === "llm_retry" &&
"message" in event &&
typeof event.message === "string" &&
event.message.trim()
) {
const e = event as { type: "llm_retry"; message: string };
toast(e.message);
}
},
onError(error) {
setOptimisticMessages([]);
setOptimisticThreadId(null);
setLiveMessagesThreadId(null);
setPendingSupersededRunIds(new Set());
setPendingSupersededMessageIds(new Set());
toast.error(getStreamErrorMessage(error));
pendingUsageBaselineMessageIdsRef.current = new Set(
messagesRef.current
.map(messageIdentity)
.filter((id): id is string => Boolean(id)),
);
if (threadIdRef.current && !isMock) {
void queryClient.invalidateQueries({
queryKey: threadTokenUsageQueryKey(threadIdRef.current),
});
}
},
onFinish(state) {
listeners.current.onFinish?.(state.values);
pendingUsageBaselineMessageIdsRef.current = new Set(
messagesRef.current
.map(messageIdentity)
.filter((id): id is string => Boolean(id)),
);
invalidateStoppedThreadCaches(queryClient, threadIdRef.current, isMock);
},
});
const stopThread = useCallback(async () => {
const stoppedThreadId =
threadIdRef.current ?? displayThreadId ?? threadId ?? null;
await stopThreadAndInvalidateCaches(
queryClient,
() => thread.stop(),
stoppedThreadId,
isMock,
);
}, [displayThreadId, isMock, queryClient, thread, threadId]);
const hasVisibleStreamState =
Boolean(threadId) || liveMessagesThreadId === currentViewThreadId;
const persistedMessages = useMemo(
() =>
hasVisibleStreamState
? thread.messages.filter(
(message) =>
!message.id || !pendingSupersededMessageIds.has(message.id),
)
: [],
[hasVisibleStreamState, pendingSupersededMessageIds, thread.messages],
);
const visibleHistory = useMemo(
() => (threadId ? history : []),
[history, threadId],
);
const humanMessageCount = persistedMessages.filter(
(m) => m.type === "human",
).length;
const latestMessageCountsRef = useRef({ humanMessageCount });
const sendInFlightRef = useRef(false);
const messagesRef = useRef<Message[]>([]);
// Synchronous bridge for messages rescued from context summarization. The
// archived-history `setState` (via appendMessages) lands on a different
// schedule than the live thread external store, so the merge reads this buffer
// to avoid dropping rescued messages in the render window before history
// catches up (#3825).
const pendingArchivedMessagesRef = useRef<Message[]>([]);
// The thread the rescue buffer belongs to, captured when onUpdateEvent fills
// it. The merge only overlays the buffer when this matches the viewed
// `threadId`, so a previous thread's rescued messages can never flash into
// another thread or the new-chat screen (#3825).
const pendingArchiveThreadIdRef = useRef<string | null>(null);
const summarizedRef = useRef<Set<string>>(null);
// Track human message count before sending to prevent clearing optimistic
// messages before the server's human message arrives (e.g. when AI messages
// from "messages-tuple" events arrive before the input human message from
// "values" events).
const prevHumanMsgCountRef = useRef(humanMessageCount);
latestMessageCountsRef.current = { humanMessageCount };
summarizedRef.current ??= new Set<string>();
// Reset thread-local pending UI state when switching between threads so
// optimistic messages and in-flight guards do not leak across chat views.
useEffect(() => {
startedRef.current = false;
sendInFlightRef.current = false;
messagesRef.current = [];
pendingArchivedMessagesRef.current = [];
pendingArchiveThreadIdRef.current = null;
summarizedRef.current = new Set<string>();
pendingUsageBaselineMessageIdsRef.current = new Set();
setPendingSupersededRunIds(new Set());
setPendingSupersededMessageIds(new Set());
prevHumanMsgCountRef.current =
latestMessageCountsRef.current.humanMessageCount;
}, [threadId]);
// Release archive-buffer entries once the canonical history state has absorbed
// them, so the synchronous bridge stays transient and never resurrects a
// message that history later filters out (e.g. a superseded run) (#3825).
useEffect(() => {
pendingArchivedMessagesRef.current = pruneConfirmedArchivedMessages(
pendingArchivedMessagesRef.current,
visibleHistory,
);
}, [visibleHistory]);
useEffect(() => {
if (optimisticThreadId && optimisticThreadId !== currentViewThreadId) {
setOptimisticMessages([]);
setOptimisticThreadId(null);
}
if (liveMessagesThreadId && liveMessagesThreadId !== currentViewThreadId) {
setLiveMessagesThreadId(null);
}
}, [currentViewThreadId, liveMessagesThreadId, optimisticThreadId]);
// When streaming starts without a baseline (e.g. reconnection, run started
// from another client, or page reload mid-stream), snapshot the current
// messages so only *new* messages are treated as "pending" for token usage.
useEffect(() => {
if (
thread.isLoading &&
pendingUsageBaselineMessageIdsRef.current.size === 0
) {
pendingUsageBaselineMessageIdsRef.current = new Set(
persistedMessages
.map(messageIdentity)
.filter((id): id is string => Boolean(id)),
);
}
}, [persistedMessages, thread.isLoading]);
// Clear optimistic when server messages arrive.
// For messages with a human optimistic message, wait until the server's
// human message has arrived to avoid clearing before the input message
// appears in the stream (the input message may arrive via "values" events
// after individual "messages-tuple" events for AI messages).
const optimisticMessageCount = optimisticMessages.length;
const hasHumanOptimistic = optimisticMessages.some((m) => m.type === "human");
useEffect(() => {
if (optimisticMessageCount === 0) return;
const newHumanMsgArrived = humanMessageCount > prevHumanMsgCountRef.current;
if (!hasHumanOptimistic || newHumanMsgArrived) {
setOptimisticMessages([]);
setOptimisticThreadId(null);
}
}, [hasHumanOptimistic, humanMessageCount, optimisticMessageCount]);
const sendMessage = useCallback(
async (
threadId: string,
message: PromptInputMessage,
extraContext?: Record<string, unknown>,
options?: SendMessageOptions,
) => {
if (sendInFlightRef.current) {
return;
}
sendInFlightRef.current = true;
// The send has genuinely proceeded past the in-flight guard, so callers
// can now run one-time cleanup that must not fire on the dropped path.
options?.onSent?.();
const text = message.text.trim();
// Capture the current human message count before showing optimistic
// messages so we can wait for the server's copy of the user input.
prevHumanMsgCountRef.current = humanMessageCount;
pendingUsageBaselineMessageIdsRef.current = new Set(
persistedMessages
.map(messageIdentity)
.filter((id): id is string => Boolean(id)),
);
// Build optimistic files list with uploading status
const optimisticFiles: FileInMessage[] = (message.files ?? []).map(
(f) => ({
filename: f.filename ?? "",
size: 0,
status: "uploading" as const,
}),
);
const hideFromUI = options?.additionalKwargs?.hide_from_ui === true;
const optimisticAdditionalKwargs = {
...options?.additionalKwargs,
...(optimisticFiles.length > 0 ? { files: optimisticFiles } : {}),
};
const newOptimistic: Message[] = [];
if (!hideFromUI) {
newOptimistic.push({
type: "human",
id: `opt-human-${Date.now()}`,
content: text ? [{ type: "text", text }] : "",
additional_kwargs: optimisticAdditionalKwargs,
});
}
if (optimisticFiles.length > 0 && !hideFromUI) {
// Mock AI message while files are being uploaded
newOptimistic.push({
type: "ai",
id: `opt-ai-${Date.now()}`,
content: t.uploads.uploadingFiles,
additional_kwargs: { element: "task" },
});
}
setOptimisticThreadId(threadId);
setLiveMessagesThreadId(threadId);
setOptimisticMessages(newOptimistic);
listeners.current.onSend?.(threadId);
let uploadedFileInfo: UploadedFileInfo[] = [];
try {
// Upload files first if any
if (message.files && message.files.length > 0) {
setIsUploading(true);
try {
const filePromises = message.files.map((fileUIPart) =>
promptInputFilePartToFile(fileUIPart),
);
const conversionResults = await Promise.all(filePromises);
const files = conversionResults.filter(
(file): file is File => file !== null,
);
const failedConversions = conversionResults.length - files.length;
if (failedConversions > 0) {
throw new Error(
`Failed to prepare ${failedConversions} attachment(s) for upload. Please retry.`,
);
}
if (!threadId) {
throw new Error("Thread is not ready for file upload.");
}
if (files.length > 0) {
const uploadResponse = await uploadFiles(threadId, files);
uploadedFileInfo = uploadResponse.files;
// Update optimistic human message with uploaded status + paths
const uploadedFiles: FileInMessage[] = uploadedFileInfo.map(
(info) => ({
filename: info.filename,
size: info.size,
path: info.virtual_path,
status: "uploaded" as const,
}),
);
setOptimisticMessages((messages) => {
if (messages.length > 1 && messages[0]) {
const humanMessage: Message = messages[0];
return [
{
...humanMessage,
additional_kwargs: { files: uploadedFiles },
},
...messages.slice(1),
];
}
return messages;
});
}
} catch (error) {
const errorMessage =
error instanceof Error
? error.message
: "Failed to upload files.";
toast.error(errorMessage);
setOptimisticMessages([]);
setOptimisticThreadId(null);
setLiveMessagesThreadId(null);
throw error;
} finally {
setIsUploading(false);
}
}
// Build files metadata for submission (included in additional_kwargs)
const filesForSubmit: FileInMessage[] = uploadedFileInfo.map(
(info) => ({
filename: info.filename,
size: info.size,
path: info.virtual_path,
status: "uploaded" as const,
}),
);
await thread.submit(
{
messages: buildThreadSubmitMessages({
text,
additionalKwargs: options?.additionalKwargs,
additionalInputMessages: options?.additionalInputMessages,
filesForSubmit,
}),
},
{
threadId: threadId,
streamSubgraphs: true,
streamResumable: true,
config: {
recursion_limit: 1000,
},
context: {
...extraContext,
...context,
thinking_enabled: context.mode !== "flash",
is_plan_mode: context.mode === "pro" || context.mode === "ultra",
subagent_enabled: context.mode === "ultra",
reasoning_effort:
context.reasoning_effort ??
(context.mode === "ultra"
? "high"
: context.mode === "pro"
? "medium"
: context.mode === "thinking"
? "low"
: undefined),
thread_id: threadId,
},
},
);
void queryClient.invalidateQueries({ queryKey: ["threads", "search"] });
void queryClient.invalidateQueries({
queryKey: INFINITE_THREADS_QUERY_KEY_PREFIX,
});
} catch (error) {
setOptimisticMessages([]);
setOptimisticThreadId(null);
setLiveMessagesThreadId(null);
setIsUploading(false);
throw error;
} finally {
sendInFlightRef.current = false;
}
},
[
thread,
t.uploads.uploadingFiles,
context,
queryClient,
humanMessageCount,
persistedMessages,
],
);
const regenerateMessage = useCallback(
async (
threadId: string,
messageId: string,
supersededMessageIds: string[] = [messageId],
) => {
if (sendInFlightRef.current || !threadId || !messageId) {
return;
}
sendInFlightRef.current = true;
prevHumanMsgCountRef.current = humanMessageCount;
pendingUsageBaselineMessageIdsRef.current = new Set(
persistedMessages
.map(messageIdentity)
.filter((id): id is string => Boolean(id)),
);
setLiveMessagesThreadId(threadId);
listeners.current.onSend?.(threadId);
let preparedSupersededRunId: string | null = null;
let preparedSupersededMessageIds: string[] = [];
try {
const response = await fetch(
`${getBackendBaseURL()}/api/threads/${encodeURIComponent(
threadId,
)}/runs/regenerate/prepare`,
{
method: "POST",
headers: {
"Content-Type": "application/json",
},
credentials: "include",
body: JSON.stringify({ message_id: messageId }),
},
);
if (!response.ok) {
throw new Error(await readResponseErrorMessage(response));
}
const prepared = (await response.json()) as RegeneratePrepareResponse;
preparedSupersededRunId = prepared.target_run_id;
preparedSupersededMessageIds = supersededMessageIds;
setPendingSupersededRunIds((current) => {
const next = new Set(current);
next.add(prepared.target_run_id);
return next;
});
setPendingSupersededMessageIds((current) => {
const next = new Set(current);
for (const id of supersededMessageIds) {
next.add(id);
}
return next;
});
await thread.submit(prepared.input, {
threadId,
checkpoint: prepared.checkpoint,
metadata: prepared.metadata,
streamSubgraphs: true,
streamResumable: true,
config: {
recursion_limit: 1000,
},
context: {
...context,
thinking_enabled: context.mode !== "flash",
is_plan_mode: context.mode === "pro" || context.mode === "ultra",
subagent_enabled: context.mode === "ultra",
reasoning_effort:
context.reasoning_effort ??
(context.mode === "ultra"
? "high"
: context.mode === "pro"
? "medium"
: context.mode === "thinking"
? "low"
: undefined),
thread_id: threadId,
},
});
void queryClient.invalidateQueries({ queryKey: ["thread", threadId] });
void queryClient.invalidateQueries({ queryKey: ["threads", "search"] });
void queryClient.invalidateQueries({
queryKey: INFINITE_THREADS_QUERY_KEY_PREFIX,
});
void queryClient.invalidateQueries({
queryKey: threadTokenUsageQueryKey(threadId),
});
} catch (error) {
setLiveMessagesThreadId(null);
if (preparedSupersededRunId) {
const supersededRunId = preparedSupersededRunId;
setPendingSupersededRunIds((current) =>
removeSetItems(current, [supersededRunId]),
);
setPendingSupersededMessageIds((current) =>
removeSetItems(current, preparedSupersededMessageIds),
);
}
toast.error(getStreamErrorMessage(error));
} finally {
sendInFlightRef.current = false;
}
},
[context, humanMessageCount, persistedMessages, queryClient, thread],
);
// Cache the latest thread messages in a ref to compare against incoming history messages for deduplication,
// and to allow access to the full message list in onUpdateEvent without causing re-renders.
if (persistedMessages.length >= messagesRef.current.length) {
messagesRef.current = persistedMessages;
}
const visibleOptimisticMessages = getVisibleOptimisticMessages(
optimisticThreadId === currentViewThreadId ? optimisticMessages : [],
prevHumanMsgCountRef.current,
humanMessageCount,
);
// Overlay the summarization rescue buffer only onto the history of the thread
// it was captured from. visibleHistory is gated on `threadId`, so comparing the
// same prop keeps the buffer from flashing into another thread or the new-chat
// screen, and reading it here (instead of clearing a ref during render) is
// concurrent-mode safe (#3825).
const rescueBuffer = pendingArchivedMessagesRef.current;
const effectiveHistory =
rescueBuffer.length > 0 && pendingArchiveThreadIdRef.current === threadId
? resolvePreservedHistory(visibleHistory, rescueBuffer)
: visibleHistory;
const mergedMessages = mergeMessages(
effectiveHistory,
persistedMessages,
visibleOptimisticMessages,
);
const pendingUsageMessages = thread.isLoading
? getMessagesAfterBaseline(
persistedMessages,
pendingUsageBaselineMessageIdsRef.current,
)
: [];
// Merge history, live stream, and optimistic messages for display
// History messages may overlap with thread.messages; thread.messages take precedence
const mergedThread = {
...thread,
stop: stopThread,
values: hasVisibleStreamState ? thread.values : EMPTY_THREAD_VALUES,
messages: mergedMessages,
} as typeof thread;
return {
thread: mergedThread,
pendingUsageMessages,
sendMessage,
regenerateMessage,
isUploading,
isHistoryLoading,
hasMoreHistory,
loadMoreHistory,
} as const;
}
type ThreadHistoryOptions = {
enabled?: boolean;
pendingSupersededRunIds?: ReadonlySet<string>;
};
export function useThreadHistory(
threadId: string,
{ enabled = true, pendingSupersededRunIds }: ThreadHistoryOptions = {},
) {
const runs = useThreadRuns(threadId, { enabled });
const threadIdRef = useRef(threadId);
const runsRef = useRef(runs.data ?? []);
const indexRef = useRef(-1);
const loadingRef = useRef(false);
const pendingLoadRef = useRef(false);
const loadingRunIdRef = useRef<string | null>(null);
const loadedRunIdsRef = useRef<Set<string>>(new Set());
const runBeforeSeqRef = useRef<Map<string, number>>(new Map());
const loadGenerationRef = useRef(0);
const [loading, setLoading] = useState(false);
const [messageRows, setMessageRows] = useState<RunMessage[]>([]);
const [appendedMessages, setAppendedMessages] = useState<Message[]>([]);
const supersededRunIds = useMemo(() => {
return getSupersededRunIds(runs.data, pendingSupersededRunIds);
}, [pendingSupersededRunIds, runs.data]);
const messages = useMemo(() => {
return buildVisibleHistoryMessages(
messageRows,
supersededRunIds,
appendedMessages,
);
}, [appendedMessages, messageRows, supersededRunIds]);
const loadMessages = useCallback(async () => {
if (!enabled) {
return;
}
const loadGeneration = loadGenerationRef.current;
if (loadingRef.current) {
const pendingRunIndex = findLatestUnloadedRunIndex(
runsRef.current,
loadedRunIdsRef.current,
);
const pendingRun = runsRef.current[pendingRunIndex];
if (pendingRun && pendingRun.run_id !== loadingRunIdRef.current) {
pendingLoadRef.current = true;
}
return;
}
if (runsRef.current.length === 0) {
return;
}
loadingRef.current = true;
setLoading(true);
try {
let consecutiveEmptyLoads = 0;
do {
pendingLoadRef.current = false;
const nextRunIndex = findLatestUnloadedRunIndex(
runsRef.current,
loadedRunIdsRef.current,
);
indexRef.current = nextRunIndex;
const run = runsRef.current[nextRunIndex];
if (!run) {
indexRef.current = -1;
return;
}
const requestThreadId = threadIdRef.current;
loadingRunIdRef.current = run.run_id;
const beforeSeq = runBeforeSeqRef.current.get(run.run_id);
const url = buildRunMessagesUrl(
getBackendBaseURL(),
requestThreadId,
run.run_id,
beforeSeq,
);
const result: RunMessagesPageResponse = await fetch(url, {
method: "GET",
headers: {
"Content-Type": "application/json",
},
credentials: "include",
}).then((res) => {
return res.json();
});
if (
loadGenerationRef.current !== loadGeneration ||
threadIdRef.current !== requestThreadId
) {
return;
}
const _messages = result.data.filter(
(m) => !m.metadata.caller?.startsWith("middleware:"),
);
setMessageRows((prev) =>
dedupeRunMessagesByIdentity([..._messages, ...prev]),
);
const nextBeforeSeq = getNextRunMessagesBeforeSeq(result);
if (typeof nextBeforeSeq === "number") {
runBeforeSeqRef.current.set(run.run_id, nextBeforeSeq);
pendingLoadRef.current = true;
} else if (nextBeforeSeq === undefined) {
console.warn(
`Run ${run.run_id} returned has_more without message seq values; leaving it pending for retry.`,
);
} else {
runBeforeSeqRef.current.delete(run.run_id);
loadedRunIdsRef.current.add(run.run_id);
if (
shouldAutoContinueOnEmptyRun(
_messages.length,
consecutiveEmptyLoads,
)
) {
consecutiveEmptyLoads += 1;
pendingLoadRef.current = true;
} else {
consecutiveEmptyLoads = 0;
}
}
indexRef.current = findLatestUnloadedRunIndex(
runsRef.current,
loadedRunIdsRef.current,
);
} while (pendingLoadRef.current);
} catch (err) {
console.error(err);
} finally {
if (loadGenerationRef.current === loadGeneration) {
loadingRef.current = false;
loadingRunIdRef.current = null;
setLoading(false);
}
}
}, [enabled]);
useEffect(() => {
const threadChanged = threadIdRef.current !== threadId;
threadIdRef.current = threadId;
if (!enabled || threadChanged) {
loadGenerationRef.current += 1;
runsRef.current = [];
indexRef.current = -1;
pendingLoadRef.current = false;
loadingRunIdRef.current = null;
loadedRunIdsRef.current = new Set();
runBeforeSeqRef.current = new Map();
loadingRef.current = false;
setLoading(false);
setMessageRows([]);
setAppendedMessages([]);
}
if (!enabled) {
return;
}
if (runs.data && runs.data.length > 0) {
runsRef.current = runs.data ?? [];
indexRef.current = findLatestUnloadedRunIndex(
runs.data,
loadedRunIdsRef.current,
);
}
loadMessages().catch(() => {
toast.error("Failed to load thread history.");
});
}, [enabled, threadId, runs.data, loadMessages]);
const appendMessages = useCallback((_messages: Message[]) => {
setAppendedMessages((prev) => {
return dedupeMessagesByIdentity([...prev, ..._messages]);
});
}, []);
const hasThreadId = Boolean(threadId);
const hasUnloadedRuns = Boolean(
runs.data?.some((run) => !loadedRunIdsRef.current.has(run.run_id)),
);
const isRunsLoading =
enabled &&
hasThreadId &&
(runs.isLoading || (runs.isFetching && !runs.data));
const isRunsUnresolved =
enabled && hasThreadId && !runs.data && !runs.isError;
const hasMore =
enabled && hasThreadId && (indexRef.current >= 0 || hasUnloadedRuns);
return {
runs: runs.data,
messages,
loading: loading || isRunsLoading || isRunsUnresolved,
appendMessages,
hasMore,
loadMore: loadMessages,
};
}
export function useThreads(
params: ThreadSearchParams = DEFAULT_THREAD_SEARCH_PARAMS,
) {
const apiClient = getAPIClient();
return useQuery<AgentThread[]>({
...buildThreadsSearchQueryOptions(apiClient, params),
});
}
export const INFINITE_THREADS_PAGE_SIZE = 50;
export const INFINITE_THREADS_QUERY_KEY_PREFIX = [
"threads",
"searchInfinite",
] as const;
const INFINITE_THREADS_NEXT_PAGE_PARAM = Symbol(
"deerflow.infiniteThreads.nextPageParam",
);
type InfiniteThreadsParams = Omit<
Parameters<ThreadsClient["search"]>[0],
"limit" | "offset"
>;
type InfiniteThreadsSearchClient = {
threads: {
search: ThreadsClient["search"];
};
};
type InfiniteThreadsPageWithNextParam = AgentThread[] & {
[INFINITE_THREADS_NEXT_PAGE_PARAM]?: number;
};
function annotateInfiniteThreadsPage(
page: AgentThread[],
nextPageParam: number | undefined,
): AgentThread[] {
if (nextPageParam !== undefined) {
Reflect.set(page, INFINITE_THREADS_NEXT_PAGE_PARAM, nextPageParam);
}
return page;
}
export async function fetchInfiniteThreadsPage(
apiClient: InfiniteThreadsSearchClient,
params: InfiniteThreadsParams,
pageParam: number,
pageSize: number = INFINITE_THREADS_PAGE_SIZE,
): Promise<AgentThread[]> {
const threads: AgentThread[] = [];
let offset = pageParam;
let nextPageParam: number | undefined;
while (threads.length < pageSize) {
const currentLimit = pageSize - threads.length;
const response = (await apiClient.threads.search<AgentThreadState>({
...params,
limit: currentLimit,
offset,
})) as AgentThread[];
threads.push(...filterThreadSearchResults(response, params));
offset += response.length;
if (response.length < currentLimit) {
nextPageParam = undefined;
break;
}
nextPageParam = offset;
}
return annotateInfiniteThreadsPage(threads, nextPageParam);
}
export function getInfiniteThreadsNextPageParam(
lastPage: AgentThread[],
allPages: AgentThread[][],
pageSize: number = INFINITE_THREADS_PAGE_SIZE,
): number | undefined {
const annotatedNextPageParam = Reflect.get(
lastPage as InfiniteThreadsPageWithNextParam,
INFINITE_THREADS_NEXT_PAGE_PARAM,
);
if (typeof annotatedNextPageParam === "number") {
return annotatedNextPageParam;
}
if (lastPage.length < pageSize) {
return undefined;
}
return allPages.reduce((sum, page) => sum + page.length, 0);
}
export function mapInfiniteThreadsCache(
oldData: InfiniteData<AgentThread[]> | undefined,
mapper: (thread: AgentThread) => AgentThread,
): InfiniteData<AgentThread[]> | undefined {
if (!oldData) {
return oldData;
}
return {
...oldData,
pages: oldData.pages.map((page) => page.map(mapper)),
};
}
export function filterInfiniteThreadsCache(
oldData: InfiniteData<AgentThread[]> | undefined,
predicate: (thread: AgentThread) => boolean,
): InfiniteData<AgentThread[]> | undefined {
if (!oldData) {
return oldData;
}
return {
...oldData,
pages: oldData.pages.map((page) => page.filter(predicate)),
};
}
export function useInfiniteThreads(
params: InfiniteThreadsParams = {
sortBy: "updated_at",
sortOrder: "desc",
select: ["thread_id", "updated_at", "values", "metadata"],
},
) {
const apiClient = getAPIClient();
return useInfiniteQuery<
AgentThread[],
Error,
InfiniteData<AgentThread[]>,
readonly unknown[],
number
>({
queryKey: [...INFINITE_THREADS_QUERY_KEY_PREFIX, params],
initialPageParam: 0,
queryFn: async ({ pageParam }) =>
fetchInfiniteThreadsPage(
apiClient,
params,
pageParam,
INFINITE_THREADS_PAGE_SIZE,
),
getNextPageParam: (lastPage, allPages) =>
getInfiniteThreadsNextPageParam(lastPage, allPages),
refetchOnWindowFocus: false,
});
}
export function useThreadRuns(
threadId?: string,
{ enabled = true }: { enabled?: boolean } = {},
) {
const apiClient = getAPIClient();
return useQuery<Run[]>({
queryKey: ["thread", threadId],
queryFn: async () => {
if (!threadId) {
return [];
}
const response = await apiClient.runs.list(threadId);
return response;
},
enabled: enabled && Boolean(threadId),
refetchOnWindowFocus: false,
});
}
export function useThreadMetadata(
threadId?: string | null,
{
enabled = true,
isMock = false,
}: { enabled?: boolean; isMock?: boolean } = {},
) {
const apiClient = getAPIClient(isMock);
return useQuery<AgentThread | null>({
queryKey: ["thread", "metadata", threadId, isMock],
queryFn: async () => {
if (!threadId) {
return null;
}
try {
const response = await apiClient.threads.get(threadId);
return response as AgentThread;
} catch (error) {
if (isThreadMissingError(error)) {
return null;
}
throw error;
}
},
enabled: enabled && Boolean(threadId),
retry: false,
refetchOnWindowFocus: false,
});
}
export function useThreadTokenUsage(
threadId?: string | null,
{ enabled = true }: { enabled?: boolean } = {},
) {
return useQuery<ThreadTokenUsageResponse | null>({
queryKey: threadTokenUsageQueryKey(threadId),
queryFn: async () => {
if (!threadId) {
return null;
}
return fetchThreadTokenUsage(threadId);
},
enabled: enabled && Boolean(threadId),
retry: false,
refetchOnWindowFocus: false,
});
}
export function useRunDetail(threadId: string, runId: string) {
const apiClient = getAPIClient();
return useQuery<Run>({
queryKey: ["thread", threadId, "run", runId],
queryFn: async () => {
const response = await apiClient.runs.get(threadId, runId);
return response;
},
refetchOnWindowFocus: false,
});
}
async function deleteLocalThreadData(threadId: string) {
const response = await fetch(
`${getBackendBaseURL()}/api/threads/${encodeURIComponent(threadId)}`,
{
method: "DELETE",
},
);
if (!response.ok) {
const error = await response
.json()
.catch(() => ({ detail: "Failed to delete local thread data." }));
throw new Error(error.detail ?? "Failed to delete local thread data.");
}
}
async function deleteThreadEverywhere(
apiClient: ThreadDeleteClient,
threadId: string,
) {
await apiClient.threads.delete(threadId);
await deleteLocalThreadData(threadId);
}
export async function findSidecarThreadIdsForParent(
apiClient: ThreadSidecarSearchClient,
parentThreadId: string,
) {
const threadIds: string[] = [];
const limit = 100;
let offset = 0;
while (true) {
const response = await apiClient.threads.search({
metadata: {
[SIDECAR_METADATA_KEY]: true,
parent_thread_id: parentThreadId,
},
limit,
offset,
sortBy: "updated_at",
sortOrder: "desc",
select: ["thread_id", "metadata"],
});
for (const thread of response) {
if (
isSidecarThread(thread) &&
thread.metadata?.parent_thread_id === parentThreadId
) {
threadIds.push(thread.thread_id);
}
}
if (response.length < limit) {
break;
}
offset += response.length;
}
return threadIds;
}
async function deleteSidecarThreadsForParent(
apiClient: ThreadDeleteClient,
parentThreadId: string,
) {
let sidecarThreadIds: string[];
try {
sidecarThreadIds = await findSidecarThreadIdsForParent(
apiClient,
parentThreadId,
);
} catch (err) {
console.warn(
`Failed to look up sidecar threads for parent ${parentThreadId}; skipping cascade cleanup. Orphaned sidecar threads may remain.`,
err,
);
return [];
}
const results = await Promise.allSettled(
sidecarThreadIds.map((threadId) =>
deleteThreadEverywhere(apiClient, threadId),
),
);
const failedDeletions = results
.map((result, index) =>
result.status === "rejected"
? { threadId: sidecarThreadIds[index], reason: result.reason }
: null,
)
.filter((entry): entry is { threadId: string; reason: unknown } =>
Boolean(entry),
);
if (failedDeletions.length > 0) {
console.warn(
`Failed to delete ${failedDeletions.length} sidecar thread(s) for parent ${parentThreadId}; orphaned sidecar threads may remain.`,
failedDeletions,
);
}
return sidecarThreadIds.filter((_, index) => {
return results[index]?.status === "fulfilled";
});
}
export function useDeleteThread() {
const queryClient = useQueryClient();
const apiClient = getAPIClient() as ThreadDeleteClient;
return useMutation({
mutationFn: async ({
threadId,
onRemoteDeleted,
}: {
threadId: string;
onRemoteDeleted?: () => void;
}) => {
const deletedSidecarThreadIds = await deleteSidecarThreadsForParent(
apiClient,
threadId,
);
await apiClient.threads.delete(threadId);
onRemoteDeleted?.();
await deleteLocalThreadData(threadId);
return deletedSidecarThreadIds;
},
onSuccess(deletedSidecarThreadIds, { threadId }) {
const deletedThreadIds = new Set([threadId, ...deletedSidecarThreadIds]);
queryClient.setQueriesData(
{
queryKey: ["threads", "search"],
exact: false,
},
(oldData: Array<AgentThread> | undefined) => {
if (oldData == null) {
return oldData;
}
return oldData.filter((t) => !deletedThreadIds.has(t.thread_id));
},
);
queryClient.setQueriesData(
{
queryKey: INFINITE_THREADS_QUERY_KEY_PREFIX,
exact: false,
},
(oldData: InfiniteData<AgentThread[]> | undefined) =>
filterInfiniteThreadsCache(
oldData,
(t) => !deletedThreadIds.has(t.thread_id),
),
);
},
onSettled() {
void queryClient.invalidateQueries({ queryKey: ["threads", "search"] });
void queryClient.invalidateQueries({
queryKey: INFINITE_THREADS_QUERY_KEY_PREFIX,
});
},
});
}
export function useRenameThread() {
const queryClient = useQueryClient();
const apiClient = getAPIClient();
return useMutation({
mutationFn: async ({
threadId,
title,
}: {
threadId: string;
title: string;
}) => {
await apiClient.threads.updateState(threadId, {
values: { title },
});
},
onSuccess(_, { threadId, title }) {
queryClient.setQueriesData(
{
queryKey: ["threads", "search"],
exact: false,
},
(oldData: Array<AgentThread>) => {
return oldData.map((t) => {
if (t.thread_id === threadId) {
return {
...t,
values: {
...t.values,
title,
},
};
}
return t;
});
},
);
queryClient.setQueriesData(
{
queryKey: INFINITE_THREADS_QUERY_KEY_PREFIX,
exact: false,
},
(oldData: InfiniteData<AgentThread[]> | undefined) =>
mapInfiniteThreadsCache(oldData, (t) =>
t.thread_id === threadId
? {
...t,
values: {
...t.values,
title,
},
}
: t,
),
);
},
});
}