From dde95bf12e51521134d0cc1d83a95520fa6899d2 Mon Sep 17 00:00:00 2001 From: ethan Date: Tue, 23 Jun 2026 13:03:22 +1000 Subject: [PATCH 1/3] Fix queued message foreground task backgrounding --- src/node/services/agentSession.ts | 10 +++-- .../services/agentSession.waitForIdle.test.ts | 23 ++++++++++ src/node/services/taskService.test.ts | 45 +++++++++++++++++++ src/node/services/taskService.ts | 5 +++ src/node/services/workspaceService.ts | 4 ++ 5 files changed, 83 insertions(+), 4 deletions(-) diff --git a/src/node/services/agentSession.ts b/src/node/services/agentSession.ts index ffd41e98f5..75e30671d4 100644 --- a/src/node/services/agentSession.ts +++ b/src/node/services/agentSession.ts @@ -3608,8 +3608,7 @@ export class AgentSession { workspaceGoalService: this.workspaceGoalService, experiments: options?.experiments, disableWorkspaceAgents: options?.disableWorkspaceAgents, - hasQueuedMessage: () => - !this.messageQueue.isEmpty() && this.messageQueue.getQueueDispatchMode() === "tool-end", + hasQueuedMessage: () => this.hasQueuedMessages("tool-end"), openaiTruncationModeOverride, }); @@ -5022,8 +5021,11 @@ export class AgentSession { this.backgroundProcessManager.setMessageQueued(this.workspaceId, false); } - hasQueuedMessages(): boolean { - return !this.messageQueue.isEmpty(); + hasQueuedMessages(dispatchMode?: "tool-end" | "turn-end"): boolean { + return ( + !this.messageQueue.isEmpty() && + (dispatchMode == null || this.messageQueue.getQueueDispatchMode() === dispatchMode) + ); } async waitForPendingStreamErrorRecoveryDecision(): Promise { diff --git a/src/node/services/agentSession.waitForIdle.test.ts b/src/node/services/agentSession.waitForIdle.test.ts index 3b5b99c11a..4e2676fff2 100644 --- a/src/node/services/agentSession.waitForIdle.test.ts +++ b/src/node/services/agentSession.waitForIdle.test.ts @@ -70,6 +70,29 @@ describe("AgentSession.waitForIdle", () => { } }); + test("tracks only tool-end queued messages for foreground wait backgrounding", async () => { + const { session, cleanup } = await createAgentSessionHarness({ + workspaceId: "wait-for-idle-tool-end-queue", + }); + + try { + session.queueMessage("later", { + model: "anthropic:claude-sonnet-4-5", + agentId: "exec", + queueDispatchMode: "turn-end", + }); + expect(session.hasQueuedMessages()).toBe(true); + expect(session.hasQueuedMessages("tool-end")).toBe(false); + + session.clearQueue(); + session.queueMessage("now"); + expect(session.hasQueuedMessages("tool-end")).toBe(true); + } finally { + session.dispose(); + await cleanup(); + } + }); + test("removes repeated aborted idle waiters during one busy turn", async () => { const { session, cleanup } = await createAgentSessionHarness({ workspaceId: "wait-for-idle-repeated-aborts", diff --git a/src/node/services/taskService.test.ts b/src/node/services/taskService.test.ts index 89806cc89d..15aa495220 100644 --- a/src/node/services/taskService.test.ts +++ b/src/node/services/taskService.test.ts @@ -363,6 +363,7 @@ function createWorkspaceServiceMocks( sendMessage: ReturnType; resumeStream: ReturnType; clearQueue: ReturnType; + hasQueuedMessages: ReturnType; hasPendingQueuedOrPreparingTurn: ReturnType; waitForPendingStreamErrorRecoveryDecision: ReturnType; remove: ReturnType; @@ -380,6 +381,7 @@ function createWorkspaceServiceMocks( sendMessage: ReturnType; resumeStream: ReturnType; clearQueue: ReturnType; + hasQueuedMessages: ReturnType; hasPendingQueuedOrPreparingTurn: ReturnType; waitForPendingStreamErrorRecoveryDecision: ReturnType; remove: ReturnType; @@ -398,6 +400,7 @@ function createWorkspaceServiceMocks( overrides?.resumeStream ?? mock((): Promise> => Promise.resolve(Ok({ started: true }))); const clearQueue = overrides?.clearQueue ?? mock((): Result => Ok(undefined)); + const hasQueuedMessages = overrides?.hasQueuedMessages ?? mock(() => false); const hasPendingQueuedOrPreparingTurn = overrides?.hasPendingQueuedOrPreparingTurn ?? mock(() => false); const waitForPendingStreamErrorRecoveryDecision = @@ -429,6 +432,7 @@ function createWorkspaceServiceMocks( sendMessage, resumeStream, clearQueue, + hasQueuedMessages, hasPendingQueuedOrPreparingTurn, waitForPendingStreamErrorRecoveryDecision, remove, @@ -444,6 +448,7 @@ function createWorkspaceServiceMocks( sendMessage, resumeStream, clearQueue, + hasQueuedMessages, hasPendingQueuedOrPreparingTurn, waitForPendingStreamErrorRecoveryDecision, remove, @@ -7449,6 +7454,46 @@ describe("TaskService", () => { expect(count2).toBe(0); }); + test("backgrounds waiters when tool-end message was already queued", async () => { + const config = await createTestConfig(rootDir); + + const parentId = "parent-ws"; + const childId = "child-task-ws"; + const projectPath = "/test/project"; + + await saveWorkspaces( + config, + projectPath, + [ + { path: `${projectPath}/parent`, id: parentId, name: "parent" }, + { + path: `${projectPath}/child`, + id: childId, + name: "agent_explore_child", + parentWorkspaceId: parentId, + agentType: "explore", + taskStatus: "running", + }, + ], + testTaskSettings(2, 3) + ); + + const hasQueuedMessages = mock(() => true); + const { workspaceService } = createWorkspaceServiceMocks({ hasQueuedMessages }); + const { taskService } = createTaskServiceHarness(config, { workspaceService }); + + const waitError = await taskService + .waitForAgentReport(childId, { + requestingWorkspaceId: parentId, + backgroundOnMessageQueued: true, + }) + .catch((error: unknown) => error); + + expect(waitError).toBeInstanceOf(ForegroundWaitBackgroundedError); + expect(hasQueuedMessages).toHaveBeenCalledWith(parentId, "tool-end"); + expect(taskService.backgroundForegroundWaitsForWorkspace(parentId)).toBe(0); + }); + test("defaults to queue-backgroundable when requestingWorkspaceId is present", async () => { const config = await createTestConfig(rootDir); diff --git a/src/node/services/taskService.ts b/src/node/services/taskService.ts index ce4d18ae52..ee8b0b8459 100644 --- a/src/node/services/taskService.ts +++ b/src/node/services/taskService.ts @@ -3769,6 +3769,11 @@ export class TaskService { this.backgroundableForegroundWaitersByWorkspaceId.set(workspaceId, set); } set.add(waiter); + // ponytail: level-trigger the edge-trigger from sendMessage, so a message + // queued before this waiter existed still backgrounds the foreground wait. + if (this.workspaceService.hasQueuedMessages(workspaceId, "tool-end")) { + this.backgroundForegroundWaitsForWorkspace(workspaceId); + } } private unregisterBackgroundableForegroundWaiter( diff --git a/src/node/services/workspaceService.ts b/src/node/services/workspaceService.ts index 441dad252f..8c918f87f4 100644 --- a/src/node/services/workspaceService.ts +++ b/src/node/services/workspaceService.ts @@ -7578,6 +7578,10 @@ export class WorkspaceService extends EventEmitter { } } + hasQueuedMessages(workspaceId: string, dispatchMode?: "tool-end" | "turn-end"): boolean { + return this.sessions.get(workspaceId.trim())?.hasQueuedMessages(dispatchMode) ?? false; + } + async waitForPendingStreamErrorRecoveryDecision(workspaceId: string): Promise { const session = this.sessions.get(workspaceId.trim()); await session?.waitForPendingStreamErrorRecoveryDecision(); From 1faac3c43f3ffcae3fcdc11ccc85fe2ab2500c86 Mon Sep 17 00:00:00 2001 From: ethan Date: Tue, 23 Jun 2026 13:07:21 +1000 Subject: [PATCH 2/3] Simplify queued message stream callback --- src/node/services/agentSession.ts | 2 +- src/node/services/aiService.ts | 6 +++--- src/node/services/streamManager.test.ts | 8 ++++---- src/node/services/streamManager.ts | 20 ++++++++++---------- 4 files changed, 18 insertions(+), 18 deletions(-) diff --git a/src/node/services/agentSession.ts b/src/node/services/agentSession.ts index 75e30671d4..c8c21ba31b 100644 --- a/src/node/services/agentSession.ts +++ b/src/node/services/agentSession.ts @@ -3608,7 +3608,7 @@ export class AgentSession { workspaceGoalService: this.workspaceGoalService, experiments: options?.experiments, disableWorkspaceAgents: options?.disableWorkspaceAgents, - hasQueuedMessage: () => this.hasQueuedMessages("tool-end"), + hasQueuedMessages: this.hasQueuedMessages.bind(this), openaiTruncationModeOverride, }); diff --git a/src/node/services/aiService.ts b/src/node/services/aiService.ts index 44084a1cf3..9e9a261ff0 100644 --- a/src/node/services/aiService.ts +++ b/src/node/services/aiService.ts @@ -250,7 +250,7 @@ export interface StreamMessageOptions { allowAgentSetGoal?: boolean; workspaceGoalService?: WorkspaceGoalService; disableWorkspaceAgents?: boolean; - hasQueuedMessage?: () => boolean; + hasQueuedMessages?: (dispatchMode?: "tool-end" | "turn-end") => boolean; muxMetadata?: MuxMessageMetadata; openaiTruncationModeOverride?: "auto" | "disabled"; } @@ -1003,7 +1003,7 @@ export class AIService extends EventEmitter { allowAgentSetGoal, workspaceGoalService, disableWorkspaceAgents, - hasQueuedMessage, + hasQueuedMessages, openaiTruncationModeOverride, muxMetadata, } = opts; @@ -2694,7 +2694,7 @@ export class AIService extends EventEmitter { maxOutputTokens, effectiveToolPolicy, streamToken, // Pass the pre-generated stream token - hasQueuedMessage, + hasQueuedMessages, metadata.name, effectiveThinkingLevel, requestHeaders, diff --git a/src/node/services/streamManager.test.ts b/src/node/services/streamManager.test.ts index c0ef40d381..38811cc35f 100644 --- a/src/node/services/streamManager.test.ts +++ b/src/node/services/streamManager.test.ts @@ -416,7 +416,7 @@ describe("StreamManager - cleanupStreamTempDir", () => { describe("StreamManager - stopWhen configuration", () => { type StopWhenCondition = (options: { steps: unknown[] }) => boolean; type BuildStopWhenCondition = (request: { - hasQueuedMessage?: () => boolean; + hasQueuedMessages?: (dispatchMode?: "tool-end" | "turn-end") => boolean; toolPolicy?: ToolPolicy; }) => StopWhenCondition[]; @@ -429,7 +429,7 @@ describe("StreamManager - stopWhen configuration", () => { function requiredToolConditionForTests(toolPolicy: ToolPolicy): StopWhenCondition { const [, , requiredToolCondition] = buildStopWhenForTests()({ - hasQueuedMessage: () => false, + hasQueuedMessages: () => false, toolPolicy, }); return requiredToolCondition; @@ -441,7 +441,7 @@ describe("StreamManager - stopWhen configuration", () => { test("returns step-cap and queued-message conditions with no policy", () => { let queued = false; - const stopWhen = buildStopWhenForTests()({ hasQueuedMessage: () => queued }); + const stopWhen = buildStopWhenForTests()({ hasQueuedMessages: () => queued }); expect(stopWhen).toHaveLength(3); const [maxStepCondition, queuedMessageCondition, requiredToolCondition] = stopWhen; @@ -666,7 +666,7 @@ describe("StreamManager - sequential tool execution", () => { headers?: Record; maxOutputTokens?: number; streamCallSettings?: Record; - hasQueuedMessage?: () => boolean; + hasQueuedMessages?: (dispatchMode?: "tool-end" | "turn-end") => boolean; toolPolicy?: ToolPolicy; toolChoice?: { type: "tool"; toolName: string }; } diff --git a/src/node/services/streamManager.ts b/src/node/services/streamManager.ts index 238e1d427b..9652aef657 100644 --- a/src/node/services/streamManager.ts +++ b/src/node/services/streamManager.ts @@ -156,7 +156,7 @@ interface StreamRequestConfig { headers?: Record; maxOutputTokens?: number; streamCallSettings?: Omit; - hasQueuedMessage?: () => boolean; + hasQueuedMessages?: (dispatchMode?: "tool-end" | "turn-end") => boolean; /** Optional hook for callers that need chunk-level visibility during streaming. */ onChunk?: StreamTextOnChunk; /** Optional hook for callers that need the live prepared step transcript. */ @@ -1392,7 +1392,7 @@ export class StreamManager extends EventEmitter { maxOutputTokens?: number, callSettingsOverrides?: ResolvedCallSettingsOverrides, toolPolicy?: ToolPolicy, - hasQueuedMessage?: () => boolean, + hasQueuedMessages?: (dispatchMode?: "tool-end" | "turn-end") => boolean, headers?: Record, anthropicCacheTtlOverride?: AnthropicCacheTtl, onChunk?: StreamTextOnChunk, @@ -1449,7 +1449,7 @@ export class StreamManager extends EventEmitter { maxOutputTokens: effectiveMaxOutputTokens, streamCallSettings: Object.keys(streamCallSettings).length > 0 ? streamCallSettings : undefined, - hasQueuedMessage, + hasQueuedMessages, onChunk, onStepMessages, toolPolicy, @@ -1457,7 +1457,7 @@ export class StreamManager extends EventEmitter { } private createStopWhenCondition( - request: Pick + request: Pick ): Array> { // Completion-tool stop check: completion/routing tools use explicit // success/ok markers (agent_report, propose_plan). @@ -1507,7 +1507,7 @@ export class StreamManager extends EventEmitter { return [ stepCountIs(100000), - () => request.hasQueuedMessage?.() ?? false, + () => request.hasQueuedMessages?.("tool-end") ?? false, hasSuccessfulRequiredToolResult, ]; } @@ -1570,7 +1570,7 @@ export class StreamManager extends EventEmitter { maxOutputTokens?: number, toolPolicy?: ToolPolicy, callSettingsOverrides?: ResolvedCallSettingsOverrides, - hasQueuedMessage?: () => boolean, + hasQueuedMessages?: (dispatchMode?: "tool-end" | "turn-end") => boolean, workspaceName?: string, thinkingLevel?: string, headers?: Record, @@ -1593,7 +1593,7 @@ export class StreamManager extends EventEmitter { maxOutputTokens, callSettingsOverrides, toolPolicy, - hasQueuedMessage, + hasQueuedMessages, headers, anthropicCacheTtlOverride, onChunk, @@ -2249,7 +2249,7 @@ export class StreamManager extends EventEmitter { fallbackState.original.maxOutputTokens, prepared.data.callSettingsOverrides, streamInfo.request.toolPolicy, - streamInfo.request.hasQueuedMessage, + streamInfo.request.hasQueuedMessages, prepared.data.headers, prepared.data.anthropicCacheTtl, streamInfo.request.onChunk, @@ -3638,7 +3638,7 @@ export class StreamManager extends EventEmitter { maxOutputTokens?: number, toolPolicy?: ToolPolicy, providedStreamToken?: StreamToken, - hasQueuedMessage?: () => boolean, + hasQueuedMessages?: (dispatchMode?: "tool-end" | "turn-end") => boolean, workspaceName?: string, thinkingLevel?: string, headers?: Record, @@ -3723,7 +3723,7 @@ export class StreamManager extends EventEmitter { maxOutputTokens, toolPolicy, callSettingsOverrides, - hasQueuedMessage, + hasQueuedMessages, workspaceName, thinkingLevel, headers, From 62d1ab092779c6e08c309e807d875098c3a821d1 Mon Sep 17 00:00:00 2001 From: ethan Date: Tue, 23 Jun 2026 13:12:27 +1000 Subject: [PATCH 3/3] Defer prequeued foreground wait backgrounding --- src/node/services/taskService.test.ts | 31 ++++++++++++++++++++++++++- src/node/services/taskService.ts | 20 ++++++++++++----- 2 files changed, 45 insertions(+), 6 deletions(-) diff --git a/src/node/services/taskService.test.ts b/src/node/services/taskService.test.ts index 15aa495220..76c768d114 100644 --- a/src/node/services/taskService.test.ts +++ b/src/node/services/taskService.test.ts @@ -526,6 +526,7 @@ describe("TaskService", () => { sendMessage?: ReturnType; remove?: ReturnType; isStreaming?: ReturnType; + hasQueuedMessages?: ReturnType; hasPendingQueuedOrPreparingTurn?: ReturnType; waitForPendingStreamErrorRecoveryDecision?: ReturnType; } = {} @@ -558,6 +559,9 @@ describe("TaskService", () => { create: createWorkspace, ...(options.sendMessage != null ? { sendMessage: options.sendMessage } : {}), ...(options.remove != null ? { remove: options.remove } : {}), + ...(options.hasQueuedMessages != null + ? { hasQueuedMessages: options.hasQueuedMessages } + : {}), ...(options.hasPendingQueuedOrPreparingTurn != null ? { hasPendingQueuedOrPreparingTurn: options.hasPendingQueuedOrPreparingTurn } : {}), @@ -1896,6 +1900,23 @@ describe("TaskService", () => { expect(taskService.backgroundForegroundWaitsForWorkspace(parentId)).toBe(0); }); + test("waitForWorkspaceTurn backgrounds when tool-end message was already queued", async () => { + const hasQueuedMessages = mock(() => true); + const { parentId, taskService } = await startWorkspaceTurnForTest({ hasQueuedMessages }); + + const waitError = await taskService + .waitForWorkspaceTurn("wst_handle", { + requestingWorkspaceId: parentId, + timeoutMs: 1_000, + backgroundOnMessageQueued: true, + }) + .catch((error: unknown) => error); + + expect(waitError).toBeInstanceOf(ForegroundWaitBackgroundedError); + expect(hasQueuedMessages).toHaveBeenCalledWith(parentId, "tool-end"); + expect(taskService.backgroundForegroundWaitsForWorkspace(parentId)).toBe(0); + }); + test("disposable workspace turns are removed after completion, error, or interruption", async () => { const completedRemove = mock((): Promise> => Promise.resolve(Ok(undefined))); const completed = await startWorkspaceTurnForTest({ @@ -7472,7 +7493,7 @@ describe("TaskService", () => { name: "agent_explore_child", parentWorkspaceId: parentId, agentType: "explore", - taskStatus: "running", + taskStatus: "queued", }, ], testTaskSettings(2, 3) @@ -7481,6 +7502,11 @@ describe("TaskService", () => { const hasQueuedMessages = mock(() => true); const { workspaceService } = createWorkspaceServiceMocks({ hasQueuedMessages }); const { taskService } = createTaskServiceHarness(config, { workspaceService }); + const internal = taskService as unknown as { + backgroundableForegroundWaitersByWorkspaceId: Map>; + pendingStartWaitersByTaskId: Map; + pendingWaitersByTaskId: Map; + }; const waitError = await taskService .waitForAgentReport(childId, { @@ -7492,6 +7518,9 @@ describe("TaskService", () => { expect(waitError).toBeInstanceOf(ForegroundWaitBackgroundedError); expect(hasQueuedMessages).toHaveBeenCalledWith(parentId, "tool-end"); expect(taskService.backgroundForegroundWaitsForWorkspace(parentId)).toBe(0); + expect(internal.backgroundableForegroundWaitersByWorkspaceId.has(parentId)).toBe(false); + expect(internal.pendingStartWaitersByTaskId.has(childId)).toBe(false); + expect(internal.pendingWaitersByTaskId.has(childId)).toBe(false); }); test("defaults to queue-backgroundable when requestingWorkspaceId is present", async () => { diff --git a/src/node/services/taskService.ts b/src/node/services/taskService.ts index ee8b0b8459..0df70952a9 100644 --- a/src/node/services/taskService.ts +++ b/src/node/services/taskService.ts @@ -3769,11 +3769,6 @@ export class TaskService { this.backgroundableForegroundWaitersByWorkspaceId.set(workspaceId, set); } set.add(waiter); - // ponytail: level-trigger the edge-trigger from sendMessage, so a message - // queued before this waiter existed still backgrounds the foreground wait. - if (this.workspaceService.hasQueuedMessages(workspaceId, "tool-end")) { - this.backgroundForegroundWaitsForWorkspace(workspaceId); - } } private unregisterBackgroundableForegroundWaiter( @@ -4035,6 +4030,13 @@ export class TaskService { timeoutMs ); + if ( + shouldBackgroundOnQueuedMessage && + this.workspaceService.hasQueuedMessages(options.requestingWorkspaceId, "tool-end") + ) { + this.backgroundForegroundWaitsForWorkspace(options.requestingWorkspaceId); + } + void (async () => { const record = await this.taskHandleStore.getWorkspaceTurn( options.requestingWorkspaceId, @@ -4380,6 +4382,14 @@ export class TaskService { }; options.abortSignal.addEventListener("abort", abortListener, { once: true }); } + + if ( + shouldBackgroundOnQueuedMessage && + requestingWorkspaceId && + this.workspaceService.hasQueuedMessages(requestingWorkspaceId, "tool-end") + ) { + this.backgroundForegroundWaitsForWorkspace(requestingWorkspaceId); + } })().catch((error: unknown) => { reject(error instanceof Error ? error : new Error(String(error))); });