From f081baebe31005d49e884ff02ccf206b2ef089fe Mon Sep 17 00:00:00 2001 From: yiliang114 Date: Thu, 23 Jul 2026 14:13:50 +0800 Subject: [PATCH] fix(cli): prevent monitor turns after task_stop --- docs/design/monitor-cancel-notification.md | 39 ++++++++ .../cli/src/nonInteractive/session.test.ts | 61 +++++++++++++ packages/cli/src/nonInteractive/session.ts | 13 +++ packages/cli/src/nonInteractiveCli.test.ts | 89 +++++++++++++++++++ packages/cli/src/nonInteractiveCli.ts | 18 +++- .../cli/src/ui/hooks/useGeminiStream.test.tsx | 47 +++++++++- packages/cli/src/ui/hooks/useGeminiStream.ts | 17 +++- packages/core/src/tools/task-stop.test.ts | 3 + packages/core/src/tools/task-stop.ts | 2 +- 9 files changed, 282 insertions(+), 7 deletions(-) create mode 100644 docs/design/monitor-cancel-notification.md diff --git a/docs/design/monitor-cancel-notification.md b/docs/design/monitor-cancel-notification.md new file mode 100644 index 00000000000..ce8933e9c15 --- /dev/null +++ b/docs/design/monitor-cancel-notification.md @@ -0,0 +1,39 @@ +# Explicit Monitor Cancellation Notifications + +## Problem + +`task_stop` already returns a synchronous tool result confirming that a monitor +was cancelled. The monitor registry also emits a terminal `cancelled` +notification, which clients record as a notification user message and submit as +a new model turn. A `running` event queued just before cancellation can cause the +same extra turn even if the terminal notification is suppressed. + +## Design + +- Cancel monitors silently when cancellation comes from `task_stop`; the tool + result remains the user- and model-visible confirmation. +- Keep the registry's default cancellation behavior unchanged for other callers. +- At drain time, discard queued `running` monitor notifications whose registry + entry is now explicitly `cancelled`. This check applies to the interactive + queue, the persistent stream-json queue, and the one-shot headless queue. +- Continue delivering natural `completed` and `failed` notifications, along with + terminal notifications emitted by non-`task_stop` cancellation paths. + +ACP already rejects `running` monitor notifications, so silent explicit +cancellation is sufficient for that client. + +Owner-routed monitor notifications stay inside an agent's input queue rather +than the user's conversation. They are outside this session-notification fix; +in the common tool-call path, any queued event is delivered alongside the +already-required `task_stop` tool result instead of creating a session turn. + +## Verification + +- `task_stop` cancels and aborts a monitor without invoking its notification + callback. +- Each client drops a queued `running` event after the monitor is explicitly + cancelled. +- Existing terminal-notification tests continue to demonstrate that natural + completion and failure are delivered. +- A real model-driven `monitor` then `task_stop` run produces no follow-up + notification turn. diff --git a/packages/cli/src/nonInteractive/session.test.ts b/packages/cli/src/nonInteractive/session.test.ts index 77aaa26f65e..af286b71ef1 100644 --- a/packages/cli/src/nonInteractive/session.test.ts +++ b/packages/cli/src/nonInteractive/session.test.ts @@ -60,6 +60,7 @@ interface ConfigOverrides { let mockMonitorRegistry: { setNotificationCallback: ReturnType; setRegisterCallback: ReturnType; + get: ReturnType; abortAll: ReturnType; }; let mockBackgroundShellRegistry: { @@ -190,6 +191,7 @@ describe('runNonInteractiveStreamJson', () => { mockMonitorRegistry = { setNotificationCallback: vi.fn(), setRegisterCallback: vi.fn(), + get: vi.fn().mockReturnValue({ status: 'running' }), abortAll: vi.fn(), }; mockBackgroundShellRegistry = { @@ -838,6 +840,65 @@ describe('runNonInteractiveStreamJson', () => { ); }); + it('drops a queued running monitor event after cancellation', async () => { + const initRequest = createControlRequest('initialize'); + const userMessage = createUserMessage('Start then stop a monitor'); + let closeInput: (() => void) | undefined; + let monitorCallback: + | (( + displayText: string, + modelText: string, + meta: { + monitorId: string; + toolUseId?: string; + status: string; + }, + ) => void) + | undefined; + let monitorStatus = 'running'; + + mockMonitorRegistry.get.mockImplementation(() => ({ + status: monitorStatus, + })); + mockMonitorRegistry.setNotificationCallback.mockImplementation((cb) => { + monitorCallback = cb; + }); + runNonInteractiveMock.mockImplementationOnce(async () => { + monitorCallback?.( + 'Monitor "logs" event #1: ready', + 'running', + { + monitorId: 'mon_1', + toolUseId: 'tool_mon_1', + status: 'running', + }, + ); + monitorStatus = 'cancelled'; + }); + + mockInputReader.read = async function* () { + yield initRequest; + yield userMessage; + await new Promise((resolve) => { + closeInput = resolve; + }); + }; + + const sessionPromise = runNonInteractiveStreamJson(config, ''); + await vi.waitFor(() => { + expect(runNonInteractiveMock).toHaveBeenCalledTimes(1); + }); + closeInput?.(); + await sessionPromise; + + expect(runNonInteractiveMock).toHaveBeenCalledTimes(1); + expect(mockOutputAdapter.emitUserMessage).not.toHaveBeenCalled(); + expect(mockOutputAdapter.emitSystemMessage).not.toHaveBeenCalledWith( + 'task_notification', + expect.anything(), + ); + }); + it('stops accepting new monitor events before EOF drain', async () => { const initRequest = createControlRequest('initialize'); const userMessage = createUserMessage('Start a monitor'); diff --git a/packages/cli/src/nonInteractive/session.ts b/packages/cli/src/nonInteractive/session.ts index f2fafd326d6..51769b83b60 100644 --- a/packages/cli/src/nonInteractive/session.ts +++ b/packages/cli/src/nonInteractive/session.ts @@ -598,6 +598,19 @@ class Session { ): Promise { await this.waitForInitialization(); + batch = batch.filter((item) => { + if (item.sdkNotification.status !== 'running') { + return true; + } + return ( + this.config.getMonitorRegistry().get(item.sdkNotification.task_id) + ?.status !== 'cancelled' + ); + }); + if (batch.length === 0) { + return; + } + for (const item of batch) { this.outputAdapter.emitUserMessage([{ text: item.displayText }]); this.outputAdapter.emitSystemMessage( diff --git a/packages/cli/src/nonInteractiveCli.test.ts b/packages/cli/src/nonInteractiveCli.test.ts index 01b68e1d656..fc15cc551e9 100644 --- a/packages/cli/src/nonInteractiveCli.test.ts +++ b/packages/cli/src/nonInteractiveCli.test.ts @@ -3283,6 +3283,95 @@ describe('runNonInteractive', () => { expect(userEnvelopes).toHaveLength(0); }); + it('drops a queued running monitor event after cancellation', async () => { + (mockConfig.getOutputFormat as Mock).mockReturnValue( + OutputFormat.STREAM_JSON, + ); + (mockConfig.getIncludePartialMessages as Mock).mockReturnValue(false); + setupMetricsMock(); + + const writes: string[] = []; + processStdoutSpy.mockImplementation((chunk: string | Uint8Array) => { + if (typeof chunk === 'string') { + writes.push(chunk); + } else { + writes.push(Buffer.from(chunk).toString('utf8')); + } + return true; + }); + + const notificationXml = + '\n' + + 'mon_1\n' + + 'monitor\n' + + 'running\n' + + 'Monitor emitted event #1.\n' + + 'ready\n' + + ''; + let monitorStatus = 'running'; + mockMonitorRegistry.get.mockImplementation(() => ({ + status: monitorStatus, + })); + mockMonitorRegistry.setNotificationCallback.mockImplementation((cb) => { + if (!cb) return; + cb('Monitor "logs" event #1: ready', notificationXml, { + monitorId: 'mon_1', + toolUseId: 'tool_mon_1', + status: 'running', + eventCount: 1, + }); + monitorStatus = 'cancelled'; + }); + mockGeminiClient.sendMessageStream.mockReturnValueOnce( + createStreamFromEvents([ + { type: GeminiEventType.Content, value: 'Monitor stopped.' }, + { + type: GeminiEventType.Finished, + value: { + reason: undefined, + usageMetadata: { totalTokenCount: 2 }, + }, + }, + ]), + ); + + await runNonInteractive( + mockConfig, + mockSettings, + 'Start then stop a monitor', + 'prompt-monitor-cancel', + ); + + expect(mockGeminiClient.sendMessageStream).toHaveBeenCalledTimes(1); + const envelopes = writes + .join('') + .split('\n') + .filter((line) => line.trim().length > 0) + .map((line) => JSON.parse(line)); + expect( + envelopes.some( + (env) => + env.type === 'user' && + Array.isArray(env.message?.content) && + env.message.content.some( + (block: unknown) => + typeof block === 'object' && + block !== null && + 'text' in block && + block.text === 'Monitor "logs" event #1: ready', + ), + ), + ).toBe(false); + expect( + envelopes.some( + (env) => + env.type === 'system' && + env.subtype === 'task_notification' && + env.data?.task_id === 'mon_1', + ), + ).toBe(false); + }); + it('does not let late monitor output keep one-shot runs alive', async () => { (mockConfig.getOutputFormat as Mock).mockReturnValue( OutputFormat.STREAM_JSON, diff --git a/packages/cli/src/nonInteractiveCli.ts b/packages/cli/src/nonInteractiveCli.ts index 3dc9c320227..60b448bf61a 100644 --- a/packages/cli/src/nonInteractiveCli.ts +++ b/packages/cli/src/nonInteractiveCli.ts @@ -456,6 +456,7 @@ export async function runNonInteractive( displayText: string; modelText: string; sendMessageType: SendMessageType; + monitorId?: string; sdkNotification?: { task_id: string; tool_use_id?: string; @@ -469,6 +470,13 @@ export async function runNonInteractive( } const localQueue: LocalQueueItem[] = []; const sdkOnlyMonitorQueue: LocalQueueItem[] = []; + const isCancelledMonitorEvent = (item: LocalQueueItem) => + Boolean( + item.monitorId && + item.sdkNotification?.status === 'running' && + config.getMonitorRegistry().get(item.monitorId)?.status === + 'cancelled', + ); const emitNotificationToSdk = (item: LocalQueueItem) => { if (item.sendMessageType !== SendMessageType.Notification) return; adapter.emitUserMessage([{ text: item.displayText }]); @@ -478,7 +486,10 @@ export async function runNonInteractive( }; const flushQueuedNotificationsToSdk = (queue: LocalQueueItem[]) => { while (queue.length > 0) { - emitNotificationToSdk(queue.shift()!); + const item = queue.shift()!; + if (!isCancelledMonitorEvent(item)) { + emitNotificationToSdk(item); + } } }; let captureMonitorTurnsInLocalQueue = true; @@ -935,6 +946,7 @@ export async function runNonInteractive( displayText, modelText, sendMessageType: SendMessageType.Notification, + monitorId: meta.monitorId, sdkNotification: { task_id: meta.monitorId, tool_use_id: meta.toolUseId, @@ -1884,7 +1896,9 @@ export async function runNonInteractive( splitIdx++; } } - const batch = localQueue.splice(0, splitIdx); + const batch = localQueue + .splice(0, splitIdx) + .filter((item) => !isCancelledMonitorEvent(item)); if (batch.length === 0) return; diff --git a/packages/cli/src/ui/hooks/useGeminiStream.test.tsx b/packages/cli/src/ui/hooks/useGeminiStream.test.tsx index f43fc041040..860d368974a 100644 --- a/packages/cli/src/ui/hooks/useGeminiStream.test.tsx +++ b/packages/cli/src/ui/hooks/useGeminiStream.test.tsx @@ -203,6 +203,10 @@ describe('useGeminiStream', () => { let mockCancelAllToolCalls: Mock; let mockMarkToolsAsSubmitted: Mock; let mockBackgroundShellRegistry: { setNotificationCallback: Mock }; + let mockMonitorRegistry: { + setNotificationCallback: Mock; + get: Mock; + }; let handleAtCommandSpy: MockInstance; beforeEach(() => { @@ -235,6 +239,10 @@ describe('useGeminiStream', () => { mockBackgroundShellRegistry = { setNotificationCallback: vi.fn(), }; + mockMonitorRegistry = { + setNotificationCallback: vi.fn(), + get: vi.fn().mockReturnValue({ status: 'running' }), + }; mockConfig = { apiKey: 'test-api-key', @@ -288,9 +296,7 @@ describe('useGeminiStream', () => { setNotificationCallback: vi.fn(), })), getBackgroundShellRegistry: vi.fn(() => mockBackgroundShellRegistry), - getMonitorRegistry: vi.fn(() => ({ - setNotificationCallback: vi.fn(), - })), + getMonitorRegistry: vi.fn(() => mockMonitorRegistry), } as unknown as Config; mockOnDebugMessage = vi.fn(); mockHandleSlashCommand = vi.fn().mockResolvedValue(false); @@ -6972,6 +6978,41 @@ describe('useGeminiStream', () => { ).toBeUndefined(); }); + it('drops a queued running monitor event after cancellation', async () => { + let monitorStatus = 'running'; + mockMonitorRegistry.get.mockImplementation(() => ({ + status: monitorStatus, + })); + renderTestHook(); + + const callback = mockMonitorRegistry.setNotificationCallback.mock + .calls[0][0] as ( + displayText: string, + modelText: string, + meta: { + monitorId: string; + status: string; + }, + ) => void; + mockSendMessageStream.mockClear(); + mockAddItem.mockClear(); + + await act(async () => { + callback( + 'Monitor "logs" event #1: ready', + 'running', + { monitorId: 'mon_1', status: 'running' }, + ); + monitorStatus = 'cancelled'; + }); + + expect(mockSendMessageStream).not.toHaveBeenCalled(); + expect(mockAddItem).not.toHaveBeenCalledWith( + expect.objectContaining({ type: 'notification' }), + expect.any(Number), + ); + }); + // Regression for #7156: progress setState calls issued from inside a // background subagent's AsyncLocalStorage frame can batch with the // notification trigger into one React commit, so the drain effect diff --git a/packages/cli/src/ui/hooks/useGeminiStream.ts b/packages/cli/src/ui/hooks/useGeminiStream.ts index 7b823b376b4..dea00419a63 100644 --- a/packages/cli/src/ui/hooks/useGeminiStream.ts +++ b/packages/cli/src/ui/hooks/useGeminiStream.ts @@ -3763,6 +3763,7 @@ export const useGeminiStream = ( displayText: string; modelText: string; sendMessageType: SendMessageType; + monitor?: { id: string; status: string }; onDelivered?: () => void; onDeliveryFailed?: () => void; }> @@ -3922,6 +3923,7 @@ export const useGeminiStream = ( displayText, modelText, sendMessageType: SendMessageType.Notification, + monitor: { id: meta.monitorId, status: meta.status }, }); setNotificationTrigger((n) => n + 1); }); @@ -3953,6 +3955,19 @@ export const useGeminiStream = ( // triggered the commit. runOutsideAgentContext(() => { const queue = notificationQueueRef.current; + const monitorRegistry = config.getMonitorRegistry(); + for (let i = queue.length - 1; i >= 0; i--) { + const monitor = queue[i]!.monitor; + if ( + monitor?.status === 'running' && + monitorRegistry.get(monitor.id)?.status === 'cancelled' + ) { + queue.splice(i, 1); + } + } + if (queue.length === 0) { + return; + } const targetType = queue[0]!.sendMessageType; // Cron prompts must run as individual turns — each needs its own @@ -3997,7 +4012,7 @@ export const useGeminiStream = ( }); }); } - }, [streamingState, submitQuery, notificationTrigger, addItem]); + }, [streamingState, submitQuery, notificationTrigger, addItem, config]); // ─── Teammate message integration ───────────────────────── // Each entry carries the full nonce-tagged envelope (`modelText`, diff --git a/packages/core/src/tools/task-stop.test.ts b/packages/core/src/tools/task-stop.test.ts index 6a6caf6f744..08b0d3bdd4a 100644 --- a/packages/core/src/tools/task-stop.test.ts +++ b/packages/core/src/tools/task-stop.test.ts @@ -233,6 +233,8 @@ describe('TaskStopTool', () => { describe('monitor support', () => { it('cancels a running monitor', async () => { const ac = new AbortController(); + const notificationCallback = vi.fn(); + monitorRegistry.setNotificationCallback(notificationCallback); monitorRegistry.register({ monitorId: 'mon_123', command: 'tail -f app.log', @@ -259,6 +261,7 @@ describe('TaskStopTool', () => { expect(result.returnDisplay).toContain('watch app log'); expect(monitorRegistry.get('mon_123')!.status).toBe('cancelled'); expect(ac.signal.aborted).toBe(true); + expect(notificationCallback).not.toHaveBeenCalled(); }); it('returns NOT_RUNNING when the monitor already completed', async () => { diff --git a/packages/core/src/tools/task-stop.ts b/packages/core/src/tools/task-stop.ts index b51e4ebbaf1..2310328f514 100644 --- a/packages/core/src/tools/task-stop.ts +++ b/packages/core/src/tools/task-stop.ts @@ -122,7 +122,7 @@ class TaskStopInvocation extends BaseToolInvocation< if (monitorEntry.status !== 'running') { return notRunningError('monitor', taskId, monitorEntry.status); } - monitorRegistry.cancel(taskId); + monitorRegistry.cancel(taskId, { notify: false }); return { llmContent: // Unlike background shells (which settle asynchronously when the