From 5dd97a2b106b0d9a3887901fe3207e0782dc3a7c Mon Sep 17 00:00:00 2001 From: Roomote Date: Fri, 21 Aug 2026 03:21:59 +0000 Subject: [PATCH 1/2] fix(task): preserve history during resume hydration --- .../__tests__/taskMessages.spec.ts | 31 ++++--- src/core/task-persistence/index.ts | 7 +- src/core/task-persistence/taskMessages.ts | 59 +++++++++----- src/core/task/Task.ts | 37 +++++---- .../task/__tests__/Task.persistence.spec.ts | 81 +++++++++++++++++++ .../Task.resume-eviction-race.spec.ts | 9 ++- src/core/webview/ClineProvider.ts | 7 +- 7 files changed, 176 insertions(+), 55 deletions(-) diff --git a/src/core/task-persistence/__tests__/taskMessages.spec.ts b/src/core/task-persistence/__tests__/taskMessages.spec.ts index c6bc360c05..6956fe667d 100644 --- a/src/core/task-persistence/__tests__/taskMessages.spec.ts +++ b/src/core/task-persistence/__tests__/taskMessages.spec.ts @@ -68,7 +68,7 @@ describe("taskMessages.saveTaskMessages", () => { }) describe("taskMessages.readTaskMessages", () => { - it("returns empty array when file contains invalid JSON", async () => { + it("rejects invalid JSON without treating it as empty history", async () => { const taskId = "task-corrupt-json" // Manually create the task directory and write corrupted JSON const taskDir = path.join(tmpBaseDir, "tasks", taskId) @@ -76,26 +76,35 @@ describe("taskMessages.readTaskMessages", () => { const filePath = path.join(taskDir, "ui_messages.json") await fs.writeFile(filePath, "{not valid json!!!", "utf8") - const result = await readTaskMessages({ - taskId, - globalStoragePath: tmpBaseDir, + await expect(readTaskMessages({ taskId, globalStoragePath: tmpBaseDir })).rejects.toMatchObject({ + kind: "invalid", }) - - expect(result).toEqual([]) }) - it("returns [] when file contains valid JSON that is not an array", async () => { + it("rejects valid non-array JSON without treating it as empty history", async () => { const taskId = "task-non-array-json" const taskDir = path.join(tmpBaseDir, "tasks", taskId) await fs.mkdir(taskDir, { recursive: true }) const filePath = path.join(taskDir, "ui_messages.json") await fs.writeFile(filePath, JSON.stringify("hello"), "utf8") - const result = await readTaskMessages({ - taskId, - globalStoragePath: tmpBaseDir, + await expect(readTaskMessages({ taskId, globalStoragePath: tmpBaseDir })).rejects.toMatchObject({ + kind: "invalid", }) + }) + + it("distinguishes a missing history file from an empty history", async () => { + await expect(readTaskMessages({ taskId: "task-missing", globalStoragePath: tmpBaseDir })).rejects.toMatchObject( + { kind: "not_found" }, + ) + }) + + it("returns an explicitly persisted empty history", async () => { + const taskId = "task-empty" + const taskDir = path.join(tmpBaseDir, "tasks", taskId) + await fs.mkdir(taskDir, { recursive: true }) + await fs.writeFile(path.join(taskDir, "ui_messages.json"), "[]", "utf8") - expect(result).toEqual([]) + await expect(readTaskMessages({ taskId, globalStoragePath: tmpBaseDir })).resolves.toEqual([]) }) }) diff --git a/src/core/task-persistence/index.ts b/src/core/task-persistence/index.ts index edc4d860b5..4c7e7fbcaa 100644 --- a/src/core/task-persistence/index.ts +++ b/src/core/task-persistence/index.ts @@ -1,4 +1,9 @@ export { type ApiMessage, readApiMessages, saveApiMessages } from "./apiMessages" -export { readTaskMessages, saveTaskMessages } from "./taskMessages" +export { + readTaskMessages, + saveTaskMessages, + TaskMessagesReadError, + type TaskMessagesReadErrorKind, +} from "./taskMessages" export { taskMetadata } from "./taskMetadata" export { TaskHistoryStore, assertValidTransition } from "./TaskHistoryStore" diff --git a/src/core/task-persistence/taskMessages.ts b/src/core/task-persistence/taskMessages.ts index cee66432d9..900335ed04 100644 --- a/src/core/task-persistence/taskMessages.ts +++ b/src/core/task-persistence/taskMessages.ts @@ -4,11 +4,22 @@ import * as fs from "fs/promises" import type { ClineMessage } from "@roo-code/types" -import { fileExistsAtPath } from "../../utils/fs" - import { GlobalFileNames } from "../../shared/globalFileNames" import { getTaskDirectoryPath } from "../../utils/storage" +export type TaskMessagesReadErrorKind = "not_found" | "invalid" | "io_error" + +export class TaskMessagesReadError extends Error { + constructor( + public readonly kind: TaskMessagesReadErrorKind, + message: string, + public readonly originalError?: unknown, + ) { + super(message) + this.name = "TaskMessagesReadError" + } +} + export type ReadTaskMessagesOptions = { taskId: string globalStoragePath: string @@ -20,27 +31,33 @@ export async function readTaskMessages({ }: ReadTaskMessagesOptions): Promise { const taskDir = await getTaskDirectoryPath(globalStoragePath, taskId) const filePath = path.join(taskDir, GlobalFileNames.uiMessages) - const fileExists = await fileExistsAtPath(filePath) - - if (fileExists) { - try { - const parsedData = JSON.parse(await fs.readFile(filePath, "utf8")) - if (!Array.isArray(parsedData)) { - console.warn( - `[readTaskMessages] Parsed data is not an array (got ${typeof parsedData}), returning empty. TaskId: ${taskId}, Path: ${filePath}`, - ) - return [] - } - return parsedData - } catch (error) { - console.warn( - `[readTaskMessages] Failed to parse ${filePath} for task ${taskId}, returning empty: ${error instanceof Error ? error.message : String(error)}`, - ) - return [] - } + + let fileContent: string + try { + fileContent = await fs.readFile(filePath, "utf8") + } catch (error) { + const kind = + typeof error === "object" && error !== null && "code" in error && error.code === "ENOENT" + ? "not_found" + : "io_error" + throw new TaskMessagesReadError(kind, `Failed to read task messages for ${taskId} at ${filePath}`, error) + } + + let parsedData: unknown + try { + parsedData = JSON.parse(fileContent) + } catch (error) { + throw new TaskMessagesReadError("invalid", `Failed to parse task messages for ${taskId} at ${filePath}`, error) + } + + if (!Array.isArray(parsedData)) { + throw new TaskMessagesReadError( + "invalid", + `Task messages for ${taskId} at ${filePath} must be an array, got ${typeof parsedData}`, + ) } - return [] + return parsedData } export type SaveTaskMessagesOptions = { diff --git a/src/core/task/Task.ts b/src/core/task/Task.ts index 4be087394e..e8c63a954d 100644 --- a/src/core/task/Task.ts +++ b/src/core/task/Task.ts @@ -1065,14 +1065,18 @@ export class Task extends EventEmitter implements TaskLike { } public async overwriteClineMessages(newMessages: ClineMessage[]) { - this.clineMessages = newMessages - restoreTodoListForTask(this) + this.hydrateClineMessages(newMessages) await this.saveClineMessages() + } + + private hydrateClineMessages(messages: ClineMessage[]) { + this.clineMessages = messages + restoreTodoListForTask(this) - // When overwriting messages (e.g., during task resume), repopulate the cloud sync tracking Set + // When hydrating or overwriting messages, repopulate the cloud sync tracking Set // with timestamps from all non-partial messages to prevent re-syncing previously synced messages this.cloudSyncedMessageTimestamps.clear() - for (const msg of newMessages) { + for (const msg of messages) { if (msg.partial !== true) { this.cloudSyncedMessageTimestamps.add(msg.ts) } @@ -1988,7 +1992,11 @@ export class Task extends EventEmitter implements TaskLike { private async resumeTaskFromHistory() { try { - const modifiedClineMessages = await this.getSavedClineMessages() + const modifiedClineMessages = [...(await this.getSavedClineMessages())] + + if (this.abort || this.abandoned) { + return + } // Remove any resume messages that may have been added before. const lastRelevantMessageIndex = findLastIndex( @@ -2000,16 +2008,6 @@ export class Task extends EventEmitter implements TaskLike { modifiedClineMessages.splice(lastRelevantMessageIndex + 1) } - // Remove any trailing reasoning-only UI messages that were not part of the persisted API conversation - while (modifiedClineMessages.length > 0) { - const last = modifiedClineMessages[modifiedClineMessages.length - 1] - if (last.type === "say" && last.say === "reasoning") { - modifiedClineMessages.pop() - } else { - break - } - } - // Since we don't use `api_req_finished` anymore, we need to check if the // last `api_req_started` has a cost value, if it doesn't and no // cancellation reason to present, then we remove it since it indicates @@ -2028,8 +2026,9 @@ export class Task extends EventEmitter implements TaskLike { } } - await this.overwriteClineMessages(modifiedClineMessages) - this.clineMessages = await this.getSavedClineMessages() + // Avoid a standalone write during hydration. The resume ask will persist only + // after all history reads succeed and the task is still active. + this.hydrateClineMessages(modifiedClineMessages) // Now present the cline messages to the user and ask if they want to // resume (NOTE: we ran into a bug before where the @@ -2039,6 +2038,10 @@ export class Task extends EventEmitter implements TaskLike { // the task first. this.apiConversationHistory = await this.getSavedApiConversationHistory() + if (this.abort || this.abandoned) { + return + } + const lastClineMessage = this.clineMessages .slice() .reverse() diff --git a/src/core/task/__tests__/Task.persistence.spec.ts b/src/core/task/__tests__/Task.persistence.spec.ts index 19bd0c7f34..f9f71b3eac 100644 --- a/src/core/task/__tests__/Task.persistence.spec.ts +++ b/src/core/task/__tests__/Task.persistence.spec.ts @@ -584,6 +584,87 @@ describe("Task persistence", () => { }) }) + describe("resumeTaskFromHistory", () => { + it.each(["not_found", "invalid", "io_error"] as const)( + "does not persist when hydration fails with %s", + async (kind) => { + mockReadTaskMessages.mockRejectedValue(Object.assign(new Error(`history ${kind}`), { kind })) + + const task = new Task({ + provider: mockProvider, + apiConfiguration: mockApiConfig, + historyItem: { + id: `issue-1279-${kind}`, + number: 1, + ts: 1, + task: "Original task", + status: "completed", + tokensIn: 10, + tokensOut: 5, + totalCost: 0.001, + }, + initialStatus: "completed", + startTask: false, + }) + const askSpy = vi.spyOn(task, "ask") + + await expect(getTaskPersistenceAccess(task).resumeTaskFromHistory()).rejects.toThrow(`history ${kind}`) + await task.abortTask(true) + + expect(askSpy).not.toHaveBeenCalled() + expect(mockSaveTaskMessages).not.toHaveBeenCalled() + expect(mockProvider.updateTaskHistory).not.toHaveBeenCalled() + }, + ) + + it("preserves finalized trailing reasoning without rewriting history during hydration", async () => { + const messages = [ + { ts: 1, type: "say" as const, say: "text" as const, text: "Original task" }, + { ts: 2, type: "say" as const, say: "completion_result" as const, text: "Initial result" }, + { ts: 3, type: "ask" as const, ask: "resume_completed_task" as const }, + { ts: 4, type: "say" as const, say: "user_feedback" as const, text: "Continue investigating" }, + { + ts: 5, + type: "say" as const, + say: "reasoning" as const, + text: "Critical current conclusion", + partial: false, + }, + ] + mockReadTaskMessages.mockResolvedValue(messages) + mockReadApiMessages.mockResolvedValue([ + { role: "user", content: [{ type: "text", text: "Continue investigating" }] }, + ]) + + const task = new Task({ + provider: mockProvider, + apiConfiguration: mockApiConfig, + historyItem: { + id: "issue-1279-current", + number: 1, + ts: 5, + task: "Original task", + status: "completed", + tokensIn: 10, + tokensOut: 5, + totalCost: 0.001, + }, + initialStatus: "completed", + startTask: false, + }) + vi.spyOn(task, "ask").mockImplementation(async (type) => { + expect(type).toBe("resume_completed_task") + expect(task.clineMessages).toContainEqual( + expect.objectContaining({ text: "Critical current conclusion", partial: false }), + ) + throw new Error("stop after hydration") + }) + + await expect(getTaskPersistenceAccess(task).resumeTaskFromHistory()).rejects.toThrow("stop after hydration") + expect(mockSaveTaskMessages).not.toHaveBeenCalled() + }) + }) + // ── flushPendingToolResultsToHistory — save failure/success ─────────── describe("flushPendingToolResultsToHistory persistence", () => { diff --git a/src/core/task/__tests__/Task.resume-eviction-race.spec.ts b/src/core/task/__tests__/Task.resume-eviction-race.spec.ts index 8766f38d5b..334fd4c02e 100644 --- a/src/core/task/__tests__/Task.resume-eviction-race.spec.ts +++ b/src/core/task/__tests__/Task.resume-eviction-race.spec.ts @@ -185,9 +185,7 @@ describe("Task resume/eviction race (Work #1 (no message) regression)", () => { // Hold the disk read open so the task is aborted while clineMessages is // still empty — the same window a user hits by navigating away quickly. const readDeferred = createDeferred() - mockReadTaskMessages - .mockReturnValueOnce(readDeferred.promise) // first read: held open to simulate the race window - .mockResolvedValue([]) // second read (resumeTaskFromHistory:2023): post-abort, safe fallback + mockReadTaskMessages.mockReturnValueOnce(readDeferred.promise) const updateTaskHistory = vi.fn().mockResolvedValue([]) const mockProvider = makeMockProvider(updateTaskHistory) @@ -222,5 +220,10 @@ describe("Task resume/eviction race (Work #1 (no message) regression)", () => { { ts: historyItem.ts + 1, type: "say", say: "completion_result", text: "Done." }, ]) await runPromise + + // The abandoned hydration must not resume and persist after its read settles. + expect(mockSaveTaskMessages).not.toHaveBeenCalled() + expect(updateTaskHistory).not.toHaveBeenCalled() + expect(mockReadTaskMessages).toHaveBeenCalledTimes(1) }) }) diff --git a/src/core/webview/ClineProvider.ts b/src/core/webview/ClineProvider.ts index 021a2fba91..000acea7c8 100644 --- a/src/core/webview/ClineProvider.ts +++ b/src/core/webview/ClineProvider.ts @@ -3953,8 +3953,11 @@ export class ClineProvider taskId: parentTaskId, globalStoragePath, }) - } catch { - parentClineMessages = [] + } catch (error) { + this.log( + `[reopenParentFromDelegation] Failed to read messages for parent ${parentTaskId}: ${error instanceof Error ? error.message : String(error)}`, + ) + return false } let parentApiMessages: any[] = [] From 3f5b2d04b585af22b0b2c9b5b261571a6c429209 Mon Sep 17 00:00:00 2001 From: Roomote Date: Sat, 22 Aug 2026 13:28:32 +0000 Subject: [PATCH 2/2] fix(task): merge concurrent history snapshots --- .../history-resume-delegation.spec.ts | 36 ++++++ .../__tests__/apiMessages.spec.ts | 35 +++++- .../__tests__/mergeMessageSnapshots.spec.ts | 106 ++++++++++++++++++ .../__tests__/taskMessages.spec.ts | 64 ++++++++++- src/core/task-persistence/apiMessages.ts | 5 +- .../task-persistence/mergeMessageSnapshots.ts | 81 +++++++++++++ src/core/task-persistence/taskMessages.ts | 41 +++++-- src/core/task/Task.ts | 10 +- .../task/__tests__/Task.persistence.spec.ts | 59 ++++++++++ src/core/webview/ClineProvider.ts | 14 ++- 10 files changed, 435 insertions(+), 16 deletions(-) create mode 100644 src/core/task-persistence/__tests__/mergeMessageSnapshots.spec.ts create mode 100644 src/core/task-persistence/mergeMessageSnapshots.ts diff --git a/src/__tests__/history-resume-delegation.spec.ts b/src/__tests__/history-resume-delegation.spec.ts index 24ea04dd8e..9c08031659 100644 --- a/src/__tests__/history-resume-delegation.spec.ts +++ b/src/__tests__/history-resume-delegation.spec.ts @@ -231,6 +231,7 @@ describe("History resume delegation - parent metadata transitions", () => { ]), taskId: "p1", globalStoragePath: "/storage", + merge: true, }), ) @@ -250,6 +251,7 @@ describe("History resume delegation - parent metadata transitions", () => { ]), taskId: "p1", globalStoragePath: "/storage", + merge: true, }), ) @@ -261,6 +263,40 @@ describe("History resume delegation - parent metadata transitions", () => { expect(apiCall.messages).toHaveLength(2) // 1 original + 1 injected }) + it("does not reopen or overwrite a parent when its UI history cannot be read", async () => { + const parentItem = { + id: "parent-read-failure", + status: "delegated", + awaitingChildId: "child-read-failure", + childIds: ["child-read-failure"], + ts: 100, + task: "Parent", + tokensIn: 0, + tokensOut: 0, + totalCost: 0, + } + const log = vi.fn() + const provider = makeProviderStub({ + contextProxy: { globalStorageUri: { fsPath: "/storage" } }, + getTaskWithId: vi.fn().mockResolvedValue({ historyItem: parentItem }), + getCurrentTask: vi.fn(() => ({ taskId: "child-read-failure" })), + taskHistoryStore: makeTaskHistoryStoreStub({ id: "child-read-failure", status: "active" }, parentItem), + log, + }) + vi.mocked(readTaskMessages).mockRejectedValue(new Error("history unavailable")) + + const result = await ClineProvider.prototype.reopenParentFromDelegation.call(provider, { + parentTaskId: "parent-read-failure", + childTaskId: "child-read-failure", + completionResultSummary: "Child done", + }) + + expect(result).toBe(false) + expect(log).toHaveBeenCalledWith(expect.stringContaining("history unavailable")) + expect(saveTaskMessages).not.toHaveBeenCalled() + expect(saveApiMessages).not.toHaveBeenCalled() + }) + it("reopenParentFromDelegation injects tool_result when new_task tool_use exists in API history", async () => { const parentItem = { id: "p-tool", diff --git a/src/core/task-persistence/__tests__/apiMessages.spec.ts b/src/core/task-persistence/__tests__/apiMessages.spec.ts index aa725f4744..b154cafaae 100644 --- a/src/core/task-persistence/__tests__/apiMessages.spec.ts +++ b/src/core/task-persistence/__tests__/apiMessages.spec.ts @@ -4,7 +4,7 @@ import * as os from "os" import * as path from "path" import * as fs from "fs/promises" -import { readApiMessages } from "../apiMessages" +import { readApiMessages, saveApiMessages } from "../apiMessages" let tmpBaseDir: string @@ -84,3 +84,36 @@ describe("apiMessages.readApiMessages", () => { expect(result).toEqual([]) }) }) + +describe("apiMessages.saveApiMessages", () => { + it("merges a concurrent disk suffix when requested", async () => { + const taskId = "task-merge-api" + const taskDir = path.join(tmpBaseDir, "tasks", taskId) + await fs.mkdir(taskDir, { recursive: true }) + const filePath = path.join(taskDir, "api_conversation_history.json") + await fs.writeFile( + filePath, + JSON.stringify([ + { role: "user", content: "disk prefix", ts: 1 }, + { role: "assistant", content: "disk suffix", ts: 3 }, + ]), + "utf8", + ) + + await saveApiMessages({ + taskId, + globalStoragePath: tmpBaseDir, + merge: true, + messages: [ + { role: "user", content: "updated prefix", ts: 1 }, + { role: "assistant", content: "incoming", ts: 2 }, + ], + }) + + expect(JSON.parse(await fs.readFile(filePath, "utf8"))).toEqual([ + expect.objectContaining({ content: "updated prefix", ts: 1 }), + expect.objectContaining({ content: "incoming", ts: 2 }), + expect.objectContaining({ content: "disk suffix", ts: 3 }), + ]) + }) +}) diff --git a/src/core/task-persistence/__tests__/mergeMessageSnapshots.spec.ts b/src/core/task-persistence/__tests__/mergeMessageSnapshots.spec.ts new file mode 100644 index 0000000000..a9b1a31093 --- /dev/null +++ b/src/core/task-persistence/__tests__/mergeMessageSnapshots.spec.ts @@ -0,0 +1,106 @@ +import { mergeApiMessageSnapshots, mergeClineMessageSnapshots } from "../mergeMessageSnapshots" + +describe("mergeClineMessageSnapshots", () => { + it("preserves disk-only messages and applies incoming updates in timestamp order", () => { + const result = mergeClineMessageSnapshots( + [ + { ts: 1, type: "say", say: "text", text: "old" }, + { ts: 3, type: "say", say: "text", text: "newer disk suffix" }, + ], + [ + { ts: 1, type: "say", say: "text", text: "updated" }, + { ts: 2, type: "say", say: "text", text: "incoming" }, + ], + ) + + expect(result).toEqual([ + expect.objectContaining({ ts: 1, text: "updated" }), + expect.objectContaining({ ts: 2, text: "incoming" }), + expect.objectContaining({ ts: 3, text: "newer disk suffix" }), + ]) + }) + + it("does not regress completed or answered message state", () => { + const result = mergeClineMessageSnapshots( + [{ ts: 1, type: "ask", ask: "tool", partial: false, isAnswered: true }], + [{ ts: 1, type: "ask", ask: "tool", partial: true, isAnswered: false }], + ) + + expect(result).toEqual([expect.objectContaining({ ts: 1, partial: false, isAnswered: true })]) + }) + + it("uses the incoming message when a timestamp is reused for a different message identity", () => { + expect( + mergeClineMessageSnapshots( + [{ ts: 1, type: "say", say: "text", text: "old" }], + [{ ts: 1, type: "ask", ask: "followup", text: "new" }], + ), + ).toEqual([{ ts: 1, type: "ask", ask: "followup", text: "new" }]) + }) + + it("returns incoming data when either snapshot is not an array", () => { + expect(mergeClineMessageSnapshots(null, [{ ts: 1 }])).toEqual([{ ts: 1 }]) + expect(mergeClineMessageSnapshots([], "invalid")).toBe("invalid") + }) +}) + +describe("mergeApiMessageSnapshots", () => { + it("retains equal-timestamp records and keeps tool calls before their results", () => { + const result = mergeApiMessageSnapshots( + [ + { role: "assistant", content: "old", ts: 1 }, + { role: "user", content: "same timestamp sibling", ts: 1 }, + { role: "user", content: [{ type: "tool_result", tool_use_id: "call-1", content: "ok" }], ts: 3 }, + ], + [ + { role: "assistant", content: "updated", ts: 1 }, + { + role: "assistant", + content: [{ type: "tool_use", id: "call-1", name: "read_file", input: {} }], + ts: 2, + }, + ], + ) + + expect(result).toEqual([ + expect.objectContaining({ role: "assistant", content: "updated", ts: 1 }), + expect.objectContaining({ role: "user", content: "same timestamp sibling", ts: 1 }), + expect.objectContaining({ role: "assistant", ts: 2 }), + expect.objectContaining({ role: "user", ts: 3 }), + ]) + }) + + it("preserves only the unmatched legacy disk tail", () => { + const result = mergeApiMessageSnapshots( + [ + { role: "user", content: "old prefix" }, + { role: "assistant", content: "disk tail" }, + ], + [{ role: "user", content: "updated prefix" }], + ) + + expect(result).toEqual([ + { role: "user", content: "updated prefix" }, + { role: "assistant", content: "disk tail" }, + ]) + }) + + it("keeps legacy prefixes ahead of newer timestamped messages", () => { + const result = mergeApiMessageSnapshots( + [ + { role: "user", content: "legacy prefix" }, + { role: "assistant", content: "disk suffix", ts: 3 }, + ], + [ + { role: "user", content: "updated legacy prefix" }, + { role: "assistant", content: "incoming", ts: 2 }, + ], + ) + + expect(result).toEqual([ + { role: "user", content: "updated legacy prefix" }, + expect.objectContaining({ content: "incoming", ts: 2 }), + expect.objectContaining({ content: "disk suffix", ts: 3 }), + ]) + }) +}) diff --git a/src/core/task-persistence/__tests__/taskMessages.spec.ts b/src/core/task-persistence/__tests__/taskMessages.spec.ts index 6956fe667d..61cecdb353 100644 --- a/src/core/task-persistence/__tests__/taskMessages.spec.ts +++ b/src/core/task-persistence/__tests__/taskMessages.spec.ts @@ -1,11 +1,18 @@ -import { describe, it, expect, vi, beforeEach } from "vitest" +import { describe, it, expect, vi, beforeEach, afterEach } from "vitest" import * as os from "os" import * as path from "path" import * as fs from "fs/promises" +import type { ClineMessage } from "@roo-code/types" + // Mocks (use hoisted to avoid initialization ordering issues) const hoisted = vi.hoisted(() => ({ safeWriteJsonMock: vi.fn().mockResolvedValue(undefined), + readFileMock: vi.fn(), +})) +vi.mock("fs/promises", async (importOriginal) => ({ + ...(await importOriginal()), + readFile: hoisted.readFileMock, })) vi.mock("../../../utils/safeWriteJson", () => ({ safeWriteJson: hoisted.safeWriteJsonMock, @@ -18,10 +25,17 @@ let tmpBaseDir: string beforeEach(async () => { hoisted.safeWriteJsonMock.mockClear() + const actualFs = await vi.importActual("fs/promises") + hoisted.readFileMock.mockReset().mockImplementation(actualFs.readFile) // Create a unique, writable temp directory to act as globalStoragePath tmpBaseDir = await fs.mkdtemp(path.join(os.tmpdir(), "roo-test-")) }) +afterEach(() => { + vi.useRealTimers() + vi.restoreAllMocks() +}) + describe("taskMessages.saveTaskMessages", () => { beforeEach(() => { hoisted.safeWriteJsonMock.mockClear() @@ -48,6 +62,7 @@ describe("taskMessages.saveTaskMessages", () => { expect(hoisted.safeWriteJsonMock).toHaveBeenCalledTimes(1) const [, persisted] = hoisted.safeWriteJsonMock.mock.calls[0] expect(persisted).toEqual(messages) + expect(hoisted.safeWriteJsonMock.mock.calls[0][2]).toBeUndefined() }) it("persists messages without modification when no metadata", async () => { @@ -65,6 +80,23 @@ describe("taskMessages.saveTaskMessages", () => { const [, persisted] = hoisted.safeWriteJsonMock.mock.calls[0] expect(persisted).toEqual(messages) }) + + it("passes the history merge callback only when requested", async () => { + const messages: ClineMessage[] = [{ ts: 2, type: "say", say: "text", text: "incoming" }] + await saveTaskMessages({ + messages, + taskId: "task-merge", + globalStoragePath: tmpBaseDir, + merge: true, + }) + + const merge = hoisted.safeWriteJsonMock.mock.calls[0][2]?.merge + expect(merge).toBeTypeOf("function") + expect(merge([{ ts: 1, type: "say", say: "text", text: "disk" }], messages)).toEqual([ + expect.objectContaining({ ts: 1, text: "disk" }), + expect.objectContaining({ ts: 2, text: "incoming" }), + ]) + }) }) describe("taskMessages.readTaskMessages", () => { @@ -107,4 +139,34 @@ describe("taskMessages.readTaskMessages", () => { await expect(readTaskMessages({ taskId, globalStoragePath: tmpBaseDir })).resolves.toEqual([]) }) + + it("retries one transient missing-file read after a jittered delay", async () => { + vi.spyOn(Math, "random").mockReturnValue(0) + const missing = Object.assign(new Error("missing"), { code: "ENOENT" }) + hoisted.readFileMock.mockRejectedValueOnce(missing).mockResolvedValueOnce("[]") + + await expect(readTaskMessages({ taskId: "task-retry", globalStoragePath: tmpBaseDir })).resolves.toEqual([]) + expect(hoisted.readFileMock).toHaveBeenCalledTimes(2) + }) + + it("throws when the missing-file retry also fails", async () => { + vi.spyOn(Math, "random").mockReturnValue(0) + const missing = Object.assign(new Error("missing"), { code: "ENOENT" }) + hoisted.readFileMock.mockRejectedValue(missing) + + await expect( + readTaskMessages({ taskId: "task-still-missing", globalStoragePath: tmpBaseDir }), + ).rejects.toMatchObject({ kind: "not_found" }) + expect(hoisted.readFileMock).toHaveBeenCalledTimes(2) + }) + + it("does not retry non-ENOENT read failures", async () => { + const denied = Object.assign(new Error("denied"), { code: "EACCES" }) + hoisted.readFileMock.mockRejectedValueOnce(denied) + + await expect(readTaskMessages({ taskId: "task-denied", globalStoragePath: tmpBaseDir })).rejects.toMatchObject({ + kind: "io_error", + }) + expect(hoisted.readFileMock).toHaveBeenCalledTimes(1) + }) }) diff --git a/src/core/task-persistence/apiMessages.ts b/src/core/task-persistence/apiMessages.ts index 7672f6f7ee..3fdcd376a5 100644 --- a/src/core/task-persistence/apiMessages.ts +++ b/src/core/task-persistence/apiMessages.ts @@ -8,6 +8,7 @@ import { fileExistsAtPath } from "../../utils/fs" import { GlobalFileNames } from "../../shared/globalFileNames" import { getTaskDirectoryPath } from "../../utils/storage" +import { mergeApiMessageSnapshots } from "./mergeMessageSnapshots" export type ApiMessage = Anthropic.MessageParam & { ts?: number @@ -110,12 +111,14 @@ export async function saveApiMessages({ messages, taskId, globalStoragePath, + merge = false, }: { messages: ApiMessage[] taskId: string globalStoragePath: string + merge?: boolean }) { const taskDir = await getTaskDirectoryPath(globalStoragePath, taskId) const filePath = path.join(taskDir, GlobalFileNames.apiConversationHistory) - await safeWriteJson(filePath, messages) + await safeWriteJson(filePath, messages, merge ? { merge: mergeApiMessageSnapshots } : undefined) } diff --git a/src/core/task-persistence/mergeMessageSnapshots.ts b/src/core/task-persistence/mergeMessageSnapshots.ts new file mode 100644 index 0000000000..1b158fecff --- /dev/null +++ b/src/core/task-persistence/mergeMessageSnapshots.ts @@ -0,0 +1,81 @@ +type MessageRecord = Record & { ts?: unknown } + +function isRecord(value: unknown): value is MessageRecord { + return typeof value === "object" && value !== null +} + +function mergeTimestampedSnapshots( + existing: unknown, + incoming: unknown, + mergeMatch: (disk: MessageRecord, next: MessageRecord) => MessageRecord, +): unknown { + if (!Array.isArray(existing) || !Array.isArray(incoming)) { + return incoming + } + + const existingGroups = new Map() + const existingLegacy: unknown[] = [] + + for (const message of existing) { + if (isRecord(message) && typeof message.ts === "number") { + const group = existingGroups.get(message.ts) ?? [] + group.push(message) + existingGroups.set(message.ts, group) + } else { + existingLegacy.push(message) + } + } + + const consumedByTimestamp = new Map() + let incomingLegacyCount = 0 + const merged = incoming.map((message) => { + if (!isRecord(message) || typeof message.ts !== "number") { + incomingLegacyCount++ + return message + } + + const consumed = consumedByTimestamp.get(message.ts) ?? 0 + consumedByTimestamp.set(message.ts, consumed + 1) + const diskMessage = existingGroups.get(message.ts)?.[consumed] + return diskMessage ? mergeMatch(diskMessage, message) : message + }) + + const diskOnlyTimestamped = [...existingGroups.entries()].flatMap(([timestamp, messages]) => + messages.slice(consumedByTimestamp.get(timestamp) ?? 0), + ) + + for (const diskMessage of diskOnlyTimestamped) { + const insertionIndex = merged.findIndex( + (message) => isRecord(message) && typeof message.ts === "number" && message.ts > (diskMessage.ts as number), + ) + if (insertionIndex === -1) { + merged.push(diskMessage) + } else { + merged.splice(insertionIndex, 0, diskMessage) + } + } + + merged.push(...existingLegacy.slice(incomingLegacyCount)) + return merged +} + +export function mergeClineMessageSnapshots(existing: unknown, incoming: unknown): unknown { + return mergeTimestampedSnapshots(existing, incoming, (disk, next) => { + if (disk.type !== next.type || disk.say !== next.say || disk.ask !== next.ask) { + return next + } + + const merged = { ...disk, ...next } + if (disk.partial === false && next.partial === true) { + merged.partial = false + } + if (disk.isAnswered === true) { + merged.isAnswered = true + } + return merged + }) +} + +export function mergeApiMessageSnapshots(existing: unknown, incoming: unknown): unknown { + return mergeTimestampedSnapshots(existing, incoming, (_disk, next) => next) +} diff --git a/src/core/task-persistence/taskMessages.ts b/src/core/task-persistence/taskMessages.ts index 900335ed04..13b0ae728f 100644 --- a/src/core/task-persistence/taskMessages.ts +++ b/src/core/task-persistence/taskMessages.ts @@ -6,6 +6,7 @@ import type { ClineMessage } from "@roo-code/types" import { GlobalFileNames } from "../../shared/globalFileNames" import { getTaskDirectoryPath } from "../../utils/storage" +import { mergeClineMessageSnapshots } from "./mergeMessageSnapshots" export type TaskMessagesReadErrorKind = "not_found" | "invalid" | "io_error" @@ -25,6 +26,29 @@ export type ReadTaskMessagesOptions = { globalStoragePath: string } +const READ_RETRY_MIN_MS = 10 +const READ_RETRY_RANGE_MS = 291 + +function getErrorCode(error: unknown): string | undefined { + return typeof error === "object" && error !== null && "code" in error && typeof error.code === "string" + ? error.code + : undefined +} + +async function readFileWithMissingRetry(filePath: string): Promise { + try { + return await fs.readFile(filePath, "utf8") + } catch (error) { + if (getErrorCode(error) !== "ENOENT") { + throw error + } + + const retryDelay = READ_RETRY_MIN_MS + Math.floor(Math.random() * READ_RETRY_RANGE_MS) + await new Promise((resolve) => setTimeout(resolve, retryDelay)) + return fs.readFile(filePath, "utf8") + } +} + export async function readTaskMessages({ taskId, globalStoragePath, @@ -34,12 +58,9 @@ export async function readTaskMessages({ let fileContent: string try { - fileContent = await fs.readFile(filePath, "utf8") + fileContent = await readFileWithMissingRetry(filePath) } catch (error) { - const kind = - typeof error === "object" && error !== null && "code" in error && error.code === "ENOENT" - ? "not_found" - : "io_error" + const kind = getErrorCode(error) === "ENOENT" ? "not_found" : "io_error" throw new TaskMessagesReadError(kind, `Failed to read task messages for ${taskId} at ${filePath}`, error) } @@ -64,10 +85,16 @@ export type SaveTaskMessagesOptions = { messages: ClineMessage[] taskId: string globalStoragePath: string + merge?: boolean } -export async function saveTaskMessages({ messages, taskId, globalStoragePath }: SaveTaskMessagesOptions) { +export async function saveTaskMessages({ + messages, + taskId, + globalStoragePath, + merge = false, +}: SaveTaskMessagesOptions) { const taskDir = await getTaskDirectoryPath(globalStoragePath, taskId) const filePath = path.join(taskDir, GlobalFileNames.uiMessages) - await safeWriteJson(filePath, messages) + await safeWriteJson(filePath, messages, merge ? { merge: mergeClineMessageSnapshots } : undefined) } diff --git a/src/core/task/Task.ts b/src/core/task/Task.ts index e8c63a954d..a3b097d92c 100644 --- a/src/core/task/Task.ts +++ b/src/core/task/Task.ts @@ -905,7 +905,7 @@ export class Task extends EventEmitter implements TaskLike { async overwriteApiConversationHistory(newHistory: ApiMessage[]) { this.apiConversationHistory = newHistory - await this.saveApiConversationHistory() + await this.saveApiConversationHistory(false) } /** @@ -987,12 +987,13 @@ export class Task extends EventEmitter implements TaskLike { return saved } - private async saveApiConversationHistory(): Promise { + private async saveApiConversationHistory(merge = true): Promise { try { await saveApiMessages({ messages: structuredClone(this.apiConversationHistory), taskId: this.taskId, globalStoragePath: this.globalStoragePath, + merge, }) return true } catch (error) { @@ -1066,7 +1067,7 @@ export class Task extends EventEmitter implements TaskLike { public async overwriteClineMessages(newMessages: ClineMessage[]) { this.hydrateClineMessages(newMessages) - await this.saveClineMessages() + await this.saveClineMessages(false) } private hydrateClineMessages(messages: ClineMessage[]) { @@ -1102,12 +1103,13 @@ export class Task extends EventEmitter implements TaskLike { } } - private async saveClineMessages(): Promise { + private async saveClineMessages(merge = true): Promise { try { await saveTaskMessages({ messages: structuredClone(this.clineMessages), taskId: this.taskId, globalStoragePath: this.globalStoragePath, + merge, }) if (this._taskApiConfigName === undefined) { diff --git a/src/core/task/__tests__/Task.persistence.spec.ts b/src/core/task/__tests__/Task.persistence.spec.ts index f9f71b3eac..71a5f82e66 100644 --- a/src/core/task/__tests__/Task.persistence.spec.ts +++ b/src/core/task/__tests__/Task.persistence.spec.ts @@ -305,6 +305,20 @@ describe("Task persistence", () => { const result = await task.retrySaveApiConversationHistory() expect(result).toBe(true) + expect(mockSaveApiMessages).toHaveBeenCalledWith(expect.objectContaining({ merge: true })) + }) + + it("uses authoritative replacement for explicit API history overwrites", async () => { + const task = new Task({ + provider: mockProvider, + apiConfiguration: mockApiConfig, + task: "test task", + startTask: false, + }) + + await task.overwriteApiConversationHistory([{ role: "user", content: "replacement" }]) + + expect(mockSaveApiMessages).toHaveBeenCalledWith(expect.objectContaining({ merge: false })) }) it("returns false on failure", async () => { @@ -398,6 +412,20 @@ describe("Task persistence", () => { const result = await (task as Record).saveClineMessages() expect(result).toBe(true) + expect(mockSaveTaskMessages).toHaveBeenCalledWith(expect.objectContaining({ merge: true })) + }) + + it("uses authoritative replacement for explicit UI history overwrites", async () => { + const task = new Task({ + provider: mockProvider, + apiConfiguration: mockApiConfig, + task: "test task", + startTask: false, + }) + + await task.overwriteClineMessages([{ ts: 1, type: "say", say: "text", text: "replacement" }]) + + expect(mockSaveTaskMessages).toHaveBeenCalledWith(expect.objectContaining({ merge: false })) }) it("returns false on failure", async () => { @@ -663,6 +691,37 @@ describe("Task persistence", () => { await expect(getTaskPersistenceAccess(task).resumeTaskFromHistory()).rejects.toThrow("stop after hydration") expect(mockSaveTaskMessages).not.toHaveBeenCalled() }) + + it("stops after API history hydration when the task is aborted", async () => { + const apiMessagesDeferred = + createDeferred }>>() + mockReadTaskMessages.mockResolvedValue([{ ts: 1, type: "say", say: "text", text: "Original task" }]) + mockReadApiMessages.mockReturnValue(apiMessagesDeferred.promise) + + const task = new Task({ + provider: mockProvider, + apiConfiguration: mockApiConfig, + historyItem: { + id: "issue-1279-api-abort", + number: 1, + ts: 1, + task: "Original task", + tokensIn: 10, + tokensOut: 5, + totalCost: 0.001, + }, + startTask: false, + }) + const askSpy = vi.spyOn(task, "ask") + const resumePromise = getTaskPersistenceAccess(task).resumeTaskFromHistory() + await vi.waitFor(() => expect(mockReadApiMessages).toHaveBeenCalled()) + + await task.abortTask(true) + apiMessagesDeferred.resolve([{ role: "user", content: [{ type: "text", text: "Original task" }] }]) + await resumePromise + + expect(askSpy).not.toHaveBeenCalled() + }) }) // ── flushPendingToolResultsToHistory — save failure/success ─────────── diff --git a/src/core/webview/ClineProvider.ts b/src/core/webview/ClineProvider.ts index 000acea7c8..a792b441f0 100644 --- a/src/core/webview/ClineProvider.ts +++ b/src/core/webview/ClineProvider.ts @@ -3991,7 +3991,12 @@ export class ClineProvider ) { parentClineMessages.push(subtaskUiMessage) } - await saveTaskMessages({ messages: parentClineMessages, taskId: parentTaskId, globalStoragePath }) + await saveTaskMessages({ + messages: parentClineMessages, + taskId: parentTaskId, + globalStoragePath, + merge: true, + }) // Find the tool_use_id from the last assistant message's new_task tool_use let toolUseId: string | undefined @@ -4076,7 +4081,12 @@ export class ClineProvider } } - await saveApiMessages({ messages: parentApiMessages as any, taskId: parentTaskId, globalStoragePath }) + await saveApiMessages({ + messages: parentApiMessages as any, + taskId: parentTaskId, + globalStoragePath, + merge: true, + }) // 4) Close child instance if still open (single-open-task invariant). // This MUST happen BEFORE marking the child "completed" because