Backend - chat.ts now routes both regular chat and QuickChat through createResolvedAgentSession instead of branching to createFnAgent for the no-runtime-hint case. This removes the divergent path where pi-ai's cleanupSessionResources(sessionId) could tear down resources the next generation depends on. - sendMessage's finally only disposes the agent if it still owns the activeGenerations slot. A newer generation that has pre-empted us cleans up its own agent in its own finally — disposing here would yank the underlying CLI process out from under it. - __setCreateFnAgent test helper now mirrors its mock into the createResolvedAgentSession slot so existing test setups still work after the unification. Frontend - Extract createChatStreamHandlers (RAF coalescing, accumulators, tool-call dedup, fallback handling) — useChat and useQuickChat were duplicating ~85 LOC each. Both now compose the shared factory. - Move shared chat types into chatTypes.ts. The hooks re-export them for backward compatibility with existing consumers. - Removed per-message Markdown/plain-text eye toggles. A single thread-level toggle in the chat header now flips every assistant bubble (including the streaming one) between rendered Markdown and plain text. - Model-only chats hide the per-message agent identity row entirely; the model name is already in the thread header. Tests - Updated ChatView tests to reflect the new render-toggle contract (single header toggle drives all bubbles) and the model-only avatar suppression. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
213 lines
7.5 KiB
TypeScript
213 lines
7.5 KiB
TypeScript
import type { ChatMessage } from "@fusion/core";
|
|
import type { Dispatch, RefObject, SetStateAction } from "react";
|
|
import type { ChatMessageInfo, FallbackInfo, ToolCallInfo } from "./chatTypes";
|
|
|
|
/**
|
|
* Inputs for the chat streaming-handler factory.
|
|
*
|
|
* The shared factory owns the per-stream accumulator state (text, thinking,
|
|
* tool calls, fallback info), the requestAnimationFrame coalescing of state
|
|
* updates, and the SSE event → state-setter wiring. Caller-specific behaviour
|
|
* for the terminal events (`onDone`, `onError`) and the optional
|
|
* `onFallbackSession` model-swap is provided through callbacks so that
|
|
* `useChat` and `useQuickChat` can plug in their own session-management
|
|
* semantics without re-implementing the streaming machinery.
|
|
*/
|
|
export interface CreateChatStreamHandlersOptions {
|
|
/** Active session id — used by `onFallbackSession` for parent-side updates. */
|
|
sessionId: string;
|
|
/** Optimistic temp id of the user message added before the stream started. */
|
|
tempUserMessageId: string;
|
|
/**
|
|
* The latest text/thinking/tool-call snapshots that are committed to React
|
|
* state. We pass setters (not values) so the factory can flush per-frame
|
|
* without rerunning the parent's effects.
|
|
*/
|
|
setStreamingText: Dispatch<SetStateAction<string>>;
|
|
setStreamingThinking: Dispatch<SetStateAction<string>>;
|
|
setStreamingToolCalls: Dispatch<SetStateAction<ToolCallInfo[]>>;
|
|
/**
|
|
* Caller-side `cancelStreamingFlushes` ref slot. The factory writes its own
|
|
* cancel function here so `stopStreaming` (in either parent hook) can call
|
|
* it to abort pending RAF flushes regardless of which sendMessage owns them.
|
|
*/
|
|
cancelStreamingFlushesRef: RefObject<(() => void) | null>;
|
|
/** Optional toast helper, used to surface fallback-model warnings + errors. */
|
|
addToast?: (message: string, level: "error" | "warning" | "success") => void;
|
|
/** Caller-supplied terminal handlers — bind in their own state setters. */
|
|
onDone: (data: {
|
|
messageId: string;
|
|
message?: ChatMessage;
|
|
accumulated: {
|
|
text: string;
|
|
thinking: string;
|
|
toolCalls: ToolCallInfo[];
|
|
fallbackInfo?: FallbackInfo;
|
|
};
|
|
}) => void;
|
|
onError: (data: string, tempUserMessageId: string) => void;
|
|
/**
|
|
* Fallback-model side effect for the parent (e.g. updating the session list
|
|
* or the active session's model fields). The factory still emits the toast.
|
|
*/
|
|
onFallbackSession?: (data: FallbackInfo, sessionId: string) => void;
|
|
}
|
|
|
|
export interface ChatStreamHandlers {
|
|
onThinking: (delta: string) => void;
|
|
onText: (delta: string) => void;
|
|
onToolStart: (data: { toolName: string; args?: Record<string, unknown> }) => void;
|
|
onToolEnd: (data: { toolName: string; isError: boolean; result?: unknown }) => void;
|
|
onFallback: (data: FallbackInfo) => void;
|
|
onDone: (data: { messageId: string; message?: ChatMessage }) => void;
|
|
onError: (data: string) => void;
|
|
}
|
|
|
|
export interface CreateChatStreamHandlersResult {
|
|
handlers: ChatStreamHandlers;
|
|
/** Cancel any pending RAF flushes for this stream. Idempotent. */
|
|
cancelFlushes: () => void;
|
|
}
|
|
|
|
/**
|
|
* Build the SSE handler bundle that `streamChatResponse` consumes. This is the
|
|
* portion of the chat send/stream flow that was identical between `useChat`
|
|
* and `useQuickChat`; extracting it keeps both hooks in sync when we tweak
|
|
* coalescing, tool-call dedup, fallback toasts, etc. The terminal events
|
|
* (`onDone`/`onError`) and parent-side fallback bookkeeping stay caller-owned
|
|
* because each hook handles message persistence and error recovery
|
|
* differently.
|
|
*
|
|
* The factory writes its `cancelFlushes` into `cancelStreamingFlushesRef.current`
|
|
* so the parent's `stopStreaming` can drain pending RAF callbacks before
|
|
* clearing transient streaming state — preventing a flushed delta from
|
|
* flashing back into the UI after a stop.
|
|
*/
|
|
export function createChatStreamHandlers(
|
|
options: CreateChatStreamHandlersOptions,
|
|
): CreateChatStreamHandlersResult {
|
|
const {
|
|
sessionId,
|
|
tempUserMessageId,
|
|
setStreamingText,
|
|
setStreamingThinking,
|
|
setStreamingToolCalls,
|
|
cancelStreamingFlushesRef,
|
|
addToast,
|
|
onDone,
|
|
onError,
|
|
onFallbackSession,
|
|
} = options;
|
|
|
|
let capturedText = "";
|
|
let capturedThinking = "";
|
|
let capturedToolCalls: ToolCallInfo[] = [];
|
|
let capturedFallbackInfo: FallbackInfo | undefined;
|
|
|
|
// Coalesce per-token state updates to one render per animation frame.
|
|
// ReactMarkdown re-parses the entire growing string on every render and
|
|
// every prior message also re-renders, so unthrottled setState here pegs
|
|
// the main thread on long replies.
|
|
let textRaf: number | null = null;
|
|
let thinkingRaf: number | null = null;
|
|
const flushText = (): void => {
|
|
textRaf = null;
|
|
setStreamingText(capturedText);
|
|
};
|
|
const flushThinking = (): void => {
|
|
thinkingRaf = null;
|
|
setStreamingThinking(capturedThinking);
|
|
};
|
|
const cancelFlushes = (): void => {
|
|
if (textRaf !== null) {
|
|
cancelAnimationFrame(textRaf);
|
|
textRaf = null;
|
|
}
|
|
if (thinkingRaf !== null) {
|
|
cancelAnimationFrame(thinkingRaf);
|
|
thinkingRaf = null;
|
|
}
|
|
};
|
|
cancelStreamingFlushesRef.current = cancelFlushes;
|
|
|
|
const handlers: ChatStreamHandlers = {
|
|
onThinking: (delta: string) => {
|
|
capturedThinking += delta;
|
|
if (thinkingRaf === null) {
|
|
thinkingRaf = requestAnimationFrame(flushThinking);
|
|
}
|
|
},
|
|
onText: (delta: string) => {
|
|
capturedText += delta;
|
|
if (textRaf === null) {
|
|
textRaf = requestAnimationFrame(flushText);
|
|
}
|
|
},
|
|
onToolStart: (data: { toolName: string; args?: Record<string, unknown> }) => {
|
|
capturedToolCalls = [
|
|
...capturedToolCalls,
|
|
{
|
|
toolName: data.toolName,
|
|
args: data.args,
|
|
isError: false,
|
|
status: "running",
|
|
},
|
|
];
|
|
setStreamingToolCalls(capturedToolCalls);
|
|
},
|
|
onToolEnd: (data: { toolName: string; isError: boolean; result?: unknown }) => {
|
|
const nextToolCalls = [...capturedToolCalls];
|
|
for (let i = nextToolCalls.length - 1; i >= 0; i--) {
|
|
const candidate = nextToolCalls[i];
|
|
if (candidate?.toolName === data.toolName && candidate.status === "running") {
|
|
nextToolCalls[i] = {
|
|
...candidate,
|
|
status: "completed",
|
|
isError: data.isError,
|
|
result: data.result,
|
|
};
|
|
capturedToolCalls = nextToolCalls;
|
|
setStreamingToolCalls(nextToolCalls);
|
|
return;
|
|
}
|
|
}
|
|
capturedToolCalls = [
|
|
...nextToolCalls,
|
|
{
|
|
toolName: data.toolName,
|
|
isError: data.isError,
|
|
result: data.result,
|
|
status: "completed",
|
|
},
|
|
];
|
|
setStreamingToolCalls(capturedToolCalls);
|
|
},
|
|
onFallback: (data: FallbackInfo) => {
|
|
capturedFallbackInfo = data;
|
|
onFallbackSession?.(data, sessionId);
|
|
addToast?.(`Primary model unavailable. Switched to fallback ${data.fallbackModel}.`, "warning");
|
|
},
|
|
onDone: (data: { messageId: string; message?: ChatMessage }) => {
|
|
cancelFlushes();
|
|
onDone({
|
|
messageId: data.messageId,
|
|
message: data.message,
|
|
accumulated: {
|
|
text: capturedText,
|
|
thinking: capturedThinking,
|
|
toolCalls: capturedToolCalls,
|
|
fallbackInfo: capturedFallbackInfo,
|
|
},
|
|
});
|
|
},
|
|
onError: (data: string) => {
|
|
cancelFlushes();
|
|
onError(data, tempUserMessageId);
|
|
},
|
|
};
|
|
|
|
return { handlers, cancelFlushes };
|
|
}
|
|
|
|
export type { ChatMessageInfo };
|