/** * @license * Copyright 2025 Vybestack LLC * SPDX-License-Identifier: Apache-2.0 */ /** * useStreamState — extracted state initialization and sanitization from * useAgentStream to keep the orchestrator under 80 lines. * * Owns: initError, abortController, turnCancelled, isResponding, * lastProfileName, thought, pendingHistoryItem, lastAgentActivityTime, * queuedSubmissions, submitQueryRef, emojiFilter, sanitizeContent, * flushPendingHistoryItem, logger, gitService. */ import { useState, useRef, useCallback, useMemo } from 'react'; import { type ThoughtSummary, EmojiFilter, type EmojiFilterMode, type ThinkingBlock, GitService, } from '@vybestack/llxprt-code-core'; import { type HistoryItemWithoutId, MessageType } from '../../types.js'; import { useStateAndRef } from '../useStateAndRef.js'; import { useLogger } from '../useLogger.js'; import { type QueuedSubmission } from './types.js'; import { useQueuedSubmissions } from './useQueuedSubmissions.js'; import type { StreamRuntime } from '../../cliUiRuntime.js'; import type { SubmissionExecutor } from './useSubmitQuery.js'; import { PendingResponseBuffer } from './pendingResponseBuffer.js'; export interface UseStreamStateReturn { initError: string | null; setInitError: React.Dispatch>; abortControllerRef: React.MutableRefObject; abortActiveStream: (reason?: unknown) => void; turnCancelledRef: React.MutableRefObject; turnCancelled: boolean; setTurnCancelled: (value: boolean) => void; drainSuppressedRef: React.MutableRefObject; isResponding: boolean; setIsResponding: React.Dispatch>; lastProfileNameRef: React.MutableRefObject; lastModelInfoRef: React.MutableRefObject; lastModelIdentityRef: React.MutableRefObject; thought: ThoughtSummary | null; setThought: React.Dispatch>; pendingHistoryItem: HistoryItemWithoutId | null; pendingHistoryItemRef: React.MutableRefObject; setPendingHistoryItem: React.Dispatch< React.SetStateAction >; lastAgentActivityTime: number; setLastAgentActivityTime: React.Dispatch>; queuedSubmissionsRef: React.MutableRefObject; queuedSubmissions: readonly QueuedSubmission[]; enqueueSubmission: (submission: QueuedSubmission) => void; enqueueSubmissionFirst: (submission: QueuedSubmission) => void; requeueSubmission: (submission: QueuedSubmission) => void; dequeueSubmission: () => QueuedSubmission | undefined; clearSubmissions: () => void; tryReserveDrain: () => boolean; releaseDrain: () => void; submitQueryRef: React.MutableRefObject; pendingResponse: PendingResponseBuffer; sanitizeContent: (text: string) => { text: string; blocked: boolean; feedback?: string; }; flushPendingHistoryItem: (timestamp: number) => void; logger: ReturnType; gitService: GitService | undefined; thinkingBlocksRef: React.MutableRefObject; } function useEmojiFilterMode(runtime: StreamRuntime): EmojiFilterMode { const rawMode = runtime.ephemeral.getEphemeralSetting('emojifilter'); return typeof rawMode === 'string' && rawMode.length > 0 ? (rawMode as EmojiFilterMode) : 'auto'; } /** * Memoised on the mode string rather than the runtime object. * * The filter now backs the stateful {@link PendingResponseBuffer}, so * recreating it mid-stream would discard the accumulated response. Keying on a * primitive means a new `runtime` identity — which callers can produce on any * re-render — cannot silently truncate a reply (issue #2852). */ function useEmojiFilter(mode: EmojiFilterMode) { return useMemo( () => (mode !== 'allowed' ? new EmojiFilter({ mode }) : undefined), [mode], ); } function useSanitizeContent(emojiFilter: EmojiFilter | undefined) { return useCallback( (text: string) => { if (!emojiFilter) { return { text, feedback: undefined as string | undefined, blocked: false, }; } const result = emojiFilter.filterText(text); if (result.blocked) { return { text: '', feedback: result.systemFeedback, blocked: true as const, }; } const sanitized = typeof result.filtered === 'string' ? result.filtered : ''; return { text: sanitized, feedback: result.systemFeedback, blocked: false as const, }; }, [emojiFilter], ); } function useFlushPendingHistoryItem( addItem: (item: HistoryItemWithoutId, timestamp: number) => number, pendingHistoryItemRef: React.MutableRefObject, pendingResponse: PendingResponseBuffer, setPendingHistoryItem: React.Dispatch< React.SetStateAction >, thinkingBlocksRef: React.MutableRefObject, ) { return useCallback( (timestamp: number) => { const pending = pendingHistoryItemRef.current; if (!pending) { return; } if (pending.type === 'gemini' || pending.type === 'gemini_content') { commitAiPendingItem( pending, timestamp, addItem, pendingResponse, thinkingBlocksRef, ); } else { addItem(pending, timestamp); } setPendingHistoryItem(null); }, [ addItem, pendingHistoryItemRef, pendingResponse, setPendingHistoryItem, thinkingBlocksRef, ], ); } /** * Commits the in-progress assistant response. The text comes from * {@link PendingResponseBuffer}, which sanitised it incrementally as it * streamed, so no whole-text pass is needed here (issue #2852). */ function commitAiPendingItem( pending: HistoryItemWithoutId, timestamp: number, addItem: (item: HistoryItemWithoutId, timestamp: number) => number, pendingResponse: PendingResponseBuffer, thinkingBlocksRef: React.MutableRefObject, ): void { const { text, feedback, blocked } = pendingResponse.materialize(); pendingResponse.reset(); const thinkingBlocks = thinkingBlocksRef.current; thinkingBlocksRef.current = []; if (blocked) { addItem( { type: MessageType.ERROR, text: '[Error: Response blocked due to emoji detection]', }, timestamp, ); if (feedback) { addItem({ type: MessageType.INFO, text: feedback }, timestamp); } return; } addItem( { ...pending, text, ...(thinkingBlocks.length > 0 ? { thinkingBlocks: [...thinkingBlocks] } : {}), }, timestamp, ); if (feedback) { addItem({ type: MessageType.INFO, text: feedback }, timestamp); } } function useBasicStreamState() { const [initError, setInitError] = useState(null); const abortControllerRef = useRef(null); const abortActiveStream = useCallback((reason?: unknown) => { abortControllerRef.current?.abort(reason); }, []); const turnCancelledRef = useRef(false); const [turnCancelled, setTurnCancelledState] = useState(false); const setTurnCancelled = useCallback((value: boolean) => { turnCancelledRef.current = value; setTurnCancelledState(value); }, []); const drainSuppressedRef = useRef(false); const [isResponding, setIsResponding] = useState(false); const lastProfileNameRef = useRef(undefined); const lastModelInfoRef = useRef(null); const lastModelIdentityRef = useRef(null); const [thought, setThought] = useState(null); const [pendingHistoryItem, pendingHistoryItemRef, setPendingHistoryItem] = useStateAndRef(null); const [lastAgentActivityTime, setLastAgentActivityTime] = useState(0); const { queuedSubmissions, queuedSubmissionsRef, enqueueSubmission, enqueueSubmissionFirst, requeueSubmission, dequeueSubmission, clearSubmissions, tryReserveDrain, releaseDrain, } = useQueuedSubmissions(); const submitQueryRef = useRef(null); const thinkingBlocksRef = useRef([]); return { initError, setInitError, abortControllerRef, abortActiveStream, turnCancelledRef, turnCancelled, setTurnCancelled, drainSuppressedRef, isResponding, setIsResponding, lastProfileNameRef, lastModelInfoRef, lastModelIdentityRef, thought, setThought, pendingHistoryItem, pendingHistoryItemRef, setPendingHistoryItem, lastAgentActivityTime, setLastAgentActivityTime, queuedSubmissions, queuedSubmissionsRef, enqueueSubmission, enqueueSubmissionFirst, requeueSubmission, dequeueSubmission, clearSubmissions, tryReserveDrain, releaseDrain, submitQueryRef, thinkingBlocksRef, }; } export function useStreamState( addItem: (item: HistoryItemWithoutId, timestamp: number) => number, runtime: StreamRuntime, ): UseStreamStateReturn { const basic = useBasicStreamState(); const storage = runtime.storage; const emojiFilter = useEmojiFilter(useEmojiFilterMode(runtime)); const sanitizeContent = useSanitizeContent(emojiFilter); const pendingResponse = useMemo( () => new PendingResponseBuffer(emojiFilter), [emojiFilter], ); const flushPendingHistoryItem = useFlushPendingHistoryItem( addItem, basic.pendingHistoryItemRef, pendingResponse, basic.setPendingHistoryItem, basic.thinkingBlocksRef, ); const logger = useLogger(storage); const gitService = useMemo(() => { const projectRoot = runtime.session.getProjectRoot(); if (projectRoot.length === 0) { return undefined; } return new GitService(projectRoot, storage); }, [runtime, storage]); return { ...basic, pendingResponse, sanitizeContent, flushPendingHistoryItem, logger, gitService, }; }