diff --git a/frontend/src/app/Main.tsx b/frontend/src/app/Main.tsx index 4e60abee..e5454440 100644 --- a/frontend/src/app/Main.tsx +++ b/frontend/src/app/Main.tsx @@ -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'; diff --git a/frontend/src/app/components/DynamicIsland.tsx b/frontend/src/app/components/DynamicIsland.tsx index bac1572f..354372e6 100644 --- a/frontend/src/app/components/DynamicIsland.tsx +++ b/frontend/src/app/components/DynamicIsland.tsx @@ -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 }> = ( /> ); +// --------------------------------------------------------------------------- +// 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 = new Map(); + +const selectDynamicIslandSessions = createSelector( + [(s: RootState) => s.agents.sessions], + (raw) => { + const out: Record = {}; + const liveIds = new Set(); + 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(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); diff --git a/frontend/src/app/components/Layout/AppShell.tsx b/frontend/src/app/components/Layout/AppShell.tsx index 00c903f3..87d53c5e 100644 --- a/frontend/src/app/components/Layout/AppShell.tsx +++ b/frontend/src/app/components/Layout/AppShell.tsx @@ -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 = () => { { + const fn = (window as any).__openswarmPrefetchRoute; + if (typeof fn === 'function') fn('/apps'); + }} data-onboarding="sidebar-apps" sx={{ borderRadius: 1.5, diff --git a/frontend/src/app/pages/AgentChat/AgentChat.tsx b/frontend/src/app/pages/AgentChat/AgentChat.tsx index 28778161..7c1b55d5 100644 --- a/frontend/src/app/pages/AgentChat/AgentChat.tsx +++ b/frontend/src/app/pages/AgentChat/AgentChat.tsx @@ -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 = ({ sessionId: sessionIdProp, onClose // so it never fires during normal streaming. const reconcileTimer = useRef | 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 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 = ({ sessionId: sessionIdProp, onClose const scrollRafRef = useRef(null); const lastScrollHeightRef = useRef(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 = ({ 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 + // 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 = ({ 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 = ({ 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 = ({ 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 = ({ sessionId: sessionIdProp, onClose ); })()} - {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 ? ( = ({ sessionId: sessionIdProp, onClose ); })} - {session.streamingMessage && ( - session.streamingMessage.role === 'tool_call' ? ( - - ) : ( - - ) + {id && ( + )} - {(awaitingResponse || (session.status === 'running' && !session.streamingMessage)) && ( + {(awaitingResponse || (session.status === 'running' && !streamingMessageId)) && ( = Object.freeze({}) as Record; + const selectBrowserSessions = createSelector( [(state: RootState) => state.agents.sessions, (_: RootState, parentSessionId: string) => parentSessionId, @@ -174,6 +181,28 @@ const BrowserAgentInlineFeed: React.FC = ({ 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 = {}; + 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 = ({ 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, ); diff --git a/frontend/src/app/pages/AgentChat/ChatInput.tsx b/frontend/src/app/pages/AgentChat/ChatInput.tsx index 3645c543..e9a229fd 100644 --- a/frontend/src/app/pages/AgentChat/ChatInput.tsx +++ b/frontend/src/app/pages/AgentChat/ChatInput.tsx @@ -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(({ onSend, disabled, mode, }), [c]); const [images, setImages] = useState([]); + // 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(null); const [isDragOver, setIsDragOver] = useState(false); const [isUploading, setIsUploading] = useState(false); @@ -773,18 +791,19 @@ const ChatInput = forwardRef(({ 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(({ 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((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(({ 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(({ 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(({ 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); diff --git a/frontend/src/app/pages/AgentChat/StreamingBubble.tsx b/frontend/src/app/pages/AgentChat/StreamingBubble.tsx new file mode 100644 index 00000000..862f3e1f --- /dev/null +++ b/frontend/src/app/pages/AgentChat/StreamingBubble.tsx @@ -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 = ({ 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(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 ( + + ); + } + + return ( + + ); +}; + +export default StreamingBubble; diff --git a/frontend/src/app/pages/Dashboard/AgentCard.tsx b/frontend/src/app/pages/Dashboard/AgentCard.tsx index de7e78a8..3fdca4c6 100644 --- a/frontend/src/app/pages/Dashboard/AgentCard.tsx +++ b/frontend/src/app/pages/Dashboard/AgentCard.tsx @@ -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 = ({ 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 = ({ 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 = ({ }; - 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 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 = ({ {session.mode} - {fmtSeconds(getAgentWorkTime(session.messages, session.status).last)} + {session.cost_usd > 0 && hasApiKey && ( diff --git a/frontend/src/app/pages/Dashboard/BrowserAgentOverlay.tsx b/frontend/src/app/pages/Dashboard/BrowserAgentOverlay.tsx index ab289d02..089086c6 100644 --- a/frontend/src/app/pages/Dashboard/BrowserAgentOverlay.tsx +++ b/frontend/src/app/pages/Dashboard/BrowserAgentOverlay.tsx @@ -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 = ({ 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 = ({ 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 = ({ 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); diff --git a/frontend/src/app/pages/Dashboard/Dashboard.tsx b/frontend/src/app/pages/Dashboard/Dashboard.tsx index d61fcbb0..121f393b 100644 --- a/frontend/src/app/pages/Dashboard/Dashboard.tsx +++ b/frontend/src/app/pages/Dashboard/Dashboard.tsx @@ -480,12 +480,28 @@ const DashboardInner: React.FC = ({ 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 = ({ 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]); diff --git a/frontend/src/app/pages/Dashboard/useCanvasControls.ts b/frontend/src/app/pages/Dashboard/useCanvasControls.ts index 944dd5d7..5d3f2894 100644 --- a/frontend/src/app/pages/Dashboard/useCanvasControls.ts +++ b/frontend/src/app/pages/Dashboard/useCanvasControls.ts @@ -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 | 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); diff --git a/frontend/src/shared/canvasInteractionState.ts b/frontend/src/shared/canvasInteractionState.ts new file mode 100644 index 00000000..e1fe8674 --- /dev/null +++ b/frontend/src/shared/canvasInteractionState.ts @@ -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); }; +} diff --git a/frontend/src/shared/state/agentsSlice.ts b/frontend/src/shared/state/agentsSlice.ts index cb7f03a1..95ef1599 100644 --- a/frontend/src/shared/state/agentsSlice.ts +++ b/frontend/src/shared/state/agentsSlice.ts @@ -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; 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; 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, diff --git a/frontend/src/shared/state/store.ts b/frontend/src/shared/state/store.ts index a9a5dbea..ada25436 100644 --- a/frontend/src/shared/state/store.ts +++ b/frontend/src/shared/state/store.ts @@ -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, diff --git a/frontend/src/shared/state/streamingSlice.ts b/frontend/src/shared/state/streamingSlice.ts new file mode 100644 index 00000000..6ff1e628 --- /dev/null +++ b/frontend/src/shared/state/streamingSlice.ts @@ -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('agents/clearSessionMessages'); +const closeSessionFromWsAction = createAction<{ id: string }>('agents/closeSessionFromWs'); +const removeSessionAction = createAction('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('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; +} + +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) { + 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); +} diff --git a/frontend/src/shared/ws/WebSocketManager.ts b/frontend/src/shared/ws/WebSocketManager.ts index 0287290c..fad8d98d 100644 --- a/frontend/src/shared/ws/WebSocketManager.ts +++ b/frontend/src/shared/ws/WebSocketManager.ts @@ -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';