diff --git a/packages/ui/src/hooks/useEventStream.ts b/packages/ui/src/hooks/useEventStream.ts index 9c779be7..fa70121b 100644 --- a/packages/ui/src/hooks/useEventStream.ts +++ b/packages/ui/src/hooks/useEventStream.ts @@ -639,7 +639,7 @@ export const useEventStream = () => { case 'session.status': { - const sessionId = typeof props.sessionID === 'string' ? props.sessionID : null; + const sessionId = readStringProp(props, ['sessionID', 'sessionId']); const statusObj = (typeof props.status === 'object' && props.status !== null) ? props.status as Record : null; const statusType = typeof statusObj?.type === 'string' ? statusObj.type : null; const statusInfo = statusObj ?? {}; @@ -664,7 +664,7 @@ export const useEventStream = () => { case 'openchamber:session-status': { - const sessionId = typeof props.sessionId === 'string' ? props.sessionId : null; + const sessionId = readStringProp(props, ['sessionId', 'sessionID']); const status = typeof props.status === 'string' ? props.status : null; const needsAttention = typeof props.needsAttention === 'boolean' ? props.needsAttention : false; const timestamp = typeof props.timestamp === 'number' ? props.timestamp : Date.now(); @@ -784,17 +784,38 @@ export const useEventStream = () => { // Fallback: if we see assistant parts but session.status hasn't arrived yet, mark busy. if (roleInfo === 'assistant') { const partType = (messagePart as { type?: unknown }).type; - const isStreamingPart = - partType === 'step-start' || - partType === 'text' || - partType === 'tool' || - partType === 'reasoning' || - partType === 'file' || - partType === 'patch'; + const partTime = (messagePart as { time?: { end?: unknown } }).time; + const partHasEnded = typeof partTime?.end === 'number'; + const toolState = (messagePart as { state?: { status?: unknown } }).state?.status; + const textContent = (messagePart as { text?: unknown }).text; + + const isStreamingPart = (() => { + if (partType === 'tool') { + return toolState === 'running' || toolState === 'pending'; + } + if (partType === 'reasoning') { + return !partHasEnded; + } + if (partType === 'text') { + const hasText = typeof textContent === 'string' && textContent.trim().length > 0; + return hasText && !partHasEnded; + } + if (partType === 'step-start') { + return true; + } + return false; + })(); if (isStreamingPart) { const currentStatus = useSessionStore.getState().sessionStatus?.get(sessionId); + const recentlyConfirmedIdle = + currentStatus?.type === 'idle' && + typeof currentStatus.confirmedAt === 'number' && + Date.now() - currentStatus.confirmedAt < 1200; if (!currentStatus || currentStatus.type === 'idle') { + if (recentlyConfirmedIdle) { + break; + } updateSessionStatus(sessionId, { type: 'busy' }, 'sse:message.part.updated'); } } @@ -1057,13 +1078,15 @@ export const useEventStream = () => { const hasParts = partsArray.length > 0; const timeObj = (messageExt as { time?: { completed?: number } }).time || {}; const completedFromServer = typeof timeObj?.completed === 'number'; - - if (!hasParts && !completedFromServer) break; - + const rawStatus = (message as { status?: unknown }).status; + const status = typeof rawStatus === 'string' ? rawStatus.toLowerCase() : null; + const hasCompletedStatus = status === 'completed' || status === 'complete'; const finishCandidate = (message as { finish?: unknown }).finish; const finish = typeof finishCandidate === 'string' ? finishCandidate : null; const eventHasStopFinish = finish === 'stop'; + if (!hasParts && !completedFromServer && !hasCompletedStatus && !eventHasStopFinish) break; + if ((messageExt as { role?: unknown }).role === 'assistant' && hasParts) { const incomingLen = computeTextLength(partsArray); const wouldShrink = existingLen > 0 && incomingLen + TEXT_SHRINK_TOLERANCE < existingLen; @@ -1111,7 +1134,7 @@ export const useEventStream = () => { const shouldFinalizeAssistantMessage = (message as { role?: string }).role === 'assistant' && - (hasCompletedTimestamp || stopMarkerPresent); + (hasCompletedTimestamp || hasCompletedStatus || stopMarkerPresent); if (shouldFinalizeAssistantMessage && (message as { role?: string }).role === 'assistant') { @@ -1741,6 +1764,7 @@ export const useEventStream = () => { clearPauseTimeout(); maybeBootstrapIfStale('visibility_restore'); + triggerSessionStatusPoll(); const isStalled = Date.now() - lastEventTimestampRef.current > 45000; if (isStalled) { @@ -1769,6 +1793,7 @@ export const useEventStream = () => { if (visibilityStateRef.current === 'visible') { clearPauseTimeout(); maybeBootstrapIfStale('window_focus'); + triggerSessionStatusPoll(); const isStalled = Date.now() - lastEventTimestampRef.current > 45000; if (isStalled) { @@ -1859,7 +1884,7 @@ export const useEventStream = () => { ); if (hasBusySessions) { - // Removed: void refreshSessionStatus(); + triggerSessionStatusPoll(); } if (now - lastEventTimestampRef.current > 45000) { Promise.resolve().then(async () => { diff --git a/packages/ui/src/hooks/useServerSessionStatus.ts b/packages/ui/src/hooks/useServerSessionStatus.ts index d0aa995b..e6038d60 100644 --- a/packages/ui/src/hooks/useServerSessionStatus.ts +++ b/packages/ui/src/hooks/useServerSessionStatus.ts @@ -25,7 +25,8 @@ interface ServerSnapshotResponse { serverTime: number; } -const IMMEDIATE_POLL_DELAY_MS = 500; // 500ms for immediate poll after notification +const IMMEDIATE_POLL_DELAY_MS = 150; +const FOLLOW_UP_POLL_DELAY_MS = 1100; // Ref to be accessed from outside (e.g., useEventStream) for triggering immediate poll let triggerImmediatePollRef: (() => void) | null = null; @@ -45,8 +46,10 @@ export const triggerSessionStatusPoll = () => { */ export function useServerSessionStatus() { const isSyncingRef = React.useRef(false); + const hasPendingImmediateSyncRef = React.useRef(false); const lastSyncAtRef = React.useRef(0); const timeoutRef = React.useRef(null); + const followUpTimeoutRef = React.useRef(null); const fetchSessionStatus = React.useCallback(async (immediate = false) => { const now = Date.now(); @@ -54,8 +57,12 @@ export function useServerSessionStatus() { return; } - // Prevent concurrent syncs + // Prevent concurrent syncs; if an immediate sync is requested while running, + // queue one more pass right after current request settles. if (isSyncingRef.current) { + if (immediate) { + hasPendingImmediateSyncRef.current = true; + } return; } @@ -65,6 +72,7 @@ export function useServerSessionStatus() { try { const snapshotResponse = await fetch('/api/sessions/snapshot', { method: 'GET', + cache: 'no-store', headers: { Accept: 'application/json' }, }); @@ -179,9 +187,37 @@ export function useServerSessionStatus() { console.warn('[useServerSessionStatus] Error fetching session status:', error); } finally { isSyncingRef.current = false; + if (hasPendingImmediateSyncRef.current) { + hasPendingImmediateSyncRef.current = false; + setTimeout(() => { + void fetchSessionStatus(true); + }, 120); + } } }, []); + // Function to trigger immediate snapshot sync from external modules + const triggerImmediatePoll = React.useCallback(() => { + // Clear any pending timeout + if (timeoutRef.current) { + clearTimeout(timeoutRef.current); + } + if (followUpTimeoutRef.current) { + clearTimeout(followUpTimeoutRef.current); + } + + // Schedule immediate sync with small delay to batch rapid calls + timeoutRef.current = setTimeout(() => { + void fetchSessionStatus(true); + }, IMMEDIATE_POLL_DELAY_MS); + + // Run one follow-up sync after short settle period to catch delayed + // server status transitions that happen right after reconnect/restore. + followUpTimeoutRef.current = setTimeout(() => { + void fetchSessionStatus(true); + }, FOLLOW_UP_POLL_DELAY_MS); + }, [fetchSessionStatus]); + // Initial snapshot sync on mount React.useEffect(() => { void fetchSessionStatus(true); @@ -190,6 +226,9 @@ export function useServerSessionStatus() { if (timeoutRef.current) { clearTimeout(timeoutRef.current); } + if (followUpTimeoutRef.current) { + clearTimeout(followUpTimeoutRef.current); + } }; }, [fetchSessionStatus]); @@ -197,10 +236,7 @@ export function useServerSessionStatus() { React.useEffect(() => { const handleVisibilityChange = () => { if (document.visibilityState === 'visible') { - // Small delay to let the browser settle - timeoutRef.current = setTimeout(() => { - void fetchSessionStatus(true); - }, 100); + triggerImmediatePoll(); } }; @@ -208,20 +244,7 @@ export function useServerSessionStatus() { return () => { document.removeEventListener('visibilitychange', handleVisibilityChange); }; - }, [fetchSessionStatus]); - - // Function to trigger immediate snapshot sync from external modules - const triggerImmediatePoll = React.useCallback(() => { - // Clear any pending timeout - if (timeoutRef.current) { - clearTimeout(timeoutRef.current); - } - - // Schedule immediate sync with small delay to batch rapid calls - timeoutRef.current = setTimeout(() => { - void fetchSessionStatus(true); - }, IMMEDIATE_POLL_DELAY_MS); - }, [fetchSessionStatus]); + }, [triggerImmediatePoll]); // Update the ref for external access React.useEffect(() => { diff --git a/packages/ui/src/lib/opencode/client.ts b/packages/ui/src/lib/opencode/client.ts index 44cce646..34428f8c 100644 --- a/packages/ui/src/lib/opencode/client.ts +++ b/packages/ui/src/lib/opencode/client.ts @@ -147,6 +147,11 @@ class OpencodeService { private globalSseListeners: Set<(event: RoutedOpencodeEvent) => void> = new Set(); private globalSseOpenListeners: Set<() => void> = new Set(); private globalSseErrorListeners: Set<(error: unknown) => void> = new Set(); + private globalSseQueue: Array = []; + private globalSseBuffer: Array = []; + private globalSseCoalesced: Map = new Map(); + private globalSseFlushTimer: ReturnType | null = null; + private globalSseLastFlushAt = 0; constructor(baseUrl: string = DEFAULT_BASE_URL) { const desktopBase = resolveDesktopBaseUrl(); @@ -1166,13 +1171,7 @@ class OpencodeService { } private emitGlobalSseEvent(event: RoutedOpencodeEvent) { - for (const listener of this.globalSseListeners) { - try { - listener(event); - } catch (error) { - console.warn('[OpencodeClient] Global SSE listener error:', error); - } - } + this.enqueueGlobalSseEvent(event); } private notifyGlobalSseOpen() { @@ -1228,6 +1227,110 @@ class OpencodeService { this.globalSseAbortController.abort(); } this.globalSseAbortController = null; + this.clearGlobalSseQueue(); + } + + private clearGlobalSseQueue() { + if (this.globalSseFlushTimer) { + clearTimeout(this.globalSseFlushTimer); + this.globalSseFlushTimer = null; + } + this.globalSseQueue.length = 0; + this.globalSseBuffer.length = 0; + this.globalSseCoalesced.clear(); + } + + private getGlobalSseCoalesceKey(event: RoutedOpencodeEvent): string | null { + const payload = event.payload as unknown as Record; + const eventType = typeof payload.type === 'string' ? payload.type : null; + if (!eventType) { + return null; + } + + const properties = + typeof payload.properties === 'object' && payload.properties !== null + ? (payload.properties as Record) + : null; + + if (eventType === 'session.status') { + const sessionId = typeof properties?.sessionID === 'string' + ? properties.sessionID + : typeof properties?.sessionId === 'string' + ? properties.sessionId + : null; + if (!sessionId) { + return null; + } + return `session.status:${event.directory}:${sessionId}`; + } + + if (eventType === 'openchamber:session-status') { + const sessionId = typeof properties?.sessionId === 'string' + ? properties.sessionId + : typeof properties?.sessionID === 'string' + ? properties.sessionID + : null; + if (!sessionId) { + return null; + } + return `openchamber:session-status:${sessionId}`; + } + + return null; + } + + private flushGlobalSseQueue = () => { + if (this.globalSseFlushTimer) { + clearTimeout(this.globalSseFlushTimer); + this.globalSseFlushTimer = null; + } + + if (this.globalSseQueue.length === 0) { + return; + } + + const events = this.globalSseQueue; + this.globalSseQueue = this.globalSseBuffer; + this.globalSseBuffer = events; + this.globalSseQueue.length = 0; + this.globalSseCoalesced.clear(); + this.globalSseLastFlushAt = Date.now(); + + for (const event of events) { + if (!event) continue; + for (const listener of this.globalSseListeners) { + try { + listener(event); + } catch (error) { + console.warn('[OpencodeClient] Global SSE listener error:', error); + } + } + } + + this.globalSseBuffer.length = 0; + }; + + private scheduleGlobalSseFlush() { + if (this.globalSseFlushTimer) { + return; + } + const elapsed = Date.now() - this.globalSseLastFlushAt; + const delay = Math.max(0, 16 - elapsed); + this.globalSseFlushTimer = setTimeout(this.flushGlobalSseQueue, delay); + } + + private enqueueGlobalSseEvent(event: RoutedOpencodeEvent) { + const key = this.getGlobalSseCoalesceKey(event); + if (key) { + const existingIndex = this.globalSseCoalesced.get(key); + if (existingIndex !== undefined) { + this.globalSseQueue[existingIndex] = undefined; + } + this.globalSseCoalesced.set(key, this.globalSseQueue.length); + } + + this.globalSseQueue.push(event); + this.scheduleGlobalSseFlush(); } private async runGlobalSseLoop(abortController: AbortController): Promise { @@ -1319,6 +1422,8 @@ class OpencodeService { const delay = Math.min(3000 * Math.pow(2, attempt), 30000); await new Promise((resolve) => setTimeout(resolve, delay)); } + + this.flushGlobalSseQueue(); } subscribeToGlobalEvents(