Download src/contexts/conversation-websocket-context.tsx from SaylorTwift/openhands: direct link, hf CLI and curl.
- Browser
- Download file 48 kB
-
https://huggingface.co/SaylorTwift/openhands/resolve/main/src/contexts/conversation-websocket-context.tsx
- Command line
-
hf download hf://SaylorTwift/openhands/src/contexts/conversation-websocket-context.tsx
-
curl -L -o conversation-websocket-context.tsx https://huggingface.co/SaylorTwift/openhands/resolve/main/src/contexts/conversation-websocket-context.tsx
48 kB
| import React, { | |
| createContext, | |
| useContext, | |
| useEffect, | |
| useLayoutEffect, | |
| useState, | |
| useCallback, | |
| useMemo, | |
| useRef, | |
| } from "react"; | |
| import { ConversationClient } from "@openhands/typescript-client/clients"; | |
| import { useQueryClient } from "@tanstack/react-query"; | |
| import { useWebSocket, WebSocketHookOptions } from "#/hooks/use-websocket"; | |
| import { SERVER_CONNECTION_ERROR_MESSAGE } from "#/constants/server-connection-error"; | |
| import { useEventStore } from "#/stores/use-event-store"; | |
| import { useErrorMessageStore } from "#/stores/error-message-store"; | |
| import { useOptimisticUserMessageStore } from "#/stores/optimistic-user-message-store"; | |
| import { useConversationStateStore } from "#/stores/conversation-state-store"; | |
| import { useCommandStore } from "#/stores/command-store"; | |
| import { useBrowserStore } from "#/stores/browser-store"; | |
| import { useGoalStore } from "#/stores/goal-store"; | |
| import { | |
| isAgentServerEvent, | |
| isAgentErrorEvent, | |
| isUserMessageEvent, | |
| isActionEvent, | |
| isConversationStateUpdateEvent, | |
| isFullStateConversationStateUpdateEvent, | |
| isAgentStatusConversationStateUpdateEvent, | |
| isStatsConversationStateUpdateEvent, | |
| isGoalConversationStateUpdateEvent, | |
| isExecuteBashActionEvent, | |
| isExecuteBashObservationEvent, | |
| isDisplayableErrorEvent, | |
| isPlanningFileEditorObservationEvent, | |
| isBrowserObservationEvent, | |
| isBrowserNavigateActionEvent, | |
| isSwitchLLMObservationEvent, | |
| isCanvasUIActionEvent, | |
| isStreamingDeltaEvent, | |
| isLaunchChildConversationActionEvent, | |
| } from "#/types/agent-server/type-guards"; | |
| import { | |
| createStreamingDeltaBatcher, | |
| StreamingDeltaBatcher, | |
| } from "#/utils/streaming-delta-batcher"; | |
| import { handleCanvasUIAction } from "#/services/canvas-ui"; | |
| import { handleLaunchChildConversationAction } from "#/services/child-conversation-launch"; | |
| import { ConversationStateUpdateEventStats } from "#/types/agent-server/core/events/conversation-state-event"; | |
| import type { | |
| ConversationErrorEvent, | |
| ServerErrorEvent, | |
| } from "#/types/agent-server/core/events/conversation-state-event"; | |
| import { handleActionEventCacheInvalidation } from "#/utils/cache-utils"; | |
| import { buildWebSocketUrl } from "#/utils/websocket-url"; | |
| import type { | |
| AppConversation, | |
| SendMessageRequest, | |
| } from "#/api/conversation-service/agent-server-conversation-service.types"; | |
| import EventService from "#/api/event-service/event-service.api"; | |
| import { getAgentServerClientOptions } from "#/api/agent-server-client-options"; | |
| import { useConversationStore } from "#/stores/conversation-store"; | |
| import { trackError } from "#/utils/error-handler"; | |
| import { useReadConversationFile } from "#/hooks/mutation/use-read-conversation-file"; | |
| import useMetricsStore, { type MetricsState } from "#/stores/metrics-store"; | |
| import { useConversationHistory } from "#/hooks/query/use-conversation-history"; | |
| import { setConversationState } from "#/utils/conversation-local-storage"; | |
| import { | |
| recordModelSwitchMessage, | |
| seedModelSwitchesFromHistory, | |
| stampActiveLlmProfile, | |
| } from "#/hooks/chat/record-model-switch-message"; | |
| import { | |
| invalidateConversationQueries, | |
| updateConversationLlmModelInCache, | |
| } from "#/hooks/mutation/conversation-mutation-utils"; | |
| import { isPlanFilePath } from "#/utils/plan-file"; | |
| export type WebSocketConnectionState = | |
| | "CONNECTING" | |
| | "OPEN" | |
| | "CLOSED" | |
| | "CLOSING"; | |
| interface SendMessageResult { | |
| queued: boolean; // true if message was queued for later delivery, false if sent immediately | |
| } | |
| interface ConversationWebSocketContextType { | |
| connectionState: WebSocketConnectionState; | |
| /** | |
| * The main connection's own state, unmerged with the planning connection — | |
| * see `useMainWebSocketStatus`. | |
| */ | |
| mainConnectionState: WebSocketConnectionState; | |
| sendMessage: (message: SendMessageRequest) => Promise<SendMessageResult>; | |
| isLoadingHistory: boolean; | |
| reconnect: () => void; | |
| } | |
| const ConversationWebSocketContext = createContext< | |
| ConversationWebSocketContextType | undefined | |
| >(undefined); | |
| /** | |
| * Extract the text body of an echoed user `MessageEvent` for matching against | |
| * the optimistic pending-message queue. The server wraps the original | |
| * `args.content` string in one or more `TextContent` entries (alongside any | |
| * `ImageContent` entries for inline images), so concatenating the `text` | |
| * fields gives us back the exact prompt we sent. | |
| */ | |
| function extractMessageEventText( | |
| event: import("#/types/agent-server/core/events/message-event").MessageEvent, | |
| ): string { | |
| return event.llm_message.content | |
| .filter( | |
| (part): part is { type: "text"; text: string } => part.type === "text", | |
| ) | |
| .map((part) => part.text) | |
| .join(""); | |
| } | |
| export function ConversationWebSocketProvider({ | |
| children, | |
| conversationId, | |
| conversationUrl, | |
| sessionApiKey, | |
| subConversations, | |
| subConversationIds, | |
| }: { | |
| children: React.ReactNode; | |
| conversationId?: string; | |
| conversationUrl?: string | null; | |
| sessionApiKey?: string | null; | |
| subConversations?: AppConversation[]; | |
| subConversationIds?: string[]; | |
| }) { | |
| // Separate connection state tracking for each WebSocket | |
| const [mainConnectionState, setMainConnectionState] = | |
| useState<WebSocketConnectionState>("CONNECTING"); | |
| const [planningConnectionState, setPlanningConnectionState] = | |
| useState<WebSocketConnectionState>("CONNECTING"); | |
| // Track if we've ever successfully connected for each connection | |
| // Don't show errors until after first successful connection | |
| const hasConnectedRefMain = React.useRef(false); | |
| const hasConnectedRefPlanning = React.useRef(false); | |
| const queryClient = useQueryClient(); | |
| const addEvent = useEventStore((state) => state.addEvent); | |
| const addEvents = useEventStore((state) => state.addEvents); | |
| const clearEventsForConversation = useEventStore( | |
| (state) => state.clearEventsForConversation, | |
| ); | |
| const { setErrorMessage, removeErrorMessage, clearConnectionError } = | |
| useErrorMessageStore(); | |
| const consumeMatchingPendingMessage = useOptimisticUserMessageStore( | |
| (state) => state.consumeMatchingPendingMessage, | |
| ); | |
| const { setExecutionStatus } = useConversationStateStore(); | |
| const { appendInput, appendOutput } = useCommandStore(); | |
| const resetBrowserStore = useBrowserStore((state) => state.reset); | |
| // Coalesce streaming deltas to ≤1 store commit/render per frame. | |
| // Separate batchers keep the main and planning streams from ever merging. | |
| const mainDeltaBatcherRef = useRef<StreamingDeltaBatcher | null>(null); | |
| if (mainDeltaBatcherRef.current === null) { | |
| mainDeltaBatcherRef.current = createStreamingDeltaBatcher((delta) => { | |
| useEventStore.getState().addEvent(delta); | |
| // A delta means connectivity recovered — mirror handleNonErrorEvent. | |
| useErrorMessageStore.getState().clearConnectionError(); | |
| }); | |
| } | |
| const planningDeltaBatcherRef = useRef<StreamingDeltaBatcher | null>(null); | |
| if (planningDeltaBatcherRef.current === null) { | |
| planningDeltaBatcherRef.current = createStreamingDeltaBatcher((delta) => { | |
| useEventStore | |
| .getState() | |
| .addEvent({ ...delta, isFromPlanningAgent: true }); | |
| useErrorMessageStore.getState().clearConnectionError(); | |
| }); | |
| } | |
| // History loading state. | |
| // - Main conversation history is now loaded via REST (`useConversationHistory`), | |
| // so its loading state mirrors the REST query state (see below). | |
| // - Planning sub-conversation history still streams over the WebSocket using | |
| // `resend_mode='all'`, so we keep the count-based detection for it. | |
| const [isLoadingHistoryPlanning, setIsLoadingHistoryPlanning] = | |
| useState(true); | |
| const [expectedEventCountPlanning, setExpectedEventCountPlanning] = useState< | |
| number | null | |
| >(null); | |
| const { setPlanContent } = useConversationStore(); | |
| useEffect(() => { | |
| setPlanContent(null); | |
| }, [conversationId, setPlanContent]); | |
| const { mutate: readConversationFile } = useReadConversationFile(); | |
| // Track planning-agent received events (still WS-driven). | |
| const receivedEventCountRefPlanning = useRef(0); | |
| // Track the latest PlanningFileEditorObservation for Plan.md during history replay | |
| const latestPlanningFileEventRef = useRef<{ | |
| path: string; | |
| conversationId: string; | |
| } | null>(null); | |
| const handleNonErrorEvent = useCallback(() => { | |
| // A normal event means connectivity recovered: clear a transient connection | |
| // error, but keep sticky conversation errors (e.g. a wrong API key). | |
| clearConnectionError(); | |
| }, [clearConnectionError]); | |
| // Helper function to update metrics from stats event | |
| const updateMetricsFromStats = useCallback( | |
| (event: ConversationStateUpdateEventStats) => { | |
| const usageToMetrics = event.value.usage_to_metrics; | |
| if (!usageToMetrics) { | |
| return; | |
| } | |
| // usage_to_metrics is keyed by arbitrary LLM usage ids ("default", | |
| // "condenser", "profile:<name>:<uuid>", …) — combine across all of | |
| // them, mirroring getCombinedMetrics on the REST path. | |
| const combined = Object.values(usageToMetrics).reduce<{ | |
| cost: number; | |
| maxBudgetPerTask: number | null; | |
| usage: MetricsState["usage"]; | |
| }>( | |
| (acc, metrics) => { | |
| acc.cost += metrics.accumulated_cost; | |
| if ( | |
| acc.maxBudgetPerTask === null && | |
| metrics.max_budget_per_task !== null | |
| ) { | |
| acc.maxBudgetPerTask = metrics.max_budget_per_task; | |
| } | |
| const tokenUsage = metrics.accumulated_token_usage; | |
| if (tokenUsage) { | |
| acc.usage = { | |
| prompt_tokens: | |
| (acc.usage?.prompt_tokens ?? 0) + tokenUsage.prompt_tokens, | |
| completion_tokens: | |
| (acc.usage?.completion_tokens ?? 0) + | |
| tokenUsage.completion_tokens, | |
| cache_read_tokens: | |
| (acc.usage?.cache_read_tokens ?? 0) + | |
| tokenUsage.cache_read_tokens, | |
| cache_write_tokens: | |
| (acc.usage?.cache_write_tokens ?? 0) + | |
| tokenUsage.cache_write_tokens, | |
| context_window: Math.max( | |
| acc.usage?.context_window ?? 0, | |
| tokenUsage.context_window, | |
| ), | |
| per_turn_token: Math.max( | |
| acc.usage?.per_turn_token ?? 0, | |
| tokenUsage.per_turn_token, | |
| ), | |
| }; | |
| } | |
| return acc; | |
| }, | |
| { cost: 0, maxBudgetPerTask: null, usage: null }, | |
| ); | |
| useMetricsStore.getState().setMetrics({ | |
| cost: combined.cost, | |
| max_budget_per_task: combined.maxBudgetPerTask, | |
| usage: combined.usage, | |
| }); | |
| }, | |
| [], | |
| ); | |
| // Initial REST history load: fetch the most recent events and seed the | |
| // store. Older events are paginated in via `useLoadOlderEvents` when the | |
| // user scrolls to the top of the chat. The WebSocket connection waits for | |
| // this query so it can subscribe with `resend_mode='since'` and avoid | |
| // re-streaming everything REST already returned. | |
| const { data: preloadedHistory, isPending: isPreloadingHistory } = | |
| useConversationHistory(conversationId); | |
| // Skeleton only on the genuine first load (no cached data yet). On return the | |
| // cached page is present, so `isPending` is false and we render the | |
| // last-known discussion immediately while the tail refetch runs in the | |
| // background — the socket gate below also keys on `isPending`, so that | |
| // refetch never drops a live socket. | |
| const isLoadingHistoryMain = !!conversationId && isPreloadingHistory; | |
| // Clear the (global, not conversation-scoped) event store when the active | |
| // conversation changes, BEFORE the preloaded-history effect below re-seeds | |
| // it. This MUST live here rather than in the route component: a parent's | |
| // passive effect runs *after* this child's layout effects, so clearing from | |
| // the route would wipe the freshly seeded history. On a conversation switch | |
| // the history page is already cached, so `preloadedHistory` is available | |
| // synchronously — without ordering the clear first, the user's already-echoed | |
| // message gets seeded then immediately wiped, leaving only the `since` | |
| // WebSocket resend (the agent's reply). Re-entering the same conversation is | |
| // a no-op, so the store survives navigating away to Settings and back. | |
| useLayoutEffect(() => { | |
| const nextId = conversationId ?? null; | |
| if (useEventStore.getState().loadedConversationId === nextId) { | |
| return; | |
| } | |
| // Single atomic action: clears the previous conversation's events and | |
| // records the new loaded id in one `set`, so no subscriber can observe a | |
| // half-applied state (events gone but the old id still reported). | |
| clearEventsForConversation(nextId); | |
| resetBrowserStore(); | |
| // The metrics store is conversation-scoped state too: without a reset the | |
| // previous conversation's usage/cost keeps rendering in the new | |
| // conversation's meter until fresh WS stats arrive — and a brand-new | |
| // conversation sends none, so the stale figure stuck indefinitely. | |
| useMetricsStore.getState().resetMetrics(); | |
| }, [conversationId, clearEventsForConversation, resetBrowserStore]); | |
| useLayoutEffect(() => { | |
| if (!preloadedHistory || preloadedHistory.events.length === 0) { | |
| return; | |
| } | |
| addEvents(preloadedHistory.events); | |
| // The first user message of a cloud start-task conversation is persisted | |
| // server-side and reaches us via this REST preload, not over the WebSocket | |
| // (which subscribes with resend_mode='since' after the latest preloaded | |
| // timestamp). Consume any matching optimistic "Sending…" bubble here too — | |
| // mirroring the WS handler — so it doesn't linger as a duplicate of the echo. | |
| if (conversationId) { | |
| // Rebuild inline "Switched to" messages from the REST-preloaded history. | |
| // The live store writers (WS handler / user action) never see preloaded | |
| // events, so without this past model switches wouldn't render on reload. | |
| // Read the post-`addEvents` `uiEvents` (actions replaced by observations, | |
| // Think/Finish observations dropped) — not the raw history — so anchors | |
| // match the ids the renderer actually mounts. | |
| seedModelSwitchesFromHistory( | |
| conversationId, | |
| useEventStore.getState().uiEvents, | |
| ); | |
| for (const event of preloadedHistory.events) { | |
| if (isUserMessageEvent(event)) { | |
| consumeMatchingPendingMessage( | |
| conversationId, | |
| extractMessageEventText(event), | |
| ); | |
| } | |
| } | |
| } | |
| }, [ | |
| preloadedHistory, | |
| addEvents, | |
| conversationId, | |
| consumeMatchingPendingMessage, | |
| ]); | |
| /** | |
| * Timestamp of the latest event we already have from REST. Used as | |
| * `after_timestamp` when opening the WebSocket so the server only resends | |
| * events strictly after this point. `null` until the first REST page lands | |
| * (the WS connection is gated on that — see `wsUrl` below). During | |
| * background refetches `preloadedHistory` keeps the last-known page, so the | |
| * anchor holds steady instead of flipping to null; reconnects read the | |
| * freshest value from the options ref at connect time. | |
| */ | |
| const initialAfterTimestamp = useMemo<string | null>(() => { | |
| const events = preloadedHistory?.events ?? []; | |
| const latest = events[events.length - 1]; | |
| if (!latest || !("timestamp" in latest) || !latest.timestamp) return null; | |
| return latest.timestamp; | |
| }, [preloadedHistory]); | |
| // Build WebSocket URL from props. | |
| // | |
| // We deliberately wait for the FIRST history load (`isPending`: no data for | |
| // this query key yet) before opening the socket, so the WS subscription can | |
| // use `resend_mode='since'` with a meaningful `after_timestamp` instead of | |
| // falling back to `resend_mode='all'`. The gate is intentionally NOT on | |
| // `isFetching`: background refetches (e.g. the `refetchOnMount` fired when | |
| // returning to a conversation) must never tear a live socket down — on a | |
| // flaky link that caused a refetch → teardown → reconnect → refetch loop | |
| // that kept the conversation stuck at "Connecting" for minutes. Connecting | |
| // during a background refetch anchors `since` to the cached tail; the | |
| // overlap with the refetched page is deduped by the event store and the | |
| // `isDuplicateEvent` guards in the message handlers. A query-key reset | |
| // (backend swap / new session key) makes `isPending` true again, so a | |
| // genuine reset still re-gates. If the initial load errors, `isPending` | |
| // flips false and we fall through to connect with `resend_mode='all'` so | |
| // the user still sees live events. | |
| const wsUrl = useMemo(() => { | |
| if (!conversationId || !conversationUrl) { | |
| return null; | |
| } | |
| if (isPreloadingHistory) { | |
| return null; | |
| } | |
| return buildWebSocketUrl(conversationId, conversationUrl); | |
| }, [conversationId, conversationUrl, isPreloadingHistory]); | |
| // Derived from `subConversationIds` (the pre-filtered, tag-verified id | |
| // list) rather than the resolved `subConversations` entry, which lands a | |
| // tick later — routing on it would leave a window where the first prompt | |
| // sent in plan mode fell through to the parent (the code agent) instead of | |
| // the planner. | |
| const planningConversationId = useMemo( | |
| () => subConversationIds?.[0] ?? null, | |
| [subConversationIds], | |
| ); | |
| const planningAgentWsUrl = useMemo(() => { | |
| if (!subConversations?.length) { | |
| return null; | |
| } | |
| // Currently, there is only one sub-conversation and it uses the planning agent. | |
| const planningAgentConversation = subConversations[0]; | |
| if ( | |
| !planningAgentConversation?.id || | |
| !planningAgentConversation.conversation_url | |
| ) { | |
| return null; | |
| } | |
| return buildWebSocketUrl( | |
| planningAgentConversation.id, | |
| planningAgentConversation.conversation_url, | |
| ); | |
| }, [subConversations]); | |
| // Merged connection state - reflects combined status of both connections | |
| const connectionState = useMemo<WebSocketConnectionState>(() => { | |
| // If planning agent connection doesn't exist, use main connection state | |
| if (!planningAgentWsUrl) { | |
| return mainConnectionState; | |
| } | |
| // If either is connecting, merged state is connecting | |
| if ( | |
| mainConnectionState === "CONNECTING" || | |
| planningConnectionState === "CONNECTING" | |
| ) { | |
| return "CONNECTING"; | |
| } | |
| // If both are open, merged state is open | |
| if (mainConnectionState === "OPEN" && planningConnectionState === "OPEN") { | |
| return "OPEN"; | |
| } | |
| // If both are closed, merged state is closed | |
| if ( | |
| mainConnectionState === "CLOSED" && | |
| planningConnectionState === "CLOSED" | |
| ) { | |
| return "CLOSED"; | |
| } | |
| // If either is closing, merged state is closing | |
| if ( | |
| mainConnectionState === "CLOSING" || | |
| planningConnectionState === "CLOSING" | |
| ) { | |
| return "CLOSING"; | |
| } | |
| // Default to closed if states don't match expected patterns | |
| return "CLOSED"; | |
| }, [mainConnectionState, planningConnectionState, planningAgentWsUrl]); | |
| useEffect(() => { | |
| if ( | |
| expectedEventCountPlanning !== null && | |
| receivedEventCountRefPlanning.current >= expectedEventCountPlanning && | |
| isLoadingHistoryPlanning | |
| ) { | |
| setIsLoadingHistoryPlanning(false); | |
| } | |
| }, [ | |
| expectedEventCountPlanning, | |
| isLoadingHistoryPlanning, | |
| receivedEventCountRefPlanning, | |
| ]); | |
| // Call API once after history loading completes if we tracked any PlanningFileEditorObservation events | |
| useEffect(() => { | |
| if (!isLoadingHistoryPlanning && latestPlanningFileEventRef.current) { | |
| const { path, conversationId: currentPlanningConversationId } = | |
| latestPlanningFileEventRef.current; | |
| readConversationFile( | |
| { | |
| conversationId: currentPlanningConversationId, | |
| filePath: path, | |
| }, | |
| { | |
| onSuccess: (fileContent) => { | |
| setPlanContent(fileContent); | |
| }, | |
| onError: (error) => { | |
| console.warn("Failed to read conversation file:", error); | |
| }, | |
| }, | |
| ); | |
| // Clear the ref after calling the API | |
| latestPlanningFileEventRef.current = null; | |
| } | |
| }, [isLoadingHistoryPlanning, readConversationFile, setPlanContent]); | |
| useEffect(() => { | |
| hasConnectedRefMain.current = false; | |
| setIsLoadingHistoryPlanning(!!subConversationIds?.length); | |
| setExpectedEventCountPlanning(null); | |
| receivedEventCountRefPlanning.current = 0; | |
| // Reset the tracked event ref when sub-conversations change | |
| latestPlanningFileEventRef.current = null; | |
| }, [subConversationIds]); | |
| // Reset hasConnected flags when the conversation changes. | |
| useEffect(() => { | |
| hasConnectedRefMain.current = false; | |
| hasConnectedRefPlanning.current = false; | |
| // Reset the tracked event ref when conversation changes | |
| latestPlanningFileEventRef.current = null; | |
| }, [conversationId]); | |
| // Drop buffered deltas on conversation switch/unmount: the store is cleared on | |
| // switch, so flushing them would leak into the next conversation. | |
| useEffect(() => { | |
| const mainBatcher = mainDeltaBatcherRef.current; | |
| const planningBatcher = planningDeltaBatcherRef.current; | |
| return () => { | |
| mainBatcher?.reset(); | |
| planningBatcher?.reset(); | |
| }; | |
| }, [conversationId]); | |
| // Merged loading history state - true if either connection is still loading | |
| const isLoadingHistory = useMemo( | |
| () => isLoadingHistoryMain || isLoadingHistoryPlanning, | |
| [isLoadingHistoryMain, isLoadingHistoryPlanning], | |
| ); | |
| // Separate message handlers for each connection | |
| const handleMainMessage = useCallback( | |
| (messageEvent: MessageEvent) => { | |
| try { | |
| const event = JSON.parse(messageEvent.data); | |
| // History loading for the main conversation is REST-driven now; | |
| // every WS message is a new event we add to the store. | |
| // Use type guard to validate v1 event structure | |
| if (isAgentServerEvent(event)) { | |
| // Buffer deltas; nothing else in this handler applies to them. | |
| if (isStreamingDeltaEvent(event)) { | |
| mainDeltaBatcherRef.current?.enqueue(event); | |
| return; | |
| } | |
| // Flush buffered deltas before this event so it can't overtake them. | |
| mainDeltaBatcherRef.current?.flush(); | |
| // A reconnect replays the backlog from a stale anchor. The store | |
| // dedups by id, but the side-effects below aren't idempotent, so skip | |
| // them for replayed events (#1656). | |
| const isDuplicateEvent = useEventStore | |
| .getState() | |
| .eventIds.has(event.id ?? ""); | |
| const switchLLMObservation = isSwitchLLMObservationEvent(event) | |
| ? event | |
| : null; | |
| addEvent(event); | |
| if (isDuplicateEvent) { | |
| return; | |
| } | |
| // Handle displayable error events - show error banner | |
| // AgentErrorEvent errors are displayed inline in the chat, not as banners | |
| if (isDisplayableErrorEvent(event)) { | |
| const errorEvent = event as | |
| | ConversationErrorEvent | |
| | ServerErrorEvent; | |
| const classification = | |
| "classification" in errorEvent ? errorEvent.classification : null; | |
| trackError({ | |
| source: "conversation", | |
| metadata: { | |
| eventId: errorEvent.id, | |
| errorCode: errorEvent.code, | |
| }, | |
| classification, | |
| }); | |
| setErrorMessage( | |
| errorEvent.detail, | |
| "conversation", | |
| errorEvent.code, | |
| classification, | |
| ); | |
| } else { | |
| handleNonErrorEvent(); | |
| } | |
| // LLM errors render inline in the chat (see ErrorEventMessage); track | |
| // them for analytics but keep them out of the banner above the chat box. | |
| if (isAgentErrorEvent(event)) { | |
| trackError({ | |
| source: "agent", | |
| metadata: { | |
| eventId: event.id, | |
| toolName: event.tool_name, | |
| toolCallId: event.tool_call_id, | |
| }, | |
| classification: event.classification, | |
| }); | |
| } | |
| // Clear optimistic user message when a user message is confirmed. | |
| // We match by the echoed text content (with FIFO fallback inside the | |
| // store), so an echo for "second" pops "second" — not whichever | |
| // pending entry happens to be oldest — protecting against any | |
| // out-of-order delivery between conversations or sub-agents. | |
| if (isUserMessageEvent(event)) { | |
| if (conversationId) { | |
| consumeMatchingPendingMessage( | |
| conversationId, | |
| extractMessageEventText(event), | |
| ); | |
| // Clear draft from localStorage - message was successfully delivered | |
| setConversationState(conversationId, { draftMessage: null }); | |
| } | |
| } | |
| // Handle cache invalidation for ActionEvent | |
| if (isActionEvent(event)) { | |
| const currentConversationId = | |
| conversationId || "test-conversation-id"; // TODO: Get from context | |
| handleActionEventCacheInvalidation( | |
| event, | |
| currentConversationId, | |
| queryClient, | |
| ); | |
| } | |
| // Handle conversation state updates | |
| // TODO: Tests | |
| if (isConversationStateUpdateEvent(event)) { | |
| if ( | |
| isFullStateConversationStateUpdateEvent(event) && | |
| conversationId | |
| ) { | |
| setExecutionStatus(conversationId, event.value.execution_status); | |
| } | |
| if ( | |
| isAgentStatusConversationStateUpdateEvent(event) && | |
| conversationId | |
| ) { | |
| setExecutionStatus(conversationId, event.value); | |
| } | |
| if (isStatsConversationStateUpdateEvent(event)) { | |
| updateMetricsFromStats(event); | |
| } | |
| // Mirror goal status into the store. Intentionally duplicated across | |
| // the main and planning WebSocket handlers (like the execution_status | |
| // and stats branches above), not a merge artifact. | |
| if (isGoalConversationStateUpdateEvent(event) && conversationId) { | |
| useGoalStore.getState().setStatus(conversationId, event.value); | |
| } | |
| } | |
| // Handle ExecuteBashAction events - add command as input to terminal | |
| if (isExecuteBashActionEvent(event)) { | |
| appendInput(event.action.command); | |
| } | |
| // Handle ExecuteBashObservation events - add output to terminal | |
| if (isExecuteBashObservationEvent(event)) { | |
| // Extract text content from the observation content array | |
| const textContent = event.observation.content | |
| .filter((c) => c.type === "text") | |
| .map((c) => c.text) | |
| .join("\n"); | |
| appendOutput(textContent); | |
| } | |
| // Handle BrowserObservation events - update browser store with screenshot | |
| if (isBrowserObservationEvent(event)) { | |
| const { screenshot_data: screenshotData } = event.observation; | |
| if (screenshotData) { | |
| const screenshotSrc = screenshotData.startsWith("data:") | |
| ? screenshotData | |
| : `data:image/png;base64,${screenshotData}`; | |
| useBrowserStore.getState().setScreenshotSrc(screenshotSrc); | |
| } | |
| } | |
| // Handle BrowserNavigateAction events - update browser store with URL | |
| if (isBrowserNavigateActionEvent(event)) { | |
| useBrowserStore.getState().setUrl(event.action.url); | |
| } | |
| if ( | |
| conversationId && | |
| switchLLMObservation && | |
| !switchLLMObservation.observation.is_error | |
| ) { | |
| recordModelSwitchMessage( | |
| conversationId, | |
| switchLLMObservation.observation.profile_name, | |
| ); | |
| // Mirror the user-driven `/model` path: persist the profile so the | |
| // chat-header switcher shows the right name after a reload, even | |
| // when several profiles share a model (#1082). Stamp with the | |
| // observation's own timestamp so a later history seed of this same | |
| // event can't roll it back (or needlessly rewrite it). | |
| stampActiveLlmProfile( | |
| conversationId, | |
| switchLLMObservation.observation.profile_name, | |
| switchLLMObservation.timestamp, | |
| ); | |
| if (switchLLMObservation.observation.active_model) { | |
| updateConversationLlmModelInCache( | |
| queryClient, | |
| conversationId, | |
| switchLLMObservation.observation.active_model, | |
| ); | |
| } | |
| invalidateConversationQueries(queryClient, conversationId); | |
| } | |
| // Handle canvas_ui ActionEvents from both the legacy Python tool and | |
| // the client-defined JSON tool. The server acknowledges immediately; | |
| // the actual UI change happens here on the client. | |
| if (isCanvasUIActionEvent(event)) { | |
| handleCanvasUIAction(event.action, conversationId ?? null); | |
| } | |
| // Same client-tool pattern, but the work is a network call: launch | |
| // the requested child conversation and post the outcome back so the | |
| // agent learns the id the server-side acknowledgement can't carry. | |
| if (conversationId && isLaunchChildConversationActionEvent(event)) { | |
| void handleLaunchChildConversationAction( | |
| event.action, | |
| conversationId, | |
| event.tool_call_id, | |
| ); | |
| } | |
| } | |
| } catch (error) { | |
| console.warn("Failed to parse WebSocket message as JSON:", error); | |
| } | |
| }, | |
| [ | |
| addEvent, | |
| setErrorMessage, | |
| consumeMatchingPendingMessage, | |
| queryClient, | |
| conversationId, | |
| setExecutionStatus, | |
| appendInput, | |
| appendOutput, | |
| updateMetricsFromStats, | |
| handleNonErrorEvent, | |
| ], | |
| ); | |
| const handlePlanningMessage = useCallback( | |
| (messageEvent: MessageEvent) => { | |
| try { | |
| const event = JSON.parse(messageEvent.data); | |
| // Track received events for history loading (count ALL events from WebSocket) | |
| // Always count when loading, even if we don't have the expected count yet | |
| if (isLoadingHistoryPlanning) { | |
| receivedEventCountRefPlanning.current += 1; | |
| if ( | |
| expectedEventCountPlanning !== null && | |
| receivedEventCountRefPlanning.current >= expectedEventCountPlanning | |
| ) { | |
| setIsLoadingHistoryPlanning(false); | |
| } | |
| } | |
| // Use type guard to validate v1 event structure | |
| if (isAgentServerEvent(event)) { | |
| // Buffer deltas (the commit re-applies the planning flag). | |
| if (isStreamingDeltaEvent(event)) { | |
| planningDeltaBatcherRef.current?.enqueue(event); | |
| return; | |
| } | |
| // Flush buffered deltas before this event so it can't overtake them. | |
| planningDeltaBatcherRef.current?.flush(); | |
| // Skip non-idempotent side-effects for replayed events, as in the | |
| // main handler (#1656). | |
| const isDuplicateEvent = useEventStore | |
| .getState() | |
| .eventIds.has(event.id ?? ""); | |
| // Mark this event as coming from the planning agent | |
| const eventWithPlanningFlag = { | |
| ...event, | |
| isFromPlanningAgent: true, | |
| }; | |
| addEvent(eventWithPlanningFlag); | |
| if (isDuplicateEvent) { | |
| return; | |
| } | |
| // Handle displayable error events - show error banner | |
| // AgentErrorEvent errors are displayed inline in the chat, not as banners | |
| if (isDisplayableErrorEvent(event)) { | |
| const errorEvent = event as | |
| | ConversationErrorEvent | |
| | ServerErrorEvent; | |
| const classification = | |
| "classification" in errorEvent ? errorEvent.classification : null; | |
| trackError({ | |
| source: "planning_conversation", | |
| metadata: { | |
| eventId: errorEvent.id, | |
| errorCode: errorEvent.code, | |
| }, | |
| classification, | |
| }); | |
| setErrorMessage( | |
| errorEvent.detail, | |
| "conversation", | |
| errorEvent.code, | |
| classification, | |
| ); | |
| } else { | |
| handleNonErrorEvent(); | |
| } | |
| // LLM errors render inline in the chat (see ErrorEventMessage); track | |
| // them for analytics but keep them out of the banner above the chat box. | |
| if (isAgentErrorEvent(event)) { | |
| trackError({ | |
| source: "planning_agent", | |
| metadata: { | |
| eventId: event.id, | |
| toolName: event.tool_name, | |
| toolCallId: event.tool_call_id, | |
| }, | |
| classification: event.classification, | |
| }); | |
| } | |
| // Clear optimistic user message when a user message is confirmed. | |
| // Always scope to the main `conversationId` (where the user types) | |
| // and match on the echoed content so the planning sub-agent's own | |
| // events can never consume a main-conversation pending entry. | |
| if (isUserMessageEvent(event)) { | |
| if (conversationId) { | |
| consumeMatchingPendingMessage( | |
| conversationId, | |
| extractMessageEventText(event), | |
| ); | |
| setConversationState(conversationId, { draftMessage: null }); | |
| } | |
| } | |
| // Handle cache invalidation for ActionEvent | |
| if (isActionEvent(event)) { | |
| const planningAgentConversation = subConversations?.[0]; | |
| const currentConversationId = | |
| planningAgentConversation?.id || "test-conversation-id"; // TODO: Get from context | |
| handleActionEventCacheInvalidation( | |
| event, | |
| currentConversationId, | |
| queryClient, | |
| ); | |
| } | |
| // Handle conversation state updates | |
| // TODO: Tests | |
| if (isConversationStateUpdateEvent(event)) { | |
| // Scope to the planning agent's own conversation id, not the main | |
| // `conversationId` — this socket reports the planning helper | |
| // conversation's run/idle transitions, which must never overwrite | |
| // the main conversation's status in the shared store. | |
| if ( | |
| isFullStateConversationStateUpdateEvent(event) && | |
| planningConversationId | |
| ) { | |
| setExecutionStatus( | |
| planningConversationId, | |
| event.value.execution_status, | |
| ); | |
| } | |
| if ( | |
| isAgentStatusConversationStateUpdateEvent(event) && | |
| planningConversationId | |
| ) { | |
| setExecutionStatus(planningConversationId, event.value); | |
| } | |
| if (isStatsConversationStateUpdateEvent(event)) { | |
| updateMetricsFromStats(event); | |
| } | |
| // Mirror goal status into the store. Intentionally duplicated across | |
| // the main and planning WebSocket handlers (like the execution_status | |
| // and stats branches above), not a merge artifact. | |
| if (isGoalConversationStateUpdateEvent(event) && conversationId) { | |
| useGoalStore.getState().setStatus(conversationId, event.value); | |
| } | |
| } | |
| // Handle ExecuteBashAction events - add command as input to terminal | |
| if (isExecuteBashActionEvent(event)) { | |
| appendInput(event.action.command); | |
| } | |
| // Handle ExecuteBashObservation events - add output to terminal | |
| if (isExecuteBashObservationEvent(event)) { | |
| // Extract text content from the observation content array | |
| const textContent = event.observation.content | |
| .filter((c) => c.type === "text") | |
| .map((c) => c.text) | |
| .join("\n"); | |
| appendOutput(textContent); | |
| } | |
| // Handle PlanningFileEditorObservation - only update plan for Plan.md | |
| if (isPlanningFileEditorObservationEvent(event)) { | |
| const { path } = event.observation; | |
| if (isPlanFilePath(path)) { | |
| const planningAgentConversation = subConversations?.[0]; | |
| const planningConversationId = planningAgentConversation?.id; | |
| if (planningConversationId && path) { | |
| if (isLoadingHistoryPlanning) { | |
| latestPlanningFileEventRef.current = { | |
| path, | |
| conversationId: planningConversationId, | |
| }; | |
| } else { | |
| readConversationFile( | |
| { | |
| conversationId: planningConversationId, | |
| filePath: path, | |
| }, | |
| { | |
| onSuccess: (fileContent) => { | |
| setPlanContent(fileContent); | |
| }, | |
| onError: (error) => { | |
| console.warn( | |
| "Failed to read conversation file:", | |
| error, | |
| ); | |
| }, | |
| }, | |
| ); | |
| } | |
| } | |
| } | |
| } | |
| } | |
| } catch (error) { | |
| console.warn("Failed to parse WebSocket message as JSON:", error); | |
| } | |
| }, | |
| [ | |
| addEvent, | |
| isLoadingHistoryPlanning, | |
| expectedEventCountPlanning, | |
| setErrorMessage, | |
| consumeMatchingPendingMessage, | |
| queryClient, | |
| subConversations, | |
| planningConversationId, | |
| conversationId, | |
| setExecutionStatus, | |
| appendInput, | |
| appendOutput, | |
| readConversationFile, | |
| setPlanContent, | |
| updateMetricsFromStats, | |
| handleNonErrorEvent, | |
| ], | |
| ); | |
| // Separate WebSocket options for main connection | |
| const mainWebsocketOptions: WebSocketHookOptions = useMemo(() => { | |
| // History was already loaded over REST (`useConversationHistory`). | |
| // Subscribe with `resend_mode='since'` so the server only resends events | |
| // strictly after the latest one we already have. If REST returned no | |
| // events at all (brand-new conversation), fall back to `'all'` so any | |
| // events that may have been written between the REST call and the WS | |
| // handshake still show up. Dedup in the event store handles overlap. | |
| const queryParams: Record<string, string | boolean> = initialAfterTimestamp | |
| ? { resend_mode: "since", after_timestamp: initialAfterTimestamp } | |
| : { resend_mode: "all" }; | |
| return { | |
| queryParams, | |
| sessionApiKey, | |
| reconnect: { enabled: true }, | |
| onOpen: () => { | |
| setMainConnectionState("OPEN"); | |
| hasConnectedRefMain.current = true; // Mark that we've successfully connected | |
| clearConnectionError(); // Clear a previous connection error; keep sticky conversation errors | |
| }, | |
| onClose: () => { | |
| setMainConnectionState("CLOSED"); | |
| }, | |
| onError: () => { | |
| setMainConnectionState("CLOSED"); | |
| // Only show error message if we've previously connected successfully | |
| if (hasConnectedRefMain.current) { | |
| setErrorMessage(SERVER_CONNECTION_ERROR_MESSAGE, "connection"); | |
| } | |
| }, | |
| onMessage: handleMainMessage, | |
| }; | |
| }, [ | |
| handleMainMessage, | |
| setErrorMessage, | |
| clearConnectionError, | |
| sessionApiKey, | |
| initialAfterTimestamp, | |
| ]); | |
| // Separate WebSocket options for planning agent connection | |
| const planningWebsocketOptions: WebSocketHookOptions = useMemo(() => { | |
| const queryParams: Record<string, string | boolean> = { | |
| resend_all: true, | |
| }; | |
| const planningAgentConversation = subConversations?.[0]; | |
| const planningApiKey = | |
| planningAgentConversation?.session_api_key ?? sessionApiKey; | |
| return { | |
| queryParams, | |
| sessionApiKey: planningApiKey, | |
| reconnect: { enabled: true }, | |
| onOpen: async () => { | |
| setPlanningConnectionState("OPEN"); | |
| hasConnectedRefPlanning.current = true; // Mark that we've successfully connected | |
| clearConnectionError(); // Clear a previous connection error; keep sticky conversation errors | |
| // Fetch expected event count for history loading detection | |
| if ( | |
| planningAgentConversation?.id && | |
| planningAgentConversation.conversation_url | |
| ) { | |
| try { | |
| const count = await EventService.getEventCount( | |
| planningAgentConversation.id, | |
| planningAgentConversation.conversation_url, | |
| planningAgentConversation.session_api_key, | |
| ); | |
| setExpectedEventCountPlanning(count); | |
| // If no events expected, mark as loaded immediately | |
| if (count === 0) { | |
| setIsLoadingHistoryPlanning(false); | |
| } | |
| } catch (error) { | |
| // Fall back to marking as loaded to avoid infinite loading state | |
| setIsLoadingHistoryPlanning(false); | |
| } | |
| } | |
| }, | |
| onClose: () => { | |
| setPlanningConnectionState("CLOSED"); | |
| }, | |
| onError: () => { | |
| setPlanningConnectionState("CLOSED"); | |
| // Only show error message if we've previously connected successfully | |
| if (hasConnectedRefPlanning.current) { | |
| setErrorMessage(SERVER_CONNECTION_ERROR_MESSAGE, "connection"); | |
| } | |
| }, | |
| onMessage: handlePlanningMessage, | |
| }; | |
| }, [ | |
| handlePlanningMessage, | |
| setErrorMessage, | |
| clearConnectionError, | |
| sessionApiKey, | |
| subConversations, | |
| ]); | |
| // Only attempt WebSocket connection when we have a valid URL | |
| // This prevents connection attempts during task polling phase | |
| const websocketUrl = wsUrl; | |
| const { socket: mainSocket, reconnect: reconnectMain } = useWebSocket( | |
| websocketUrl || "", | |
| mainWebsocketOptions, | |
| ); | |
| const { socket: planningAgentSocket, reconnect: reconnectPlanning } = | |
| useWebSocket(planningAgentWsUrl || "", planningWebsocketOptions); | |
| const reconnect = useCallback(() => { | |
| removeErrorMessage(); | |
| const currentMode = useConversationStore.getState().conversationMode; | |
| if (currentMode === "plan" && planningAgentWsUrl) { | |
| reconnectPlanning(); | |
| return; | |
| } | |
| reconnectMain(); | |
| }, [ | |
| planningAgentWsUrl, | |
| reconnectMain, | |
| reconnectPlanning, | |
| removeErrorMessage, | |
| ]); | |
| // V1 send message function via WebSocket | |
| // Falls back to REST API queue when WebSocket is not connected | |
| const sendMessage = useCallback( | |
| async (message: SendMessageRequest): Promise<SendMessageResult> => { | |
| const currentMode = useConversationStore.getState().conversationMode; | |
| const currentSocket = | |
| currentMode === "plan" ? planningAgentSocket : mainSocket; | |
| const targetConversationId = | |
| currentMode === "plan" ? planningConversationId : conversationId; | |
| if (currentSocket?.readyState !== WebSocket.OPEN) { | |
| // WebSocket not connected - queue message via REST API | |
| // Message will be delivered automatically when conversation becomes ready | |
| if (!targetConversationId) { | |
| // Never fall back to the parent in plan mode: without a planner | |
| // target the message would run in the code agent, which is exactly | |
| // the boundary plan mode exists to enforce. | |
| const error = new Error( | |
| currentMode === "plan" | |
| ? "Planning conversation is not ready yet" | |
| : "No conversation ID available", | |
| ); | |
| setErrorMessage(error.message); | |
| throw error; | |
| } | |
| try { | |
| await new ConversationClient(getAgentServerClientOptions()).sendEvent( | |
| targetConversationId, | |
| { | |
| role: "user", | |
| content: message.content, | |
| }, | |
| { run: true }, | |
| ); | |
| // Message queued successfully - it will be delivered when ready | |
| // Return queued: true so caller knows not to show optimistic UI | |
| return { queued: true }; | |
| } catch (error) { | |
| const errorMessage = | |
| error instanceof Error | |
| ? error.message | |
| : "Failed to queue message for delivery"; | |
| setErrorMessage(errorMessage); | |
| throw error; | |
| } | |
| } | |
| try { | |
| // Send message through WebSocket as JSON with run: true so the | |
| // agent loop starts automatically in async mode. | |
| currentSocket.send(JSON.stringify({ ...message, run: true })); | |
| return { queued: false }; | |
| } catch (error) { | |
| const errorMessage = | |
| error instanceof Error ? error.message : "Failed to send message"; | |
| setErrorMessage(errorMessage); | |
| throw error; | |
| } | |
| }, | |
| [ | |
| mainSocket, | |
| planningAgentSocket, | |
| setErrorMessage, | |
| conversationId, | |
| planningConversationId, | |
| ], | |
| ); | |
| // Track main socket state changes | |
| useEffect(() => { | |
| // Only process socket updates if we have a valid URL and socket | |
| if (mainSocket && wsUrl) { | |
| // Update state based on socket readyState | |
| const updateState = () => { | |
| switch (mainSocket.readyState) { | |
| case WebSocket.CONNECTING: | |
| setMainConnectionState("CONNECTING"); | |
| break; | |
| case WebSocket.OPEN: | |
| setMainConnectionState("OPEN"); | |
| break; | |
| case WebSocket.CLOSING: | |
| setMainConnectionState("CLOSING"); | |
| break; | |
| case WebSocket.CLOSED: | |
| setMainConnectionState("CLOSED"); | |
| break; | |
| default: | |
| setMainConnectionState("CLOSED"); | |
| break; | |
| } | |
| }; | |
| updateState(); | |
| } | |
| }, [mainSocket, wsUrl]); | |
| // Track planning agent socket state changes | |
| useEffect(() => { | |
| // Only process socket updates if we have a valid URL and socket | |
| if (planningAgentSocket && planningAgentWsUrl) { | |
| // Update state based on socket readyState | |
| const updateState = () => { | |
| switch (planningAgentSocket.readyState) { | |
| case WebSocket.CONNECTING: | |
| setPlanningConnectionState("CONNECTING"); | |
| break; | |
| case WebSocket.OPEN: | |
| setPlanningConnectionState("OPEN"); | |
| break; | |
| case WebSocket.CLOSING: | |
| setPlanningConnectionState("CLOSING"); | |
| break; | |
| case WebSocket.CLOSED: | |
| setPlanningConnectionState("CLOSED"); | |
| break; | |
| default: | |
| setPlanningConnectionState("CLOSED"); | |
| break; | |
| } | |
| }; | |
| updateState(); | |
| } | |
| }, [planningAgentSocket, planningAgentWsUrl]); | |
| const contextValue = useMemo( | |
| () => ({ | |
| connectionState, | |
| mainConnectionState, | |
| sendMessage, | |
| isLoadingHistory, | |
| reconnect, | |
| }), | |
| [ | |
| connectionState, | |
| mainConnectionState, | |
| sendMessage, | |
| isLoadingHistory, | |
| reconnect, | |
| ], | |
| ); | |
| return ( | |
| <ConversationWebSocketContext.Provider value={contextValue}> | |
| {children} | |
| </ConversationWebSocketContext.Provider> | |
| ); | |
| } | |
| export const useConversationWebSocket = | |
| (): ConversationWebSocketContextType | null => { | |
| const context = useContext(ConversationWebSocketContext); | |
| // Return null instead of throwing when not in provider | |
| // This allows the hook to be called conditionally based on conversation version | |
| return context || null; | |
| }; | |