[eric] frontend perf: streaming slice + canvas coalescing + lazy paths

This commit is contained in:
ciregenz
2026-05-15 16:11:58 -07:00
parent 1c3aab08dd
commit 058cb0c52f
16 changed files with 690 additions and 177 deletions
+26 -2
View File
@@ -40,6 +40,24 @@ const SignInGate = lazy(() => import('./components/SignInGate'));
// callback to avoid all six firing at once and contending for network
// + parse time during first paint.
if (typeof window !== 'undefined') {
// Map sidebar paths to their dynamic imports so a hover/mouseenter on
// the sidebar can preload the chunk before the click. By the time the
// user actually clicks (~100-300ms after hover), the chunk is parsed
// and React.lazy resolves instantly. Exposed on window so AppShell
// can call it without prop-drilling. Each entry is idempotent;
// webpack dedupes repeated dynamic imports.
(window as any).__openswarmPrefetchRoute = (path: string) => {
switch (path) {
case '/skills': void import('./pages/Skills/Skills'); return;
case '/actions':
case '/tools': void import('./pages/Tools/Tools'); return;
case '/modes': void import('./pages/Modes/Modes'); return;
case '/views':
case '/apps': void import('./pages/Views/Views'); return;
case '/customization': void import('./pages/Customization/Customization'); return;
case '/analytics': void import('./pages/Analytics/Analytics'); return;
}
};
const prefetchAll = () => {
void import('./pages/Views/Views');
void import('./pages/Skills/Skills');
@@ -48,11 +66,17 @@ if (typeof window !== 'undefined') {
void import('./pages/Customization/Customization');
void import('./pages/Analytics/Analytics');
};
// Tighter idle deadline (was 4000ms): we WANT these chunks loaded
// before the user's first click, so don't let the browser defer them
// indefinitely. Fallback timeout reduced from 2000ms to 500ms for the
// same reason. The cost during initial render is small (one chunk
// parse per route, deferred); the cost of paying it on first click
// is a multi-hundred-ms freeze.
const ric = (window as any).requestIdleCallback as
| ((cb: () => void, opts?: { timeout?: number }) => number)
| undefined;
if (ric) ric(prefetchAll, { timeout: 4000 });
else window.setTimeout(prefetchAll, 2000);
if (ric) ric(prefetchAll, { timeout: 1500 });
else window.setTimeout(prefetchAll, 500);
}
import { report, getSessionTraceState, getRecentActions } from '@/shared/serviceClient';
import { useRouteTracker } from '@/shared/hooks/useRouteTracker';
+77 -3
View File
@@ -14,6 +14,9 @@ import PanToolOutlinedIcon from '@mui/icons-material/PanToolOutlined';
import { motion, AnimatePresence } from 'framer-motion';
import { useNavigate } from 'react-router-dom';
import { useAppDispatch, useAppSelector } from '@/shared/hooks';
import { shallowEqual } from 'react-redux';
import { createSelector } from '@reduxjs/toolkit';
import type { RootState } from '@/shared/state/store';
import {
handleApproval,
stopAgent,
@@ -190,6 +193,67 @@ const ActivityIndicator: React.FC<{ c: ReturnType<typeof useClaudeTokens> }> = (
/>
);
// ---------------------------------------------------------------------------
// Memoized session projection
// ---------------------------------------------------------------------------
//
// DynamicIsland only reads name / status / dashboard_id / pending_approvals
// per session. We project to a stable shape so identity persists across
// streamingMessage deltas (which mutate state.streaming, not state.agents,
// but still trigger Immer to swap the agents root reference any time
// agentsSlice runs (fine in theory, but selector consumers re-fire).
//
// Per-session cache: when a session's relevant fields haven't moved,
// return the SAME inner object reference, so the outer dict can be
// dropped on shallowEqual if its key set + per-session refs match.
type DiSession = {
id: string;
name: string;
status: string;
dashboard_id?: string;
pending_approvals: AgentSession['pending_approvals'];
};
const _diSessionCache: Map<string, DiSession> = new Map();
const selectDynamicIslandSessions = createSelector(
[(s: RootState) => s.agents.sessions],
(raw) => {
const out: Record<string, DiSession> = {};
const liveIds = new Set<string>();
for (const [sid, s] of Object.entries(raw)) {
liveIds.add(sid);
const prev = _diSessionCache.get(sid);
if (
prev
&& prev.name === s.name
&& prev.status === s.status
&& prev.dashboard_id === s.dashboard_id
&& prev.pending_approvals === s.pending_approvals
) {
out[sid] = prev;
} else {
const next: DiSession = {
id: sid,
name: s.name,
status: s.status,
dashboard_id: s.dashboard_id,
pending_approvals: s.pending_approvals,
};
_diSessionCache.set(sid, next);
out[sid] = next;
}
}
// Evict cache entries for sessions that disappeared. Without this,
// long sessions of dashboard switching slowly accumulate dead refs.
for (const cached of _diSessionCache.keys()) {
if (!liveIds.has(cached)) _diSessionCache.delete(cached);
}
return out;
},
);
// ---------------------------------------------------------------------------
// Main component
// ---------------------------------------------------------------------------
@@ -200,9 +264,19 @@ const DynamicIsland: React.FC = () => {
const navigate = useNavigate();
const islandRef = useRef<HTMLDivElement>(null);
const sessions = useAppSelector((state) => state.agents.sessions);
const history = useAppSelector((state) => state.agents.history);
const trackedIds = useAppSelector((state) => state.agents.trackedNotificationIds);
// Read the whole sessions dict, but memoize its projection so the
// useSelector only emits a new value when one of the four fields we
// actually consume (name/status/dashboard_id/pending_approvals)
// changes for SOME session. createSelector caches both the inner
// per-session shape AND the outer dict, so re-runs return the same
// reference when nothing relevant moved, even though Immer flips
// the top-level dict ref on every streamed character elsewhere.
// shallowEqual: createSelector returns a fresh outer dict object on
// each re-run, but the inner refs are cached so when nothing relevant
// moved, key-by-key comparison short-circuits the re-render.
const sessions = useAppSelector(selectDynamicIslandSessions, shallowEqual);
const history = useAppSelector((state) => state.agents.history, shallowEqual);
const trackedIds = useAppSelector((state) => state.agents.trackedNotificationIds, shallowEqual);
const [userExpanded, setUserExpanded] = useState(false);
const [searchOpen, setSearchOpen] = useState(false);
@@ -38,6 +38,7 @@ import Dashboard from '@/app/pages/Dashboard/Dashboard';
import DashboardHost from '@/app/components/Layout/DashboardHost';
import { useLastDashboardId } from '@/shared/hooks/useLastDashboardId';
import { useAppDispatch, useAppSelector } from '@/shared/hooks';
import { shallowEqual } from 'react-redux';
import { fetchDashboards, createDashboard, renameDashboard } from '@/shared/state/dashboardsSlice';
import { setPendingFocusAgentId } from '@/shared/state/tempStateSlice';
import { addBrowserCard, addBrowserTab } from '@/shared/state/dashboardLayoutSlice';
@@ -163,13 +164,29 @@ const AppShell: React.FC = () => {
(window as any).openswarm?.installUpdate();
}, [installing, dispatch]);
const dashboardItems = useAppSelector((state) => state.dashboards.items);
const dashboardList = Object.values(dashboardItems).sort(
(a, b) => new Date(b.updated_at).getTime() - new Date(a.updated_at).getTime(),
// Whole-dict subscriptions are deceptively expensive: `state.dashboards.items`
// and `state.outputs.items` are top-level dicts that get a NEW reference
// on any nested mutation (RTK/Immer behavior). With default referential
// equality, AppShell re-rendered on every dashboard rename, every output
// bump, every settings refresh that touched these slices, even though
// the dict CONTENTS were structurally identical from AppShell's POV.
// shallowEqual compares one level deep (key set + each value's identity),
// so AppShell now only re-renders on real structural changes.
const dashboardItems = useAppSelector(
(state) => state.dashboards.items,
shallowEqual,
);
const dashboardList = React.useMemo(
() => Object.values(dashboardItems).sort(
(a, b) => new Date(b.updated_at).getTime() - new Date(a.updated_at).getTime(),
),
[dashboardItems],
);
const outputItems = useAppSelector((state) => state.outputs.items);
// memo so the sort doesn't re-run on every AppShell re-render.
const outputItems = useAppSelector(
(state) => state.outputs.items,
shallowEqual,
);
const appsList = React.useMemo(
() => Object.values(outputItems).sort(
(a, b) => new Date(b.updated_at).getTime() - new Date(a.updated_at).getTime(),
@@ -862,6 +879,13 @@ const AppShell: React.FC = () => {
key={item.path}
data-onboarding={item.onboarding}
onClick={() => navigate(item.path)}
onMouseEnter={() => {
// Hover-prefetch the lazy chunk so the click pays
// ~0ms instead of the multi-hundred-ms chunk parse.
// See Main.tsx for the path → import map.
const fn = (window as any).__openswarmPrefetchRoute;
if (typeof fn === 'function') fn(item.path);
}}
sx={{
display: 'flex',
alignItems: 'center',
@@ -911,6 +935,10 @@ const AppShell: React.FC = () => {
<Box sx={{ px: 1, mb: 0.25 }}>
<ListItemButton
onClick={handleAppsClick}
onMouseEnter={() => {
const fn = (window as any).__openswarmPrefetchRoute;
if (typeof fn === 'function') fn('/apps');
}}
data-onboarding="sidebar-apps"
sx={{
borderRadius: 1.5,
+63 -66
View File
@@ -40,6 +40,7 @@ import {
} from '@/shared/state/agentsSlice';
import { fetchModes } from '@/shared/state/modesSlice';
import { createSessionWs } from '@/shared/ws/WebSocketManager';
import StreamingBubble from './StreamingBubble';
import MessageBubble from './MessageBubble';
import CompactionMarker from './CompactionMarker';
import MessageActionBar from './MessageActionBar';
@@ -351,7 +352,13 @@ const AgentChat: React.FC<AgentChatProps> = ({ sessionId: sessionIdProp, onClose
// so it never fires during normal streaming.
const reconcileTimer = useRef<ReturnType<typeof setTimeout> | null>(null);
const messageCount = session?.messages?.length ?? 0;
const hasStreaming = !!session?.streamingMessage;
// Subscribe only to the streaming MESSAGE ID (stable across the 30Hz
// delta updates), never to the content. The actual streaming text
// renders inside the leaf <StreamingBubble> below, which subscribes to
// the content itself. This keeps AgentChat's render and useEffects
// dormant during streaming; only the bubble updates per delta.
const streamingMessageId = useAppSelector((s) => id ? s.streaming.bySession[id]?.id ?? null : null);
const hasStreaming = !!streamingMessageId;
useEffect(() => {
if (reconcileTimer.current) {
@@ -417,7 +424,11 @@ const AgentChat: React.FC<AgentChatProps> = ({ sessionId: sessionIdProp, onClose
const scrollRafRef = useRef<number | null>(null);
const lastScrollHeightRef = useRef<number>(0);
useEffect(() => {
// Shared scroll-stick routine. Used both by the structural-events
// useEffect below (new message lands / stream starts/ends) and by
// StreamingBubble's onStreamGrew callback (per-delta growth). RAF +
// height-grew gate ensures we only set scrollTop when needed.
const stickToBottomIfNeeded = useCallback(() => {
if (!isAtBottomRef.current) return;
if (scrollRafRef.current != null) return;
scrollRafRef.current = requestAnimationFrame(() => {
@@ -425,20 +436,19 @@ const AgentChat: React.FC<AgentChatProps> = ({ sessionId: sessionIdProp, onClose
if (!isAtBottomRef.current) return;
const el = scrollContainerRef.current;
if (!el) return;
// Only set scrollTop when the scrollable height actually grew.
// Otherwise we're forcing a paint for nothing — and on a
// streaming turn we get one of these per delta, which thrashes
// the compositor for zero visible benefit. The native
// overflow-anchor on the container already keeps the viewport
// pinned to the bottom; this JS fallback only needs to handle
// the rare case where anchoring misses (legacy WebKit,
// virtualized children, dynamic-height inserts).
const newHeight = el.scrollHeight;
if (newHeight === lastScrollHeightRef.current) return;
lastScrollHeightRef.current = newHeight;
el.scrollTop = newHeight;
});
}, [session?.messages.length, session?.streamingMessage?.content]);
}, []);
useEffect(() => {
stickToBottomIfNeeded();
// Structural triggers only: a new message lands or a stream
// starts/ends. Streaming content updates trigger this via
// <StreamingBubble onStreamGrew={stickToBottomIfNeeded} /> instead
// so AgentChat stays dormant during the 30Hz delta storm.
}, [session?.messages.length, streamingMessageId, stickToBottomIfNeeded]);
useEffect(() => () => {
if (scrollRafRef.current != null) {
@@ -447,19 +457,31 @@ const AgentChat: React.FC<AgentChatProps> = ({ sessionId: sessionIdProp, onClose
}
}, []);
const handleSend = (prompt: string, images?: Array<{ data: string; media_type: string }>, contextPaths?: Array<{ path: string; type: 'file' | 'directory' }>, forcedTools?: string[], attachedSkills?: Array<{ id: string; name: string; content: string }>, selectedBrowserIds?: string[]) => {
if (!id) return;
// Sending a message is a clear intent signal: the user wants to see
// the response. Force-scroll to bottom regardless of isAtBottomRef.
scrollToBottom();
const msg: QueuedMessage = { prompt, images, contextPaths, forcedTools, attachedSkills, selectedBrowserIds };
if (agentBusy) {
messageQueueRef.current.push(msg);
setQueueLength(messageQueueRef.current.length);
return;
}
dispatchMessage(msg);
};
// useCallback so ChatInput's memo equality holds across AgentChat
// re-renders driven by unrelated session state. Captures agentBusy
// through the dependency so a stale "busy" closure doesn't ever route
// a message past the queue.
const handleSend = useCallback(
(
prompt: string,
images?: Array<{ data: string; media_type: string }>,
contextPaths?: Array<{ path: string; type: 'file' | 'directory' }>,
forcedTools?: string[],
attachedSkills?: Array<{ id: string; name: string; content: string }>,
selectedBrowserIds?: string[],
) => {
if (!id) return;
scrollToBottom();
const msg: QueuedMessage = { prompt, images, contextPaths, forcedTools, attachedSkills, selectedBrowserIds };
if (agentBusy) {
messageQueueRef.current.push(msg);
setQueueLength(messageQueueRef.current.length);
return;
}
dispatchMessage(msg);
},
[id, scrollToBottom, agentBusy, dispatchMessage],
);
const handleModeChange = useCallback((newMode: string) => {
setMode(newMode);
@@ -485,10 +507,10 @@ const AgentChat: React.FC<AgentChatProps> = ({ sessionId: sessionIdProp, onClose
dispatch(handleApproval({ requestId, behavior: 'deny', message }));
};
const handleStop = () => {
const handleStop = useCallback(() => {
if (!id) return;
dispatch(stopAgent({ sessionId: id }));
};
}, [id, dispatch]);
const handleResume = useCallback(() => {
if (!id) return;
@@ -609,15 +631,14 @@ const AgentChat: React.FC<AgentChatProps> = ({ sessionId: sessionIdProp, onClose
for (const msg of activeBranchMessages) {
totalChars += stringifyContent(msg.content).length;
}
if (session?.streamingMessage) {
totalChars += (session.streamingMessage.content || '').length;
}
const used = Math.round(totalChars / 4);
return { used, limit };
// Depending on streamingMessage.id (not .content) recomputes once per
// turn instead of per painted character. The header pct gauge would
// otherwise re-run a full-message length sum every animation frame.
}, [activeBranchMessages, session?.system_prompt, session?.streamingMessage?.id, model, modelsByProvider]);
// Streaming content's contribution to the context estimate is no
// longer included here: we'd have to subscribe to the streaming
// text and re-run this sum on every painted character, defeating
// the whole point of isolating AgentChat from delta updates. The
// header gauge will catch up when stream_end commits the message.
}, [activeBranchMessages, session?.system_prompt, streamingMessageId, model, modelsByProvider]);
const sessionRunning = session?.status === 'running' || session?.status === 'waiting_approval';
@@ -1119,7 +1140,7 @@ const AgentChat: React.FC<AgentChatProps> = ({ sessionId: sessionIdProp, onClose
</Box>
);
})()}
{renderItems.filter((item) => !session.streamingMessage || item.id !== session.streamingMessage.id).map((item) => {
{renderItems.filter((item) => !streamingMessageId || item.id !== streamingMessageId).map((item) => {
const isCompactionAnchor = !!session.compacted_through_msg_id && item.id === session.compacted_through_msg_id;
const compactionChip = isCompactionAnchor ? (
<CompactionMarker
@@ -1209,39 +1230,15 @@ const AgentChat: React.FC<AgentChatProps> = ({ sessionId: sessionIdProp, onClose
</Box>
);
})}
{session.streamingMessage && (
session.streamingMessage.role === 'tool_call' ? (
<ToolCallBubble
key={`streaming-${session.streamingMessage.id}`}
isStreaming
isPending
sessionId={session.id}
call={{
id: session.streamingMessage.id,
role: 'tool_call',
content: { tool: session.streamingMessage.tool_name || '', input: session.streamingMessage.content },
timestamp: new Date().toISOString(),
branch_id: session.active_branch_id || 'main',
parent_id: null,
}}
/>
) : (
<MessageBubble
key={`streaming-${session.streamingMessage.id}`}
isStreaming
dynamicTurnLabel={session.turn_label?.label}
message={{
id: session.streamingMessage.id,
role: session.streamingMessage.role,
content: session.streamingMessage.content,
timestamp: new Date().toISOString(),
branch_id: session.active_branch_id || 'main',
parent_id: null,
}}
/>
)
{id && (
<StreamingBubble
sessionId={id}
activeBranchId={session.active_branch_id || 'main'}
turnLabel={session.turn_label?.label}
onStreamGrew={stickToBottomIfNeeded}
/>
)}
{(awaitingResponse || (session.status === 'running' && !session.streamingMessage)) && (
{(awaitingResponse || (session.status === 'running' && !streamingMessageId)) && (
<ThinkingBubble
label={session.turn_label?.label}
seedKey={`${session.id}:${session.messages?.length ?? 0}`}
@@ -20,8 +20,10 @@ import AccountTreeOutlinedIcon from '@mui/icons-material/AccountTreeOutlined';
import CodeOutlinedIcon from '@mui/icons-material/CodeOutlined';
import BuildOutlinedIcon from '@mui/icons-material/BuildOutlined';
import { createSelector } from '@reduxjs/toolkit';
import { shallowEqual } from 'react-redux';
import { useAppSelector, useAppDispatch } from '@/shared/hooks';
import { AgentMessage, AgentSession, fetchBrowserAgentChildren, handleApproval } from '@/shared/state/agentsSlice';
import type { StreamingMessage } from '@/shared/state/streamingSlice';
import { useClaudeTokens, useThemeMode } from '@/shared/styles/ThemeContext';
import type { RootState } from '@/shared/state/store';
@@ -150,6 +152,11 @@ const lightFeedColors: FeedColors = {
scrollThumb: '#ccc9c0',
};
// Stable empty-object reference for the streaming selector to return
// when there are no browser sessions yet; keeps shallowEqual happy
// across renders so we don't churn on an "empty" dict literal.
const EMPTY_STREAMING: Record<string, StreamingMessage> = Object.freeze({}) as Record<string, StreamingMessage>;
const selectBrowserSessions = createSelector(
[(state: RootState) => state.agents.sessions,
(_: RootState, parentSessionId: string) => parentSessionId,
@@ -174,6 +181,28 @@ const BrowserAgentInlineFeed: React.FC<Props> = ({ parentSessionId, browserId })
const browserSessions = useAppSelector((state) =>
selectBrowserSessions(state, parentSessionId, browserId),
);
// Subscribe to only the streaming entries that belong to THIS feed's
// browser sessions. Previously this read the full bySession dict,
// which re-rendered the feed on every streamed character from every
// agent on the dashboard, which was the "glitching when agent is
// using the browser" experience. With shallowEqual we only re-render
// when one of our specific browser sessions actually gets a delta.
const browserSessionIds = useMemo(
() => browserSessions.map((s) => s.id).sort().join(','),
[browserSessions],
);
const streamingBySession = useAppSelector(
(state) => {
if (!browserSessionIds) return EMPTY_STREAMING;
const out: Record<string, StreamingMessage> = {};
for (const id of browserSessionIds.split(',')) {
const entry = state.streaming.bySession[id];
if (entry) out[id] = entry;
}
return out;
},
shallowEqual,
);
useEffect(() => {
if (browserSessions.length === 0 && fetchedForSession.current !== parentSessionId) {
@@ -191,15 +220,16 @@ const BrowserAgentInlineFeed: React.FC<Props> = ({ parentSessionId, browserId })
const entry = formatMessage(msg);
if (entry) entries.push(entry);
}
if (session.streamingMessage?.role === 'assistant' && session.streamingMessage.content) {
entries.push({ type: 'thought', text: session.streamingMessage.content });
const stream: StreamingMessage | undefined = streamingBySession[session.id];
if (stream?.role === 'assistant' && stream.content) {
entries.push({ type: 'thought', text: stream.content });
}
return { session, entries };
});
}, [browserSessions]);
}, [browserSessions, streamingBySession]);
const totalMessages = browserSessions.reduce(
(n, s) => n + s.messages.length + (s.streamingMessage ? 1 : 0),
(n, s) => n + s.messages.length + (streamingBySession[s.id] ? 1 : 0),
0,
);
+76 -15
View File
@@ -90,6 +90,12 @@ export interface AttachedImage {
data: string;
media_type: string;
preview: string;
// Set by addImageFiles when we use URL.createObjectURL for the preview
// instead of a data URL. handleSend reads this with FileReader at
// send time so we don't carry the full base64 in memory between
// attach and send. Falls back to base64 conversion of the data URL
// if the file is missing (e.g. paste-from-clipboard with raw base64).
_file?: File;
}
export interface ForcedToolGroup {
@@ -663,6 +669,18 @@ const ChatInput = forwardRef<ChatInputHandle, Props>(({ onSend, disabled, mode,
}), [c]);
const [images, setImages] = useState<AttachedImage[]>([]);
// Track current images via ref so the unmount cleanup sees the latest
// list (not the empty-array snapshot from the effect's first run) and
// can revoke any outstanding blob: preview URLs.
const imagesRef = useRef(images);
imagesRef.current = images;
useEffect(() => () => {
for (const img of imagesRef.current) {
if (img.preview?.startsWith('blob:')) {
try { URL.revokeObjectURL(img.preview); } catch { /* nothing */ }
}
}
}, []);
const [lightboxSrc, setLightboxSrc] = useState<string | null>(null);
const [isDragOver, setIsDragOver] = useState(false);
const [isUploading, setIsUploading] = useState(false);
@@ -773,18 +791,19 @@ const ChatInput = forwardRef<ChatInputHandle, Props>(({ onSend, disabled, mode,
}, []);
const addImageFiles = useCallback((files: FileList | File[]) => {
// Two-stage: thumbnail / lightbox preview comes from a blob: URL
// (kept by the browser as a binary file handle, not a JS string),
// and the base64 only materializes asynchronously for the actual
// send payload. Holding only the blob URL keeps a ~2MB screenshot
// attachment from also costing ~2.7MB of JS heap as a data URL.
// The base64 promise is awaited inside handleSend.
Array.from(files).forEach((file) => {
if (!file.type.startsWith('image/')) return;
const reader = new FileReader();
reader.onload = () => {
const result = reader.result as string;
const base64 = result.split(',')[1];
setImages((prev) => [
...prev,
{ data: base64, media_type: file.type, preview: result },
]);
};
reader.readAsDataURL(file);
const previewUrl = URL.createObjectURL(file);
setImages((prev) => [
...prev,
{ data: '', media_type: file.type, preview: previewUrl, _file: file } as AttachedImage,
]);
});
}, []);
@@ -847,9 +866,31 @@ const ChatInput = forwardRef<ChatInputHandle, Props>(({ onSend, disabled, mode,
}
const selectedEls = elementSelection?.elementsByOwner?.[ownerId] ?? [];
let allImages = images.length > 0
? images.map(({ data, media_type }) => ({ data, media_type }))
: [];
// Materialize image base64 at send time so we don't keep ~2.7MB
// strings in component state for every attached screenshot. Images
// added via addImageFiles carry a File reference (_file) and read
// their bytes on demand; legacy paste flows that wrote `data`
// directly still work. FileReader is async, so this is a Promise.all.
let allImages: Array<{ data: string; media_type: string }> = [];
if (images.length > 0) {
allImages = await Promise.all(images.map(async (img) => {
if (img.data) return { data: img.data, media_type: img.media_type };
if (img._file) {
const base64 = await new Promise<string>((resolve, reject) => {
const reader = new FileReader();
reader.onload = () => {
const r = reader.result as string;
resolve(r.split(',')[1] ?? '');
};
reader.onerror = () => reject(reader.error || new Error('FileReader failed'));
reader.readAsDataURL(img._file!);
});
return { data: base64, media_type: img.media_type };
}
return { data: '', media_type: img.media_type };
}));
allImages = allImages.filter((i) => i.data);
}
if (selectedEls.length > 0) {
const lines: string[] = ['\n\n---\nSelected UI Elements:\n'];
@@ -923,6 +964,13 @@ const ChatInput = forwardRef<ChatInputHandle, Props>(({ onSend, disabled, mode,
);
editor.innerHTML = '';
_draftStore.delete(ownerId);
// Revoke any blob: URLs we minted for previews so the underlying
// bytes can be freed by the browser.
for (const img of images) {
if (img.preview?.startsWith('blob:')) {
try { URL.revokeObjectURL(img.preview); } catch { /* nothing */ }
}
}
setImages([]);
setContextPaths([]);
setForcedTools([]);
@@ -1150,7 +1198,16 @@ const ChatInput = forwardRef<ChatInputHandle, Props>(({ onSend, disabled, mode,
}, [addImageFiles, uploadAndAttachFiles]);
const removeImage = useCallback((idx: number) => {
setImages((prev) => prev.filter((_, i) => i !== idx));
setImages((prev) => {
const removed = prev[idx];
// Revoke the blob URL we minted in addImageFiles so the browser
// can free the underlying bytes; data URLs have no resource to
// revoke so the `blob:` check is sufficient.
if (removed?.preview?.startsWith('blob:')) {
try { URL.revokeObjectURL(removed.preview); } catch { /* nothing */ }
}
return prev.filter((_, i) => i !== idx);
});
}, []);
const menuPaperProps = {
@@ -2423,4 +2480,8 @@ const ChatInput = forwardRef<ChatInputHandle, Props>(({ onSend, disabled, mode,
ChatInput.displayName = 'ChatInput';
export default ChatInput;
// Memoize across parent re-renders driven by unrelated state (most
// commonly AgentChat re-rendering because its session-local data
// updated). The parent passes callbacks via useCallback and primitive
// props, so the default shallow comparison is correct.
export default React.memo(ChatInput);
@@ -0,0 +1,88 @@
import React, { useEffect, useRef } from 'react';
import { useStreamingMessage } from '@/shared/state/streamingSlice';
import MessageBubble from './MessageBubble';
import ToolCallBubble from './ToolCallBubble';
interface Props {
sessionId: string;
activeBranchId: string;
turnLabel?: string | null;
// Called when the streamed content grows so the host scroll container
// can stick to the bottom. We pass a callback instead of doing the
// scroll math here so AgentChat keeps ownership of its scroll state
// (isAtBottomRef etc.). The callback is invoked from a RAF, so it's
// safe to do DOM reads/writes inside.
onStreamGrew?: () => void;
}
// Leaf component that subscribes to the streaming entry for a single
// session and renders the appropriate bubble. Isolating this in its own
// component is what keeps AgentChat from re-rendering on every painted
// character. AgentChat only knows whether a stream exists (boolean
// selector elsewhere), not the per-character content. StreamingBubble
// itself does re-render at the streaming rate, but it has no children
// beyond a MessageBubble/ToolCallBubble, so React reconciliation stays
// local and cheap.
const StreamingBubble: React.FC<Props> = ({ sessionId, activeBranchId, turnLabel, onStreamGrew }) => {
const streamingMessage = useStreamingMessage(sessionId);
// Fire onStreamGrew once per render (i.e. per delta) on a RAF so the
// host can scroll if it wants to. RAF coalesces multiple deltas in
// the same frame into one host call. The ref-callback keeps the
// useEffect dep array minimal: we don't want to re-run effects on
// every callback identity change from the parent.
const onGrewRef = useRef(onStreamGrew);
onGrewRef.current = onStreamGrew;
const rafRef = useRef<number | null>(null);
useEffect(() => {
if (!streamingMessage) return;
if (rafRef.current != null) return;
rafRef.current = requestAnimationFrame(() => {
rafRef.current = null;
onGrewRef.current?.();
});
return () => {
if (rafRef.current != null) {
cancelAnimationFrame(rafRef.current);
rafRef.current = null;
}
};
});
if (!streamingMessage) return null;
if (streamingMessage.role === 'tool_call') {
return (
<ToolCallBubble
key={`streaming-${streamingMessage.id}`}
isStreaming
isPending
sessionId={sessionId}
call={{
id: streamingMessage.id,
role: 'tool_call',
content: { tool: streamingMessage.tool_name || '', input: streamingMessage.content },
timestamp: new Date().toISOString(),
branch_id: activeBranchId,
parent_id: null,
}}
/>
);
}
return (
<MessageBubble
key={`streaming-${streamingMessage.id}`}
isStreaming
dynamicTurnLabel={turnLabel}
message={{
id: streamingMessage.id,
role: streamingMessage.role,
content: streamingMessage.content,
timestamp: new Date().toISOString(),
branch_id: activeBranchId,
parent_id: null,
}}
/>
);
};
export default StreamingBubble;
+51 -13
View File
@@ -32,6 +32,8 @@ import { parseMcpToolName, getMcpShortAction } from '@/app/pages/AgentChat/ToolC
import { useClaudeTokens } from '@/shared/styles/ThemeContext';
import { useDashboardActive } from '@/shared/hooks/useDashboardActive';
import { useOverlayScrollPassthrough } from './useOverlayScrollPassthrough';
import { useStreamingMessage } from '@/shared/state/streamingSlice';
import { isCanvasInteractionActive, onCanvasInteractionEnd } from '@/shared/canvasInteractionState';
// ---------------------------------------------------------------------------
// Helper components & functions (unchanged)
@@ -80,6 +82,23 @@ function fmtSeconds(seconds: number): string {
return `${hours}h ${minutes % 60}m`;
}
// Self-ticking elapsed-time renderer. Owns its own 1Hz interval so only
// this leaf re-renders per second while a session is active; the rest
// of AgentCard stays put. Memoized on `status` + `messages` so it
// doesn't re-tick after the session goes terminal.
const ElapsedTimer: React.FC<{
messages: Array<{ role: string; timestamp: string; elapsed_ms?: number; hidden?: boolean }>;
status: string;
}> = React.memo(({ messages, status }) => {
const [, setTick] = React.useState(0);
React.useEffect(() => {
if (status !== 'running' && status !== 'waiting_approval') return;
const id = setInterval(() => setTick((t) => (t + 1) & 0xffff), 1000);
return () => clearInterval(id);
}, [status]);
return <>{fmtSeconds(getAgentWorkTime(messages, status).last)}</>;
});
function getAgentWorkTime(
messages: Array<{ role: string; timestamp: string; elapsed_ms?: number; hidden?: boolean }>,
status: string,
@@ -322,16 +341,35 @@ const AgentCard: React.FC<Props> = ({
useEffect(() => {
const el = cardBoxRef.current;
if (!el || !onMeasuredHeight) return;
// Remember the most recent height seen during a suppressed window
// (pan/drag/zoom in progress). When the interaction ends, fire it
// through so the layout reconciles to the truth right then.
let suppressedHeight: number | null = null;
const ro = new ResizeObserver((entries) => {
// Short-circuit when dashboard is hidden — observer stays attached so
// the next resize after returning to the dashboard fires correctly.
if (!isDashboardActiveRef.current) return;
// Short-circuit during active canvas interaction (pan/drag/wheel).
// During those gestures we don't care about millimeter-precise card
// heights; re-measuring on every streamed character was forcing
// Dashboard re-renders mid-pan via setMeasuredHeightsTick. Stash
// the latest height instead and flush on gesture end.
if (isCanvasInteractionActive()) {
for (const entry of entries) suppressedHeight = entry.contentRect.height;
return;
}
for (const entry of entries) {
onMeasuredHeight(session.id, entry.contentRect.height);
}
});
ro.observe(el);
return () => ro.disconnect();
const unsub = onCanvasInteractionEnd(() => {
if (suppressedHeight != null && isDashboardActiveRef.current) {
onMeasuredHeight(session.id, suppressedHeight);
}
suppressedHeight = null;
});
return () => { ro.disconnect(); unsub(); };
}, [session.id, onMeasuredHeight]);
// ---- Glow state (for branched cards) ----
@@ -364,7 +402,6 @@ const AgentCard: React.FC<Props> = ({
draft: { color: c.accent.primary, bg: c.bg.secondary },
};
const [, setTick] = useState(0);
const isDraft = session.status === 'draft';
// ---- Drag via header (pointer events) ----
@@ -557,19 +594,20 @@ const AgentCard: React.FC<Props> = ({
};
useEffect(() => {
if (session.status === 'running' || session.status === 'waiting_approval') {
const interval = setInterval(() => setTick((t) => t + 1), 1000);
return () => clearInterval(interval);
}
}, [session.status]);
// Elapsed-time display owns its own 1Hz tick via <ElapsedTimer/> below;
// we don't force-re-render the whole 1000+ line AgentCard every second
// anymore (each card running × 1Hz = wasted reconciliation budget).
const lastMessage = session.messages[session.messages.length - 1];
const isStreaming = !!session.streamingMessage;
// Subscribe to this card's own streaming entry from the streaming
// slice. Per-character mutations no longer churn the sessions dict,
// so other cards stay stable while this one streams.
const streamingMessage = useStreamingMessage(session.id);
const isStreaming = !!streamingMessage;
const previewContent = isStreaming
? (session.streamingMessage!.role === 'tool_call'
? `[${getToolDisplayName(session.streamingMessage!.tool_name || '')}] ${session.streamingMessage!.content}`
: session.streamingMessage!.content
? (streamingMessage!.role === 'tool_call'
? `[${getToolDisplayName(streamingMessage!.tool_name || '')}] ${streamingMessage!.content}`
: streamingMessage!.content
).slice(0, 120)
: lastMessage && typeof lastMessage.content === 'string'
? lastMessage.content.slice(0, 120)
@@ -985,7 +1023,7 @@ const AgentCard: React.FC<Props> = ({
{session.mode}
</Typography>
<Typography variant="caption" sx={{ color: c.text.tertiary }}>
{fmtSeconds(getAgentWorkTime(session.messages, session.status).last)}
<ElapsedTimer messages={session.messages} status={session.status} />
</Typography>
{session.cost_usd > 0 && hasApiKey && (
<Typography variant="caption" sx={{ color: c.accent.primary }}>
@@ -14,6 +14,7 @@ import CloseFullscreenIcon from '@mui/icons-material/CloseFullscreen';
import CheckCircleOutlineIcon from '@mui/icons-material/CheckCircleOutline';
import ErrorOutlineIcon from '@mui/icons-material/ErrorOutline';
import { AgentSession, AgentMessage, stopAgent, handleApproval } from '@/shared/state/agentsSlice';
import { useStreamingMessage } from '@/shared/state/streamingSlice';
import { useAppDispatch, useAppSelector } from '@/shared/hooks';
import { useClaudeTokens } from '@/shared/styles/ThemeContext';
@@ -81,6 +82,8 @@ const BrowserAgentOverlay: React.FC<Props> = ({ session, browserWidth, browserHe
const isRunning = session.status === 'running' || session.status === 'waiting_approval';
const browserDone = session.status === 'completed' || session.status === 'error' || session.status === 'stopped';
// Streaming message lives in its own slice; see streamingSlice.ts.
const streamingMessage = useStreamingMessage(session.id);
// Only truly "done" (fade + hide) when the parent is also finished.
// While the parent is still active, the overlay stays visible in a
// "waiting for next task" state between sub-tasks.
@@ -130,7 +133,7 @@ const BrowserAgentOverlay: React.FC<Props> = ({ session, browserWidth, browserHe
if (scrollRef.current) {
scrollRef.current.scrollTop = scrollRef.current.scrollHeight;
}
}, [session.messages.length, session.streamingMessage]);
}, [session.messages.length, streamingMessage]);
const handleStop = useCallback(() => {
if (!confirmStop) {
@@ -153,9 +156,8 @@ const BrowserAgentOverlay: React.FC<Props> = ({ session, browserWidth, browserHe
.map(summarizeMessage)
.filter((e) => e.type !== 'skip' && e.type !== 'result');
const streamingMsg = session.streamingMessage;
if (streamingMsg && streamingMsg.role === 'assistant' && streamingMsg.content) {
entries.push({ type: 'thought', text: streamingMsg.content });
if (streamingMessage && streamingMessage.role === 'assistant' && streamingMessage.content) {
entries.push({ type: 'thought', text: streamingMessage.content });
}
const collapsedW = Math.min(300, browserWidth - 24);
+26 -3
View File
@@ -480,12 +480,28 @@ const DashboardInner: React.FC<DashboardProps> = ({ dashboardId, isActive = true
hasFittedRef.current = false;
restoredExpandedRef.current = false;
dispatch(resetLayout());
// CRITICAL path: these populate the cards the user expects to see
// on first paint. Don't defer.
dispatch(fetchSessions({ dashboardId }));
dispatch(fetchHistory({ dashboardId }));
dispatch(fetchLayout(dashboardId));
dispatch(fetchOutputs());
dashboardWs.connect();
const cleanupBrowserHandler = initBrowserCommandHandler();
// DEFERRABLE: history list (for the search palette) and outputs
// (for the apps panel) aren't on the first-paint path. Same for the
// dashboard WS connection (it carries cross-session events; opens
// ~100ms later costs nothing). Pushing these into the post-paint
// window measurably improves LCP because the initial render
// pipeline isn't competing with their thunks/network setup.
const idleHandle = (typeof window !== 'undefined' && (window as any).requestIdleCallback)
? (window as any).requestIdleCallback(() => {
dispatch(fetchHistory({ dashboardId }));
dispatch(fetchOutputs());
dashboardWs.connect();
}, { timeout: 2000 })
: window.setTimeout(() => {
dispatch(fetchHistory({ dashboardId }));
dispatch(fetchOutputs());
dashboardWs.connect();
}, 200);
// Pre-warm Anthropic's prompt cache for sessions on this dashboard
// ~250ms after mount (debounced; AbortController cancels on
@@ -523,6 +539,13 @@ const DashboardInner: React.FC<DashboardProps> = ({ dashboardId, isActive = true
warmAbort.abort();
cleanupBrowserHandler();
dashboardWs.disconnect();
// Cancel any not-yet-fired idle work; the cleanup handler can't
// run partially if the dashboard switches before idle fired.
if (typeof window !== 'undefined') {
const cancelIdle = (window as any).cancelIdleCallback;
if (cancelIdle && typeof idleHandle === 'number') cancelIdle(idleHandle);
else if (typeof idleHandle === 'number') clearTimeout(idleHandle);
}
};
}, [dispatch, dashboardId]);
@@ -1,4 +1,5 @@
import { useState, useCallback, useRef, useEffect, useMemo, RefObject } from 'react';
import { setCanvasInteractionActive } from '@/shared/canvasInteractionState';
const MIN_ZOOM = 0.15;
const MAX_ZOOM = 3.0;
@@ -195,6 +196,13 @@ export function useCanvasControls(zoomSensitivity: number = 50, contentBounds?:
let pendingZoomDy = 0;
let pendingZoomCenter: { cx: number; cy: number } | null = null;
let wheelRafId: number | null = null;
// Trackpad wheel gestures don't have a "gestureend" event; we infer
// it from idle time. ~140ms after the last wheel event we declare
// the gesture over and unset the interaction flag, which un-pauses
// ResizeObservers etc. The 140ms window is short enough to feel
// responsive on re-engage and long enough to absorb the inter-burst
// gaps inside a continuous swipe.
let wheelIdleTimer: ReturnType<typeof setTimeout> | null = null;
const flushWheel = () => {
wheelRafId = null;
@@ -224,6 +232,15 @@ export function useCanvasControls(zoomSensitivity: number = 50, contentBounds?:
};
const scheduleWheelFlush = () => {
// Mark the canvas as actively-interacting and (re)arm the idle
// timer. Any ResizeObserver / streaming reconciler that checks the
// flag will bail until the user's gesture goes quiet for ~140ms.
setCanvasInteractionActive(true);
if (wheelIdleTimer != null) clearTimeout(wheelIdleTimer);
wheelIdleTimer = setTimeout(() => {
wheelIdleTimer = null;
setCanvasInteractionActive(false);
}, 140);
if (wheelRafId != null) return;
wheelRafId = requestAnimationFrame(flushWheel);
};
@@ -347,6 +364,9 @@ export function useCanvasControls(zoomSensitivity: number = 50, contentBounds?:
el.removeEventListener('wheel', onWheel);
window.removeEventListener('openswarm:canvas-wheel-zoom', onForwardedZoom);
if (wheelRafId != null) cancelAnimationFrame(wheelRafId);
if (wheelIdleTimer != null) clearTimeout(wheelIdleTimer);
// Don't leave the flag stuck on if the canvas unmounts mid-gesture.
setCanvasInteractionActive(false);
};
}, [enabled]);
@@ -355,6 +375,7 @@ export function useCanvasControls(zoomSensitivity: number = 50, contentBounds?:
cancelAnimation();
cancelInertia();
setIsPanning(true);
setCanvasInteractionActive(true);
velocityHistoryRef.current = [{ x: e.clientX, y: e.clientY, t: performance.now() }];
panStartRef.current = {
x: e.clientX,
@@ -434,6 +455,7 @@ export function useCanvasControls(zoomSensitivity: number = 50, contentBounds?:
}
panStartRef.current = null;
setIsPanning(false);
setCanvasInteractionActive(false);
// Only spring back if we were actually panning (not on simple clicks)
if (wasPanning && !didInertia) {
springBackIfNeeded();
@@ -446,6 +468,7 @@ export function useCanvasControls(zoomSensitivity: number = 50, contentBounds?:
if (panStartRef.current) {
panStartRef.current = null;
setIsPanning(false);
setCanvasInteractionActive(false);
}
};
window.addEventListener('mouseup', onUp);
@@ -0,0 +1,36 @@
// Plain-JS shared ref (NOT React state) for "is the user currently
// interacting with the canvas" (pan/drag/wheel/zoom). Read on hot paths
// like AgentCard's ResizeObserver to suppress expensive work during the
// gesture. Setting/clearing the ref does NOT trigger any React re-renders.
//
// Why this pattern instead of Redux or context: ResizeObserver callbacks
// fire dozens of times per second during streaming. We want them to bail
// in O(1) without a subscription that itself has overhead. A module-level
// mutable holder + a one-shot "interaction ended" event meets both.
let _isPanning = false;
const listeners: Set<() => void> = new Set();
export function isCanvasInteractionActive(): boolean {
return _isPanning;
}
export function setCanvasInteractionActive(active: boolean) {
if (_isPanning === active) return;
const wasActive = _isPanning;
_isPanning = active;
// Fire the end-of-interaction notification so listeners can flush work
// that was suppressed during the gesture (re-measure heights, dispatch
// pending state updates, etc.).
if (wasActive && !active) {
for (const fn of listeners) {
try { fn(); } catch (e) { console.warn('[canvas-interaction] listener threw', e); }
}
}
}
export function onCanvasInteractionEnd(fn: () => void): () => void {
listeners.add(fn);
return () => { listeners.delete(fn); };
}
+23 -59
View File
@@ -54,12 +54,9 @@ export interface MessageBranch {
created_at: string;
}
export interface StreamingMessage {
id: string;
role: 'assistant' | 'tool_call' | 'thinking';
content: string;
tool_name?: string;
}
// StreamingMessage type moved to streamingSlice. Import from there if you
// need the shape directly.
export type { StreamingMessage } from './streamingSlice';
export interface ToolGroupMeta {
id: string;
@@ -89,7 +86,10 @@ export interface AgentSession {
pending_approvals: ApprovalRequest[];
branches: Record<string, MessageBranch>;
active_branch_id: string;
streamingMessage: StreamingMessage | null;
// streamingMessage lives in `state.streaming.bySession[id]` now;
// read it via the selectors in streamingSlice. Kept off this type so
// that any reader still trying to access it gets a compile error and
// is migrated to the new location.
target_directory?: string | null;
tool_group_meta: Record<string, ToolGroupMeta>;
dashboard_id?: string;
@@ -543,7 +543,6 @@ const agentsSlice = createSlice({
pending_approvals: [],
branches: { main: { id: 'main', parent_branch_id: null, fork_point_message_id: null, created_at: new Date().toISOString() } },
active_branch_id: 'main',
streamingMessage: null,
target_directory: targetDirectory || null,
tool_group_meta: {},
thinking_level: thinkingLevel,
@@ -660,7 +659,6 @@ const agentsSlice = createSlice({
state.sessions[action.payload.id] = {
...action.payload,
pending_approvals: mergedApprovals,
streamingMessage: existing?.streamingMessage ?? action.payload.streamingMessage ?? null,
tool_group_meta: { ...existing?.tool_group_meta, ...action.payload.tool_group_meta },
};
if (action.payload.status === 'running' && !state.trackedNotificationIds.includes(action.payload.id)) {
@@ -712,9 +710,8 @@ const agentsSlice = createSlice({
);
if (optIdx >= 0) {
session.messages[optIdx] = { ...incoming, optimistic_status: undefined };
if (session.streamingMessage?.id === incoming.id) {
session.streamingMessage = null;
}
// streamingMessage cleanup is handled by streamingSlice's
// extraReducers listening to this action.
return;
}
}
@@ -724,9 +721,7 @@ const agentsSlice = createSlice({
} else {
session.messages.push(incoming);
}
if (session.streamingMessage?.id === incoming.id) {
session.streamingMessage = null;
}
// streamingMessage cleanup is handled by streamingSlice's extraReducers.
},
// Synchronous "you sent a message" bubble dispatched from the
@@ -810,40 +805,12 @@ const agentsSlice = createSlice({
session.turn_label = null;
},
streamStart(
state,
action: PayloadAction<{ sessionId: string; messageId: string; role: 'assistant' | 'tool_call' | 'thinking'; toolName?: string }>
) {
const session = state.sessions[action.payload.sessionId];
if (session) {
session.streamingMessage = {
id: action.payload.messageId,
role: action.payload.role,
content: '',
tool_name: action.payload.toolName,
};
}
},
streamDelta(
state,
action: PayloadAction<{ sessionId: string; messageId: string; delta: string }>
) {
const session = state.sessions[action.payload.sessionId];
if (session?.streamingMessage?.id === action.payload.messageId) {
session.streamingMessage.content += action.payload.delta;
}
},
streamEnd(
state,
action: PayloadAction<{ sessionId: string; messageId: string }>
) {
const session = state.sessions[action.payload.sessionId];
if (session?.streamingMessage?.id === action.payload.messageId) {
session.streamingMessage = null;
}
},
// streamStart / streamDelta / streamEnd live in streamingSlice now.
// Mutating per-character on `session.streamingMessage` previously
// changed the top-level `sessions` dict reference 30Hz × N agents,
// forcing Dashboard (subscribed to sessions) to re-render at the same
// rate. Keeping the streaming text in a separate slice keeps the
// sessions dict stable during streaming.
addApprovalRequest(
state,
@@ -1095,7 +1062,6 @@ const agentsSlice = createSlice({
pending_approvals: existing?.pending_approvals?.length
? existing.pending_approvals
: s.pending_approvals ?? [],
streamingMessage: existing?.streamingMessage ?? s.streamingMessage ?? null,
tool_group_meta: { ...existing?.tool_group_meta, ...s.tool_group_meta },
};
if (activeStatuses.has(s.status) && !state.trackedNotificationIds.includes(s.id)) {
@@ -1107,7 +1073,7 @@ const agentsSlice = createSlice({
state.loading = false;
})
.addCase(launchAgent.fulfilled, (state, action) => {
state.sessions[action.payload.id] = { ...action.payload, streamingMessage: null, tool_group_meta: action.payload.tool_group_meta ?? {} };
state.sessions[action.payload.id] = { ...action.payload, tool_group_meta: action.payload.tool_group_meta ?? {} };
state.activeSessionId = action.payload.id;
if (!state.expandedSessionIds.includes(action.payload.id)) {
state.expandedSessionIds.push(action.payload.id);
@@ -1120,7 +1086,7 @@ const agentsSlice = createSlice({
const { draftId, session } = action.payload;
const shouldExpand = action.meta.arg.expand !== false;
delete state.sessions[draftId];
state.sessions[session.id] = { ...session, streamingMessage: null, tool_group_meta: session.tool_group_meta ?? {} };
state.sessions[session.id] = { ...session, tool_group_meta: session.tool_group_meta ?? {} };
state.activeSessionId = session.id;
state.draftLaunchMap[draftId] = session.id;
state.expandedSessionIds = state.expandedSessionIds.map((id) => (id === draftId ? session.id : id));
@@ -1170,8 +1136,11 @@ const agentsSlice = createSlice({
const session = state.sessions[action.payload];
if (session) {
session.status = 'stopped';
session.streamingMessage = null;
session.pending_approvals = [];
// streamingMessage cleanup is handled by streamingSlice via
// clearStreamingForSession. We dispatch it explicitly here
// because stopAgent.fulfilled isn't one of the action types
// we listen for in streamingSlice's extraReducers.
}
})
.addCase(handleApproval.fulfilled, (state, action) => {
@@ -1261,7 +1230,7 @@ const agentsSlice = createSlice({
})
.addCase(resumeSession.fulfilled, (state, action) => {
const session = action.payload;
state.sessions[session.id] = { ...session, streamingMessage: null, tool_group_meta: session.tool_group_meta ?? {} };
state.sessions[session.id] = { ...session, tool_group_meta: session.tool_group_meta ?? {} };
delete state.history[session.id];
state.activeSessionId = session.id;
if (!state.expandedSessionIds.includes(session.id)) {
@@ -1282,7 +1251,6 @@ const agentsSlice = createSlice({
state.sessions[session.id] = {
...session,
pending_approvals: session.pending_approvals ?? existing?.pending_approvals ?? [],
streamingMessage: existing?.streamingMessage ?? null,
tool_group_meta: session.tool_group_meta ?? existing?.tool_group_meta ?? {},
};
})
@@ -1310,7 +1278,6 @@ const agentsSlice = createSlice({
if (!state.sessions[session.id]) {
state.sessions[session.id] = {
...session,
streamingMessage: null,
tool_group_meta: session.tool_group_meta ?? {},
};
}
@@ -1358,9 +1325,6 @@ export const {
recordCompaction,
setTurnLabel,
clearTurnLabel,
streamStart,
streamDelta,
streamEnd,
addApprovalRequest,
removeApprovalRequest,
updateSessionCost,
+2
View File
@@ -1,6 +1,7 @@
import { configureStore } from '@reduxjs/toolkit';
import tempStateReducer from './tempStateSlice';
import agentsReducer from './agentsSlice';
import streamingReducer from './streamingSlice';
import skillsReducer from './skillsSlice';
import toolsReducer from './toolsSlice';
import modesReducer from './modesSlice';
@@ -20,6 +21,7 @@ export const store = configureStore({
reducer: {
tempState: tempStateReducer,
agents: agentsReducer,
streaming: streamingReducer,
skills: skillsReducer,
tools: toolsReducer,
modes: modesReducer,
+125
View File
@@ -0,0 +1,125 @@
import { createSlice, PayloadAction, createAction } from '@reduxjs/toolkit';
// Action type strings for cross-slice listening: we react to agentsSlice
// events (addMessage, editMessage, etc.) by clearing the streaming entry,
// matching the old in-place behavior. Using createAction with the same
// name lets streamingSlice's extraReducers catch the dispatch even though
// the action itself is owned by agentsSlice.
const addMessageAction = createAction<{ sessionId: string; message: { id: string } }>('agents/addMessage');
const editMessageFulfilled = createAction<{ sessionId: string }>('agents/editMessage/fulfilled');
const clearSessionMessagesAction = createAction<string>('agents/clearSessionMessages');
const closeSessionFromWsAction = createAction<{ id: string }>('agents/closeSessionFromWs');
const removeSessionAction = createAction<string>('agents/removeSession');
// stopAgent thunk's fulfilled action carries the sessionId as payload.
// Listening here lets us drop the streaming entry the moment the user
// stops a running agent, matching the previous in-place behavior.
const stopAgentFulfilledAction = createAction<string>('agents/stopAgent/fulfilled');
// Streaming-message state lives in its own slice (separate from agents/
// sessions) so that the high-frequency mutation of streamingMessage.content
// on every painted character doesn't bubble up through the sessions dict
// reference. Previously each painted character changed `state.agents.sessions`
// via Immer, causing every component subscribed to `state.agents.sessions`
// (Dashboard.tsx in particular: 30 useEffects, many selectors) to
// re-render at 30Hz × N streaming agents. Moving this out keeps the
// sessions dict stable during streaming; only structural events (start,
// end, status change, new message) mutate it now.
export interface StreamingMessage {
id: string;
role: 'assistant' | 'tool_call' | 'thinking';
content: string;
tool_name?: string;
}
interface StreamingState {
// Keyed by sessionId. Map semantics: an entry exists iff that session
// currently has an in-flight streaming message; it's removed on
// stream_end. Use selectStreamingMessage(sessionId) to read.
bySession: Record<string, StreamingMessage>;
}
const initialState: StreamingState = {
bySession: {},
};
const streamingSlice = createSlice({
name: 'streaming',
initialState,
reducers: {
streamStart(
state,
action: PayloadAction<{ sessionId: string; messageId: string; role: StreamingMessage['role']; toolName?: string }>,
) {
state.bySession[action.payload.sessionId] = {
id: action.payload.messageId,
role: action.payload.role,
content: '',
tool_name: action.payload.toolName,
};
},
streamDelta(
state,
action: PayloadAction<{ sessionId: string; messageId: string; delta: string }>,
) {
const entry = state.bySession[action.payload.sessionId];
if (entry && entry.id === action.payload.messageId) {
entry.content += action.payload.delta;
}
},
streamEnd(
state,
action: PayloadAction<{ sessionId: string; messageId: string }>,
) {
const entry = state.bySession[action.payload.sessionId];
if (entry && entry.id === action.payload.messageId) {
delete state.bySession[action.payload.sessionId];
}
},
// Used when a session is fully closed/removed so we don't leak a
// stuck streaming entry for a session that no longer exists.
clearStreamingForSession(state, action: PayloadAction<string>) {
delete state.bySession[action.payload];
},
},
extraReducers: (builder) => {
// When a final message lands for a session, clear the streaming
// entry if it matches (the streaming bubble in the UI was acting as
// a placeholder; now the real message takes over). Matches the
// original behavior that lived in agentsSlice.addMessage.
builder.addCase(addMessageAction, (state, action) => {
const entry = state.bySession[action.payload.sessionId];
if (entry && entry.id === action.payload.message.id) {
delete state.bySession[action.payload.sessionId];
}
});
// Edit / clear / close / remove all wipe any in-flight streaming
// bubble regardless of id match: the session's been mutated.
builder.addCase(editMessageFulfilled, (state, action) => {
delete state.bySession[action.payload.sessionId];
});
builder.addCase(clearSessionMessagesAction, (state, action) => {
delete state.bySession[action.payload];
});
builder.addCase(closeSessionFromWsAction, (state, action) => {
delete state.bySession[action.payload.id];
});
builder.addCase(removeSessionAction, (state, action) => {
delete state.bySession[action.payload];
});
builder.addCase(stopAgentFulfilledAction, (state, action) => {
delete state.bySession[action.payload];
});
},
});
export const { streamStart, streamDelta, streamEnd, clearStreamingForSession } = streamingSlice.actions;
export default streamingSlice.reducer;
// Reader hook. Each call subscribes only to that one session's streaming
// entry, so unrelated agents' deltas don't trigger re-renders. Returns
// null when no stream is active for the session.
import { useAppSelector } from '@/shared/hooks';
export function useStreamingMessage(sessionId: string | null | undefined) {
return useAppSelector((s) => sessionId ? s.streaming.bySession[sessionId] ?? null : null);
}
+1 -3
View File
@@ -5,9 +5,6 @@ import {
updateSessionName,
updateGroupMeta,
addMessage,
streamStart,
streamDelta,
streamEnd,
addApprovalRequest,
removeApprovalRequest,
updateSessionStatus,
@@ -25,6 +22,7 @@ import {
setTurnLabel,
clearTurnLabel,
} from '../state/agentsSlice';
import { streamStart, streamDelta, streamEnd } from '../state/streamingSlice';
import { addBrowserCardFromBackend, removeBrowserCard, setBrowserCardPosition, setGlowingBrowserCards, GRID_GAP } from '../state/dashboardLayoutSlice';
import { getAuthToken } from '../config';
import { notifyAgentCompletion } from '../notifications';