diff --git a/packages/ui/src/hooks/useQueuedMessageAutoSend.test.ts b/packages/ui/src/hooks/useQueuedMessageAutoSend.test.ts index 70c5eb43..f5a6a343 100644 --- a/packages/ui/src/hooks/useQueuedMessageAutoSend.test.ts +++ b/packages/ui/src/hooks/useQueuedMessageAutoSend.test.ts @@ -1,6 +1,8 @@ import { beforeEach, describe, expect, mock, test } from 'bun:test'; -import type { Agent } from '@opencode-ai/sdk/v2'; +import type { Agent, Message } from '@opencode-ai/sdk/v2'; import type { QueuedMessage } from '../stores/messageQueueStore'; +import { ChildStoreManager } from '@/sync/child-store'; +import { setSyncRefs } from '@/sync/sync-refs'; let visibleAgents: Agent[] = []; const sendMessageCalls: unknown[][] = []; @@ -32,6 +34,7 @@ import { createQueuedAutoSendRetryScheduler, getQueuedAutoSendRetryDelayMs, isQueuedAutoSendBackedOff, + resolveQueuedSessionStatusType, sendQueuedAutoSendPayload, shouldDispatchQueuedAutoSend, } from './useQueuedMessageAutoSend'; @@ -119,6 +122,59 @@ describe('queued auto-send retry backoff', () => { }); }); +describe('resolveQueuedSessionStatusType', () => { + const DIRECTORY = '/repo'; + + const assistantMessage = (id: string, completed?: number): Message => ({ + id, + role: 'assistant', + sessionID: 'ses_1', + time: { created: 1, ...(completed !== undefined ? { completed } : {}) }, + } as Message); + + let childStores: ChildStoreManager; + + beforeEach(() => { + childStores = new ChildStoreManager(); + const store = childStores.ensureChild(DIRECTORY, { bootstrap: false }); + store.setState({ status: 'complete', session_status: {}, message: {} }); + setSyncRefs({} as never, childStores, DIRECTORY); + }); + + test('treats a session with an in-flight assistant turn as busy even when the status entry is missing', () => { + // The server status map only lists busy/retry sessions, so a missed busy + // event leaves NO status entry while the turn is still streaming. The + // queue gate must not read that absence as idle: queued prompts would be + // dispatched into the running turn and merged into one model response. + childStores.ensureChild(DIRECTORY, { bootstrap: false }).setState({ + message: { ses_1: [assistantMessage('msg_streaming')] }, + }); + + expect(resolveQueuedSessionStatusType('ses_1', DIRECTORY)).toBe('busy'); + }); + + test('resolves an explicit busy or retry status entry', () => { + const store = childStores.ensureChild(DIRECTORY, { bootstrap: false }); + store.setState({ session_status: { ses_1: { type: 'busy' } } }); + expect(resolveQueuedSessionStatusType('ses_1', DIRECTORY)).toBe('busy'); + store.setState({ session_status: { ses_1: { type: 'retry', attempt: 2, message: 'boom', next: 30 } } }); + expect(resolveQueuedSessionStatusType('ses_1', DIRECTORY)).toBe('retry'); + }); + + test('resolves idle when the trailing assistant message has completed', () => { + const store = childStores.ensureChild(DIRECTORY, { bootstrap: false }); + store.setState({ message: { ses_1: [assistantMessage('msg_done', 5)] } }); + expect(resolveQueuedSessionStatusType('ses_1', DIRECTORY)).toBe('idle'); + }); + + test('resolves an explicit idle entry and unknown sessions as idle', () => { + const store = childStores.ensureChild(DIRECTORY, { bootstrap: false }); + store.setState({ session_status: { ses_1: { type: 'idle' } } }); + expect(resolveQueuedSessionStatusType('ses_1', DIRECTORY)).toBe('idle'); + expect(resolveQueuedSessionStatusType('ses_unknown', DIRECTORY)).toBe('idle'); + }); +}); + describe('buildQueuedAutoSendPayload', () => { beforeEach(() => { visibleAgents = []; diff --git a/packages/ui/src/hooks/useQueuedMessageAutoSend.ts b/packages/ui/src/hooks/useQueuedMessageAutoSend.ts index e7b39c97..38e95713 100644 --- a/packages/ui/src/hooks/useQueuedMessageAutoSend.ts +++ b/packages/ui/src/hooks/useQueuedMessageAutoSend.ts @@ -173,11 +173,50 @@ export const shouldDispatchQueuedAutoSend = ( && currentStatusType === 'idle'; }; +/** + * Resolve the live status the queue gate should honor for a session. + * + * The server's `/session/status` map only lists busy/retry sessions — idle + * sessions are absent — so a missing entry means "idle per the snapshot", not + * "no information". A missed busy event therefore leaves no entry while a turn + * is still streaming. The trailing in-flight assistant message is the live + * evidence of that running turn: treat it as busy so the queue never dispatches + * into it (mirrors `useSessionActivity`'s fallback). The entry becomes idle the + * moment the message completes or an idle status event lands. This reads the + * directory child store directly so both the effect-loop gate and the + * dispatch-time re-check agree. + */ +export const resolveQueuedSessionStatusType = ( + sessionId: string, + directory: string, +): SessionStatusType => { + const state = getDirectoryState(directory); + const statusType = state?.session_status?.[sessionId]?.type; + if (statusType === 'busy' || statusType === 'retry') { + return statusType; + } + const sessionMessages = state?.message?.[sessionId]; + const lastMessage = sessionMessages && sessionMessages.length > 0 + ? sessionMessages[sessionMessages.length - 1] + : undefined; + if ( + lastMessage?.role === 'assistant' + && typeof (lastMessage as { time?: { completed?: number } }).time?.completed !== 'number' + ) { + return 'busy'; + } + return 'idle'; +}; + export function useQueuedMessageAutoSend(enabledOrOptions?: boolean | { enabled?: boolean }) { const enabled = typeof enabledOrOptions === 'boolean' ? enabledOrOptions : (enabledOrOptions?.enabled ?? true); const queuedMessages = useMessageQueueStore((state) => state.queuedMessages); const autoReviewRuns = useAutoReviewStore((state) => state.runsByOriginalSessionID); const sessionStatusRecord = useDirectorySync((state) => state.session_status); + // Message completion clears the in-flight fallback in + // resolveQueuedSessionStatusType; subscribe so the queue drains the moment + // the trailing assistant message completes even if status events were missed. + const sessionMessages = useDirectorySync((state) => state.message); const currentDirectory = useDirectoryStore((state) => state.currentDirectory); const inFlightSessionsRef = React.useRef>(new Set()); @@ -216,7 +255,7 @@ export function useQueuedMessageAutoSend(enabledOrOptions?: boolean | { enabled? return; } - const currentStatus = getDirectoryState(target.directory)?.session_status?.[sessionId]?.type ?? 'idle'; + const currentStatus = resolveQueuedSessionStatusType(sessionId, target.directory); if (currentStatus !== 'idle') { return; } @@ -294,7 +333,7 @@ export function useQueuedMessageAutoSend(enabledOrOptions?: boolean | { enabled? const target = parseMessageQueueKey(key); if (!target || target.runtimeKey !== getRuntimeKey() || target.directory !== currentDirectory) return; const { sessionId } = target; - const currentStatusType = (statusRecord[sessionId]?.type ?? 'idle') as SessionStatusType; + const currentStatusType = resolveQueuedSessionStatusType(sessionId, target.directory); const previousStatusType = previousStatusRef.current.get(sessionId); const wasAutoReviewBlocked = autoReviewBlockedSessionsRef.current.has(sessionId); const isAutoReviewRunning = useAutoReviewStore.getState().isRunningForSession(sessionId); @@ -315,5 +354,5 @@ export function useQueuedMessageAutoSend(enabledOrOptions?: boolean | { enabled? }); previousStatusRef.current = nextStatusMap; - }, [enabled, queuedMessages, sessionStatusRecord, autoReviewRuns, currentDirectory, retryTick, retryScheduler]); + }, [enabled, queuedMessages, sessionStatusRecord, sessionMessages, autoReviewRuns, currentDirectory, retryTick, retryScheduler]); }