diff --git a/docs/architecture/task-lifecycle-model.md b/docs/architecture/task-lifecycle-model.md index 355b08aa7e..744ac8eb81 100644 --- a/docs/architecture/task-lifecycle-model.md +++ b/docs/architecture/task-lifecycle-model.md @@ -41,18 +41,21 @@ TLA+/PlusCal or Quint with TLC becomes a better fit when the lifecycle needs tem ## Production mapping -| Model concept | Production concept | -| ------------------------- | ------------------------------------------------------------------------------------ | -| Task record and status | `HistoryItem` persisted by `TaskHistoryStore` | -| `delegate(parent, child)` | `ClineProvider.delegateParentAndOpenChild` | -| `interrupt(child)` | cancellation or eviction through `markDelegatedChildInterrupted` | -| `complete(child)` | `ClineProvider.reopenParentFromDelegation` | -| `abandon(child)` | `ClineProvider.abandonSubtask` | -| Atomic event step | `atomicReadAndUpdate`, `atomicUpdatePair`, and per-parent delegation transition lock | -| Event interleaving | Competing completion, cancellation, abandonment, and new delegation calls | +| Model concept | Production concept | +| ------------------------- | -------------------------------------------------------------------------------------------------------------------- | +| Task record and status | `HistoryItem` persisted by `TaskHistoryStore` | +| `delegate(parent, child)` | `ClineProvider.delegateParentAndOpenChild` | +| `interrupt(child)` | cancellation or eviction through `markDelegatedChildInterrupted` | +| `complete(child)` | `ClineProvider.reopenParentFromDelegation` | +| `abandon(child)` | `ClineProvider.abandonSubtask` | +| Pending-action settlement | `TaskHistoryStore.clearPendingActionIfMatching` compare-and-clear in the rejected-delegation settlement path (#1714) | +| Atomic event step | `atomicReadAndUpdate`, `atomicUpdatePair`, and per-parent delegation transition lock | +| Event interleaving | Competing completion, cancellation, abandonment, and new delegation calls | The model has three fixed task slots, enough to cover competing siblings and a nested parent-child-grandchild chain. It explores every reachable interleaving through depth 12, deduplicating canonical states. Representative checks also exercise rejected operations that do not create a new state: a second concurrent delegation while the first child is active, stale completion after re-delegation, late completion after abandonment, completion after interruption, and nested completion. Named semantic landmarks require the graph to retain interrupted-child re-delegation and nested delegation even when the raw state total changes. +Each task slot can also hold one of two pending `create_subtask` actions. A `stage` action mirrors `setPendingTaskAction` overwrite semantics, delegation clears the action its request carried, completion clears the child's action only when its event carries the matching action ID, and a `settle-rejected` action models the settlement that follows an authoritative delegation rejection (#1714). Production settles through the typed `LifecycleTransitionError` from the shared guards: the provider calls the disk-authoritative `TaskHistoryStore.clearPendingActionIfMatching` compare-and-clear under the per-file lock, then propagates the original rejection. Six named witnesses must remain reachable: settlement from an interrupted record after rejection, settlement through a successful active delegation, unrelated-action preservation during completion, stale-action protection where a settlement targeting one action ID leaves a replacement action intact, matching-ID completion clearing, and replacement-ID completion preservation. A mismatched pending-action request keeps its production behavior: the atomic update throws before any transition, and no settlement runs. + Production completion also accepts a recovery-compatible `active` parent that still awaits the returning child, then clears the stale pointers. Normal model transitions never create that intermediate state, so it is covered by a focused reducer test rather than admitted as a generally valid reachable state. ## Shared-store concurrency model @@ -63,6 +66,7 @@ The same `pnpm lifecycle:model-check` command also runs a second bounded explore - store read/update operations hold the host mutex, while live-task snapshots used by completion and message saves may outlive it; - a write delta is computed relative to that host's cache; - revalidation under the per-file disk lock checks only status-transition legality; +- the rejected-delegation settlement compare-and-clear decides inside the disk merge callback, so a replacement action persisted by another host is never cleared; - fields absent from the delta preserve the current disk value, `childIds` are unioned, and other same-field conflicts are last-writer-wins; - `atomicUpdatePair` commits its files in order, with another host able to act between file commits; - successful pair-operation cache entries publish together after both file writes; if the second write fails, the cache publishes only the first committed record; @@ -79,7 +83,7 @@ CI fails if either exact causal witness or violation class changes, a witness di The known-unsafe witnesses currently compare exact shortest action sequences. This is intentionally simple and reviewable, but brittle to harmless action renames or serialization refactors. A causal partial-order comparator would reduce that brittleness but would add a second trace-equivalence protocol to maintain. Until that complexity is justified, update an exact witness only after confirming the terminal violation class and required causal ordering are unchanged. -`TaskHistoryStore.realConcurrency.spec.ts` complements the abstract interleavings with one synchronized integration smoke check through the real `proper-lockfile` and filesystem rename path; broader VS Code E2E remains reserved for restart and extension-host behavior. +`TaskHistoryStore.realConcurrency.spec.ts` complements the abstract interleavings with real-filesystem checks through the real `proper-lockfile` and filesystem rename path, including stale-settlement compare-and-clear. History-file deletion holds the same per-file advisory lock as `safeWriteJson` around each unlink, and the real-filesystem suite covers a deletion that targets the settlement read-to-commit window plus schema and task-ID validation of the locked disk record before settlement (#1726). That lock serialization is best-effort: a lock acquisition or unlink failure is swallowed, so a deletion can proceed without the lock and the two sides serialize only when both acquire it. Restart recovery is covered by `Task.persistence.spec.ts`: it refreshes the persisted record first, treats a refreshed interrupted `create_subtask` action as rejected, settles that exact action before replay, adopts the authoritative pending action, and a failed refresh, lookup, or settlement stops replay rather than creating another child. Broader VS Code E2E remains reserved for other restart and extension-host behavior. ## Task cleanup protocol model @@ -134,6 +138,8 @@ The task delegation checker currently enforces: 5. Parent-child lineage is acyclic. 6. Completed task records cannot be changed by later lifecycle events. 7. Active-child re-delegation, stale completion after ownership moves to another child, duplicate/late completion, and abandonment of a live child are rejected by the shared production guards. +8. A rejected delegation settles only the exact matching pending `create_subtask` action. Settlement preserves status, lineage, and accounting, never clears a replacement or different-kind action, and never mutates a completed record. The settlement compare-and-clear reads the persisted record under the per-file lock, so it never clears from a stale host cache. +9. A completion clears the completing child's pending action only when the completion event carries the exact matching action ID. A completion with no action ID or a different ID preserves the pending action. The completion persistence checker additionally enforces: @@ -147,15 +153,15 @@ These are safety claims within the documented bounds. The checks do not claim li ## Coverage audit -| Protocol area | Coverage status | Production/model relationship | Explicit limits and open points | -| ------------------------------ | ---------------------------------------------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------ | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | -| Delegation lifecycle | Production-backed bounded universal | The explorer calls the four production reducers for three task slots through depth 12. | Excludes provider instances, persistence failures, scheduler state, most live `Task` behavior, and generation identity for delayed pre-interruption completion. Recovery-compatible active-parent completion is test-only. | -| Shared-store concurrency | Production-backed bounded scenarios plus known-unsafe witnesses | The explorer imports production delta/merge functions and reducers; a real-filesystem test is a smoke check. | Does not prove crash safety, filesystem/lock semantics, arbitrary processes, or loss-free same-field merging. #1469 and #1021 remain unsafe. | -| Provider handoff and scheduler | Mixed: production-backed reducers/selector plus abstract bounded protocol | Commits use production reducers; provider ownership, publication, transition locks, and permits are model abstractions through depth 15. | Selector correctness does not refine all downstream readers. Scheduler tests cover concrete permit behavior separately. | -| Optional fan-out scope | Planned-only abstract bounded protocol outside baseline CI | The model has two sibling slots and two abstract permits and imports no production fan-out transition. | Excluded from baseline closure; production fan-out remains separately scoped future functionality. | -| Cleanup | Abstract bounded universal plus adapter tests | Abort, disposal, settlement, rejection, and provider shutdown are modeled as protocol/environment actions. | No direct execution of all production cleanup methods, filesystem/editor promises, timing liveness, fairness, or arbitrary task counts. | -| Parser request scope | Production-backed bounded schedule replay | The checker executes production parser APIs across 924 order-preserving schedules for two scopes. | Assumes callers stop invoking a finalized scope; transport behavior, arbitrary request counts, indices, and malformed histories are outside the claim. | -| Completion persistence | Abstract bounded universal plus production tests and one fresh-host E2E path | The model abstracts persistence as a durable phase with at most two write starts; production guards and retry paths are tested separately. | Production permits more retries; no power-loss/filesystem proof, fairness, arbitrary retry count, complete delegated fallback, provider status metadata, or downstream event-consumer model. | +| Protocol area | Coverage status | Production/model relationship | Explicit limits and open points | +| ------------------------------ | ---------------------------------------------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------ | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | +| Delegation lifecycle | Production-backed bounded universal | The explorer calls the production lifecycle and settlement reducers for three task slots through depth 12. | Excludes provider instances, persistence failures, scheduler state, most live `Task` behavior, and generation identity for delayed pre-interruption completion. Recovery-compatible active-parent completion is test-only. Settlement models the matching-action request only; a mismatched action aborts the update unmodeled. | +| Shared-store concurrency | Production-backed bounded scenarios plus known-unsafe witnesses | The explorer imports production delta/merge functions and reducers; a real-filesystem test is a smoke check. | Does not prove crash safety, filesystem/lock semantics, arbitrary processes, or loss-free same-field merging. #1469 and #1021 remain unsafe. | +| Provider handoff and scheduler | Mixed: production-backed reducers/selector plus abstract bounded protocol | Commits use production reducers; provider ownership, publication, transition locks, and permits are model abstractions through depth 15. | Selector correctness does not refine all downstream readers. Scheduler tests cover concrete permit behavior separately. | +| Optional fan-out scope | Planned-only abstract bounded protocol outside baseline CI | The model has two sibling slots and two abstract permits and imports no production fan-out transition. | Excluded from baseline closure; production fan-out remains separately scoped future functionality. | +| Cleanup | Abstract bounded universal plus adapter tests | Abort, disposal, settlement, rejection, and provider shutdown are modeled as protocol/environment actions. | No direct execution of all production cleanup methods, filesystem/editor promises, timing liveness, fairness, or arbitrary task counts. | +| Parser request scope | Production-backed bounded schedule replay | The checker executes production parser APIs across 924 order-preserving schedules for two scopes. | Assumes callers stop invoking a finalized scope; transport behavior, arbitrary request counts, indices, and malformed histories are outside the claim. | +| Completion persistence | Abstract bounded universal plus production tests and one fresh-host E2E path | The model abstracts persistence as a durable phase with at most two write starts; production guards and retry paths are tested separately. | Production permits more retries; no power-loss/filesystem proof, fairness, arbitrary retry count, complete delegated fallback, provider status metadata, or downstream event-consumer model. | The production mapping above names primary lifecycle transitions, not every mutation or consumer. Generic store upserts, reconciliation, repair replay, migrations, tool entry points, webview/public API abandonment, provider status updates, and public `TaskCompleted` re-emission remain outside the persisted reducer graph unless explicitly named by a submodel or focused test. For task-local mode, the consumers named under #921/#1623 are a confirmed set, not an exhaustive repository-wide inventory; other mode-sensitive tools must be audited before claiming universal reader isolation. @@ -197,6 +203,7 @@ The [Task lifecycle verification GAP report](./task-lifecycle-gap-report.md) is - Shared status vocabulary: [#612](https://github.com/Zoo-Code-Org/Zoo-Code/issues/612). - Completion visibility history: [#1453](https://github.com/Zoo-Code-Org/Zoo-Code/issues/1453) and [#1279](https://github.com/Zoo-Code-Org/Zoo-Code/issues/1279). - Cross-instance history preservation: [#920](https://github.com/Zoo-Code-Org/Zoo-Code/issues/920). +- Rejected-delegation pending-action settlement: [#1714](https://github.com/Zoo-Code-Org/Zoo-Code/issues/1714). - Parser request scoping: [#1468](https://github.com/Zoo-Code-Org/Zoo-Code/issues/1468). ## Extending the model diff --git a/scripts/check-task-lifecycle.ts b/scripts/check-task-lifecycle.ts index 73e9078366..bf05ad68e9 100644 --- a/scripts/check-task-lifecycle.ts +++ b/scripts/check-task-lifecycle.ts @@ -1,12 +1,13 @@ import assert from "node:assert/strict" -import type { HistoryItem } from "../packages/types/src/history" +import type { HistoryItem, PendingTaskAction } from "../packages/types/src/history" import { abandonDelegatedChild, completeDelegatedChild, delegateTaskToChild, interruptDelegatedChild, + settleRejectedCreateSubtaskAction, } from "../src/core/task-persistence/taskLifecycle" const taskIds = ["parent", "child-a", "child-b"] as const @@ -16,6 +17,9 @@ type ModelState = Record interface Transition { name: string next: ModelState + delegation?: { parentId: TaskId } + completion?: { childId: TaskId; pendingActionId?: string } + settlement?: { taskId: TaskId; actionId: string } } interface TraceStep { @@ -23,9 +27,16 @@ interface TraceStep { state: ModelState } +interface WitnessContext { + prev: ModelState + next: ModelState + transition: Transition +} + const MAX_DEPTH = 12 const MAX_STATES = 10_000 -const expectedActions = ["delegate", "interrupt", "complete", "abandon"] as const +const actionIds = ["action-1", "action-2"] as const +const expectedActions = ["delegate", "interrupt", "complete", "abandon", "stage", "settle-rejected"] as const const semanticLandmarks = { "interrupted-child-redelegation": (state: ModelState) => state.parent?.status === "delegated" && @@ -37,6 +48,58 @@ const semanticLandmarks = { state["child-a"]?.status === "delegated" && state["child-a"].awaitingChildId === "child-b", } satisfies Record boolean> +const semanticWitnesses = { + "interrupted-pending-delegation-settled": ({ prev, next, transition }: WitnessContext) => + transition.settlement !== undefined && + prev[transition.settlement.taskId]?.status === "interrupted" && + prev[transition.settlement.taskId]?.pendingAction?.actionId === transition.settlement.actionId && + next[transition.settlement.taskId]?.pendingAction === undefined, + "successful-active-delegation-settlement": ({ prev, next, transition }: WitnessContext) => + transition.delegation !== undefined && + prev[transition.delegation.parentId]?.status === "active" && + prev[transition.delegation.parentId]?.pendingAction?.kind === "create_subtask" && + next[transition.delegation.parentId]?.pendingAction === undefined, + "completion-preserves-unrelated-pending-action": ({ prev, next, transition }: WitnessContext) => { + if (transition.completion === undefined) return false + const beforeAction = prev[transition.completion.childId]?.pendingAction + return ( + beforeAction !== undefined && + next[transition.completion.childId]?.status === "completed" && + canonicalTask(next[transition.completion.childId]?.pendingAction) === canonicalTask(beforeAction) + ) + }, + "stale-settlement-rejected": ({ prev, next, transition }: WitnessContext) => { + if (transition.settlement === undefined) return false + const beforeAction = prev[transition.settlement.taskId]?.pendingAction + return ( + beforeAction?.kind === "create_subtask" && + beforeAction.actionId !== transition.settlement.actionId && + next[transition.settlement.taskId]?.pendingAction?.actionId === beforeAction.actionId + ) + }, + "matching-completion-clears-pending-action": ({ prev, next, transition }: WitnessContext) => { + if (transition.completion?.pendingActionId === undefined) return false + const beforeAction = prev[transition.completion.childId]?.pendingAction + const after = next[transition.completion.childId] + return ( + beforeAction?.kind === "create_subtask" && + beforeAction.actionId === transition.completion.pendingActionId && + after?.status === "completed" && + after.pendingAction === undefined + ) + }, + "replacement-completion-preserves-pending-action": ({ prev, next, transition }: WitnessContext) => { + if (transition.completion?.pendingActionId === undefined) return false + const beforeAction = prev[transition.completion.childId]?.pendingAction + const after = next[transition.completion.childId] + return ( + beforeAction?.kind === "create_subtask" && + beforeAction.actionId !== transition.completion.pendingActionId && + after?.status === "completed" && + canonicalTask(after.pendingAction) === canonicalTask(beforeAction) + ) + }, +} satisfies Record boolean> function task(id: TaskId, parentTaskId?: TaskId): HistoryItem { return { @@ -54,6 +117,17 @@ function task(id: TaskId, parentTaskId?: TaskId): HistoryItem { } } +function createSubtaskAction(actionId: (typeof actionIds)[number]): PendingTaskAction { + return { + kind: "create_subtask", + actionId, + approvalText: "{}", + mode: "code", + message: `message for ${actionId}`, + todos: [], + } +} + function initialState(): ModelState { return { parent: task("parent"), "child-a": undefined, "child-b": undefined } } @@ -70,18 +144,38 @@ function transitions(state: ModelState): Transition[] { const parent = state[parentId] if (!parent) continue + const awaitedStatus = parent.awaitingChildId ? state[parent.awaitingChildId as TaskId]?.status : undefined + const delegationValid = + parent.status === "active" || (parent.status === "delegated" && awaitedStatus === "interrupted") + for (const childId of taskIds) { if (childId === parentId || state[childId]) continue - const awaitedStatus = parent.awaitingChildId ? state[parent.awaitingChildId as TaskId]?.status : undefined - if (parent.status !== "active" && !(parent.status === "delegated" && awaitedStatus === "interrupted")) { - continue - } - const delegated = delegateTaskToChild(parent, childId, awaitedStatus) + if (!delegationValid) continue + const delegated = { ...delegateTaskToChild(parent, childId, awaitedStatus), pendingAction: undefined } result.push({ name: `delegate(${parentId}, ${childId})`, next: replace(state, delegated, task(childId, parentId)), + delegation: { parentId }, }) } + + for (const actionId of actionIds) { + if (parent.status === "active" || parent.status === "interrupted") { + result.push({ + name: `stage(${parentId}, ${actionId})`, + next: replace(state, { ...parent, pendingAction: createSubtaskAction(actionId) }), + }) + } + + const pending = parent.pendingAction + if (!delegationValid && parent.status !== "completed" && pending?.kind === "create_subtask") { + result.push({ + name: `settle-rejected(${parentId}, ${actionId})`, + next: replace(state, settleRejectedCreateSubtaskAction(parent, actionId)), + settlement: { taskId: parentId, actionId }, + }) + } + } } for (const childId of taskIds) { @@ -104,7 +198,19 @@ function transitions(state: ModelState): Transition[] { result.push({ name: `complete(${childId})`, next: replace(state, completed.parent, completed.child), + completion: { childId }, }) + for (const actionId of actionIds) { + const completedChild: HistoryItem = + child.pendingAction?.actionId === actionId + ? { ...completed.child, pendingAction: undefined } + : completed.child + result.push({ + name: `complete(${childId}, ${actionId})`, + next: replace(state, completed.parent, completedChild), + completion: { childId, pendingActionId: actionId }, + }) + } } if (parent.status === "delegated" && parent.awaitingChildId === child.id && child.status === "interrupted") { @@ -189,10 +295,49 @@ function checkTransitionInvariants(previous: ModelState, transition: Transition) violations.push(`${id}: completed task changed after ${transition.name}`) } } + + const settlement = transition.settlement + if (settlement) { + const before = previous[settlement.taskId] + const after = transition.next[settlement.taskId] + const beforeAction = before?.pendingAction + const afterAction = after?.pendingAction + + if (afterAction && canonicalTask(afterAction) !== canonicalTask(beforeAction)) { + violations.push(`${settlement.taskId}: settlement after ${transition.name} modified a replacement action`) + } + if ( + !afterAction && + beforeAction && + !(beforeAction.kind === "create_subtask" && beforeAction.actionId === settlement.actionId) + ) { + violations.push(`${settlement.taskId}: settlement after ${transition.name} cleared a non-matching action`) + } + const beforeRest = { ...before, pendingAction: undefined } + const afterRest = { ...after, pendingAction: undefined } + if (canonicalTask(beforeRest) !== canonicalTask(afterRest)) { + violations.push( + `${settlement.taskId}: settlement after ${transition.name} changed status, lineage, or accounting`, + ) + } + } + + const completion = transition.completion + if (completion) { + const beforeAction = previous[completion.childId]?.pendingAction + const afterAction = transition.next[completion.childId]?.pendingAction + + if (afterAction && canonicalTask(afterAction) !== canonicalTask(beforeAction)) { + violations.push(`${completion.childId}: completion after ${transition.name} replaced a pending action`) + } + if (!afterAction && beforeAction && beforeAction.actionId !== completion.pendingActionId) { + violations.push(`${completion.childId}: completion after ${transition.name} cleared a non-matching action`) + } + } return violations } -function canonicalTask(value: HistoryItem | undefined): string { +function canonicalTask(value: unknown): string { return JSON.stringify(value ?? null) } @@ -204,6 +349,7 @@ function runModelCheck(): number { const visited = new Set([canonical(start)]) const reachedActions = new Set() const reachedLandmarks = new Set() + const reachedWitnesses = new Set() const frontier: ModelState[] = [] for (let index = 0; index < queue.length; index++) { @@ -220,6 +366,9 @@ function runModelCheck(): number { for (const transition of transitions(node.state)) { reachedActions.add(transition.name.slice(0, transition.name.indexOf("("))) + for (const [name, matches] of Object.entries(semanticWitnesses)) { + if (matches({ prev: node.state, next: transition.next, transition })) reachedWitnesses.add(name) + } const transitionViolations = checkTransitionInvariants(node.state, transition) const trace = [...node.trace, { action: transition.name, state: transition.next }] if (transitionViolations.length) { @@ -244,6 +393,10 @@ function runModelCheck(): number { if (missingLandmarks.length) { throw new Error(`Task lifecycle model has unreachable semantic landmarks: ${missingLandmarks.join(", ")}`) } + const missingWitnesses = Object.keys(semanticWitnesses).filter((name) => !reachedWitnesses.has(name)) + if (missingWitnesses.length) { + throw new Error(`Task lifecycle model has unreachable semantic witnesses: ${missingWitnesses.join(", ")}`) + } const unexploredSuccessor = frontier .flatMap((state) => transitions(state)) .find((transition) => !visited.has(canonical(transition.next))) @@ -278,10 +431,29 @@ function runRepresentativeScenarios(): void { const interruptedCompletion = completeDelegatedChild(delegated, interruptedA, "resumed result") assert.equal(interruptedCompletion.child.status, "completed") assert.equal(interruptedCompletion.parent.status, "active") + + const rejectedParent: HistoryItem = { + ...parent, + status: "interrupted", + pendingAction: createSubtaskAction("action-1"), + } + assert.throws(() => delegateTaskToChild(rejectedParent, childA.id), /Invalid task status transition/) + const settled = settleRejectedCreateSubtaskAction(rejectedParent, "action-1") + assert.equal(settled.status, "interrupted") + assert.equal(settled.pendingAction, undefined) + assert.equal(settled.childIds, rejectedParent.childIds) + + const replacement = { ...rejectedParent, pendingAction: createSubtaskAction("action-2") } + assert.equal(settleRejectedCreateSubtaskAction(replacement, "action-1"), replacement) + assert.equal( + settleRejectedCreateSubtaskAction({ ...rejectedParent, status: "completed" }, "action-1").pendingAction + ?.actionId, + "action-1", + ) } runRepresentativeScenarios() const checkedStates = runModelCheck() console.log( - `Task lifecycle model check passed: ${checkedStates} reachable states, ${expectedActions.length}/${expectedActions.length} actions reachable, ${Object.keys(semanticLandmarks).length}/${Object.keys(semanticLandmarks).length} landmarks reached, depth <= ${MAX_DEPTH}, ${taskIds.length} task slots`, + `Task lifecycle model check passed: ${checkedStates} reachable states, ${expectedActions.length}/${expectedActions.length} actions reachable, ${Object.keys(semanticLandmarks).length}/${Object.keys(semanticLandmarks).length} landmarks reached, ${Object.keys(semanticWitnesses).length}/${Object.keys(semanticWitnesses).length} semantic witnesses reached, depth <= ${MAX_DEPTH}, ${taskIds.length} task slots`, ) diff --git a/src/__tests__/ClineProvider.delegation.spec.ts b/src/__tests__/ClineProvider.delegation.spec.ts index 422c264e2c..68f39c6517 100644 --- a/src/__tests__/ClineProvider.delegation.spec.ts +++ b/src/__tests__/ClineProvider.delegation.spec.ts @@ -1,10 +1,15 @@ // npx vitest run __tests__/provider-delegation.spec.ts +import * as fs from "fs/promises" +import * as os from "os" +import * as path from "path" + import { describe, it, expect, vi } from "vitest" import type { HistoryItem } from "@roo-code/types" import { providerIdentifiers, RooCodeEventName } from "@roo-code/types" import { ClineProvider } from "../core/webview/ClineProvider" import { TaskScheduler } from "../core/task/TaskScheduler" +import { LifecycleTransitionError } from "../core/task-persistence" const parentHistoryItem: HistoryItem = { id: "parent-1", @@ -753,4 +758,493 @@ describe("ClineProvider.delegateParentAndOpenChild()", () => { expect(deleteTaskWithId).toHaveBeenCalledWith("child-1", false) expect(createTaskWithHistoryItem).toHaveBeenCalledWith(parentHistoryItem) }) + + it("settles an interrupted parent's pending action when delegation is rejected", async () => { + const pendingAction = { + kind: "create_subtask" as const, + actionId: "create-action", + approvalText: "{}", + mode: "code", + message: "Do something", + todos: [], + } + let current: HistoryItem = { + ...parentHistoryItem, + status: "interrupted", + pendingAction, + } + const parentTask = makeParentTask() + const child = { taskId: "child-1", run: vi.fn().mockResolvedValue(undefined) } + const getCurrentTask = vi.fn().mockReturnValue(parentTask) + const createTask = vi.fn(async () => { + getCurrentTask.mockReturnValue(child) + return child + }) + const taskHistoryStore = { + invalidate: vi.fn().mockResolvedValue(undefined), + get: vi.fn(() => current), + atomicReadAndUpdate: vi.fn(async (_taskId: string, updater: (item: HistoryItem) => HistoryItem) => { + current = updater(current) + return [current] + }), + clearPendingActionIfMatching: vi.fn(async (_taskId: string, actionId: string) => { + if (current.pendingAction?.kind === "create_subtask" && current.pendingAction.actionId === actionId) { + current = { ...current, pendingAction: undefined } + } + return current + }), + } + const provider = { + taskScheduler: new TaskScheduler(), + recentTasksCache: [parentHistoryItem], + emit: vi.fn(), + getCurrentTask, + removeClineFromStack: vi.fn().mockResolvedValue(undefined), + createTask, + getTaskWithId: vi.fn().mockImplementation(async () => ({ historyItem: current })), + handleModeSwitch: vi.fn().mockResolvedValue(undefined), + deleteTaskWithId: vi.fn().mockResolvedValue(undefined), + createTaskWithHistoryItem: vi.fn().mockResolvedValue(undefined), + log: vi.fn(), + isViewLaunched: false, + taskHistoryStore, + } as unknown as ClineProvider + + await expect( + ClineProvider.prototype.delegateParentAndOpenChild.call(provider, { + parentTaskId: "parent-1", + message: pendingAction.message, + initialTodos: pendingAction.todos, + mode: pendingAction.mode, + pendingActionId: pendingAction.actionId, + }), + ).rejects.toThrow("Invalid task status transition: interrupted → delegated") + + expect(current.pendingAction).toBeUndefined() + expect(current.status).toBe("interrupted") + expect((provider as unknown as { recentTasksCache?: HistoryItem[] }).recentTasksCache).toBeUndefined() + expect(provider.deleteTaskWithId).toHaveBeenCalledWith("child-1", false) + expect(provider.createTaskWithHistoryItem).toHaveBeenCalledWith( + expect.objectContaining({ + status: "interrupted", + pendingAction: undefined, + }), + ) + }) + + it("rolls back a typed lifecycle rejection without a pending action ID", async () => { + const interruptedParent: HistoryItem = { + ...parentHistoryItem, + status: "interrupted", + } + const parentTask = makeParentTask() + const child = { taskId: "child-1", run: vi.fn().mockResolvedValue(undefined) } + const getCurrentTask = vi.fn().mockReturnValue(parentTask) + const createTask = vi.fn(async () => { + getCurrentTask.mockReturnValue(child) + return child + }) + let transitionError: LifecycleTransitionError | undefined + const atomicReadAndUpdate = vi.fn(async (_taskId: string, updater: (item: HistoryItem) => HistoryItem) => { + try { + updater(interruptedParent) + } catch (error) { + if (error instanceof LifecycleTransitionError) transitionError = error + throw error + } + return [] + }) + const clearPendingActionIfMatching = vi.fn() + const deleteTaskWithId = vi.fn().mockResolvedValue(undefined) + const createTaskWithHistoryItem = vi.fn().mockResolvedValue(undefined) + const provider = { + taskScheduler: new TaskScheduler(), + emit: vi.fn(), + getCurrentTask, + removeClineFromStack: vi.fn().mockResolvedValue(undefined), + createTask, + getTaskWithId: vi.fn().mockResolvedValue({ historyItem: interruptedParent }), + handleModeSwitch: vi.fn().mockResolvedValue(undefined), + deleteTaskWithId, + createTaskWithHistoryItem, + log: vi.fn(), + isViewLaunched: false, + taskHistoryStore: { + invalidate: vi.fn().mockResolvedValue(undefined), + get: vi.fn(() => interruptedParent), + atomicReadAndUpdate, + clearPendingActionIfMatching, + }, + } as unknown as ClineProvider + + let rejection: unknown + try { + await ClineProvider.prototype.delegateParentAndOpenChild.call(provider, { + parentTaskId: "parent-1", + message: "Do something", + initialTodos: [], + mode: "code", + }) + } catch (error) { + rejection = error + } + + expect(rejection).toBe(transitionError) + expect(transitionError).toBeInstanceOf(LifecycleTransitionError) + expect(clearPendingActionIfMatching).not.toHaveBeenCalled() + expect(deleteTaskWithId).toHaveBeenCalledWith("child-1", false) + expect(createTaskWithHistoryItem).toHaveBeenCalledWith(interruptedParent) + }) + + it("does not settle a pending action after an unrelated persistence failure", async () => { + const persistenceError = new Error("parent persistence failed") + const pendingAction = { + kind: "create_subtask" as const, + actionId: "create-action", + approvalText: "{}", + mode: "code", + message: "Do something", + todos: [], + } + const interruptedParent: HistoryItem = { + ...parentHistoryItem, + status: "interrupted", + pendingAction, + } + const parentTask = makeParentTask() + const child = { taskId: "child-1", run: vi.fn().mockResolvedValue(undefined) } + const getCurrentTask = vi.fn().mockReturnValue(parentTask) + const createTask = vi.fn(async () => { + getCurrentTask.mockReturnValue(child) + return child + }) + const atomicReadAndUpdate = vi.fn().mockRejectedValue(persistenceError) + const clearPendingActionIfMatching = vi.fn().mockResolvedValue(interruptedParent) + const createTaskWithHistoryItem = vi.fn().mockResolvedValue(undefined) + const provider = { + taskScheduler: new TaskScheduler(), + emit: vi.fn(), + getCurrentTask, + removeClineFromStack: vi.fn().mockResolvedValue(undefined), + createTask, + getTaskWithId: vi.fn().mockResolvedValue({ historyItem: interruptedParent }), + handleModeSwitch: vi.fn().mockResolvedValue(undefined), + deleteTaskWithId: vi.fn().mockResolvedValue(undefined), + createTaskWithHistoryItem, + log: vi.fn(), + isViewLaunched: false, + taskHistoryStore: { + invalidate: vi.fn().mockResolvedValue(undefined), + get: vi.fn(() => interruptedParent), + atomicReadAndUpdate, + clearPendingActionIfMatching, + }, + } as unknown as ClineProvider + + await expect( + ClineProvider.prototype.delegateParentAndOpenChild.call(provider, { + parentTaskId: "parent-1", + message: pendingAction.message, + initialTodos: pendingAction.todos, + mode: pendingAction.mode, + pendingActionId: pendingAction.actionId, + }), + ).rejects.toThrow(persistenceError) + + expect(atomicReadAndUpdate).toHaveBeenCalledTimes(1) + expect(clearPendingActionIfMatching).not.toHaveBeenCalled() + expect(interruptedParent.pendingAction).toBe(pendingAction) + expect(provider.deleteTaskWithId).toHaveBeenCalledWith("child-1", false) + expect(createTaskWithHistoryItem).toHaveBeenCalledWith(interruptedParent) + }) + + it("does not restore a rejected pending action when its settlement write fails", async () => { + const settlementError = new Error("pending-action settlement failed") + const pendingAction = { + kind: "create_subtask" as const, + actionId: "create-action", + approvalText: "{}", + mode: "code", + message: "Do something", + todos: [], + } + const interruptedParent: HistoryItem = { + ...parentHistoryItem, + status: "interrupted", + pendingAction, + } + const parentTask = makeParentTask() + const child = { taskId: "child-1", run: vi.fn().mockResolvedValue(undefined) } + const getCurrentTask = vi.fn().mockReturnValue(parentTask) + const createTask = vi.fn(async () => { + getCurrentTask.mockReturnValue(child) + return child + }) + const atomicReadAndUpdate = vi.fn(async (_taskId: string, updater: (item: HistoryItem) => HistoryItem) => { + updater(interruptedParent) + return [] + }) + const clearPendingActionIfMatching = vi.fn().mockRejectedValue(settlementError) + const createTaskWithHistoryItem = vi.fn().mockResolvedValue(undefined) + const provider = { + taskScheduler: new TaskScheduler(), + emit: vi.fn(), + getCurrentTask, + removeClineFromStack: vi.fn().mockResolvedValue(undefined), + createTask, + getTaskWithId: vi.fn().mockResolvedValue({ historyItem: interruptedParent }), + handleModeSwitch: vi.fn().mockResolvedValue(undefined), + deleteTaskWithId: vi.fn().mockResolvedValue(undefined), + createTaskWithHistoryItem, + log: vi.fn(), + isViewLaunched: false, + taskHistoryStore: { + invalidate: vi.fn().mockResolvedValue(undefined), + get: vi.fn(() => interruptedParent), + atomicReadAndUpdate, + clearPendingActionIfMatching, + }, + } as unknown as ClineProvider + + await expect( + ClineProvider.prototype.delegateParentAndOpenChild.call(provider, { + parentTaskId: "parent-1", + message: pendingAction.message, + initialTodos: pendingAction.todos, + mode: pendingAction.mode, + pendingActionId: pendingAction.actionId, + }), + ).rejects.toThrow("Invalid task status transition: interrupted → delegated") + + expect(atomicReadAndUpdate).toHaveBeenCalledTimes(1) + expect(clearPendingActionIfMatching).toHaveBeenCalledTimes(1) + expect(provider.deleteTaskWithId).toHaveBeenCalledWith("child-1", false) + expect(createTaskWithHistoryItem).not.toHaveBeenCalled() + }) + + it("does not restore a completed parent when settlement preserves the rejected action", async () => { + const pendingAction = { + kind: "create_subtask" as const, + actionId: "create-action", + approvalText: "{}", + mode: "code", + message: "Do something", + todos: [], + } + const interruptedParent: HistoryItem = { + ...parentHistoryItem, + status: "interrupted", + pendingAction, + } + const completedParent: HistoryItem = { + ...parentHistoryItem, + status: "completed", + pendingAction, + } + const parentTask = makeParentTask() + const child = { taskId: "child-1", run: vi.fn().mockResolvedValue(undefined) } + const getCurrentTask = vi.fn().mockReturnValue(parentTask) + const createTask = vi.fn(async () => { + getCurrentTask.mockReturnValue(child) + return child + }) + const clearPendingActionIfMatching = vi.fn().mockResolvedValue(completedParent) + const createTaskWithHistoryItem = vi.fn().mockResolvedValue(undefined) + const provider = { + taskScheduler: new TaskScheduler(), + emit: vi.fn(), + getCurrentTask, + removeClineFromStack: vi.fn().mockResolvedValue(undefined), + createTask, + getTaskWithId: vi.fn().mockResolvedValue({ historyItem: completedParent }), + handleModeSwitch: vi.fn().mockResolvedValue(undefined), + deleteTaskWithId: vi.fn().mockResolvedValue(undefined), + createTaskWithHistoryItem, + log: vi.fn(), + isViewLaunched: false, + taskHistoryStore: { + invalidate: vi.fn().mockResolvedValue(undefined), + get: vi.fn(() => interruptedParent), + atomicReadAndUpdate: vi.fn(async (_taskId: string, updater: (item: HistoryItem) => HistoryItem) => { + updater(interruptedParent) + return [] + }), + clearPendingActionIfMatching, + }, + } as unknown as ClineProvider + + await expect( + ClineProvider.prototype.delegateParentAndOpenChild.call(provider, { + parentTaskId: "parent-1", + message: pendingAction.message, + initialTodos: pendingAction.todos, + mode: pendingAction.mode, + pendingActionId: pendingAction.actionId, + }), + ).rejects.toThrow("Invalid task status transition: interrupted → delegated") + + expect(clearPendingActionIfMatching).toHaveBeenCalledWith("parent-1", pendingAction.actionId) + expect(provider.deleteTaskWithId).toHaveBeenCalledWith("child-1", false) + expect(createTaskWithHistoryItem).not.toHaveBeenCalled() + }) + + it("restores the authoritative parent record when settlement preserves a replacement action", async () => { + const pendingAction = { + kind: "create_subtask" as const, + actionId: "create-action", + approvalText: "{}", + mode: "code", + message: "Do something", + todos: [], + } + const replacementAction = { ...pendingAction, actionId: "replacement-action" } + const cachedParent: HistoryItem = { + ...parentHistoryItem, + status: "interrupted", + pendingAction, + } + const replacedParent: HistoryItem = { + ...parentHistoryItem, + status: "interrupted", + pendingAction: replacementAction, + } + const parentTask = makeParentTask() + const child = { taskId: "child-1", run: vi.fn().mockResolvedValue(undefined) } + const getCurrentTask = vi.fn().mockReturnValue(parentTask) + const createTask = vi.fn(async () => { + getCurrentTask.mockReturnValue(child) + return child + }) + const clearPendingActionIfMatching = vi.fn(async (_taskId: string, _actionId: string) => replacedParent) + const createTaskWithHistoryItem = vi.fn().mockResolvedValue(undefined) + const provider = { + taskScheduler: new TaskScheduler(), + emit: vi.fn(), + getCurrentTask, + removeClineFromStack: vi.fn().mockResolvedValue(undefined), + createTask, + getTaskWithId: vi.fn().mockResolvedValue({ historyItem: replacedParent }), + handleModeSwitch: vi.fn().mockResolvedValue(undefined), + deleteTaskWithId: vi.fn().mockResolvedValue(undefined), + createTaskWithHistoryItem, + log: vi.fn(), + isViewLaunched: false, + taskHistoryStore: { + invalidate: vi.fn().mockResolvedValue(undefined), + get: vi.fn(() => cachedParent), + atomicReadAndUpdate: vi.fn(async (_taskId: string, updater: (item: HistoryItem) => HistoryItem) => { + updater(cachedParent) + return [] + }), + clearPendingActionIfMatching, + }, + } as unknown as ClineProvider + + await expect( + ClineProvider.prototype.delegateParentAndOpenChild.call(provider, { + parentTaskId: "parent-1", + message: pendingAction.message, + initialTodos: pendingAction.todos, + mode: pendingAction.mode, + pendingActionId: pendingAction.actionId, + }), + ).rejects.toThrow("Invalid task status transition: interrupted → delegated") + + expect(clearPendingActionIfMatching).toHaveBeenCalledTimes(1) + expect(replacedParent.pendingAction).toEqual(replacementAction) + expect(provider.deleteTaskWithId).toHaveBeenCalledWith("child-1", false) + expect(createTaskWithHistoryItem).toHaveBeenCalledWith(replacedParent) + }) + + it("keeps directory cleanup and parent restoration when the child history lock failure is swallowed", async () => { + const pendingAction = { + kind: "create_subtask" as const, + actionId: "create-action", + approvalText: "{}", + mode: "code", + message: "Do something", + todos: [], + } + const interruptedParent: HistoryItem = { + ...parentHistoryItem, + status: "interrupted", + pendingAction, + } + const settledParent: HistoryItem = { ...interruptedParent, pendingAction: undefined } + const childItem: HistoryItem = { ...parentHistoryItem, id: "child-1", task: "Child" } + + const globalStorageDir = await fs.mkdtemp(path.join(os.tmpdir(), "delegation-rollback-cleanup-")) + const childDir = path.join(globalStorageDir, "tasks", "child-1") + await fs.mkdir(childDir, { recursive: true }) + await fs.writeFile(path.join(childDir, "ui_messages.json"), "[]") + + const parentTask = makeParentTask() + const child = { taskId: "child-1", run: vi.fn().mockResolvedValue(undefined) } + const getCurrentTask = vi.fn().mockReturnValue(parentTask) + const createTask = vi.fn(async () => { + getCurrentTask.mockReturnValue(child) + return child + }) + const createTaskWithHistoryItem = vi.fn().mockResolvedValue(undefined) + const getTaskWithId = vi.fn(async (taskId: string) => { + if (taskId === "parent-1") { + return { historyItem: settledParent } + } + return { taskDirPath: childDir, historyItem: childItem } + }) + const clearPendingActionIfMatching = vi.fn(async () => settledParent) + + const provider = { + taskScheduler: new TaskScheduler(), + recentTasksCache: [parentHistoryItem], + emit: vi.fn(), + getCurrentTask, + removeClineFromStack: vi.fn().mockResolvedValue(undefined), + createTask, + getTaskWithId, + handleModeSwitch: vi.fn().mockResolvedValue(undefined), + // The real deletion path: the store swallows a per-file history + // lock failure, so deleteMany resolves and cleanup continues. + deleteTaskWithId: ClineProvider.prototype.deleteTaskWithId, + createTaskWithHistoryItem, + log: vi.fn(), + postStateToWebview: vi.fn().mockResolvedValue(undefined), + isViewLaunched: false, + contextProxy: { globalStorageUri: { fsPath: globalStorageDir } }, + cwd: globalStorageDir, + taskHistoryStore: { + invalidate: vi.fn().mockResolvedValue(undefined), + get: vi.fn(() => interruptedParent), + atomicReadAndUpdate: vi.fn(async (_taskId: string, updater: (item: HistoryItem) => HistoryItem) => { + updater(interruptedParent) + return [] + }), + clearPendingActionIfMatching, + deleteMany: vi.fn().mockResolvedValue(undefined), + }, + } as unknown as ClineProvider + + try { + await expect( + ClineProvider.prototype.delegateParentAndOpenChild.call(provider, { + parentTaskId: "parent-1", + message: pendingAction.message, + initialTodos: pendingAction.todos, + mode: pendingAction.mode, + pendingActionId: pendingAction.actionId, + }), + ).rejects.toThrow("Invalid task status transition: interrupted → delegated") + + // The child task directory was removed even though the child's + // history lock failed and was swallowed, and the settled parent + // was restored before the original rejection surfaced. + await expect(fs.access(childDir)).rejects.toMatchObject({ code: "ENOENT" }) + expect(clearPendingActionIfMatching).toHaveBeenCalledWith("parent-1", "create-action") + expect(createTaskWithHistoryItem).toHaveBeenCalledWith( + expect.objectContaining({ status: "interrupted", pendingAction: undefined }), + ) + } finally { + await fs.rm(globalStorageDir, { recursive: true, force: true }).catch(() => {}) + } + }) }) diff --git a/src/__tests__/removeClineFromStack-delegation.spec.ts b/src/__tests__/removeClineFromStack-delegation.spec.ts index e211a8f4ab..bc5b8a4426 100644 --- a/src/__tests__/removeClineFromStack-delegation.spec.ts +++ b/src/__tests__/removeClineFromStack-delegation.spec.ts @@ -3,7 +3,7 @@ import { describe, it, expect, vi, type MockedFunction } from "vitest" import { ClineProvider } from "../core/webview/ClineProvider" import { TaskRegistry } from "../core/task/TaskRegistry" -import { type Task } from "../core/task/Task" +import { PendingActionSettlementError, type Task } from "../core/task/Task" import { makeProviderStub } from "./helpers/provider-stub" type MockTask = Pick & @@ -14,6 +14,7 @@ type MockTask = Pick & type PrivateClineProviderMethods = { removeClineFromStack: (this: ClineProvider) => ReturnType + cleanupFailedHistoryTask: (this: ClineProvider, task: Task, error: unknown) => Promise markDelegatedChildInterrupted: ( this: ClineProvider, ...args: Parameters @@ -176,6 +177,64 @@ describe("ClineProvider.removeClineFromStack() — pure lifecycle, no delegation }) }) +describe("ClineProvider failed history restoration cleanup", () => { + it("removes the failed task, its listeners, and its resources without saving stale history", async () => { + const cleanupListener = vi.fn() + const task = { + taskId: "failed-history-task", + instanceId: "inst-1", + emit: vi.fn(), + dispose: vi.fn().mockResolvedValue(undefined), + } as unknown as Task + const taskRegistry = new TaskRegistry() + taskRegistry.push(task) + const taskEventListeners = new Map([[task, [cleanupListener]]]) + const provider = { + taskRegistry, + taskEventListeners, + log: vi.fn(), + } as unknown as ClineProvider + + await privateClineProvider.cleanupFailedHistoryTask.call( + provider, + task, + new PendingActionSettlementError("settlement failed"), + ) + + expect(taskRegistry.getById(task.taskId)).toBeUndefined() + expect(taskRegistry.current).toBeUndefined() + expect(cleanupListener).toHaveBeenCalledOnce() + expect(taskEventListeners.has(task)).toBe(false) + expect(task.dispose).toHaveBeenCalledOnce() + }) + + it("keeps the task active after an unrelated history resume failure", async () => { + const cleanupListener = vi.fn() + const task = { + taskId: "failed-history-task", + instanceId: "inst-1", + emit: vi.fn(), + dispose: vi.fn().mockResolvedValue(undefined), + } as unknown as Task + const taskRegistry = new TaskRegistry() + taskRegistry.push(task) + const taskEventListeners = new Map([[task, [cleanupListener]]]) + const provider = { + taskRegistry, + taskEventListeners, + log: vi.fn(), + } as unknown as ClineProvider + + await privateClineProvider.cleanupFailedHistoryTask.call(provider, task, new Error("history read failed")) + + expect(taskRegistry.getById(task.taskId)).toBe(task) + expect(taskRegistry.current).toBe(task) + expect(cleanupListener).not.toHaveBeenCalled() + expect(taskEventListeners.has(task)).toBe(true) + expect(task.dispose).not.toHaveBeenCalled() + }) +}) + describe("ClineProvider.markDelegatedChildInterrupted() — live eviction path", () => { it("marks an active delegated child interrupted and leaves parent delegated", async () => { const childTaskId = "child-1" diff --git a/src/api/providers/__tests__/gemini.spec.ts b/src/api/providers/__tests__/gemini.spec.ts index 2f19028eb7..0920042586 100644 --- a/src/api/providers/__tests__/gemini.spec.ts +++ b/src/api/providers/__tests__/gemini.spec.ts @@ -332,6 +332,53 @@ describe("GeminiHandler", () => { }) }) + it("generates request-unique tool call IDs across requests (#1714)", async () => { + const metadata = { + taskId: "test-task", + tools: [{ type: "function", function: { name: "new_task", description: "", parameters: {} } }], + } satisfies ApiHandlerCreateMessageMetadata + const messages: Anthropic.Messages.MessageParam[] = [{ role: "user", content: "Delegate" }] + + // Each request restarts the tool-call counter at zero, so the + // synthesized ID must carry a request-unique component to keep + // persisted pending-action IDs distinct across requests. + const firstCallIds: string[] = [] + for (let request = 0; request < 2; request++) { + mockGenerateContentStream.mockResolvedValueOnce( + asyncStreamFrom([ + { + candidates: [{ content: { parts: [{ functionCall: { name: "new_task", args: {} } }] } }], + }, + ]), + ) + const chunks = await collectStream(handler.createMessage(systemPrompt, messages, metadata)) + const partials = chunks + .filter((chunk) => chunk.type === "tool_call_partial") + .map((chunk) => chunk as { id: string; name?: string; arguments?: string }) + // The handler emits one name partial and one arguments partial + // for the synthesized call. Assert both semantic halves so a + // duplicate name partial cannot satisfy the ID comparison. + expect(partials).toHaveLength(2) + expect( + partials.filter(({ name, arguments: args }) => name === "new_task" && args === undefined), + ).toHaveLength(1) + expect( + partials.filter(({ name, arguments: args }) => name === undefined && args === "{}"), + ).toHaveLength(1) + + const partialIds = partials.map(({ id }) => id) + expect(new Set(partialIds).size).toBe(1) + firstCallIds.push(partialIds[0]) + } + + // The first call of each request must not collide. + expect(firstCallIds[0]).not.toBe(firstCallIds[1]) + for (const id of firstCallIds) { + expect(id).toMatch(/^new_task-.+-0$/) + expect(id).not.toBe("new_task-0") + } + }) + it("should handle API errors", async () => { const mockError = new Error("Gemini API error") ;(handler["client"].models.generateContentStream as any).mockRejectedValue(mockError) diff --git a/src/api/providers/gemini.ts b/src/api/providers/gemini.ts index ec0d14e4c9..7674521d13 100644 --- a/src/api/providers/gemini.ts +++ b/src/api/providers/gemini.ts @@ -354,6 +354,10 @@ export class GeminiHandler extends BaseProvider implements SingleCompletionHandl let finalResponse: { responseId?: string } | undefined let finishReason: string | undefined + // Gemini provides no call ID, so one is synthesized here. The + // request-unique component keeps persisted pending-action IDs (#1714) + // distinct across requests even though the counter restarts at zero. + const toolCallRequestId = crypto.randomUUID() let toolCallCounter = 0 let hasContent = false let hasReasoning = false @@ -398,7 +402,7 @@ export class GeminiHandler extends BaseProvider implements SingleCompletionHandl hasContent = true // Gemini sends complete function calls in a single chunk // Emit as partial chunks for consistent handling with NativeToolCallParser - const callId = `${part.functionCall.name}-${toolCallCounter}` + const callId = `${part.functionCall.name}-${toolCallRequestId}-${toolCallCounter}` const args = JSON.stringify(part.functionCall.args) // Emit name first diff --git a/src/core/task-persistence/TaskHistoryStore.ts b/src/core/task-persistence/TaskHistoryStore.ts index 3d4cc47604..346d697d9a 100644 --- a/src/core/task-persistence/TaskHistoryStore.ts +++ b/src/core/task-persistence/TaskHistoryStore.ts @@ -4,12 +4,13 @@ import * as path from "path" import crypto from "crypto" import deepEqual from "fast-deep-equal" -import type { HistoryItem } from "@roo-code/types" +import { historyItemSchema, type HistoryItem } from "@roo-code/types" import { GlobalFileNames } from "../../shared/globalFileNames" -import { LOCK_STALE_MS, safeWriteJson } from "../../utils/safeWriteJson" +import { LOCK_STALE_MS, withFileLock } from "../../utils/fileLock" +import { safeWriteJson } from "../../utils/safeWriteJson" import { getStorageBasePath } from "../../utils/storage" -import { assertValidTransition, type HistoryItemStatus } from "./taskLifecycle" +import { assertValidTransition, settleRejectedCreateSubtaskAction, type HistoryItemStatus } from "./taskLifecycle" import { computeHistoryDelta, DeltaRejectedError, mergeHistoryDelta } from "./taskStoreConcurrency" export { assertValidTransition, type HistoryItemStatus } from "./taskLifecycle" @@ -263,6 +264,13 @@ export class TaskHistoryStore { /** * Delete a single task's history item. + * + * Deletion is best-effort: the unlink runs under the same per-file + * advisory lock as `safeWriteJson`, so a locked read-modify-write (for + * example the settlement in `clearPendingActionIfMatching`) cannot + * interleave with it. A lock or unlink failure is swallowed because the + * file may already be deleted; the in-memory eviction and the write + * through still complete. */ async delete(taskId: string): Promise { return this.withLock(async () => { @@ -272,7 +280,7 @@ export class TaskHistoryStore { // Remove per-task file (best-effort) try { const filePath = await this.getTaskFilePath(taskId) - await fs.unlink(filePath) + await withFileLock(filePath, (absoluteFilePath) => fs.unlink(absoluteFilePath)) } catch { // File may already be deleted } @@ -286,6 +294,10 @@ export class TaskHistoryStore { /** * Delete multiple tasks' history items in a batch. + * + * Every item follows the `delete` semantics and is attempted even when an + * earlier unlink fails. The single write-through runs once after the + * whole batch. */ async deleteMany(taskIds: string[]): Promise { return this.withLock(async () => { @@ -293,9 +305,10 @@ export class TaskHistoryStore { this.cache.delete(taskId) this.taskFileMtimes.delete(taskId) + // Remove per-task file (best-effort) try { const filePath = await this.getTaskFilePath(taskId) - await fs.unlink(filePath) + await withFileLock(filePath, (absoluteFilePath) => fs.unlink(absoluteFilePath)) } catch { // File may already be deleted } @@ -1060,6 +1073,93 @@ export class TaskHistoryStore { }) } + /** + * Disk-authoritative compare-and-clear for a rejected `create_subtask` + * pending action (#1714). The comparison runs inside the per-file + * advisory lock's merge callback, so the decision reads the persisted + * record rather than this store's possibly stale cache. Settlement goes + * through the shared `settleRejectedCreateSubtaskAction` reducer, so a + * completed record is never mutated and a missing, different-kind, or + * replacement pending action is preserved unchanged. The authoritative + * record is written back, the store cache is refreshed with it, and it + * is returned to the caller. + * + * Deletion by another host is authoritative (#1726): when no persisted + * record exists, the merge callback removes the stale cache entry and + * throws instead of writing the cached record back to disk. A persisted + * record must also match the canonical task-history schema and carry the + * requested task ID; malformed or mismatched records drop the stale + * cache entry and fail settlement closed without rewriting the record. + * + * @throws If the task ID is not present in the cache, no persisted record remains on disk, + * or the persisted record is invalid or belongs to a different task ID. + */ + public async clearPendingActionIfMatching(taskId: string, expectedActionId: string): Promise { + return this.withLock(async () => { + const cached = this.cache.get(taskId) + if (!cached) { + throw new Error(`[TaskHistoryStore] clearPendingActionIfMatching: task ${taskId} not found in cache`) + } + const filePath = await this.getTaskFilePath(taskId) + let authoritative: HistoryItem = cached + let missingDiskRecord = false + try { + await safeWriteJson(filePath, cached, { + merge: (existing) => { + if (existing === null || existing === undefined) { + // Writing the cached record back would recreate a task + // another host deleted, so drop the stale entry first. + // A throwing merge writes nothing, so the deleted + // record stays deleted. + missingDiskRecord = true + this.cache.delete(taskId) + this.taskFileMtimes.delete(taskId) + throw new Error( + `[TaskHistoryStore] clearPendingActionIfMatching: task ${taskId} not found in cache`, + ) + } + // Validate the locked disk record with the canonical + // task-history schema (#1726). Settlement must not + // rewrite malformed data and must not clear the action + // on a record persisted under a different task ID, so + // both cases drop the stale cache entry and fail + // closed without touching the disk record. + const parsed = historyItemSchema.safeParse(existing) + if (!parsed.success) { + this.cache.delete(taskId) + this.taskFileMtimes.delete(taskId) + throw new Error( + `[TaskHistoryStore] clearPendingActionIfMatching: task ${taskId} has an invalid disk record`, + ) + } + if (parsed.data.id !== taskId) { + this.cache.delete(taskId) + this.taskFileMtimes.delete(taskId) + throw new Error( + `[TaskHistoryStore] clearPendingActionIfMatching: task ${taskId} has a disk record with mismatched id ${parsed.data.id}`, + ) + } + const disk = existing as HistoryItem + authoritative = settleRejectedCreateSubtaskAction(disk, expectedActionId) + return authoritative + }, + }) + } catch (error) { + if (missingDiskRecord) { + throw new Error( + `[TaskHistoryStore] clearPendingActionIfMatching: task ${taskId} not found in cache`, + ) + } + throw error + } + this.cache.set(taskId, authoritative) + if (this.onWrite) { + await this.onWrite(this.getAll()) + } + return authoritative + }) + } + // ────────────────────────────── Private: Write lock ────────────────────────────── /** diff --git a/src/core/task-persistence/__tests__/TaskHistoryStore.crossInstance.spec.ts b/src/core/task-persistence/__tests__/TaskHistoryStore.crossInstance.spec.ts index cef9874e5f..5b85c9b4e2 100644 --- a/src/core/task-persistence/__tests__/TaskHistoryStore.crossInstance.spec.ts +++ b/src/core/task-persistence/__tests__/TaskHistoryStore.crossInstance.spec.ts @@ -204,6 +204,118 @@ describe("TaskHistoryStore cross-instance safety", () => { expect(storeB.getAll().length).toBe(10) }) + it("preserves the cached record when rejected-action settlement fails with an unrelated lock error", async () => { + await storeA.initialize() + const pendingAction = { + kind: "create_subtask" as const, + actionId: "action-a", + approvalText: "{}", + mode: "code", + message: "action A", + todos: [], + } + const cached = makeHistoryItem({ id: "settlement-error-task", pendingAction }) + await storeA.upsert(cached) + + const { safeWriteJson } = await import("../../../utils/safeWriteJson") + const lockError = Object.assign(new Error("lock acquisition failed"), { code: "ELOCKED" }) + vi.mocked(safeWriteJson).mockRejectedValueOnce(lockError) + + await expect(storeA.clearPendingActionIfMatching(cached.id, pendingAction.actionId)).rejects.toBe(lockError) + expect(storeA.get(cached.id)).toEqual(cached) + }) + + it("evicts stale cache state when the disk record is invalid without recreating it", async () => { + const pendingAction = { + kind: "create_subtask" as const, + actionId: "action-a", + approvalText: "{}", + mode: "code", + message: "action A", + todos: [], + } + const cached = makeHistoryItem({ id: "invalid-settlement-task", pendingAction }) + const filePath = path.join(tmpDir, "tasks", cached.id, GlobalFileNames.historyItem) + await fs.mkdir(path.dirname(filePath), { recursive: true }) + await fs.writeFile(filePath, "{invalid", "utf8") + const storeState = storeA as unknown as { + cache: Map + taskFileMtimes: Map + } + storeState.cache.set(cached.id, cached) + storeState.taskFileMtimes.set(cached.id, Date.now()) + + await expect(storeA.clearPendingActionIfMatching(cached.id, pendingAction.actionId)).rejects.toThrow( + `task ${cached.id} not found in cache`, + ) + expect(storeA.get(cached.id)).toBeUndefined() + expect(storeState.taskFileMtimes.has(cached.id)).toBe(false) + expect(await fs.readFile(filePath, "utf8")).toBe("{invalid") + }) + + it("fails settlement closed for a malformed disk record without rewriting it", async () => { + const pendingAction = { + kind: "create_subtask" as const, + actionId: "action-a", + approvalText: "{}", + mode: "code", + message: "action A", + todos: [], + } + const cached = makeHistoryItem({ id: "malformed-settlement-task", pendingAction }) + const filePath = path.join(tmpDir, "tasks", cached.id, GlobalFileNames.historyItem) + await fs.mkdir(path.dirname(filePath), { recursive: true }) + // Valid JSON that fails the canonical task-history schema: id must + // be a string and the required history fields are absent. + const malformed = JSON.stringify({ id: 123, pendingAction }) + await fs.writeFile(filePath, malformed, "utf8") + const storeState = storeA as unknown as { + cache: Map + taskFileMtimes: Map + } + storeState.cache.set(cached.id, cached) + storeState.taskFileMtimes.set(cached.id, Date.now()) + + await expect(storeA.clearPendingActionIfMatching(cached.id, pendingAction.actionId)).rejects.toThrow( + `task ${cached.id} has an invalid disk record`, + ) + expect(storeA.get(cached.id)).toBeUndefined() + expect(storeState.taskFileMtimes.has(cached.id)).toBe(false) + // The malformed record stays on disk untouched. + expect(await fs.readFile(filePath, "utf8")).toBe(malformed) + }) + + it("fails settlement closed when the disk record carries a different task id", async () => { + const pendingAction = { + kind: "create_subtask" as const, + actionId: "action-a", + approvalText: "{}", + mode: "code", + message: "action A", + todos: [], + } + const cached = makeHistoryItem({ id: "mismatch-settlement-task", pendingAction }) + const other = makeHistoryItem({ id: "other-task" }) + const filePath = path.join(tmpDir, "tasks", cached.id, GlobalFileNames.historyItem) + await fs.mkdir(path.dirname(filePath), { recursive: true }) + const mismatched = JSON.stringify(other, null, "\t") + await fs.writeFile(filePath, mismatched, "utf8") + const storeState = storeA as unknown as { + cache: Map + taskFileMtimes: Map + } + storeState.cache.set(cached.id, cached) + storeState.taskFileMtimes.set(cached.id, Date.now()) + + await expect(storeA.clearPendingActionIfMatching(cached.id, pendingAction.actionId)).rejects.toThrow( + `task ${cached.id} has a disk record with mismatched id other-task`, + ) + expect(storeA.get(cached.id)).toBeUndefined() + expect(storeState.taskFileMtimes.has(cached.id)).toBe(false) + // The other task's record stays on disk untouched. + expect(await fs.readFile(filePath, "utf8")).toBe(mismatched) + }) + /** * Host B completes a task on disk while host A's cache still has it * active. Host A's next save updates only totalCost (a full-object diff --git a/src/core/task-persistence/__tests__/TaskHistoryStore.deleteSemantics.spec.ts b/src/core/task-persistence/__tests__/TaskHistoryStore.deleteSemantics.spec.ts new file mode 100644 index 0000000000..5071c9de73 --- /dev/null +++ b/src/core/task-persistence/__tests__/TaskHistoryStore.deleteSemantics.spec.ts @@ -0,0 +1,243 @@ +// pnpm --filter roo-cline test core/task-persistence/__tests__/TaskHistoryStore.deleteSemantics.spec.ts +// +// Best-effort deletion semantics for `delete()` and `deleteMany()`. +// Lock, unlink, and fs behavior stay real by default; individual tests force +// one lock or unlink failure through the wrappers below. + +import * as fs from "fs/promises" +import * as os from "os" +import * as path from "path" + +import type { HistoryItem } from "@roo-code/types" + +import { TaskHistoryStore } from "../TaskHistoryStore" +import { withFileLock } from "../../../utils/fileLock" +import { GlobalFileNames } from "../../../shared/globalFileNames" + +vi.mock("../../../utils/storage", () => ({ + getStorageBasePath: vi.fn().mockImplementation((defaultPath: string) => defaultPath), +})) + +// The default implementation stays the real one, so `safeWriteJson` writes +// and per-file locking behave exactly as in production unless a test forces +// a failure. +vi.mock("../../../utils/fileLock", async () => { + const actual = await vi.importActual("../../../utils/fileLock") + return { ...actual, withFileLock: vi.fn(actual.withFileLock) } +}) + +vi.mock("fs/promises", async () => { + const actual = await vi.importActual("fs/promises") + return { ...actual, unlink: vi.fn(actual.unlink) } +}) + +const actualFs = await vi.importActual("fs/promises") +const actualFileLock = await vi.importActual("../../../utils/fileLock") + +function makeHistoryItem(overrides: Partial = {}): HistoryItem { + return { + id: `task-${Date.now()}-${Math.random().toString(36).substring(2, 8)}`, + number: 1, + ts: Date.now(), + task: "Test task", + tokensIn: 100, + tokensOut: 50, + totalCost: 0.01, + workspace: "/test/workspace", + ...overrides, + } +} + +function historyFilePath(storagePath: string, taskId: string): string { + return path.join(storagePath, "tasks", taskId, GlobalFileNames.historyItem) +} + +function storeInternals(store: TaskHistoryStore): { + cache: Map + taskFileMtimes: Map +} { + return { + cache: store["cache"], + taskFileMtimes: store["taskFileMtimes"], + } +} + +describe("TaskHistoryStore best-effort deletion semantics", () => { + let storagePath: string + let stores: TaskHistoryStore[] + let onWrite: ReturnType + + beforeEach(async () => { + storagePath = await fs.mkdtemp(path.join(os.tmpdir(), "task-history-delete-semantics-")) + stores = [] + onWrite = vi.fn().mockResolvedValue(undefined) + vi.mocked(withFileLock).mockImplementation(actualFileLock.withFileLock) + vi.mocked(fs.unlink).mockImplementation(actualFs.unlink) + }) + + afterEach(async () => { + for (const store of stores) { + store.dispose() + } + await fs.rm(storagePath, { recursive: true, force: true }).catch(() => {}) + }) + + function createStore(): TaskHistoryStore { + const store = new TaskHistoryStore(storagePath, { + onWrite: onWrite as (items: HistoryItem[]) => Promise, + }) + stores.push(store) + return store + } + + describe("delete()", () => { + it("unlinks under the shared per-file lock, evicts cache and mtime, and writes through once", async () => { + const store = createStore() + await store.initialize() + await store.upsert(makeHistoryItem({ id: "locked-delete" })) + onWrite.mockClear() + + await expect(store.delete("locked-delete")).resolves.toBeUndefined() + + expect(vi.mocked(withFileLock)).toHaveBeenCalledWith( + historyFilePath(storagePath, "locked-delete"), + expect.any(Function), + ) + await expect(fs.access(historyFilePath(storagePath, "locked-delete"))).rejects.toMatchObject({ + code: "ENOENT", + }) + + const { cache, taskFileMtimes } = storeInternals(store) + expect(cache.has("locked-delete")).toBe(false) + expect(taskFileMtimes.has("locked-delete")).toBe(false) + expect(store.get("locked-delete")).toBeUndefined() + + expect(onWrite).toHaveBeenCalledTimes(1) + const writtenIds = (onWrite.mock.calls[0][0] as HistoryItem[]).map((item) => item.id) + expect(writtenIds).not.toContain("locked-delete") + }) + + it("swallows a lock acquisition failure, evicts cache and mtime, and still writes through once", async () => { + const store = createStore() + await store.initialize() + await store.upsert(makeHistoryItem({ id: "lock-fail" })) + onWrite.mockClear() + + vi.mocked(withFileLock).mockRejectedValueOnce(new Error("lock acquisition timed out")) + + await expect(store.delete("lock-fail")).resolves.toBeUndefined() + + // The unlink never ran, but the in-memory eviction and the single + // write-through still completed. + await expect(fs.access(historyFilePath(storagePath, "lock-fail"))).resolves.toBeUndefined() + const { cache, taskFileMtimes } = storeInternals(store) + expect(cache.has("lock-fail")).toBe(false) + expect(taskFileMtimes.has("lock-fail")).toBe(false) + expect(store.get("lock-fail")).toBeUndefined() + expect(onWrite).toHaveBeenCalledTimes(1) + }) + + it("swallows a non-ENOENT unlink failure, evicts cache and mtime, and stays deletable afterwards", async () => { + const store = createStore() + await store.initialize() + await store.upsert(makeHistoryItem({ id: "perm-fail" })) + onWrite.mockClear() + + vi.mocked(fs.unlink).mockRejectedValueOnce( + Object.assign(new Error("EACCES: permission denied, unlink"), { code: "EACCES" }), + ) + + await expect(store.delete("perm-fail")).resolves.toBeUndefined() + + await expect(fs.access(historyFilePath(storagePath, "perm-fail"))).resolves.toBeUndefined() + const { cache, taskFileMtimes } = storeInternals(store) + expect(cache.has("perm-fail")).toBe(false) + expect(taskFileMtimes.has("perm-fail")).toBe(false) + expect(onWrite).toHaveBeenCalledTimes(1) + + // The per-file lock was released: a retry without the injected + // failure deletes the file and writes through again. + await expect(store.delete("perm-fail")).resolves.toBeUndefined() + await expect(fs.access(historyFilePath(storagePath, "perm-fail"))).rejects.toMatchObject({ + code: "ENOENT", + }) + expect(onWrite).toHaveBeenCalledTimes(2) + }) + + it("treats a missing file as a completed deletion", async () => { + const store = createStore() + await store.initialize() + onWrite.mockClear() + + await expect(store.delete("never-existed")).resolves.toBeUndefined() + expect(store.get("never-existed")).toBeUndefined() + expect(onWrite).toHaveBeenCalledTimes(1) + }) + }) + + describe("deleteMany()", () => { + it("continues the batch after a failed unlink, evicts every entry, and writes through exactly once", async () => { + const store = createStore() + await store.initialize() + await store.upsert(makeHistoryItem({ id: "batch-a", ts: 1000 })) + await store.upsert(makeHistoryItem({ id: "batch-b", ts: 2000 })) + await store.upsert(makeHistoryItem({ id: "batch-c", ts: 3000 })) + onWrite.mockClear() + + vi.mocked(fs.unlink).mockImplementation(async (p) => { + if (p === historyFilePath(storagePath, "batch-b")) { + throw Object.assign(new Error("EACCES: permission denied, unlink"), { code: "EACCES" }) + } + return actualFs.unlink(p) + }) + + await expect(store.deleteMany(["batch-a", "batch-b", "batch-c"])).resolves.toBeUndefined() + + // batch-b failed but the batch continued around it. + await expect(fs.access(historyFilePath(storagePath, "batch-a"))).rejects.toMatchObject({ code: "ENOENT" }) + await expect(fs.access(historyFilePath(storagePath, "batch-b"))).resolves.toBeUndefined() + await expect(fs.access(historyFilePath(storagePath, "batch-c"))).rejects.toMatchObject({ code: "ENOENT" }) + + const { cache, taskFileMtimes } = storeInternals(store) + expect(cache.has("batch-a")).toBe(false) + expect(cache.has("batch-b")).toBe(false) + expect(cache.has("batch-c")).toBe(false) + expect(taskFileMtimes.has("batch-a")).toBe(false) + expect(taskFileMtimes.has("batch-b")).toBe(false) + expect(taskFileMtimes.has("batch-c")).toBe(false) + expect(store.get("batch-b")).toBeUndefined() + + // One write-through, after the whole batch, seeing the final cache. + expect(onWrite).toHaveBeenCalledTimes(1) + const writtenIds = (onWrite.mock.calls[0][0] as HistoryItem[]).map((item) => item.id) + expect(writtenIds).toEqual([]) + }) + + it("swallows a lock failure for one item and continues the batch with one write-through", async () => { + const store = createStore() + await store.initialize() + await store.upsert(makeHistoryItem({ id: "lock-a", ts: 1000 })) + await store.upsert(makeHistoryItem({ id: "lock-b", ts: 2000 })) + onWrite.mockClear() + + vi.mocked(withFileLock).mockRejectedValueOnce(new Error("lock acquisition timed out")) + + await expect(store.deleteMany(["lock-a", "lock-b"])).resolves.toBeUndefined() + + // The first item's lock failed before its unlink; the second item + // still completed. + await expect(fs.access(historyFilePath(storagePath, "lock-a"))).resolves.toBeUndefined() + await expect(fs.access(historyFilePath(storagePath, "lock-b"))).rejects.toMatchObject({ code: "ENOENT" }) + + const { cache, taskFileMtimes } = storeInternals(store) + expect(cache.has("lock-a")).toBe(false) + expect(cache.has("lock-b")).toBe(false) + expect(taskFileMtimes.has("lock-a")).toBe(false) + expect(taskFileMtimes.has("lock-b")).toBe(false) + + expect(onWrite).toHaveBeenCalledTimes(1) + const writtenIds = (onWrite.mock.calls[0][0] as HistoryItem[]).map((item) => item.id) + expect(writtenIds).toEqual([]) + }) + }) +}) diff --git a/src/core/task-persistence/__tests__/TaskHistoryStore.realConcurrency.spec.ts b/src/core/task-persistence/__tests__/TaskHistoryStore.realConcurrency.spec.ts index d94ca8f782..d35dfc49cf 100644 --- a/src/core/task-persistence/__tests__/TaskHistoryStore.realConcurrency.spec.ts +++ b/src/core/task-persistence/__tests__/TaskHistoryStore.realConcurrency.spec.ts @@ -6,6 +6,13 @@ import type { HistoryItem } from "@roo-code/types" import { TaskHistoryStore } from "../TaskHistoryStore" +// Wrap the real fs/promises so a test can pause inside the per-file lock +// while every call still runs against the real filesystem. +vi.mock("fs/promises", async () => { + const actual = await vi.importActual("fs/promises") + return { ...actual, readFile: vi.fn(actual.readFile) } +}) + type WriteTaskFile = (item: HistoryItem, delta?: Partial) => Promise interface WriteBarrier { @@ -72,6 +79,17 @@ function item(id: string): HistoryItem { } } +function createAction(actionId: string, message: string) { + return { + kind: "create_subtask" as const, + actionId, + approvalText: "{}", + mode: "code", + message, + todos: [], + } +} + describe("TaskHistoryStore real cross-host locking", () => { it("preserves independent stale-cache deltas through the real per-file lock", async () => { const storagePath = await fs.mkdtemp(path.join(os.tmpdir(), "task-history-real-lock-")) @@ -124,4 +142,265 @@ describe("TaskHistoryStore real cross-host locking", () => { await fs.rm(storagePath, { recursive: true, force: true }) } }) + + it("preserves a replacement pending action when a stale store settles the prior action", async () => { + const storagePath = await fs.mkdtemp(path.join(os.tmpdir(), "task-history-stale-settlement-")) + const storeA = new TaskHistoryStore(storagePath) + const storeB = new TaskHistoryStore(storagePath) + const actionA = createAction("action-a", "action A") + const actionB = { ...actionA, actionId: "action-b", message: "action B" } + + try { + await storeA.initialize() + await storeA.upsert({ ...item("shared-task"), pendingAction: actionA }) + await storeB.initialize() + + await storeB.atomicReadAndUpdate("shared-task", (current) => ({ ...current, pendingAction: actionB })) + expect(storeA.get("shared-task")?.pendingAction).toEqual(actionA) + + const authoritative = await storeA.clearPendingActionIfMatching("shared-task", actionA.actionId) + expect(authoritative.pendingAction).toEqual(actionB) + expect(storeA.get("shared-task")?.pendingAction).toEqual(actionB) + await storeB.invalidate("shared-task") + + expect(storeB.get("shared-task")?.pendingAction).toEqual(actionB) + } finally { + storeA.dispose() + storeB.dispose() + await fs.rm(storagePath, { recursive: true, force: true }) + } + }) + + it("clears a matching create_subtask action from disk and refreshes the cache", async () => { + const storagePath = await fs.mkdtemp(path.join(os.tmpdir(), "task-history-compare-clear-")) + const storeA = new TaskHistoryStore(storagePath) + const storeB = new TaskHistoryStore(storagePath) + const actionA = createAction("action-a", "action A") + + try { + await storeA.initialize() + await storeA.upsert({ ...item("shared-task"), pendingAction: actionA }) + await storeB.initialize() + + const authoritative = await storeA.clearPendingActionIfMatching("shared-task", actionA.actionId) + + expect(authoritative.pendingAction).toBeUndefined() + expect(storeA.get("shared-task")?.pendingAction).toBeUndefined() + await storeB.invalidate("shared-task") + expect(storeB.get("shared-task")?.pendingAction).toBeUndefined() + expect(storeB.get("shared-task")).toMatchObject({ id: "shared-task", status: "active" }) + } finally { + storeA.dispose() + storeB.dispose() + await fs.rm(storagePath, { recursive: true, force: true }) + } + }) + + it("preserves a matching create_subtask action on a completed record", async () => { + const storagePath = await fs.mkdtemp(path.join(os.tmpdir(), "task-history-completed-settlement-")) + const store = new TaskHistoryStore(storagePath) + const pendingAction = createAction("action-a", "action A") + const completed: HistoryItem = { ...item("shared-task"), status: "completed", pendingAction } + const filePath = path.join(storagePath, "tasks", completed.id, "history_item.json") + + try { + await store.initialize() + await store.upsert(completed) + + const beforeDisk = JSON.parse(await fs.readFile(filePath, "utf8")) as HistoryItem + expect({ cache: store.get(completed.id), disk: beforeDisk }).toEqual({ cache: completed, disk: completed }) + + const returned = await store.clearPendingActionIfMatching(completed.id, pendingAction.actionId) + const afterDisk = JSON.parse(await fs.readFile(filePath, "utf8")) as HistoryItem + + expect({ returned, disk: afterDisk, cache: store.get(completed.id) }).toEqual({ + returned: completed, + disk: completed, + cache: completed, + }) + } finally { + store.dispose() + await fs.rm(storagePath, { recursive: true, force: true }) + } + }) + + it("clears a matching action from disk even when the calling store cache is stale", async () => { + const storagePath = await fs.mkdtemp(path.join(os.tmpdir(), "task-history-disk-compare-clear-")) + const storeA = new TaskHistoryStore(storagePath) + const storeB = new TaskHistoryStore(storagePath) + const actionA = createAction("action-a", "action A") + + try { + await storeA.initialize() + await storeA.upsert(item("shared-task")) + await storeB.initialize() + + await storeB.atomicReadAndUpdate("shared-task", (current) => ({ ...current, pendingAction: actionA })) + + const authoritative = await storeA.clearPendingActionIfMatching("shared-task", actionA.actionId) + + expect(authoritative.pendingAction).toBeUndefined() + await storeB.invalidate("shared-task") + expect(storeB.get("shared-task")?.pendingAction).toBeUndefined() + } finally { + storeA.dispose() + storeB.dispose() + await fs.rm(storagePath, { recursive: true, force: true }) + } + }) + + it("preserves a different-kind pending action with the same action ID", async () => { + const storagePath = await fs.mkdtemp(path.join(os.tmpdir(), "task-history-different-kind-")) + const storeA = new TaskHistoryStore(storagePath) + const finishAction = { + kind: "finish_subtask" as const, + actionId: "action-a", + approvalText: "{}", + parentTaskId: "parent-1", + result: "done", + } + + try { + await storeA.initialize() + await storeA.upsert({ ...item("shared-task"), pendingAction: finishAction }) + + const authoritative = await storeA.clearPendingActionIfMatching("shared-task", finishAction.actionId) + + expect(authoritative.pendingAction).toEqual(finishAction) + expect(storeA.get("shared-task")?.pendingAction).toEqual(finishAction) + } finally { + storeA.dispose() + await fs.rm(storagePath, { recursive: true, force: true }) + } + }) + + it("preserves the record when no pending action is persisted", async () => { + const storagePath = await fs.mkdtemp(path.join(os.tmpdir(), "task-history-no-action-")) + const storeA = new TaskHistoryStore(storagePath) + + try { + await storeA.initialize() + await storeA.upsert(item("shared-task")) + + const authoritative = await storeA.clearPendingActionIfMatching("shared-task", "action-a") + + expect(authoritative.pendingAction).toBeUndefined() + expect(authoritative).toMatchObject({ id: "shared-task", status: "active" }) + } finally { + storeA.dispose() + await fs.rm(storagePath, { recursive: true, force: true }) + } + }) + + it("does not recreate a task deleted by another host before settlement", async () => { + const storagePath = await fs.mkdtemp(path.join(os.tmpdir(), "task-history-deleted-settlement-")) + const storeA = new TaskHistoryStore(storagePath) + const storeB = new TaskHistoryStore(storagePath) + const actionA = createAction("action-a", "action A") + + try { + await storeA.initialize() + await storeA.upsert({ ...item("shared-task"), pendingAction: actionA }) + await storeB.initialize() + + await storeB.delete("shared-task") + expect(storeA.get("shared-task")?.pendingAction).toEqual(actionA) + + await expect(storeA.clearPendingActionIfMatching("shared-task", actionA.actionId)).rejects.toThrow( + "task shared-task not found", + ) + expect(storeA.get("shared-task")).toBeUndefined() + await storeB.invalidate("shared-task") + expect(storeB.get("shared-task")).toBeUndefined() + } finally { + storeA.dispose() + storeB.dispose() + await fs.rm(storagePath, { recursive: true, force: true }) + } + }) + + it("serializes deletion after settlement reads disk without recreating the record", async () => { + const storagePath = await fs.mkdtemp(path.join(os.tmpdir(), "task-history-delete-during-settlement-")) + const storeA = new TaskHistoryStore(storagePath) + const storeB = new TaskHistoryStore(storagePath) + const action = createAction("action-a", "action A") + const filePath = path.join(storagePath, "tasks", "shared-task", "history_item.json") + const actualFs = await vi.importActual("fs/promises") + let signalReadComplete!: () => void + const readComplete = new Promise((resolve) => { + signalReadComplete = resolve + }) + let releaseSettlement!: () => void + const settlementCanContinue = new Promise((resolve) => { + releaseSettlement = resolve + }) + + try { + await storeA.initialize() + await storeA.upsert({ ...item("shared-task"), pendingAction: action }) + await storeB.initialize() + + // Pause settlement inside its locked disk read, then start a + // deletion from another store so it targets the settlement's + // read-to-commit window. + vi.mocked(fs.readFile).mockImplementation(async (...args: Parameters) => { + const result = await actualFs.readFile(...args) + if (args[0] === filePath) { + signalReadComplete() + await settlementCanContinue + } + return result + }) + + const settlement = storeA.clearPendingActionIfMatching("shared-task", action.actionId) + await readComplete + const deletion = storeB.delete("shared-task") + let deletionSettled = false + void deletion.finally(() => { + deletionSettled = true + }) + await new Promise((resolve) => setTimeout(resolve, 25)) + // The deletion must stay blocked while settlement holds the + // per-file lock across its read-to-commit window. + expect(deletionSettled).toBe(false) + + releaseSettlement() + await expect(settlement).resolves.toMatchObject({ id: "shared-task", pendingAction: undefined }) + await deletion + + // The settlement committed first and the locked deletion + // removed the file afterwards, so the record stays deleted + // instead of being resurrected. + await expect(fs.access(filePath)).rejects.toMatchObject({ code: "ENOENT" }) + await expect(storeA.clearPendingActionIfMatching("shared-task", action.actionId)).rejects.toThrow( + "task shared-task not found", + ) + expect(storeA.get("shared-task")).toBeUndefined() + expect(storeB.get("shared-task")).toBeUndefined() + } finally { + vi.mocked(fs.readFile).mockImplementation(actualFs.readFile) + storeA.dispose() + storeB.dispose() + await fs.rm(storagePath, { recursive: true, force: true }) + } + }) + + it("rejects settlement for a task absent from the cache without creating it", async () => { + const storagePath = await fs.mkdtemp(path.join(os.tmpdir(), "task-history-cache-miss-settlement-")) + const store = new TaskHistoryStore(storagePath) + const filePath = path.join(storagePath, "tasks", "missing-task", "history_item.json") + + try { + await store.initialize() + + await expect(store.clearPendingActionIfMatching("missing-task", "action-a")).rejects.toThrow( + "task missing-task not found in cache", + ) + expect(store.get("missing-task")).toBeUndefined() + await expect(fs.access(filePath)).rejects.toMatchObject({ code: "ENOENT" }) + } finally { + store.dispose() + await fs.rm(storagePath, { recursive: true, force: true }) + } + }) }) diff --git a/src/core/task-persistence/__tests__/taskLifecycle.spec.ts b/src/core/task-persistence/__tests__/taskLifecycle.spec.ts index fe415f09f8..5d3129e023 100644 --- a/src/core/task-persistence/__tests__/taskLifecycle.spec.ts +++ b/src/core/task-persistence/__tests__/taskLifecycle.spec.ts @@ -5,6 +5,8 @@ import { completeDelegatedChild, delegateTaskToChild, interruptDelegatedChild, + LifecycleTransitionError, + settleRejectedCreateSubtaskAction, } from "../taskLifecycle" function item(id: string, overrides: Partial = {}): HistoryItem { @@ -102,3 +104,84 @@ describe("task lifecycle transitions", () => { expect(abandoned.child).toMatchObject({ parentTaskId: undefined, rootTaskId: undefined }) }) }) + +describe("settleRejectedCreateSubtaskAction", () => { + const createSubtaskAction = { + kind: "create_subtask" as const, + actionId: "create-action", + approvalText: "{}", + mode: "code", + message: "Do something", + todos: [], + } + + it("clears only the matching pending create_subtask action", () => { + const parent = item("parent", { + status: "interrupted", + parentTaskId: "root", + rootTaskId: "root", + awaitingChildId: undefined, + tokensIn: 12, + totalCost: 0.5, + pendingAction: createSubtaskAction, + }) + + const settled = settleRejectedCreateSubtaskAction(parent, "create-action") + + expect(settled).toEqual({ + ...parent, + pendingAction: undefined, + }) + expect(settled).toMatchObject({ + status: "interrupted", + parentTaskId: "root", + rootTaskId: "root", + tokensIn: 12, + totalCost: 0.5, + }) + }) + + it("never clears a replacement action with a different ID", () => { + const parent = item("parent", { + status: "interrupted", + pendingAction: { ...createSubtaskAction, actionId: "replacement-action" }, + }) + + expect(settleRejectedCreateSubtaskAction(parent, "stale-action")).toBe(parent) + }) + + it("never clears a pending action of a different kind", () => { + const parent = item("parent", { + pendingAction: { + kind: "finish_subtask", + actionId: "create-action", + approvalText: "{}", + parentTaskId: "root", + result: "done", + }, + }) + + expect(settleRejectedCreateSubtaskAction(parent, "create-action")).toBe(parent) + }) + + it("leaves a record without a pending action unchanged", () => { + const parent = item("parent", { status: "interrupted" }) + + expect(settleRejectedCreateSubtaskAction(parent, "create-action")).toBe(parent) + }) + + it("never mutates a completed record", () => { + const parent = item("parent", { status: "completed", pendingAction: createSubtaskAction }) + + expect(settleRejectedCreateSubtaskAction(parent, "create-action")).toBe(parent) + }) + + it("rejects an interrupted parent's delegation with a typed transition error", () => { + const parent = item("parent", { status: "interrupted", pendingAction: createSubtaskAction }) + + expect(() => delegateTaskToChild(parent, "child")).toThrow(LifecycleTransitionError) + expect(() => delegateTaskToChild(parent, "child")).toThrow( + "Invalid task status transition: interrupted → delegated", + ) + }) +}) diff --git a/src/core/task-persistence/index.ts b/src/core/task-persistence/index.ts index 14adeedc68..0baf9ce57e 100644 --- a/src/core/task-persistence/index.ts +++ b/src/core/task-persistence/index.ts @@ -21,6 +21,7 @@ export { delegateTaskToChild, interruptDelegatedChild, LifecycleTransitionError, + settleRejectedCreateSubtaskAction, type HistoryItemStatus, VALID_TASK_STATUS_TRANSITIONS, } from "./taskLifecycle" diff --git a/src/core/task-persistence/taskLifecycle.ts b/src/core/task-persistence/taskLifecycle.ts index efd2e1148f..b8a5e571f0 100644 --- a/src/core/task-persistence/taskLifecycle.ts +++ b/src/core/task-persistence/taskLifecycle.ts @@ -20,10 +20,27 @@ export class LifecycleTransitionError extends Error { export function assertValidTransition(from: HistoryItemStatus | undefined, to: HistoryItemStatus): void { const fromStatus: HistoryItemStatus = from ?? "active" if (!VALID_TASK_STATUS_TRANSITIONS[fromStatus].includes(to)) { - throw new Error(`Invalid task status transition: ${fromStatus} → ${to}`) + throw new LifecycleTransitionError(`Invalid task status transition: ${fromStatus} → ${to}`) } } +/** + * Settles the pending create_subtask action whose delegation the authoritative + * parent record rejected (#1714). Only the exact matching action ID is cleared; + * status, lineage, accounting, and unrelated fields are preserved. A replaced + * or different-kind pending action is never cleared. + */ +export function settleRejectedCreateSubtaskAction(parent: HistoryItem, pendingActionId: string): HistoryItem { + const pending = parent.pendingAction + if (parent.status === "completed") { + return parent + } + if (pending?.kind !== "create_subtask" || pending.actionId !== pendingActionId) { + return parent + } + return { ...parent, pendingAction: undefined } +} + export function delegateTaskToChild( parent: HistoryItem, childId: string, diff --git a/src/core/task/Task.ts b/src/core/task/Task.ts index d5313f68cf..143c6b886c 100644 --- a/src/core/task/Task.ts +++ b/src/core/task/Task.ts @@ -217,6 +217,13 @@ type AssistantMessagePersistenceCancellation = { resolve: () => void } +export class PendingActionSettlementError extends Error { + constructor(message: string, options?: ErrorOptions) { + super(message, options) + this.name = "PendingActionSettlementError" + } +} + export class Task extends EventEmitter implements TaskLike { readonly taskId: string readonly rootTaskId?: string @@ -610,7 +617,7 @@ export class Task extends EventEmitter implements TaskLike { this.parentTask = parentTask this.taskNumber = taskNumber - this.initialStatus = initialStatus + this.initialStatus = initialStatus ?? historyItem?.status this.pendingAction = historyItem?.pendingAction // Store the task's mode and API config name when it's created. @@ -995,6 +1002,60 @@ export class Task extends EventEmitter implements TaskLike { } } + /** + * An interrupted task cannot legally delegate, so a staged create-subtask + * action is a durable rejection marker rather than replayable work. The + * constructor-injected history item can be stale, so the persisted record + * is refreshed first and the refreshed action is the one settled. A failed + * refresh, lookup, or settlement stops replay instead of risking another + * doomed child. + */ + private async settleInterruptedCreateSubtaskBeforeReplay(): Promise { + if (this.initialStatus !== "interrupted") { + return + } + + const provider = this.providerRef.deref() + if (!provider) { + throw new PendingActionSettlementError( + `[Task#settleInterruptedCreateSubtaskBeforeReplay] Provider unavailable for task ${this.taskId}`, + ) + } + + try { + await provider.taskHistoryStore.reconcile({ forceRefresh: true }) + } catch (error) { + throw new PendingActionSettlementError( + `[Task#settleInterruptedCreateSubtaskBeforeReplay] Failed to refresh task history for task ${this.taskId}`, + { cause: error }, + ) + } + + const refreshedItem = provider.taskHistoryStore.get(this.taskId) + if (!refreshedItem) { + throw new PendingActionSettlementError( + `[Task#settleInterruptedCreateSubtaskBeforeReplay] Task ${this.taskId} not found in refreshed task history`, + ) + } + this.pendingAction = refreshedItem.pendingAction + + const action = refreshedItem.pendingAction + if (action?.kind !== "create_subtask") { + return + } + + let authoritative: HistoryItem + try { + authoritative = await provider.taskHistoryStore.clearPendingActionIfMatching(this.taskId, action.actionId) + } catch (error) { + throw new PendingActionSettlementError( + `[Task#settleInterruptedCreateSubtaskBeforeReplay] Failed to settle rejected action for task ${this.taskId}`, + { cause: error }, + ) + } + this.pendingAction = authoritative.pendingAction + } + private handleQueuedAskResponse(message: QueuedMessage, resolution: QueuedAskResolution): string | undefined { this.handleWebviewAskResponse(resolution.response, message.text, message.images) if (resolution.requiresDurableAck) { @@ -2373,6 +2434,7 @@ export class Task extends EventEmitter implements TaskLike { // This is important in case the user deletes messages without resuming // the task first. this.hydrateApiConversationHistory(savedApiConversationHistory) + await this.settleInterruptedCreateSubtaskBeforeReplay() if ( this.pendingAction && this.apiConversationHistory.some( diff --git a/src/core/task/__tests__/Task.persistence.spec.ts b/src/core/task/__tests__/Task.persistence.spec.ts index 8d3314a9a6..77cb813e51 100644 --- a/src/core/task/__tests__/Task.persistence.spec.ts +++ b/src/core/task/__tests__/Task.persistence.spec.ts @@ -14,7 +14,7 @@ import { import { TelemetryService } from "@roo-code/telemetry" import type { Anthropic } from "@anthropic-ai/sdk" -import { Task } from "../Task" +import { PendingActionSettlementError, Task } from "../Task" import { ClineProvider } from "../../webview/ClineProvider" import { ContextProxy } from "../../config/ContextProxy" import { providerIdentifiers } from "@roo-code/types/provider-identifiers" @@ -1244,6 +1244,13 @@ describe("Task persistence", () => { }, ]) + // Restart settlement needs the refreshed record to exist, and a + // record without a pending action keeps the replay path unchanged. + mockProvider.taskHistoryStore.get = vi.fn().mockReturnValue({ + id: "interrupted-subtask", + status: "interrupted", + }) + await getTaskPersistenceAccess(task).resumeTaskFromHistory() expect(initiateTaskLoopSpy).toHaveBeenCalledTimes(1) @@ -1309,6 +1316,13 @@ describe("Task persistence", () => { }, ]) + // Restart settlement needs the refreshed record to exist, and a + // record without a pending action keeps the replay path unchanged. + mockProvider.taskHistoryStore.get = vi.fn().mockReturnValue({ + id: "interrupted-subtask-2", + status: "interrupted", + }) + await getTaskPersistenceAccess(task).resumeTaskFromHistory() expect(initiateTaskLoopSpy).toHaveBeenCalledTimes(1) @@ -1340,6 +1354,14 @@ describe("Task persistence", () => { parentTaskId: "parent-1", result: "Done", } + const createSubtaskAction: PendingTaskAction = { + kind: "create_subtask", + actionId: "create-action", + approvalText: JSON.stringify({ tool: "newTask" }), + mode: "code", + message: "Child task", + todos: [], + } it("replays an unresolved pending action instead of a generic resume ask", async () => { const messages: ClineMessage[] = [ @@ -1383,6 +1405,404 @@ describe("Task persistence", () => { expect(mockSaveTaskMessages).not.toHaveBeenCalled() }) + it("replays a create-subtask action for an active historical task without restart settlement", async () => { + mockReadTaskMessages.mockResolvedValue([ + { ts: 1, type: "ask", ask: "tool", text: createSubtaskAction.approvalText }, + ]) + mockReadApiMessages.mockResolvedValue([{ role: "assistant", content: "Previous response" }]) + const clearRejectedAction = vi.fn() + mockProvider.taskHistoryStore.clearPendingActionIfMatching = clearRejectedAction + const task = new Task({ + provider: mockProvider, + apiConfiguration: mockApiConfig, + historyItem: { + id: "parent-1", + number: 1, + ts: 1, + task: "Parent", + tokensIn: 0, + tokensOut: 0, + totalCost: 0, + status: "active", + pendingAction: createSubtaskAction, + }, + startTask: false, + }) + const replay = vi + .spyOn(getTaskPersistenceAccess(task), "resumePendingTaskAction") + .mockResolvedValue(undefined) + + await getTaskPersistenceAccess(task).resumeTaskFromHistory() + + expect(clearRejectedAction).not.toHaveBeenCalled() + expect(mockProvider.taskHistoryStore.reconcile).not.toHaveBeenCalled() + expect(replay).toHaveBeenCalledWith(createSubtaskAction) + }) + + it("settles an interrupted create-subtask action before restart replay", async () => { + mockReadTaskMessages.mockResolvedValue([ + { ts: 1, type: "ask", ask: "tool", text: createSubtaskAction.approvalText }, + ]) + mockReadApiMessages.mockResolvedValue([{ role: "assistant", content: "Previous response" }]) + mockProvider.taskHistoryStore.reconcile = vi.fn().mockResolvedValue(undefined) + mockProvider.taskHistoryStore.get = vi.fn().mockReturnValue({ + id: "parent-1", + status: "interrupted", + pendingAction: createSubtaskAction, + }) + const clearRejectedAction = vi.fn().mockResolvedValue({ + id: "parent-1", + status: "interrupted", + pendingAction: undefined, + }) + mockProvider.taskHistoryStore.clearPendingActionIfMatching = clearRejectedAction + const task = new Task({ + provider: mockProvider, + apiConfiguration: mockApiConfig, + historyItem: { + id: "parent-1", + number: 1, + ts: 1, + task: "Parent", + tokensIn: 0, + tokensOut: 0, + totalCost: 0, + status: "interrupted", + pendingAction: createSubtaskAction, + }, + startTask: false, + }) + vi.spyOn(task, "ask").mockResolvedValue({ response: "noButtonClicked" }) + vi.spyOn(getTaskPersistenceAccess(task), "initiateTaskLoop").mockResolvedValue(undefined) + const replay = vi.spyOn(getTaskPersistenceAccess(task), "resumePendingTaskAction") + + await getTaskPersistenceAccess(task).resumeTaskFromHistory() + + expect(mockProvider.taskHistoryStore.reconcile).toHaveBeenCalledWith({ forceRefresh: true }) + expect(mockProvider.taskHistoryStore.get).toHaveBeenCalledWith("parent-1") + expect(clearRejectedAction).toHaveBeenCalledWith("parent-1", "create-action") + expect(replay).not.toHaveBeenCalled() + expect(task.ask).toHaveBeenCalledWith("resume_task") + }) + + it("does not replay an interrupted create-subtask action when restart settlement fails", async () => { + const settlementError = new Error("settlement unavailable") + mockReadTaskMessages.mockResolvedValue([]) + mockReadApiMessages.mockResolvedValue([]) + mockProvider.taskHistoryStore.reconcile = vi.fn().mockResolvedValue(undefined) + mockProvider.taskHistoryStore.get = vi.fn().mockReturnValue({ + id: "parent-1", + status: "interrupted", + pendingAction: createSubtaskAction, + }) + mockProvider.taskHistoryStore.clearPendingActionIfMatching = vi.fn().mockRejectedValue(settlementError) + const task = new Task({ + provider: mockProvider, + apiConfiguration: mockApiConfig, + historyItem: { + id: "parent-1", + number: 1, + ts: 1, + task: "Parent", + tokensIn: 0, + tokensOut: 0, + totalCost: 0, + status: "interrupted", + pendingAction: createSubtaskAction, + }, + startTask: false, + }) + const ask = vi.spyOn(task, "ask") + const replay = vi.spyOn(getTaskPersistenceAccess(task), "resumePendingTaskAction") + + const resumeError = await getTaskPersistenceAccess(task) + .resumeTaskFromHistory() + .catch((error: unknown) => error) + + expect(resumeError).toMatchObject({ + name: "PendingActionSettlementError", + cause: settlementError, + }) + expect(resumeError).toBeInstanceOf(PendingActionSettlementError) + + expect(replay).not.toHaveBeenCalled() + expect(ask).not.toHaveBeenCalled() + }) + + it("adopts the authoritative replacement pending action returned by settlement", async () => { + const replacementAction = { + ...createSubtaskAction, + actionId: "create-action-b", + message: "Replacement child", + } + mockReadTaskMessages.mockResolvedValue([ + { ts: 1, type: "ask", ask: "tool", text: createSubtaskAction.approvalText }, + ]) + mockReadApiMessages.mockResolvedValue([{ role: "assistant", content: "Previous response" }]) + mockProvider.taskHistoryStore.reconcile = vi.fn().mockResolvedValue(undefined) + mockProvider.taskHistoryStore.get = vi.fn().mockReturnValue({ + id: "parent-1", + status: "interrupted", + pendingAction: createSubtaskAction, + }) + const clearRejectedAction = vi.fn().mockResolvedValue({ + id: "parent-1", + status: "interrupted", + pendingAction: replacementAction, + }) + mockProvider.taskHistoryStore.clearPendingActionIfMatching = clearRejectedAction + const task = new Task({ + provider: mockProvider, + apiConfiguration: mockApiConfig, + historyItem: { + id: "parent-1", + number: 1, + ts: 1, + task: "Parent", + tokensIn: 0, + tokensOut: 0, + totalCost: 0, + status: "interrupted", + pendingAction: createSubtaskAction, + }, + startTask: false, + }) + const taskState = task as unknown as { pendingAction?: PendingTaskAction } + const ask = vi.spyOn(task, "ask") + const replay = vi + .spyOn(getTaskPersistenceAccess(task), "resumePendingTaskAction") + .mockResolvedValue(undefined) + + await getTaskPersistenceAccess(task).resumeTaskFromHistory() + + expect(clearRejectedAction).toHaveBeenCalledWith("parent-1", "create-action") + expect(taskState.pendingAction).toEqual(replacementAction) + expect(replay).toHaveBeenCalledWith(replacementAction) + expect(ask).not.toHaveBeenCalled() + }) + + it("settles the refreshed action when the staged action was replaced before restart", async () => { + const refreshedAction = { + ...createSubtaskAction, + actionId: "create-action-b", + message: "Refreshed child", + } + mockReadTaskMessages.mockResolvedValue([ + { ts: 1, type: "ask", ask: "tool", text: createSubtaskAction.approvalText }, + ]) + mockReadApiMessages.mockResolvedValue([{ role: "assistant", content: "Previous response" }]) + mockProvider.taskHistoryStore.reconcile = vi.fn().mockResolvedValue(undefined) + mockProvider.taskHistoryStore.get = vi.fn().mockReturnValue({ + id: "parent-1", + status: "interrupted", + pendingAction: refreshedAction, + }) + const clearRejectedAction = vi.fn().mockResolvedValue({ + id: "parent-1", + status: "interrupted", + pendingAction: undefined, + }) + mockProvider.taskHistoryStore.clearPendingActionIfMatching = clearRejectedAction + const task = new Task({ + provider: mockProvider, + apiConfiguration: mockApiConfig, + historyItem: { + id: "parent-1", + number: 1, + ts: 1, + task: "Parent", + tokensIn: 0, + tokensOut: 0, + totalCost: 0, + status: "interrupted", + pendingAction: createSubtaskAction, + }, + startTask: false, + }) + const taskState = task as unknown as { pendingAction?: PendingTaskAction } + const ask = vi.spyOn(task, "ask") + const replay = vi.spyOn(getTaskPersistenceAccess(task), "resumePendingTaskAction") + + await getTaskPersistenceAccess(task).resumeTaskFromHistory() + + // The stale staged action never reaches settlement; the refreshed + // one does, and the cleared authoritative result ends the replay. + expect(clearRejectedAction).toHaveBeenCalledTimes(1) + expect(clearRejectedAction).toHaveBeenCalledWith("parent-1", "create-action-b") + expect(taskState.pendingAction).toBeUndefined() + expect(replay).not.toHaveBeenCalled() + expect(ask).toHaveBeenCalledWith("resume_task") + }) + + it("replays without settlement when the refreshed record has no pending action", async () => { + mockReadTaskMessages.mockResolvedValue([ + { ts: 1, type: "ask", ask: "tool", text: createSubtaskAction.approvalText }, + ]) + mockReadApiMessages.mockResolvedValue([{ role: "assistant", content: "Previous response" }]) + mockProvider.taskHistoryStore.reconcile = vi.fn().mockResolvedValue(undefined) + mockProvider.taskHistoryStore.get = vi.fn().mockReturnValue({ + id: "parent-1", + status: "interrupted", + pendingAction: undefined, + }) + const clearRejectedAction = vi.fn() + mockProvider.taskHistoryStore.clearPendingActionIfMatching = clearRejectedAction + const task = new Task({ + provider: mockProvider, + apiConfiguration: mockApiConfig, + historyItem: { + id: "parent-1", + number: 1, + ts: 1, + task: "Parent", + tokensIn: 0, + tokensOut: 0, + totalCost: 0, + status: "interrupted", + pendingAction: createSubtaskAction, + }, + startTask: false, + }) + const taskState = task as unknown as { pendingAction?: PendingTaskAction } + const ask = vi.spyOn(task, "ask") + const replay = vi.spyOn(getTaskPersistenceAccess(task), "resumePendingTaskAction") + + await getTaskPersistenceAccess(task).resumeTaskFromHistory() + + expect(mockProvider.taskHistoryStore.reconcile).toHaveBeenCalledWith({ forceRefresh: true }) + expect(clearRejectedAction).not.toHaveBeenCalled() + expect(taskState.pendingAction).toBeUndefined() + expect(replay).not.toHaveBeenCalled() + expect(ask).toHaveBeenCalledWith("resume_task") + }) + + it("replays a refreshed non-create-subtask action without settlement", async () => { + const refreshedFinishAction: PendingTaskAction = { + ...pendingAction, + actionId: "finish-action-b", + } + mockReadTaskMessages.mockResolvedValue([ + { ts: 1, type: "ask", ask: "tool", text: createSubtaskAction.approvalText }, + ]) + mockReadApiMessages.mockResolvedValue([{ role: "assistant", content: "Previous response" }]) + mockProvider.taskHistoryStore.reconcile = vi.fn().mockResolvedValue(undefined) + mockProvider.taskHistoryStore.get = vi.fn().mockReturnValue({ + id: "parent-1", + status: "interrupted", + pendingAction: refreshedFinishAction, + }) + const clearRejectedAction = vi.fn() + mockProvider.taskHistoryStore.clearPendingActionIfMatching = clearRejectedAction + const task = new Task({ + provider: mockProvider, + apiConfiguration: mockApiConfig, + historyItem: { + id: "parent-1", + number: 1, + ts: 1, + task: "Parent", + tokensIn: 0, + tokensOut: 0, + totalCost: 0, + status: "interrupted", + pendingAction: createSubtaskAction, + }, + startTask: false, + }) + const taskState = task as unknown as { pendingAction?: PendingTaskAction } + const ask = vi.spyOn(task, "ask") + const replay = vi + .spyOn(getTaskPersistenceAccess(task), "resumePendingTaskAction") + .mockResolvedValue(undefined) + + await getTaskPersistenceAccess(task).resumeTaskFromHistory() + + expect(clearRejectedAction).not.toHaveBeenCalled() + expect(taskState.pendingAction).toEqual(refreshedFinishAction) + expect(replay).toHaveBeenCalledWith(refreshedFinishAction) + expect(ask).not.toHaveBeenCalled() + }) + + it("stops replay when the refreshed task history cannot be read", async () => { + const refreshError = new Error("reconcile unavailable") + mockReadTaskMessages.mockResolvedValue([]) + mockReadApiMessages.mockResolvedValue([]) + mockProvider.taskHistoryStore.reconcile = vi.fn().mockRejectedValue(refreshError) + const clearRejectedAction = vi.fn() + mockProvider.taskHistoryStore.clearPendingActionIfMatching = clearRejectedAction + const task = new Task({ + provider: mockProvider, + apiConfiguration: mockApiConfig, + historyItem: { + id: "parent-1", + number: 1, + ts: 1, + task: "Parent", + tokensIn: 0, + tokensOut: 0, + totalCost: 0, + status: "interrupted", + pendingAction: createSubtaskAction, + }, + startTask: false, + }) + const ask = vi.spyOn(task, "ask") + const replay = vi.spyOn(getTaskPersistenceAccess(task), "resumePendingTaskAction") + + const resumeError = await getTaskPersistenceAccess(task) + .resumeTaskFromHistory() + .catch((error: unknown) => error) + + expect(resumeError).toBeInstanceOf(PendingActionSettlementError) + expect(resumeError).toMatchObject({ + name: "PendingActionSettlementError", + cause: refreshError, + }) + expect(clearRejectedAction).not.toHaveBeenCalled() + expect(replay).not.toHaveBeenCalled() + expect(ask).not.toHaveBeenCalled() + }) + + it("stops replay when the task is missing from the refreshed task history", async () => { + mockReadTaskMessages.mockResolvedValue([]) + mockReadApiMessages.mockResolvedValue([]) + mockProvider.taskHistoryStore.reconcile = vi.fn().mockResolvedValue(undefined) + mockProvider.taskHistoryStore.get = vi.fn().mockReturnValue(undefined) + const clearRejectedAction = vi.fn() + mockProvider.taskHistoryStore.clearPendingActionIfMatching = clearRejectedAction + const task = new Task({ + provider: mockProvider, + apiConfiguration: mockApiConfig, + historyItem: { + id: "parent-1", + number: 1, + ts: 1, + task: "Parent", + tokensIn: 0, + tokensOut: 0, + totalCost: 0, + status: "interrupted", + pendingAction: createSubtaskAction, + }, + startTask: false, + }) + const ask = vi.spyOn(task, "ask") + const replay = vi.spyOn(getTaskPersistenceAccess(task), "resumePendingTaskAction") + + const resumeError = await getTaskPersistenceAccess(task) + .resumeTaskFromHistory() + .catch((error: unknown) => error) + + expect(resumeError).toBeInstanceOf(PendingActionSettlementError) + expect(resumeError).toMatchObject({ + name: "PendingActionSettlementError", + message: expect.stringContaining("not found in refreshed task history"), + }) + expect(clearRejectedAction).not.toHaveBeenCalled() + expect(replay).not.toHaveBeenCalled() + expect(ask).not.toHaveBeenCalled() + }) + it("reconciles an already-persisted tool result before generic resume", async () => { mockReadTaskMessages.mockResolvedValue([{ ts: 1, type: "say", say: "text", text: "Child" }]) mockReadApiMessages.mockResolvedValue([ diff --git a/src/core/webview/ClineProvider.ts b/src/core/webview/ClineProvider.ts index 6fac386383..39f19130d7 100644 --- a/src/core/webview/ClineProvider.ts +++ b/src/core/webview/ClineProvider.ts @@ -111,7 +111,7 @@ import { forceFullModelDetailsLoad, hasLoadedFullDetails } from "../../api/provi import { ContextProxy } from "../config/ContextProxy" import { ProviderSettingsManager } from "../config/ProviderSettingsManager" import { CustomModesManager } from "../config/CustomModesManager" -import { Task } from "../task/Task" +import { PendingActionSettlementError, Task } from "../task/Task" import { webviewMessageHandler } from "./webviewMessageHandler" import type { ClineMessage, TodoItem } from "@roo-code/types" @@ -125,6 +125,7 @@ import { completeDelegatedChild, delegateTaskToChild, interruptDelegatedChild, + LifecycleTransitionError, } from "../task-persistence" import { readTaskMessages } from "../task-persistence/taskMessages" import { getNonce } from "./getNonce" @@ -175,10 +176,16 @@ function scheduleTask( task: Task, source: string, run: () => Promise = () => task.run(), + onError?: (error: unknown) => void | Promise, ): void { - void scheduler - .schedule(task, run) - .catch((error) => console.error(`[${source}] taskScheduler.schedule failed:`, error)) + void scheduler.schedule(task, run).catch(async (error) => { + console.error(`[${source}] taskScheduler.schedule failed:`, error) + try { + await onError?.(error) + } catch (cleanupError) { + console.error(`[${source}] task failure cleanup failed:`, cleanupError) + } + }) } type GetStateOptions = { @@ -610,6 +617,33 @@ export class ClineProvider } } + private async cleanupFailedHistoryTask(task: Task, error: unknown): Promise { + if (!(error instanceof PendingActionSettlementError)) { + return + } + + if (this.taskRegistry.getById(task.taskId) !== task) { + return + } + + this.taskRegistry.remove(task.taskId) + task.emit(RooCodeEventName.TaskUnfocused) + + const cleanupFunctions = this.taskEventListeners.get(task) + if (cleanupFunctions) { + cleanupFunctions.forEach((cleanup) => cleanup()) + this.taskEventListeners.delete(task) + } + + try { + await task.dispose() + } catch (error) { + this.log( + `[cleanupFailedHistoryTask] dispose() failed for ${task.taskId}.${task.instanceId}: ${error instanceof Error ? error.message : String(error)}`, + ) + } + } + /** * Evicts the current task from the stack and, if it was an active delegated child, * marks it interrupted so the parent stays delegated (rather than silently losing the link). @@ -1361,7 +1395,9 @@ export class ClineProvider ) if (options?.startTask !== false) { - scheduleTask(this.taskScheduler, task, "createTaskWithHistoryItem") + scheduleTask(this.taskScheduler, task, "createTaskWithHistoryItem", undefined, (error) => + this.cleanupFailedHistoryTask(task, error), + ) } } else { await this.addClineToStack(task) @@ -1371,7 +1407,9 @@ export class ClineProvider ) if (options?.startTask !== false) { - scheduleTask(this.taskScheduler, task, "createTaskWithHistoryItem") + scheduleTask(this.taskScheduler, task, "createTaskWithHistoryItem", undefined, (error) => + this.cleanupFailedHistoryTask(task, error), + ) } } @@ -3932,6 +3970,31 @@ export class ClineProvider (err as Error)?.message ?? String(err) }`, ) + // The authoritative parent record rejected this delegation (#1714). + // Settle the matching pending create_subtask action durably through + // the disk-authoritative compare-and-clear so a retry cannot replay + // a rejected action and a replacement action from another host is + // never cleared, then propagate the original error. + let settlementFailed = false + if (pendingActionId && err instanceof LifecycleTransitionError) { + try { + const authoritative = await this.taskHistoryStore.clearPendingActionIfMatching( + parentTaskId, + pendingActionId, + ) + settlementFailed = + authoritative.pendingAction?.kind === "create_subtask" && + authoritative.pendingAction.actionId === pendingActionId + this.recentTasksCache = undefined + } catch (settlementError) { + settlementFailed = true + this.log( + `[delegateParentAndOpenChild] Failed to settle pending action ${pendingActionId} for parent ${parentTaskId}: ${ + (settlementError as Error)?.message ?? String(settlementError) + }`, + ) + } + } try { // Only pop the stack if the child we just created is still on top. // A concurrent delegation could have pushed another child since we created ours. @@ -3955,8 +4018,14 @@ export class ClineProvider ) } try { - const { historyItem: parentHistory } = await this.getTaskWithId(parentTaskId) - await this.createTaskWithHistoryItem(parentHistory) + // A failed settlement write leaves the rejected pending action in + // durable storage. Restoring the stored parent would replay it in + // this process, so leave the parent unrestored. Restart recovery also + // settles interrupted create-subtask actions before allowing replay. + if (!settlementFailed) { + const { historyItem: parentHistory } = await this.getTaskWithId(parentTaskId) + await this.createTaskWithHistoryItem(parentHistory) + } } catch (rollbackError) { this.log( `[delegateParentAndOpenChild] Failed to restore parent ${parentTaskId} during rollback: ${ diff --git a/src/core/webview/__tests__/ClineProvider.spec.ts b/src/core/webview/__tests__/ClineProvider.spec.ts index feef95d870..8ba9144f1c 100644 --- a/src/core/webview/__tests__/ClineProvider.spec.ts +++ b/src/core/webview/__tests__/ClineProvider.spec.ts @@ -87,6 +87,9 @@ vi.mock("../../../utils/storage", () => ({ getSettingsDirectoryPath: vi.fn().mockResolvedValue("/test/settings/path"), getTaskDirectoryPath: vi.fn().mockResolvedValue("/test/task/path"), getGlobalStoragePath: vi.fn().mockResolvedValue("/test/storage/path"), + // Deletion resolves the tasks directory before it removes a history + // file, so the harness must provide the passthrough base path. + getStorageBasePath: vi.fn().mockImplementation((defaultPath: string) => defaultPath), })) vi.mock("@modelcontextprotocol/sdk/types.js", () => ({ diff --git a/src/utils/fileLock.ts b/src/utils/fileLock.ts new file mode 100644 index 0000000000..9f7cad7653 --- /dev/null +++ b/src/utils/fileLock.ts @@ -0,0 +1,78 @@ +import * as path from "path" +import * as lockfile from "proper-lockfile" + +/** + * Shared staleness window for per-file advisory locks. This module owns the + * single advisory lock protocol used by `safeWriteJson` and by callers that + * must serialize with it, such as task-history deletion. + */ +export const LOCK_STALE_MS = 31_000 + +/** + * Acquire the advisory lock for one file path using the exact protocol + * `safeWriteJson` uses, so operations that hold this lock serialize with + * every `safeWriteJson` write to the same path. Callers must release the + * returned function exactly once and must not acquire the same lock again + * while holding it. + */ +export async function acquireFileLock(filePath: string): Promise<() => Promise> { + const absoluteFilePath = path.resolve(filePath) + try { + return await lockfile.lock(absoluteFilePath, { + stale: LOCK_STALE_MS, + update: 10000, // Update mtime every 10 seconds to prevent staleness if operation is long + realpath: false, // the file may not exist yet, which is acceptable + retries: { + // Configuration for retrying lock acquisition + retries: 5, // Number of retries after the initial attempt + factor: 2, // Exponential backoff factor (e.g., 100ms, 200ms, 400ms, ...) + minTimeout: 100, // Minimum time to wait before the first retry (in ms) + maxTimeout: 1000, // Maximum time to wait for any single retry (in ms) + }, + onCompromised: (err) => { + console.error(`Lock at ${absoluteFilePath} was compromised:`, err) + throw err + }, + }) + } catch (lockError) { + console.error(`Failed to acquire lock for ${absoluteFilePath}:`, lockError) + throw lockError + } +} + +/** + * Run one operation while holding the advisory lock for a file path. + * Callers must not acquire this lock again from inside `operation`, and + * they must keep the documented lock order when combining this helper + * with other locks to prevent deadlock. + */ +export async function withFileLock( + filePath: string, + operation: (absoluteFilePath: string) => Promise, +): Promise { + const absoluteFilePath = path.resolve(filePath) + const releaseLock = await acquireFileLock(absoluteFilePath) + + let result: T + try { + result = await operation(absoluteFilePath) + } catch (operationError) { + // The operation error is the primary failure. Release without + // reporting a secondary release error over it. + try { + await releaseLock() + } catch (releaseError) { + console.error(`Failed to release lock for ${absoluteFilePath}:`, releaseError) + } + throw operationError + } + + try { + await releaseLock() + } catch (releaseError) { + // The operation already succeeded, so a release failure is only + // logged, matching how `safeWriteJson` handles release failures. + console.error(`Failed to release lock for ${absoluteFilePath}:`, releaseError) + } + return result +} diff --git a/src/utils/safeWriteJson.ts b/src/utils/safeWriteJson.ts index 957a0bb20f..7da68b2a7a 100644 --- a/src/utils/safeWriteJson.ts +++ b/src/utils/safeWriteJson.ts @@ -1,9 +1,10 @@ import * as fs from "fs/promises" import * as fsSync from "fs" import * as path from "path" -import * as lockfile from "proper-lockfile" import { JsonStreamStringify } from "json-stream-stringify" +import { acquireFileLock } from "./fileLock" + /** * Options for safeWriteJson function */ @@ -61,32 +62,14 @@ async function safeWriteJson(filePath: string, data: any, options?: SafeWriteJso throw dirError } - // Acquire the lock before any file operations - try { - releaseLock = await lockfile.lock(absoluteFilePath, { - stale: LOCK_STALE_MS, - update: 10000, // Update mtime every 10 seconds to prevent staleness if operation is long - realpath: false, // the file may not exist yet, which is acceptable - retries: { - // Configuration for retrying lock acquisition - retries: 5, // Number of retries after the initial attempt - factor: 2, // Exponential backoff factor (e.g., 100ms, 200ms, 400ms, ...) - minTimeout: 100, // Minimum time to wait before the first retry (in ms) - maxTimeout: 1000, // Maximum time to wait for any single retry (in ms) - }, - onCompromised: (err) => { - console.error(`Lock at ${absoluteFilePath} was compromised:`, err) - throw err - }, - }) - } catch (lockError) { - // If lock acquisition fails, we throw immediately. - // The releaseLock remains a no-op, so the finally block in the main file operations - // try-catch-finally won't try to release an unacquired lock if this path is taken. - console.error(`Failed to acquire lock for ${absoluteFilePath}:`, lockError) - // Propagate the lock acquisition error - throw lockError - } + // Acquire the lock before any file operations. `acquireFileLock` owns the + // shared advisory lock protocol, so callers that lock the same path with + // it (for example task-history deletion) serialize with this write. + // If lock acquisition fails, it throws immediately. The releaseLock + // remains a no-op, so the finally block in the main file operations + // try-catch-finally won't try to release an unacquired lock if this + // path is taken. + releaseLock = await acquireFileLock(absoluteFilePath) // Variables to hold the actual paths of temp files if they are created. let actualTempNewFilePath: string | null = null @@ -247,6 +230,4 @@ async function _streamDataToFile(targetPath: string, data: any, prettyPrint = fa }) } -export const LOCK_STALE_MS = 31_000 - export { safeWriteJson }