From 3445ca1ecdb196de14a3073fec78df3bee93ee40 Mon Sep 17 00:00:00 2001 From: Toray Altas Date: Sun, 18 Jan 2026 19:03:12 -0500 Subject: [PATCH] fix: await subtask usage persistence before roll-up Ensure delegation roll-up reads final persisted cost/tokens by awaiting background usage collection before reopening parent. Add regression test covering late usage persistence ordering. --- .../history-resume-delegation.spec.ts | 116 +++++++++++++++++- src/core/task/Task.ts | 41 ++++++- src/core/tools/AttemptCompletionTool.ts | 3 + 3 files changed, 152 insertions(+), 8 deletions(-) diff --git a/src/__tests__/history-resume-delegation.spec.ts b/src/__tests__/history-resume-delegation.spec.ts index 93bff631f6..233e7ece15 100644 --- a/src/__tests__/history-resume-delegation.spec.ts +++ b/src/__tests__/history-resume-delegation.spec.ts @@ -3,6 +3,13 @@ import { describe, it, expect, vi, beforeEach } from "vitest" import { RooCodeEventName } from "@roo-code/types" +// Keep AttemptCompletionTool tests deterministic (TelemetryService can be undefined in unit test env) +vi.mock("@roo-code/telemetry", () => ({ + TelemetryService: { + instance: { captureTaskCompleted: vi.fn() }, + }, +})) + /* vscode mock for Task/Provider imports */ vi.mock("vscode", () => { const window = { @@ -374,13 +381,15 @@ describe("History resume delegation - parent metadata transitions", () => { }) // Verify both events emitted - const eventNames = emitSpy.mock.calls.map((c) => c[0]) + const eventNames = emitSpy.mock.calls.map((c: any[]) => c[0]) expect(eventNames).toContain(RooCodeEventName.TaskDelegationCompleted) expect(eventNames).toContain(RooCodeEventName.TaskDelegationResumed) // CRITICAL: verify ordering (TaskDelegationCompleted before TaskDelegationResumed) - const completedIdx = emitSpy.mock.calls.findIndex((c) => c[0] === RooCodeEventName.TaskDelegationCompleted) - const resumedIdx = emitSpy.mock.calls.findIndex((c) => c[0] === RooCodeEventName.TaskDelegationResumed) + const completedIdx = emitSpy.mock.calls.findIndex( + (c: any[]) => c[0] === RooCodeEventName.TaskDelegationCompleted, + ) + const resumedIdx = emitSpy.mock.calls.findIndex((c: any[]) => c[0] === RooCodeEventName.TaskDelegationResumed) expect(completedIdx).toBeGreaterThanOrEqual(0) expect(resumedIdx).toBeGreaterThan(completedIdx) }) @@ -424,7 +433,7 @@ describe("History resume delegation - parent metadata transitions", () => { }) // CRITICAL: verify legacy pause/unpause events NOT emitted - const eventNames = emitSpy.mock.calls.map((c) => c[0]) + const eventNames = emitSpy.mock.calls.map((c: any[]) => c[0]) expect(eventNames).not.toContain(RooCodeEventName.TaskPaused) expect(eventNames).not.toContain(RooCodeEventName.TaskUnpaused) expect(eventNames).not.toContain(RooCodeEventName.TaskSpawned) @@ -666,4 +675,103 @@ describe("History resume delegation - parent metadata transitions", () => { expect(orphanLinked).toBeUndefined() } }) + + it("subtask completion awaits late usage persistence before delegating (parent sees final cost)", async () => { + const parentTaskId = "p-late-cost" + const childTaskId = "c-late-cost" + + // Seed parent messages with a todo linked to the child + const todos = [{ id: "t1", content: "do subtask", status: "in_progress", subtaskId: childTaskId }] + const seededParentMessages = [ + { type: "say", say: "system_update_todos", text: JSON.stringify({ tool: "updateTodoList", todos }), ts: 1 }, + ] as any + vi.mocked(readTaskMessages).mockResolvedValue(seededParentMessages) + vi.mocked(readApiMessages).mockResolvedValue([] as any) + + // Parent history confirms relationship + const parentHistory = { + id: parentTaskId, + status: "delegated", + awaitingChildId: childTaskId, + childIds: [childTaskId], + ts: 1, + task: "Parent", + tokensIn: 0, + tokensOut: 0, + totalCost: 0, + } as any + + // Child history initially stale cost + let childHistory = { + id: childTaskId, + status: "active", + tokensIn: 10, + tokensOut: 5, + totalCost: 0, + ts: 2, + task: "Child", + } as any + + const reopenSpy = vi.fn().mockResolvedValue(undefined) + const provider: any = { + contextProxy: { globalStorageUri: { fsPath: "/storage" } }, + getTaskWithId: vi.fn(async (id: string) => { + if (id === parentTaskId) return { historyItem: parentHistory } + if (id === childTaskId) return { historyItem: childHistory } + throw new Error("unknown") + }), + reopenParentFromDelegation: reopenSpy, + } + + const childTask: any = { + parentTaskId, + taskId: childTaskId, + providerRef: { deref: () => provider }, + didToolFailInCurrentTurn: false, + todoList: undefined, + consecutiveMistakeCount: 0, + recordToolError: vi.fn(), + say: vi.fn(), + emitFinalTokenUsageUpdate: vi.fn(), + getTokenUsage: () => ({}) as any, + toolUsage: {}, + emit: vi.fn(), + ask: vi.fn().mockResolvedValue({ response: "yesButtonClicked" }), + waitForPendingUsageCollection: vi.fn(async () => { + childHistory = { ...childHistory, totalCost: 1.23 } + }), + } + + const { attemptCompletionTool } = await import("../core/tools/AttemptCompletionTool") + const askFinishSubTaskApproval = vi.fn().mockResolvedValue(true) + const handleError = vi.fn((_context: string, error: Error) => { + // Fail loudly if AttemptCompletionTool hits its catch block + throw error + }) + await attemptCompletionTool.execute({ result: "done" }, childTask, { + askApproval: vi.fn(), + handleError, + pushToolResult: vi.fn(), + removeClosingTag: vi.fn((_: any, s: any) => s), + askFinishSubTaskApproval, + toolDescription: vi.fn(), + toolProtocol: "native", + } as any) + + // Ensure we actually awaited and hit the delegation decision point + expect(childTask.waitForPendingUsageCollection).toHaveBeenCalled() + expect(provider.getTaskWithId).toHaveBeenCalledWith(childTaskId) + expect(askFinishSubTaskApproval).toHaveBeenCalled() + expect(reopenSpy).toHaveBeenCalledOnce() + + // Critical ordering: waitForPendingUsageCollection must run before delegation. + const waitCall = childTask.waitForPendingUsageCollection.mock.invocationCallOrder[0] + const reopenCall = reopenSpy.mock.invocationCallOrder[0] + expect(waitCall).toBeLessThan(reopenCall) + + // Parent roll-up reads the child's persisted history; this ensures cost was finalized + // before delegation begins (the bug fix). + expect(childHistory.totalCost).toBe(1.23) + expect(childHistory.tokensIn + childHistory.tokensOut).toBe(15) + }) }) diff --git a/src/core/task/Task.ts b/src/core/task/Task.ts index af7aed86d5..0c1d71b2a4 100644 --- a/src/core/task/Task.ts +++ b/src/core/task/Task.ts @@ -384,6 +384,15 @@ export class Task extends EventEmitter implements TaskLike { presentAssistantMessageHasPendingUpdates = false userMessageContent: (Anthropic.TextBlockParam | Anthropic.ImageBlockParam | Anthropic.ToolResultBlockParam)[] = [] userMessageContentReady = false + /** + * When an LLM stream ends, some providers may emit usage/cost information in trailing chunks. + * We drain the iterator in the background to capture those, which can complete *after* a + * subtask calls attempt_completion. + * + * This promise tracks the currently-running background usage collection so callers (notably + * delegation/roll-up logic) can await final persisted cost/tokens before reading history. + */ + private pendingUsageCollectionPromise?: Promise /** * Push a tool_result block to userMessageContent, preventing duplicates. @@ -1224,6 +1233,23 @@ export class Task extends EventEmitter implements TaskLike { } } + /** + * Best-effort wait for any in-flight background usage collection to finish. + * + * This is critical when completing subtasks: the parent roll-up reads cost/tokens + * from the child's persisted history item, which is updated by `saveClineMessages()`. + */ + public async waitForPendingUsageCollection(timeoutMs: number = DEFAULT_USAGE_COLLECTION_TIMEOUT_MS): Promise { + const pending = this.pendingUsageCollectionPromise + if (!pending) return + + try { + await Promise.race([pending, delay(timeoutMs)]) + } catch { + // Non-fatal: background usage collection already logs errors. + } + } + private findMessageByTimestamp(ts: number): ClineMessage | undefined { for (let i = this.clineMessages.length - 1; i >= 0; i--) { if (this.clineMessages[i].ts === ts) { @@ -3256,10 +3282,17 @@ export class Task extends EventEmitter implements TaskLike { } } - // Start the background task and handle any errors - drainStreamInBackgroundToFindAllUsage(lastApiReqIndex).catch((error) => { - console.error("Background usage collection failed:", error) - }) + // Start the background task and handle any errors. + // IMPORTANT: keep a reference so completion/delegation can await final persisted usage/cost. + this.pendingUsageCollectionPromise = drainStreamInBackgroundToFindAllUsage(lastApiReqIndex) + .catch((error) => { + console.error("Background usage collection failed:", error) + }) + .finally(() => { + if (this.pendingUsageCollectionPromise) { + this.pendingUsageCollectionPromise = undefined + } + }) } catch (error) { // Abandoned happens when extension is no longer waiting for the // Cline instance to finish aborting (error is thrown here when diff --git a/src/core/tools/AttemptCompletionTool.ts b/src/core/tools/AttemptCompletionTool.ts index 039036f829..cb3f8df77b 100644 --- a/src/core/tools/AttemptCompletionTool.ts +++ b/src/core/tools/AttemptCompletionTool.ts @@ -91,6 +91,9 @@ export class AttemptCompletionTool extends BaseTool<"attempt_completion"> { // This ensures the most recent stats are captured regardless of throttle timer // and properly updates the snapshot to prevent redundant emissions task.emitFinalTokenUsageUpdate() + // Ensure any trailing usage/cost emitted after stream completion is persisted + // before delegation/roll-up reads the child's history item. + await task.waitForPendingUsageCollection() TelemetryService.instance.captureTaskCompleted(task.taskId) task.emit(RooCodeEventName.TaskCompleted, task.taskId, task.getTokenUsage(), task.toolUsage)