From 673c0f4dfc1f29a8fd60d1ad0571f67021cd6358 Mon Sep 17 00:00:00 2001 From: wenshao Date: Fri, 3 Apr 2026 19:42:03 +0800 Subject: [PATCH 1/6] feat(core): implement mid-turn queue drain for agent execution Inject queued user messages between tool execution steps within a single turn, so the model sees them immediately instead of waiting for the entire round to complete. - Add `dequeueAll()` to AsyncMessageQueue - Add `midTurnDrain` callback to ReasoningLoopOptions - Drain queue after processFunctionCalls, inject as text parts - AgentComposer always enqueues directly (no local buffering) - Add QUEUE_MESSAGES_CONSUMED event for UI sync Co-Authored-By: Claude Opus 4.6 (1M context) --- .../components/agent-view/AgentComposer.tsx | 43 +++-- .../core/src/agents/runtime/agent-core.ts | 19 ++ .../core/src/agents/runtime/agent-events.ts | 14 +- .../agents/runtime/agent-interactive.test.ts | 171 +++++++++++++++++- .../src/agents/runtime/agent-interactive.ts | 30 +++ .../core/src/utils/asyncMessageQueue.test.ts | 28 +++ packages/core/src/utils/asyncMessageQueue.ts | 5 + 7 files changed, 288 insertions(+), 22 deletions(-) diff --git a/packages/cli/src/ui/components/agent-view/AgentComposer.tsx b/packages/cli/src/ui/components/agent-view/AgentComposer.tsx index d26d5db2fe6..c8d7d717df2 100644 --- a/packages/cli/src/ui/components/agent-view/AgentComposer.tsx +++ b/packages/cli/src/ui/components/agent-view/AgentComposer.tsx @@ -21,9 +21,9 @@ import { Box, Text, useStdin } from 'ink'; import { useCallback, useEffect, useMemo, useState } from 'react'; import { AgentStatus, - isTerminalStatus, ApprovalMode, APPROVAL_MODES, + AgentEventType, } from '@qwen-code/qwen-code-core'; import { useAgentViewState, @@ -184,32 +184,35 @@ export const AgentComposer: React.FC = ({ agentId }) => { [buffer, agentTabBarFocused, setAgentTabBarFocused], ); - // ── Message queue (accumulate while streaming, flush as one prompt on idle) ── + // ── Message queue (mid-turn drain: always enqueue, track display locally) ── - const [messageQueue, setMessageQueue] = useState([]); + const [pendingMessages, setPendingMessages] = useState([]); - // When agent becomes idle (and not terminal), flush queued messages. + // Clear display when messages are consumed mid-turn or agent goes idle. useEffect(() => { - if ( - streamingState === StreamingState.Idle && - messageQueue.length > 0 && - status !== undefined && - !isTerminalStatus(status) - ) { - const combined = messageQueue.join('\n'); - setMessageQueue([]); - interactiveAgent?.enqueueMessage(combined); - } - }, [streamingState, messageQueue, interactiveAgent, status]); + const emitter = interactiveAgent?.getEventEmitter(); + if (!emitter) return; + + const clearDisplay = () => setPendingMessages([]); + emitter.on(AgentEventType.QUEUE_MESSAGES_CONSUMED, clearDisplay); + emitter.on(AgentEventType.STATUS_CHANGE, clearDisplay); + + return () => { + emitter.off(AgentEventType.QUEUE_MESSAGES_CONSUMED, clearDisplay); + emitter.off(AgentEventType.STATUS_CHANGE, clearDisplay); + }; + }, [interactiveAgent]); const handleSubmit = useCallback( (text: string) => { const trimmed = text.trim(); if (!trimmed || !interactiveAgent) return; - if (streamingState === StreamingState.Idle) { - interactiveAgent.enqueueMessage(trimmed); - } else { - setMessageQueue((prev) => [...prev, trimmed]); + // Always enqueue directly — mid-turn drain picks up messages + // during tool execution rather than waiting for idle. + interactiveAgent.enqueueMessage(trimmed); + // Track for display when agent is busy + if (streamingState !== StreamingState.Idle) { + setPendingMessages((prev) => [...prev, trimmed]); } }, [interactiveAgent, streamingState], @@ -279,7 +282,7 @@ export const AgentComposer: React.FC = ({ agentId }) => { )} - + {/* Input prompt — always visible, like the main Composer */} string[]; } /** @@ -490,6 +496,19 @@ export class AgentCore { toolsList, currentResponseId, ); + + // Mid-turn queue drain: inject queued user messages alongside tool + // results so the model sees them in the next API call. + if (options?.midTurnDrain) { + const drained = options.midTurnDrain(); + if (drained.length > 0 && currentMessages[0]?.parts) { + for (const msg of drained) { + currentMessages[0].parts.push({ + text: `\n[User message received during tool execution]: ${msg}`, + }); + } + } + } } else { // No tool calls — treat this as the model's final answer. if (roundText && roundText.trim().length > 0) { diff --git a/packages/core/src/agents/runtime/agent-events.ts b/packages/core/src/agents/runtime/agent-events.ts index 4626bb0cd3e..578f9ebdd56 100644 --- a/packages/core/src/agents/runtime/agent-events.ts +++ b/packages/core/src/agents/runtime/agent-events.ts @@ -37,7 +37,8 @@ export type AgentEvent = | 'usage_metadata' | 'finish' | 'error' - | 'status_change'; + | 'status_change' + | 'queue_messages_consumed'; export enum AgentEventType { START = 'start', @@ -54,6 +55,7 @@ export enum AgentEventType { FINISH = 'finish', ERROR = 'error', STATUS_CHANGE = 'status_change', + QUEUE_MESSAGES_CONSUMED = 'queue_messages_consumed', } // ─── Event Payloads ───────────────────────────────────────── @@ -172,6 +174,15 @@ export interface AgentErrorEvent { timestamp: number; } +export interface AgentQueueDrainEvent { + subagentId: string; + /** The messages that were consumed from the queue. */ + messages: string[]; + /** Number of messages consumed. */ + count: number; + timestamp: number; +} + export interface AgentStatusChangeEvent { agentId: string; previousStatus: AgentStatus; @@ -200,6 +211,7 @@ export interface AgentEventMap { [AgentEventType.FINISH]: AgentFinishEvent; [AgentEventType.ERROR]: AgentErrorEvent; [AgentEventType.STATUS_CHANGE]: AgentStatusChangeEvent; + [AgentEventType.QUEUE_MESSAGES_CONSUMED]: AgentQueueDrainEvent; } // ─── Event Emitter ────────────────────────────────────────── diff --git a/packages/core/src/agents/runtime/agent-interactive.test.ts b/packages/core/src/agents/runtime/agent-interactive.test.ts index 5560b665f71..1e82bd5b196 100644 --- a/packages/core/src/agents/runtime/agent-interactive.test.ts +++ b/packages/core/src/agents/runtime/agent-interactive.test.ts @@ -6,7 +6,7 @@ import { describe, it, expect, vi, beforeEach } from 'vitest'; import { AgentInteractive } from './agent-interactive.js'; -import type { AgentCore } from './agent-core.js'; +import type { AgentCore, ReasoningLoopOptions } from './agent-core.js'; import { AgentEventEmitter, AgentEventType } from './agent-events.js'; import { ContextState } from './agent-headless.js'; import type { AgentInteractiveConfig } from './agent-types.js'; @@ -600,6 +600,175 @@ describe('AgentInteractive', () => { await agent.shutdown(); }); + // ─── Mid-Turn Queue Drain ─────────────────────────────────── + + it('should pass midTurnDrain callback to runReasoningLoop', async () => { + const { core } = createMockCore(); + + (core.runReasoningLoop as ReturnType).mockImplementation( + async ( + _chat: unknown, + _msgs: unknown, + _tools: unknown, + _abort: unknown, + options?: ReasoningLoopOptions, + ) => { + expect(options?.midTurnDrain).toBeTypeOf('function'); + return { text: 'Done', terminateMode: null, turnsUsed: 1 }; + }, + ); + + const config = createConfig({ initialTask: 'test' }); + const agent = new AgentInteractive(config, core); + await agent.start(context); + await vi.waitFor(() => expect(agent.getStatus()).toBe('idle')); + + const callArgs = (core.runReasoningLoop as ReturnType).mock + .calls[0]; + expect(callArgs[4]).toHaveProperty('midTurnDrain'); + + await agent.shutdown(); + }); + + it('should drain queued messages mid-turn and record them', async () => { + const { core } = createMockCore(); + + let drainCallback: (() => string[]) | undefined; + let resolveLoop: (() => void) | undefined; + + (core.runReasoningLoop as ReturnType).mockImplementation( + async ( + _chat: unknown, + _msgs: unknown, + _tools: unknown, + _abort: unknown, + options?: ReasoningLoopOptions, + ) => { + drainCallback = options?.midTurnDrain; + await new Promise((r) => { + resolveLoop = r; + }); + return { text: 'Done', terminateMode: null, turnsUsed: 1 }; + }, + ); + + const config = createConfig(); + const agent = new AgentInteractive(config, core); + await agent.start(context); + + // Start a message processing + agent.enqueueMessage('primary task'); + await vi.waitFor(() => expect(agent.getStatus()).toBe('running')); + + // Enqueue another message while "running" + agent.enqueueMessage('extra context'); + expect(agent.getQueueSize()).toBe(1); + + // Mid-turn drain should pick it up + expect(drainCallback).toBeDefined(); + const drained = drainCallback!(); + expect(drained).toEqual(['extra context']); + expect(agent.getQueueSize()).toBe(0); + + // Verify recorded as user message with metadata + const messages = agent.getMessages(); + const midTurnMsg = messages.find( + (m) => + m.role === 'user' && + m.content === 'extra context' && + m.metadata?.['midTurnInjected'] === true, + ); + expect(midTurnMsg).toBeDefined(); + + resolveLoop!(); + await vi.waitFor(() => expect(agent.getStatus()).toBe('idle')); + await agent.shutdown(); + }); + + it('should emit QUEUE_MESSAGES_CONSUMED event on drain', async () => { + const { core, emitter } = createMockCore(); + + let drainCallback: (() => string[]) | undefined; + let resolveLoop: (() => void) | undefined; + + (core.runReasoningLoop as ReturnType).mockImplementation( + async ( + _chat: unknown, + _msgs: unknown, + _tools: unknown, + _abort: unknown, + options?: ReasoningLoopOptions, + ) => { + drainCallback = options?.midTurnDrain; + await new Promise((r) => { + resolveLoop = r; + }); + return { text: 'Done', terminateMode: null, turnsUsed: 1 }; + }, + ); + + const config = createConfig(); + const agent = new AgentInteractive(config, core); + await agent.start(context); + + agent.enqueueMessage('task'); + await vi.waitFor(() => expect(agent.getStatus()).toBe('running')); + + agent.enqueueMessage('update'); + + const drainEvents: unknown[] = []; + emitter.on(AgentEventType.QUEUE_MESSAGES_CONSUMED, (e: unknown) => + drainEvents.push(e), + ); + + drainCallback!(); + expect(drainEvents).toHaveLength(1); + expect(drainEvents[0]).toMatchObject({ + messages: ['update'], + count: 1, + }); + + resolveLoop!(); + await vi.waitFor(() => expect(agent.getStatus()).toBe('idle')); + await agent.shutdown(); + }); + + it('should not emit event when drain finds empty queue', async () => { + const { core, emitter } = createMockCore(); + + let drainCallback: (() => string[]) | undefined; + + (core.runReasoningLoop as ReturnType).mockImplementation( + async ( + _chat: unknown, + _msgs: unknown, + _tools: unknown, + _abort: unknown, + options?: ReasoningLoopOptions, + ) => { + drainCallback = options?.midTurnDrain; + return { text: 'Done', terminateMode: null, turnsUsed: 1 }; + }, + ); + + const config = createConfig({ initialTask: 'test' }); + const agent = new AgentInteractive(config, core); + await agent.start(context); + await vi.waitFor(() => expect(agent.getStatus()).toBe('idle')); + + const drainEvents: unknown[] = []; + emitter.on(AgentEventType.QUEUE_MESSAGES_CONSUMED, (e: unknown) => + drainEvents.push(e), + ); + + // Drain on empty queue + const drained = drainCallback!(); + expect(drained).toEqual([]); + expect(drainEvents).toHaveLength(0); + + await agent.shutdown(); + }); + // ─── Events ──────────────────────────────────────────────── it('should emit status_change events', async () => { diff --git a/packages/core/src/agents/runtime/agent-interactive.ts b/packages/core/src/agents/runtime/agent-interactive.ts index 42e9dedce12..f67ef831869 100644 --- a/packages/core/src/agents/runtime/agent-interactive.ts +++ b/packages/core/src/agents/runtime/agent-interactive.ts @@ -186,6 +186,7 @@ export class AgentInteractive { { maxTurns: this.config.maxTurnsPerMessage, maxTimeMinutes: this.config.maxTimeMinutesPerMessage, + midTurnDrain: () => this.drainQueuedMessages(), }, ); @@ -270,6 +271,11 @@ export class AgentInteractive { } } + /** Number of messages currently waiting in the queue. */ + getQueueSize(): number { + return this.queue.size; + } + // ─── State Accessors ─────────────────────────────────────── getMessages(): readonly AgentMessage[] { @@ -342,6 +348,30 @@ export class AgentInteractive { } } + // ─── Mid-Turn Queue Drain ─────────────────────────────────── + + /** + * Drain all pending messages from the queue for mid-turn injection. + * Called by the midTurnDrain callback during the reasoning loop. + */ + private drainQueuedMessages(): string[] { + const messages = this.queue.dequeueAll(); + if (messages.length > 0) { + for (const msg of messages) { + this.addMessage('user', msg, { + metadata: { midTurnInjected: true }, + }); + } + this.core.eventEmitter?.emit(AgentEventType.QUEUE_MESSAGES_CONSUMED, { + subagentId: this.core.subagentId, + messages, + count: messages.length, + timestamp: Date.now(), + }); + } + return messages; + } + // ─── Private Helpers ─────────────────────────────────────── /** diff --git a/packages/core/src/utils/asyncMessageQueue.test.ts b/packages/core/src/utils/asyncMessageQueue.test.ts index fe54210333c..08d80437d3d 100644 --- a/packages/core/src/utils/asyncMessageQueue.test.ts +++ b/packages/core/src/utils/asyncMessageQueue.test.ts @@ -72,4 +72,32 @@ describe('AsyncMessageQueue', () => { expect(queue.dequeue()).toBe(i); } }); + + // ── dequeueAll ── + + it('should dequeueAll items and leave the queue empty', () => { + const queue = new AsyncMessageQueue(); + queue.enqueue('a'); + queue.enqueue('b'); + queue.enqueue('c'); + + const all = queue.dequeueAll(); + expect(all).toEqual(['a', 'b', 'c']); + expect(queue.size).toBe(0); + expect(queue.dequeue()).toBeNull(); + }); + + it('should return empty array when dequeueAll on empty queue', () => { + const queue = new AsyncMessageQueue(); + expect(queue.dequeueAll()).toEqual([]); + }); + + it('should not affect drained state when dequeueAll', () => { + const queue = new AsyncMessageQueue(); + queue.enqueue('x'); + queue.drain(); + const all = queue.dequeueAll(); + expect(all).toEqual(['x']); + expect(queue.isDrained).toBe(true); + }); }); diff --git a/packages/core/src/utils/asyncMessageQueue.ts b/packages/core/src/utils/asyncMessageQueue.ts index 3268718ef00..cbf0c0b8d4d 100644 --- a/packages/core/src/utils/asyncMessageQueue.ts +++ b/packages/core/src/utils/asyncMessageQueue.ts @@ -37,6 +37,11 @@ export class AsyncMessageQueue { return null; } + /** Remove and return all items currently in the queue. */ + dequeueAll(): T[] { + return this.items.splice(0); + } + /** Signal that no more items will be enqueued. */ drain(): void { this.drained = true; From e916934c805556564e48a94b45266672552c514a Mon Sep 17 00:00:00 2001 From: wenshao Date: Fri, 3 Apr 2026 20:11:40 +0800 Subject: [PATCH 2/6] feat(cli): add mid-turn queue drain to main session Extend mid-turn queue drain to the main session's tool execution path (useGeminiStream). Previously only agent tabs had this feature. - Add midTurnDrainRef parameter to useGeminiStream - Inject queued messages in handleCompletedTools before submitQuery - Bridge useMessageQueue to drain ref in AppContainer via ref pattern Co-Authored-By: Claude Opus 4.6 (1M context) --- packages/cli/src/ui/AppContainer.tsx | 12 ++++++++++++ packages/cli/src/ui/hooks/useGeminiStream.ts | 13 +++++++++++++ 2 files changed, 25 insertions(+) diff --git a/packages/cli/src/ui/AppContainer.tsx b/packages/cli/src/ui/AppContainer.tsx index 37dc325180f..b408c59d6ee 100644 --- a/packages/cli/src/ui/AppContainer.tsx +++ b/packages/cli/src/ui/AppContainer.tsx @@ -697,6 +697,7 @@ export const AppContainer = (props: AppContainerProps) => { }, [config, historyManager, settings.merged]); const cancelHandlerRef = useRef<() => void>(() => {}); + const midTurnDrainRef = useRef<(() => string[]) | null>(null); const { streamingState, @@ -728,6 +729,7 @@ export const AppContainer = (props: AppContainerProps) => { setEmbeddedShellFocused, terminalWidth, terminalHeight, + midTurnDrainRef, ); // Track whether suggestions are visible for Tab key handling @@ -751,6 +753,16 @@ export const AppContainer = (props: AppContainerProps) => { submitQuery, }); + // Bridge message queue to mid-turn drain via ref (avoids stale closure). + const messageQueueRef = useRef(messageQueue); + messageQueueRef.current = messageQueue; + midTurnDrainRef.current = () => { + const queue = messageQueueRef.current; + if (queue.length === 0) return []; + clearQueue(); + return [...queue]; + }; + // Callback for handling final submit (must be after addMessage from useMessageQueue) const handleFinalSubmit = useCallback( (submittedValue: string) => { diff --git a/packages/cli/src/ui/hooks/useGeminiStream.ts b/packages/cli/src/ui/hooks/useGeminiStream.ts index 4dd2e952193..998689c561f 100644 --- a/packages/cli/src/ui/hooks/useGeminiStream.ts +++ b/packages/cli/src/ui/hooks/useGeminiStream.ts @@ -172,6 +172,7 @@ export const useGeminiStream = ( setShellInputFocused: (value: boolean) => void, terminalWidth: number, terminalHeight: number, + midTurnDrainRef?: React.RefObject<(() => string[]) | null>, ) => { const [initError, setInitError] = useState(null); const abortControllerRef = useRef(null); @@ -1561,6 +1562,17 @@ export const useGeminiStream = ( return; } + // Mid-turn queue drain: inject queued user messages alongside tool + // results so the model sees them in the next API call. + const drained = midTurnDrainRef?.current?.() ?? []; + if (drained.length > 0) { + for (const msg of drained) { + responsesToSend.push({ + text: `\n[User message received during tool execution]: ${msg}`, + }); + } + } + submitQuery(responsesToSend, SendMessageType.ToolResult, prompt_ids[0]); }, [ @@ -1571,6 +1583,7 @@ export const useGeminiStream = ( performMemoryRefresh, modelSwitchedFromQuotaError, config, + midTurnDrainRef, ], ); From ef55bd512f95eff9452a2bc727596455b48350e5 Mon Sep 17 00:00:00 2001 From: wenshao Date: Fri, 3 Apr 2026 20:27:18 +0800 Subject: [PATCH 3/6] fix: address Copilot review feedback on mid-turn drain - Guard midTurnDrain with abort check to prevent message loss on cancel - Synchronously clear messageQueueRef to prevent duplicate drains - Only clear pending display on IDLE status, not all status changes Co-Authored-By: Claude Opus 4.6 (1M context) --- packages/cli/src/ui/AppContainer.tsx | 3 +++ .../cli/src/ui/components/agent-view/AgentComposer.tsx | 9 +++++++-- packages/core/src/agents/runtime/agent-core.ts | 3 ++- 3 files changed, 12 insertions(+), 3 deletions(-) diff --git a/packages/cli/src/ui/AppContainer.tsx b/packages/cli/src/ui/AppContainer.tsx index b408c59d6ee..6842a8de4b5 100644 --- a/packages/cli/src/ui/AppContainer.tsx +++ b/packages/cli/src/ui/AppContainer.tsx @@ -759,6 +759,9 @@ export const AppContainer = (props: AppContainerProps) => { midTurnDrainRef.current = () => { const queue = messageQueueRef.current; if (queue.length === 0) return []; + // Synchronously clear the ref to prevent duplicate drains if called + // again before React re-renders. + messageQueueRef.current = []; clearQueue(); return [...queue]; }; diff --git a/packages/cli/src/ui/components/agent-view/AgentComposer.tsx b/packages/cli/src/ui/components/agent-view/AgentComposer.tsx index c8d7d717df2..7612b020be6 100644 --- a/packages/cli/src/ui/components/agent-view/AgentComposer.tsx +++ b/packages/cli/src/ui/components/agent-view/AgentComposer.tsx @@ -194,12 +194,17 @@ export const AgentComposer: React.FC = ({ agentId }) => { if (!emitter) return; const clearDisplay = () => setPendingMessages([]); + const clearOnIdle = (event: { newStatus?: AgentStatus }) => { + if (event.newStatus === AgentStatus.IDLE) { + setPendingMessages([]); + } + }; emitter.on(AgentEventType.QUEUE_MESSAGES_CONSUMED, clearDisplay); - emitter.on(AgentEventType.STATUS_CHANGE, clearDisplay); + emitter.on(AgentEventType.STATUS_CHANGE, clearOnIdle); return () => { emitter.off(AgentEventType.QUEUE_MESSAGES_CONSUMED, clearDisplay); - emitter.off(AgentEventType.STATUS_CHANGE, clearDisplay); + emitter.off(AgentEventType.STATUS_CHANGE, clearOnIdle); }; }, [interactiveAgent]); diff --git a/packages/core/src/agents/runtime/agent-core.ts b/packages/core/src/agents/runtime/agent-core.ts index f0c62e9f94f..06602b23a7e 100644 --- a/packages/core/src/agents/runtime/agent-core.ts +++ b/packages/core/src/agents/runtime/agent-core.ts @@ -499,7 +499,8 @@ export class AgentCore { // Mid-turn queue drain: inject queued user messages alongside tool // results so the model sees them in the next API call. - if (options?.midTurnDrain) { + // Skip if aborted — messages stay in queue for the next round. + if (options?.midTurnDrain && !abortController.signal.aborted) { const drained = options.midTurnDrain(); if (drained.length > 0 && currentMessages[0]?.parts) { for (const msg of drained) { From 3c8a5bdb72305b40875a104e5e0d9793e42e56e6 Mon Sep 17 00:00:00 2001 From: wenshao Date: Fri, 3 Apr 2026 20:46:11 +0800 Subject: [PATCH 4/6] refactor: scope mid-turn drain to main session only Revert subagent-path changes (AgentCore, AgentInteractive, AgentComposer, AsyncMessageQueue, agent-events) to keep the PR focused on the main session, which is easier to test and validate. Subagent mid-turn drain can be added in a follow-up PR. Co-Authored-By: Claude Opus 4.6 (1M context) --- .../components/agent-view/AgentComposer.tsx | 48 ++--- .../core/src/agents/runtime/agent-core.ts | 20 -- .../core/src/agents/runtime/agent-events.ts | 14 +- .../agents/runtime/agent-interactive.test.ts | 171 +----------------- .../src/agents/runtime/agent-interactive.ts | 30 --- .../core/src/utils/asyncMessageQueue.test.ts | 28 --- packages/core/src/utils/asyncMessageQueue.ts | 5 - 7 files changed, 22 insertions(+), 294 deletions(-) diff --git a/packages/cli/src/ui/components/agent-view/AgentComposer.tsx b/packages/cli/src/ui/components/agent-view/AgentComposer.tsx index 7612b020be6..d26d5db2fe6 100644 --- a/packages/cli/src/ui/components/agent-view/AgentComposer.tsx +++ b/packages/cli/src/ui/components/agent-view/AgentComposer.tsx @@ -21,9 +21,9 @@ import { Box, Text, useStdin } from 'ink'; import { useCallback, useEffect, useMemo, useState } from 'react'; import { AgentStatus, + isTerminalStatus, ApprovalMode, APPROVAL_MODES, - AgentEventType, } from '@qwen-code/qwen-code-core'; import { useAgentViewState, @@ -184,40 +184,32 @@ export const AgentComposer: React.FC = ({ agentId }) => { [buffer, agentTabBarFocused, setAgentTabBarFocused], ); - // ── Message queue (mid-turn drain: always enqueue, track display locally) ── + // ── Message queue (accumulate while streaming, flush as one prompt on idle) ── - const [pendingMessages, setPendingMessages] = useState([]); + const [messageQueue, setMessageQueue] = useState([]); - // Clear display when messages are consumed mid-turn or agent goes idle. + // When agent becomes idle (and not terminal), flush queued messages. useEffect(() => { - const emitter = interactiveAgent?.getEventEmitter(); - if (!emitter) return; - - const clearDisplay = () => setPendingMessages([]); - const clearOnIdle = (event: { newStatus?: AgentStatus }) => { - if (event.newStatus === AgentStatus.IDLE) { - setPendingMessages([]); - } - }; - emitter.on(AgentEventType.QUEUE_MESSAGES_CONSUMED, clearDisplay); - emitter.on(AgentEventType.STATUS_CHANGE, clearOnIdle); - - return () => { - emitter.off(AgentEventType.QUEUE_MESSAGES_CONSUMED, clearDisplay); - emitter.off(AgentEventType.STATUS_CHANGE, clearOnIdle); - }; - }, [interactiveAgent]); + if ( + streamingState === StreamingState.Idle && + messageQueue.length > 0 && + status !== undefined && + !isTerminalStatus(status) + ) { + const combined = messageQueue.join('\n'); + setMessageQueue([]); + interactiveAgent?.enqueueMessage(combined); + } + }, [streamingState, messageQueue, interactiveAgent, status]); const handleSubmit = useCallback( (text: string) => { const trimmed = text.trim(); if (!trimmed || !interactiveAgent) return; - // Always enqueue directly — mid-turn drain picks up messages - // during tool execution rather than waiting for idle. - interactiveAgent.enqueueMessage(trimmed); - // Track for display when agent is busy - if (streamingState !== StreamingState.Idle) { - setPendingMessages((prev) => [...prev, trimmed]); + if (streamingState === StreamingState.Idle) { + interactiveAgent.enqueueMessage(trimmed); + } else { + setMessageQueue((prev) => [...prev, trimmed]); } }, [interactiveAgent, streamingState], @@ -287,7 +279,7 @@ export const AgentComposer: React.FC = ({ agentId }) => { )} - + {/* Input prompt — always visible, like the main Composer */} string[]; } /** @@ -496,20 +490,6 @@ export class AgentCore { toolsList, currentResponseId, ); - - // Mid-turn queue drain: inject queued user messages alongside tool - // results so the model sees them in the next API call. - // Skip if aborted — messages stay in queue for the next round. - if (options?.midTurnDrain && !abortController.signal.aborted) { - const drained = options.midTurnDrain(); - if (drained.length > 0 && currentMessages[0]?.parts) { - for (const msg of drained) { - currentMessages[0].parts.push({ - text: `\n[User message received during tool execution]: ${msg}`, - }); - } - } - } } else { // No tool calls — treat this as the model's final answer. if (roundText && roundText.trim().length > 0) { diff --git a/packages/core/src/agents/runtime/agent-events.ts b/packages/core/src/agents/runtime/agent-events.ts index 578f9ebdd56..4626bb0cd3e 100644 --- a/packages/core/src/agents/runtime/agent-events.ts +++ b/packages/core/src/agents/runtime/agent-events.ts @@ -37,8 +37,7 @@ export type AgentEvent = | 'usage_metadata' | 'finish' | 'error' - | 'status_change' - | 'queue_messages_consumed'; + | 'status_change'; export enum AgentEventType { START = 'start', @@ -55,7 +54,6 @@ export enum AgentEventType { FINISH = 'finish', ERROR = 'error', STATUS_CHANGE = 'status_change', - QUEUE_MESSAGES_CONSUMED = 'queue_messages_consumed', } // ─── Event Payloads ───────────────────────────────────────── @@ -174,15 +172,6 @@ export interface AgentErrorEvent { timestamp: number; } -export interface AgentQueueDrainEvent { - subagentId: string; - /** The messages that were consumed from the queue. */ - messages: string[]; - /** Number of messages consumed. */ - count: number; - timestamp: number; -} - export interface AgentStatusChangeEvent { agentId: string; previousStatus: AgentStatus; @@ -211,7 +200,6 @@ export interface AgentEventMap { [AgentEventType.FINISH]: AgentFinishEvent; [AgentEventType.ERROR]: AgentErrorEvent; [AgentEventType.STATUS_CHANGE]: AgentStatusChangeEvent; - [AgentEventType.QUEUE_MESSAGES_CONSUMED]: AgentQueueDrainEvent; } // ─── Event Emitter ────────────────────────────────────────── diff --git a/packages/core/src/agents/runtime/agent-interactive.test.ts b/packages/core/src/agents/runtime/agent-interactive.test.ts index 1e82bd5b196..5560b665f71 100644 --- a/packages/core/src/agents/runtime/agent-interactive.test.ts +++ b/packages/core/src/agents/runtime/agent-interactive.test.ts @@ -6,7 +6,7 @@ import { describe, it, expect, vi, beforeEach } from 'vitest'; import { AgentInteractive } from './agent-interactive.js'; -import type { AgentCore, ReasoningLoopOptions } from './agent-core.js'; +import type { AgentCore } from './agent-core.js'; import { AgentEventEmitter, AgentEventType } from './agent-events.js'; import { ContextState } from './agent-headless.js'; import type { AgentInteractiveConfig } from './agent-types.js'; @@ -600,175 +600,6 @@ describe('AgentInteractive', () => { await agent.shutdown(); }); - // ─── Mid-Turn Queue Drain ─────────────────────────────────── - - it('should pass midTurnDrain callback to runReasoningLoop', async () => { - const { core } = createMockCore(); - - (core.runReasoningLoop as ReturnType).mockImplementation( - async ( - _chat: unknown, - _msgs: unknown, - _tools: unknown, - _abort: unknown, - options?: ReasoningLoopOptions, - ) => { - expect(options?.midTurnDrain).toBeTypeOf('function'); - return { text: 'Done', terminateMode: null, turnsUsed: 1 }; - }, - ); - - const config = createConfig({ initialTask: 'test' }); - const agent = new AgentInteractive(config, core); - await agent.start(context); - await vi.waitFor(() => expect(agent.getStatus()).toBe('idle')); - - const callArgs = (core.runReasoningLoop as ReturnType).mock - .calls[0]; - expect(callArgs[4]).toHaveProperty('midTurnDrain'); - - await agent.shutdown(); - }); - - it('should drain queued messages mid-turn and record them', async () => { - const { core } = createMockCore(); - - let drainCallback: (() => string[]) | undefined; - let resolveLoop: (() => void) | undefined; - - (core.runReasoningLoop as ReturnType).mockImplementation( - async ( - _chat: unknown, - _msgs: unknown, - _tools: unknown, - _abort: unknown, - options?: ReasoningLoopOptions, - ) => { - drainCallback = options?.midTurnDrain; - await new Promise((r) => { - resolveLoop = r; - }); - return { text: 'Done', terminateMode: null, turnsUsed: 1 }; - }, - ); - - const config = createConfig(); - const agent = new AgentInteractive(config, core); - await agent.start(context); - - // Start a message processing - agent.enqueueMessage('primary task'); - await vi.waitFor(() => expect(agent.getStatus()).toBe('running')); - - // Enqueue another message while "running" - agent.enqueueMessage('extra context'); - expect(agent.getQueueSize()).toBe(1); - - // Mid-turn drain should pick it up - expect(drainCallback).toBeDefined(); - const drained = drainCallback!(); - expect(drained).toEqual(['extra context']); - expect(agent.getQueueSize()).toBe(0); - - // Verify recorded as user message with metadata - const messages = agent.getMessages(); - const midTurnMsg = messages.find( - (m) => - m.role === 'user' && - m.content === 'extra context' && - m.metadata?.['midTurnInjected'] === true, - ); - expect(midTurnMsg).toBeDefined(); - - resolveLoop!(); - await vi.waitFor(() => expect(agent.getStatus()).toBe('idle')); - await agent.shutdown(); - }); - - it('should emit QUEUE_MESSAGES_CONSUMED event on drain', async () => { - const { core, emitter } = createMockCore(); - - let drainCallback: (() => string[]) | undefined; - let resolveLoop: (() => void) | undefined; - - (core.runReasoningLoop as ReturnType).mockImplementation( - async ( - _chat: unknown, - _msgs: unknown, - _tools: unknown, - _abort: unknown, - options?: ReasoningLoopOptions, - ) => { - drainCallback = options?.midTurnDrain; - await new Promise((r) => { - resolveLoop = r; - }); - return { text: 'Done', terminateMode: null, turnsUsed: 1 }; - }, - ); - - const config = createConfig(); - const agent = new AgentInteractive(config, core); - await agent.start(context); - - agent.enqueueMessage('task'); - await vi.waitFor(() => expect(agent.getStatus()).toBe('running')); - - agent.enqueueMessage('update'); - - const drainEvents: unknown[] = []; - emitter.on(AgentEventType.QUEUE_MESSAGES_CONSUMED, (e: unknown) => - drainEvents.push(e), - ); - - drainCallback!(); - expect(drainEvents).toHaveLength(1); - expect(drainEvents[0]).toMatchObject({ - messages: ['update'], - count: 1, - }); - - resolveLoop!(); - await vi.waitFor(() => expect(agent.getStatus()).toBe('idle')); - await agent.shutdown(); - }); - - it('should not emit event when drain finds empty queue', async () => { - const { core, emitter } = createMockCore(); - - let drainCallback: (() => string[]) | undefined; - - (core.runReasoningLoop as ReturnType).mockImplementation( - async ( - _chat: unknown, - _msgs: unknown, - _tools: unknown, - _abort: unknown, - options?: ReasoningLoopOptions, - ) => { - drainCallback = options?.midTurnDrain; - return { text: 'Done', terminateMode: null, turnsUsed: 1 }; - }, - ); - - const config = createConfig({ initialTask: 'test' }); - const agent = new AgentInteractive(config, core); - await agent.start(context); - await vi.waitFor(() => expect(agent.getStatus()).toBe('idle')); - - const drainEvents: unknown[] = []; - emitter.on(AgentEventType.QUEUE_MESSAGES_CONSUMED, (e: unknown) => - drainEvents.push(e), - ); - - // Drain on empty queue - const drained = drainCallback!(); - expect(drained).toEqual([]); - expect(drainEvents).toHaveLength(0); - - await agent.shutdown(); - }); - // ─── Events ──────────────────────────────────────────────── it('should emit status_change events', async () => { diff --git a/packages/core/src/agents/runtime/agent-interactive.ts b/packages/core/src/agents/runtime/agent-interactive.ts index f67ef831869..42e9dedce12 100644 --- a/packages/core/src/agents/runtime/agent-interactive.ts +++ b/packages/core/src/agents/runtime/agent-interactive.ts @@ -186,7 +186,6 @@ export class AgentInteractive { { maxTurns: this.config.maxTurnsPerMessage, maxTimeMinutes: this.config.maxTimeMinutesPerMessage, - midTurnDrain: () => this.drainQueuedMessages(), }, ); @@ -271,11 +270,6 @@ export class AgentInteractive { } } - /** Number of messages currently waiting in the queue. */ - getQueueSize(): number { - return this.queue.size; - } - // ─── State Accessors ─────────────────────────────────────── getMessages(): readonly AgentMessage[] { @@ -348,30 +342,6 @@ export class AgentInteractive { } } - // ─── Mid-Turn Queue Drain ─────────────────────────────────── - - /** - * Drain all pending messages from the queue for mid-turn injection. - * Called by the midTurnDrain callback during the reasoning loop. - */ - private drainQueuedMessages(): string[] { - const messages = this.queue.dequeueAll(); - if (messages.length > 0) { - for (const msg of messages) { - this.addMessage('user', msg, { - metadata: { midTurnInjected: true }, - }); - } - this.core.eventEmitter?.emit(AgentEventType.QUEUE_MESSAGES_CONSUMED, { - subagentId: this.core.subagentId, - messages, - count: messages.length, - timestamp: Date.now(), - }); - } - return messages; - } - // ─── Private Helpers ─────────────────────────────────────── /** diff --git a/packages/core/src/utils/asyncMessageQueue.test.ts b/packages/core/src/utils/asyncMessageQueue.test.ts index 08d80437d3d..fe54210333c 100644 --- a/packages/core/src/utils/asyncMessageQueue.test.ts +++ b/packages/core/src/utils/asyncMessageQueue.test.ts @@ -72,32 +72,4 @@ describe('AsyncMessageQueue', () => { expect(queue.dequeue()).toBe(i); } }); - - // ── dequeueAll ── - - it('should dequeueAll items and leave the queue empty', () => { - const queue = new AsyncMessageQueue(); - queue.enqueue('a'); - queue.enqueue('b'); - queue.enqueue('c'); - - const all = queue.dequeueAll(); - expect(all).toEqual(['a', 'b', 'c']); - expect(queue.size).toBe(0); - expect(queue.dequeue()).toBeNull(); - }); - - it('should return empty array when dequeueAll on empty queue', () => { - const queue = new AsyncMessageQueue(); - expect(queue.dequeueAll()).toEqual([]); - }); - - it('should not affect drained state when dequeueAll', () => { - const queue = new AsyncMessageQueue(); - queue.enqueue('x'); - queue.drain(); - const all = queue.dequeueAll(); - expect(all).toEqual(['x']); - expect(queue.isDrained).toBe(true); - }); }); diff --git a/packages/core/src/utils/asyncMessageQueue.ts b/packages/core/src/utils/asyncMessageQueue.ts index cbf0c0b8d4d..3268718ef00 100644 --- a/packages/core/src/utils/asyncMessageQueue.ts +++ b/packages/core/src/utils/asyncMessageQueue.ts @@ -37,11 +37,6 @@ export class AsyncMessageQueue { return null; } - /** Remove and return all items currently in the queue. */ - dequeueAll(): T[] { - return this.items.splice(0); - } - /** Signal that no more items will be enqueued. */ drain(): void { this.drained = true; From 2a8e578dcfff0e377d6554afe3a0d1b8f36ccd33 Mon Sep 17 00:00:00 2001 From: wenshao Date: Fri, 3 Apr 2026 21:04:02 +0800 Subject: [PATCH 5/6] fix: address Copilot review on main session mid-turn drain - Move synchronous queue ref into useMessageQueue itself, expose drainQueue() for atomic drain (fixes race between addMessage and drain) - Record drained messages as USER history items so the transcript stays complete - Simplify AppContainer bridge to just midTurnDrainRef.current = drainQueue Co-Authored-By: Claude Opus 4.6 (1M context) --- packages/cli/src/ui/AppContainer.tsx | 31 ++++++++--------- packages/cli/src/ui/hooks/useGeminiStream.ts | 3 ++ packages/cli/src/ui/hooks/useMessageQueue.ts | 35 +++++++++++++++++--- 3 files changed, 47 insertions(+), 22 deletions(-) diff --git a/packages/cli/src/ui/AppContainer.tsx b/packages/cli/src/ui/AppContainer.tsx index 6842a8de4b5..4604a9c378b 100644 --- a/packages/cli/src/ui/AppContainer.tsx +++ b/packages/cli/src/ui/AppContainer.tsx @@ -746,25 +746,20 @@ export const AppContainer = (props: AppContainerProps) => { disabled: agentViewState.activeView !== 'main', }); - const { messageQueue, addMessage, clearQueue, getQueuedMessagesText } = - useMessageQueue({ - isConfigInitialized, - streamingState, - submitQuery, - }); + const { + messageQueue, + addMessage, + clearQueue, + getQueuedMessagesText, + drainQueue, + } = useMessageQueue({ + isConfigInitialized, + streamingState, + submitQuery, + }); - // Bridge message queue to mid-turn drain via ref (avoids stale closure). - const messageQueueRef = useRef(messageQueue); - messageQueueRef.current = messageQueue; - midTurnDrainRef.current = () => { - const queue = messageQueueRef.current; - if (queue.length === 0) return []; - // Synchronously clear the ref to prevent duplicate drains if called - // again before React re-renders. - messageQueueRef.current = []; - clearQueue(); - return [...queue]; - }; + // Bridge message queue to mid-turn drain ref for useGeminiStream. + midTurnDrainRef.current = drainQueue; // Callback for handling final submit (must be after addMessage from useMessageQueue) const handleFinalSubmit = useCallback( diff --git a/packages/cli/src/ui/hooks/useGeminiStream.ts b/packages/cli/src/ui/hooks/useGeminiStream.ts index 998689c561f..ee6acb17719 100644 --- a/packages/cli/src/ui/hooks/useGeminiStream.ts +++ b/packages/cli/src/ui/hooks/useGeminiStream.ts @@ -1570,6 +1570,8 @@ export const useGeminiStream = ( responsesToSend.push({ text: `\n[User message received during tool execution]: ${msg}`, }); + // Record in UI history so the transcript stays complete. + addItem({ type: MessageType.USER, text: msg }, Date.now()); } } @@ -1584,6 +1586,7 @@ export const useGeminiStream = ( modelSwitchedFromQuotaError, config, midTurnDrainRef, + addItem, ], ); diff --git a/packages/cli/src/ui/hooks/useMessageQueue.ts b/packages/cli/src/ui/hooks/useMessageQueue.ts index 517040fecfa..00af564b38f 100644 --- a/packages/cli/src/ui/hooks/useMessageQueue.ts +++ b/packages/cli/src/ui/hooks/useMessageQueue.ts @@ -4,7 +4,7 @@ * SPDX-License-Identifier: Apache-2.0 */ -import { useCallback, useEffect, useState } from 'react'; +import { useCallback, useEffect, useRef, useState } from 'react'; import { StreamingState } from '../types.js'; export interface UseMessageQueueOptions { @@ -18,6 +18,12 @@ export interface UseMessageQueueReturn { addMessage: (message: string) => void; clearQueue: () => void; getQueuedMessagesText: () => string; + /** + * Atomically drain all queued messages. Returns the drained messages + * and clears both the synchronous ref and React state. Safe to call + * from non-React contexts (e.g., tool completion callbacks). + */ + drainQueue: () => string[]; } /** @@ -31,17 +37,22 @@ export function useMessageQueue({ submitQuery, }: UseMessageQueueOptions): UseMessageQueueReturn { const [messageQueue, setMessageQueue] = useState([]); + // Synchronous ref mirrors React state so non-React callbacks (e.g., + // mid-turn drain in handleCompletedTools) always see the latest queue. + const queueRef = useRef([]); // Add a message to the queue const addMessage = useCallback((message: string) => { const trimmedMessage = message.trim(); if (trimmedMessage.length > 0) { - setMessageQueue((prev) => [...prev, trimmedMessage]); + queueRef.current = [...queueRef.current, trimmedMessage]; + setMessageQueue(queueRef.current); } }, []); // Clear the entire queue const clearQueue = useCallback(() => { + queueRef.current = []; setMessageQueue([]); }, []); @@ -51,6 +62,15 @@ export function useMessageQueue({ return messageQueue.join('\n\n'); }, [messageQueue]); + // Atomically drain all queued messages (synchronous, safe from callbacks). + const drainQueue = useCallback((): string[] => { + const drained = queueRef.current; + if (drained.length === 0) return []; + queueRef.current = []; + setMessageQueue([]); + return drained; + }, []); + // Process queued messages when streaming becomes idle useEffect(() => { if ( @@ -61,15 +81,22 @@ export function useMessageQueue({ // Combine all messages with double newlines for clarity const combinedMessage = messageQueue.join('\n\n'); // Clear the queue and submit - setMessageQueue([]); + clearQueue(); submitQuery(combinedMessage); } - }, [isConfigInitialized, streamingState, messageQueue, submitQuery]); + }, [ + isConfigInitialized, + streamingState, + messageQueue, + submitQuery, + clearQueue, + ]); return { messageQueue, addMessage, clearQueue, getQueuedMessagesText, + drainQueue, }; } From 56373257bc39e86f21d150d61d3f9e1e588592da Mon Sep 17 00:00:00 2001 From: wenshao Date: Fri, 3 Apr 2026 21:51:02 +0800 Subject: [PATCH 6/6] fix: guard mid-turn drain against cancelled turns - Skip drain when turnCancelledRef or abortController signal is set, so queued messages stay for the next turn instead of being lost - Restore ref-based queue bridge (drainQueue removed from useMessageQueue) - Keep synchronous ref clear to prevent duplicate drains Co-Authored-By: Claude Opus 4.6 (1M context) --- packages/cli/src/ui/AppContainer.tsx | 30 +++++++++++--------- packages/cli/src/ui/hooks/useGeminiStream.ts | 6 +++- 2 files changed, 22 insertions(+), 14 deletions(-) diff --git a/packages/cli/src/ui/AppContainer.tsx b/packages/cli/src/ui/AppContainer.tsx index 4604a9c378b..1465a176399 100644 --- a/packages/cli/src/ui/AppContainer.tsx +++ b/packages/cli/src/ui/AppContainer.tsx @@ -746,20 +746,24 @@ export const AppContainer = (props: AppContainerProps) => { disabled: agentViewState.activeView !== 'main', }); - const { - messageQueue, - addMessage, - clearQueue, - getQueuedMessagesText, - drainQueue, - } = useMessageQueue({ - isConfigInitialized, - streamingState, - submitQuery, - }); + const { messageQueue, addMessage, clearQueue, getQueuedMessagesText } = + useMessageQueue({ + isConfigInitialized, + streamingState, + submitQuery, + }); - // Bridge message queue to mid-turn drain ref for useGeminiStream. - midTurnDrainRef.current = drainQueue; + // Bridge message queue to mid-turn drain via ref. + // Sync ref on every render so the drain callback always reads latest state. + const messageQueueRef = useRef(messageQueue); + messageQueueRef.current = messageQueue; + midTurnDrainRef.current = () => { + const queue = messageQueueRef.current; + if (queue.length === 0) return []; + messageQueueRef.current = []; + clearQueue(); + return [...queue]; + }; // Callback for handling final submit (must be after addMessage from useMessageQueue) const handleFinalSubmit = useCallback( diff --git a/packages/cli/src/ui/hooks/useGeminiStream.ts b/packages/cli/src/ui/hooks/useGeminiStream.ts index ee6acb17719..979f84748e4 100644 --- a/packages/cli/src/ui/hooks/useGeminiStream.ts +++ b/packages/cli/src/ui/hooks/useGeminiStream.ts @@ -1564,7 +1564,11 @@ export const useGeminiStream = ( // Mid-turn queue drain: inject queued user messages alongside tool // results so the model sees them in the next API call. - const drained = midTurnDrainRef?.current?.() ?? []; + // Skip if the turn was cancelled — messages stay in queue for next turn. + const drained = + turnCancelledRef.current || abortControllerRef.current?.signal.aborted + ? [] + : (midTurnDrainRef?.current?.() ?? []); if (drained.length > 0) { for (const msg of drained) { responsesToSend.push({