/** * @license * Copyright 2025 Vybestack LLC * SPDX-License-Identifier: Apache-2.0 */ /** * Extracted stream event handler hooks from useAgentStream. * Contains the per-event handler useCallbacks consumed by the AgentEvent * dispatcher (agentEventDispatcher.dispatchAgentEvent), and displayUserMessage. * Multi-turn continuation is owned by the AgenticLoop in * @vybestack/llxprt-code-agents, not by this module. */ import type React from 'react'; import { useCallback, useMemo } from 'react'; import { type ServerErrorEvent as ErrorEvent, type ServerChatCompressedEvent, type MessageSenderType, type ToolCallRequestInfo, parseAndFormatApiError, type ThinkingBlock, type ThoughtSummary, type ServerContentEvent, type AgentRequestInput, } from '@vybestack/llxprt-code-core'; import { logUserPrompt, type UserPromptEvent, } from '@vybestack/llxprt-code-telemetry'; import type { Agent } from '@vybestack/llxprt-code-agents'; import { type LoadedSettings } from '../../../config/settings.js'; import { type HistoryItemWithoutId, type HistoryItemToolGroup, MessageType, ToolCallStatus, type SlashCommandProcessorResult, } from '../../types.js'; import { type RemoveHistoryItems, type UseHistoryManagerReturn, } from '../useHistoryManager.js'; import { showCitations, getCurrentProfileName, buildApiErrorInfo, } from './streamUtils.js'; import { getActiveProviderNameForApiError, getErrorFallbackModel, } from '../../../utils/apiErrorFormatting.js'; import { resolveContentPrefixIdentity, createCliModelIdentityRuntime, } from '../../utils/modelIdentity.js'; /** * Shared content-prefix identity resolver. Reads fresh runtime state at call * time (no dependencies on component state/props), so a single stable reference * can be reused by both ContentEventDeps and StreamEventDeps without risk of * staleness (issue #2263). */ function defaultGetContentPrefixIdentity(): string | null { try { return resolveContentPrefixIdentity(createCliModelIdentityRuntime()); } catch { return null; } } import { processContentEvent, type ContentEventDeps, } from './contentEventProcessor.js'; import type { PendingResponseBuffer } from './pendingResponseBuffer.js'; import { prepareQueryForAgent as prepareQueryImpl, type PrepareQueryDeps, } from './queryPreparer.js'; import { getTokenLimitForConfiguredContext } from './contextLimit.js'; import type { StreamRuntime, UiSubagentManager } from '../../cliUiRuntime.js'; interface StreamEventHandlersResult { handleContentEvent: ( eventValue: ServerContentEvent['value'], currentAgentMessageBuffer: string, userMessageTimestamp: number, ) => string; handleUserCancelledEvent: (userMessageTimestamp: number) => void; handleErrorEvent: ( eventValue: ErrorEvent['value'], userMessageTimestamp: number, options?: { clearQueue?: boolean }, ) => void; handleCitationEvent: (text: string, userMessageTimestamp: number) => void; /** * Finished-notice handler for the public AgentEvent done{stop|refusal} path. * Receives the pre-computed message (or null for no notice). */ handleFinishedNotice: ( message: string | null, userMessageTimestamp: number, ) => void; handleChatCompressionEvent: ( eventValue: ServerChatCompressedEvent['value'], userMessageTimestamp: number, ) => void; handleMaxSessionTurnsEvent: () => void; handleContextWindowWillOverflowEvent: ( estimatedRequestTokenCount: number, remainingTokenCount: number, ) => void; /** * Discards the abandoned attempt's render state (issue #3048): nulls the * pending AI item without flushing, resets buffer/thinking state, and * retracts committed stable segments. The turn stays responding. */ handleStreamAttemptDiscarded: () => void; handleLoopDetectedEvent: () => void; displayUserMessage: ( trimmedQuery: string, userMessageTimestamp: number, ) => void; prepareQueryForAgent: ( query: AgentRequestInput, userMessageTimestamp: number, abortSignal: AbortSignal, promptId: string, ) => Promise<{ queryToSend: AgentRequestInput | null; shouldProceed: boolean; }>; } interface StreamEventHandlerDeps { runtime: StreamRuntime; // @plan:ISSUE-2376 — the Agent surface supplies named-tool lookup // (agent.tools.get) for @file processing; threaded alongside the #2384 // StreamRuntime rather than through getToolRegistry. agent: Agent; settings: LoadedSettings; addItem: UseHistoryManagerReturn['addItem']; removeItems?: RemoveHistoryItems; onDebugMessage: (message: string) => void; onCancelSubmit: (shouldRestorePrompt?: boolean) => void; sanitizeContent: (text: string) => { text: string; blocked: boolean; feedback?: string; }; flushPendingHistoryItem: (timestamp: number) => void; pendingResponse: PendingResponseBuffer; pendingHistoryItemRef: React.MutableRefObject; thinkingBlocksRef: React.MutableRefObject; turnCancelledRef: React.MutableRefObject; clearSubmissions: () => void; setPendingHistoryItem: React.Dispatch< React.SetStateAction >; setIsResponding: React.Dispatch>; setThought: React.Dispatch>; setLastAgentActivityTime: React.Dispatch>; scheduleToolCalls: ( requests: ToolCallRequestInfo[], signal: AbortSignal, ) => Promise; abortActiveStream: (reason?: unknown) => void; handleShellCommand: (query: string, signal: AbortSignal) => boolean; handleSlashCommand: ( cmd: AgentRequestInput, ) => Promise; logger: | { logMessage: (sender: MessageSenderType, text: string) => Promise } | null | undefined; shellModeActive: boolean; loopDetectedRef: React.MutableRefObject; lastProfileNameRef: React.MutableRefObject; lastModelInfoRef: React.MutableRefObject; lastModelIdentityRef: React.MutableRefObject; // Optional: sourced from the app runtime when available. StreamRuntime does // not yet expose getSubagentManager, so this is undefined in the streaming // path until a runtime slice is added; @subagent mentions degrade gracefully. subagentManager?: UiSubagentManager; } export function useStreamEventHandlers( deps: StreamEventHandlerDeps, ): StreamEventHandlersResult { const handleContentEvent = useContentEventHandler(deps); const handleStreamAttemptDiscarded = useStreamAttemptDiscardedHandler(deps); const handleLoopDetectedEvent = useCallback( () => deps.addItem( { type: 'info', text: 'A potential loop was detected. This can happen due to repetitive tool calls or other model behavior. The request has been halted.', }, Date.now(), ), [deps], ); const handlers = useStreamHandlers( deps, handleContentEvent, handleStreamAttemptDiscarded, handleLoopDetectedEvent, ); const displayUserMessage = useDisplayUserMessage(deps); const prepareQueryForAgent = usePrepareQueryForAgent(deps); return { ...handlers, displayUserMessage, prepareQueryForAgent, }; } function useContentEventHandler(deps: StreamEventHandlerDeps) { const contentEventDeps = useContentEventDeps(deps, deps.pendingResponse); return useCallback( ( eventValue: ServerContentEvent['value'], currentAgentMessageBuffer: string, userMessageTimestamp: number, ): string => processContentEvent( eventValue, currentAgentMessageBuffer, userMessageTimestamp, contentEventDeps, ), [contentEventDeps], ); } /** * REQ-3048-008/009: resets the uncommitted render state of the abandoned * attempt and retracts stable segments it already committed. */ function useStreamAttemptDiscardedHandler(deps: StreamEventHandlerDeps) { const { pendingHistoryItemRef, setPendingHistoryItem, pendingResponse, thinkingBlocksRef, setThought, removeItems, } = deps; return useCallback(() => { const pending = pendingHistoryItemRef.current; if ( pending && (pending.type === 'gemini' || pending.type === 'gemini_content') ) { setPendingHistoryItem(null); } pendingResponse.reset(); // Resolve the retraction capability before draining the ledger so a // missing removeItems wiring fails fast without losing committed segment // IDs. reset() does not touch the ledger, so an early throw preserves the // ids for a later, correctly-wired retraction. const retract = requireHistoryRetraction(removeItems); const retracted = pendingResponse.drainCommittedSegments(); if (retracted.length > 0) { retract(retracted); } thinkingBlocksRef.current = []; setThought(null); }, [ pendingHistoryItemRef, setPendingHistoryItem, pendingResponse, removeItems, thinkingBlocksRef, setThought, ]); } function requireHistoryRetraction( removeItems: RemoveHistoryItems | undefined, ): RemoveHistoryItems { if (!removeItems) { throw new Error('History retraction is required after streamed content'); } return removeItems; } function useStreamHandlers( deps: StreamEventHandlerDeps, handleContentEvent: ( eventValue: ServerContentEvent['value'], currentAgentMessageBuffer: string, userMessageTimestamp: number, ) => string, handleStreamAttemptDiscarded: () => void, handleLoopDetectedEvent: () => void, ): HandlerMap { return { handleContentEvent, handleStreamAttemptDiscarded, handleUserCancelledEvent: useUserCancelledHandler(deps), handleErrorEvent: useErrorEventHandler(deps), handleCitationEvent: useCitationEventHandler(deps), handleFinishedNotice: useFinishedNoticeHandler(deps), handleChatCompressionEvent: useChatCompressionHandler(deps), handleMaxSessionTurnsEvent: useMaxSessionTurnsHandler(deps), handleContextWindowWillOverflowEvent: useContextOverflowHandler(deps), handleLoopDetectedEvent, }; } function useUserCancelledHandler(deps: StreamEventHandlerDeps) { const { addItem, flushPendingHistoryItem, pendingHistoryItemRef, setIsResponding, setPendingHistoryItem, setThought, turnCancelledRef, pendingResponse, } = deps; return useCallback( (userMessageTimestamp: number) => { if (turnCancelledRef.current) return; if (pendingHistoryItemRef.current) { if (pendingHistoryItemRef.current.type === 'tool_group') { const pendingItem: HistoryItemToolGroup = { ...pendingHistoryItemRef.current, tools: pendingHistoryItemRef.current.tools.map((tool) => tool.status === ToolCallStatus.Pending || tool.status === ToolCallStatus.Confirming || tool.status === ToolCallStatus.Executing ? { ...tool, status: ToolCallStatus.Canceled } : tool, ), }; addItem(pendingItem, userMessageTimestamp); } else { flushPendingHistoryItem(userMessageTimestamp); } pendingResponse.endCommittedSegments(); setPendingHistoryItem(null); } addItem( { type: MessageType.INFO, text: 'User cancelled the request.' }, userMessageTimestamp, ); setIsResponding(false); setThought(null); }, [ addItem, flushPendingHistoryItem, pendingHistoryItemRef, pendingResponse, setIsResponding, setPendingHistoryItem, setThought, turnCancelledRef, ], ); } function useErrorEventHandler(deps: StreamEventHandlerDeps) { const { addItem, runtime, flushPendingHistoryItem, pendingHistoryItemRef, clearSubmissions, setPendingHistoryItem, setThought, pendingResponse, } = deps; return useCallback( ( eventValue: ErrorEvent['value'], userMessageTimestamp: number, options?: { clearQueue?: boolean }, ) => { if (options?.clearQueue ?? true) clearSubmissions(); setThought(null); if (pendingHistoryItemRef.current) { flushPendingHistoryItem(userMessageTimestamp); pendingResponse.endCommittedSegments(); setPendingHistoryItem(null); } const apiErrorInfo = buildApiErrorInfo(runtime); const providerName = getActiveProviderNameForApiError(apiErrorInfo); const fallbackModel = getErrorFallbackModel(apiErrorInfo, providerName); addItem( { type: MessageType.ERROR, text: parseAndFormatApiError( eventValue.error, undefined, fallbackModel, providerName, ), }, userMessageTimestamp, ); }, [ addItem, runtime, flushPendingHistoryItem, pendingResponse, pendingHistoryItemRef, clearSubmissions, setPendingHistoryItem, setThought, ], ); } function useCitationEventHandler(deps: StreamEventHandlerDeps) { const { addItem, runtime, flushPendingHistoryItem, pendingHistoryItemRef, setPendingHistoryItem, pendingResponse, settings, } = deps; return useCallback( (text: string, userMessageTimestamp: number) => { if (!showCitations(settings, runtime)) return; if (pendingHistoryItemRef.current) { flushPendingHistoryItem(userMessageTimestamp); pendingResponse.endCommittedSegments(); setPendingHistoryItem(null); } addItem({ type: MessageType.INFO, text }, userMessageTimestamp); }, [ addItem, runtime, flushPendingHistoryItem, pendingHistoryItemRef, pendingResponse, setPendingHistoryItem, settings, ], ); } /** * Finished-notice handler for the public AgentEvent path. The agentEventDispatcher * computes the message (from a FinishedValue) and passes it here; this handler * only renders the WARNING item. Null message → no item (parity). */ function useFinishedNoticeHandler(deps: StreamEventHandlerDeps) { const { addItem } = deps; return useCallback( (message: string | null, userMessageTimestamp: number) => { if (!message) return; addItem( { type: 'info', text: `WARNING: ${message}` }, userMessageTimestamp, ); }, [addItem], ); } function useChatCompressionHandler(deps: StreamEventHandlerDeps) { const { addItem, runtime, pendingHistoryItemRef, pendingResponse, setPendingHistoryItem, } = deps; return useCallback( ( eventValue: ServerChatCompressedEvent['value'], userMessageTimestamp: number, ) => { if (pendingHistoryItemRef.current) { addItem(pendingHistoryItemRef.current, userMessageTimestamp); pendingResponse.endCommittedSegments(); setPendingHistoryItem(null); } return addItem( { type: 'info', text: `IMPORTANT: This conversation approached the input token limit for ${runtime.model.getModel()}. ` + `A compressed context will be sent for future messages (compressed from: ` + `${eventValue?.originalTokenCount ?? 'unknown'} to ` + `${eventValue?.newTokenCount ?? 'unknown'} tokens).`, }, Date.now(), ); }, [ addItem, runtime, pendingHistoryItemRef, pendingResponse, setPendingHistoryItem, ], ); } function useMaxSessionTurnsHandler(deps: StreamEventHandlerDeps) { const { addItem, runtime } = deps; return useCallback( () => addItem( { type: 'info', text: `The session has reached the maximum number of turns: ${runtime.sessionLimits.getMaxSessionTurns()}. Please update this limit in your setting.json file.`, }, Date.now(), ), [addItem, runtime], ); } function useContextOverflowHandler(deps: StreamEventHandlerDeps) { const { addItem, runtime, onCancelSubmit } = deps; return useCallback( (estimatedRequestTokenCount: number, remainingTokenCount: number) => { onCancelSubmit(true); const limit = getTokenLimitForConfiguredContext(runtime); const isLessThan75Percent = limit > 0 && remainingTokenCount < limit * 0.75; let text = `Sending this message (${estimatedRequestTokenCount} tokens) might exceed the remaining context window limit (${remainingTokenCount} tokens).`; if (isLessThan75Percent) text += ' Please try reducing the size of your message or use the `/compress` command to compress the chat history.'; addItem({ type: 'info', text }, Date.now()); }, [addItem, runtime, onCancelSubmit], ); } function usePrepareQueryForAgent(deps: StreamEventHandlerDeps) { const prepareQueryDeps = usePrepareQueryDeps(deps); return useCallback( async ( query: AgentRequestInput, userMessageTimestamp: number, abortSignal: AbortSignal, prompt_id: string, ) => prepareQueryImpl( query, userMessageTimestamp, abortSignal, prompt_id, prepareQueryDeps, ), [prepareQueryDeps], ); } function useContentEventDeps( deps: StreamEventHandlerDeps, pendingResponse: PendingResponseBuffer, ): ContentEventDeps { return useMemo( () => ({ addItem: deps.addItem, pendingResponse, sanitizeContent: deps.sanitizeContent, flushPendingHistoryItem: deps.flushPendingHistoryItem, pendingHistoryItemRef: deps.pendingHistoryItemRef, thinkingBlocksRef: deps.thinkingBlocksRef, turnCancelledRef: deps.turnCancelledRef, setPendingHistoryItem: deps.setPendingHistoryItem, getContentPrefixIdentity: defaultGetContentPrefixIdentity, }), [ deps.addItem, pendingResponse, deps.sanitizeContent, deps.flushPendingHistoryItem, deps.pendingHistoryItemRef, deps.thinkingBlocksRef, deps.turnCancelledRef, deps.setPendingHistoryItem, ], ); } function usePrepareQueryDeps(deps: StreamEventHandlerDeps): PrepareQueryDeps { return useMemo( () => ({ runtime: deps.runtime, getToolHandle: (name: string) => deps.agent.tools.get(name), logUserPrompt: (event: UserPromptEvent) => logUserPrompt( { getSessionId: () => deps.runtime.session.getSessionId(), getTelemetryLogPromptsEnabled: () => deps.runtime.settings.getTelemetryLogPromptsEnabled(), }, event, ), addItem: deps.addItem, onDebugMessage: deps.onDebugMessage, handleShellCommand: deps.handleShellCommand, handleSlashCommand: deps.handleSlashCommand, logger: deps.logger, shellModeActive: deps.shellModeActive, scheduleToolCalls: deps.scheduleToolCalls, subagentManager: deps.subagentManager, }), [ deps.runtime, deps.agent, deps.addItem, deps.onDebugMessage, deps.handleShellCommand, deps.handleSlashCommand, deps.logger, deps.shellModeActive, deps.scheduleToolCalls, deps.subagentManager, ], ); } function useDisplayUserMessage(deps: StreamEventHandlerDeps) { const { addItem, runtime, lastProfileNameRef } = deps; return useCallback( (trimmedQuery: string, userMessageTimestamp: number) => { addItem( { type: MessageType.USER, text: trimmedQuery }, userMessageTimestamp, ); // Inline profile_change notifications are now owned exclusively by the // ModelInfo event path (agentEventDispatcher.handleModelInfoEvent). // We still track lastProfileNameRef for backward-compatible diagnostics. const liveProfileName = getCurrentProfileName(runtime); lastProfileNameRef.current = liveProfileName ?? undefined; }, [addItem, runtime, lastProfileNameRef], ); } type HandlerMap = Pick< StreamEventHandlersResult, | 'handleContentEvent' | 'handleStreamAttemptDiscarded' | 'handleUserCancelledEvent' | 'handleErrorEvent' | 'handleChatCompressionEvent' | 'handleFinishedNotice' | 'handleMaxSessionTurnsEvent' | 'handleContextWindowWillOverflowEvent' | 'handleCitationEvent' | 'handleLoopDetectedEvent' >;