diff --git a/docs/architecture/task-lifecycle-model.md b/docs/architecture/task-lifecycle-model.md index 355b08aa7e..55611d5166 100644 --- a/docs/architecture/task-lifecycle-model.md +++ b/docs/architecture/task-lifecycle-model.md @@ -51,10 +51,14 @@ TLA+/PlusCal or Quint with TLC becomes a better fit when the lifecycle needs tem | 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. +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. Runtime ownership is modeled separately from persisted status: an `owner-loss` fault can remove the live owner of an active or delegated child without changing its history record, matching process termination or session skip. Recovery is enabled only when the awaited delegation chain has no live owner and ends in an interrupted/completed task, a delegated task with no `awaitingChildId`, or a missing awaited record whose task ID also has no live owner. 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, nested delegation, and owner-loss recovery even when the raw state total changes. 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. +The `recover-active` action models startup repair after an active child loses its owner. `recoverDelegationParent` clears the parent's child pointers but interrupts it when an ancestor still awaits it; a top-level parent becomes active. The repair journal accepts both target statuses so a crash between the child and parent writes cannot break that ancestor link. Filesystem tests cover journal replay and subsequent re-delegation. Provider tests separately cover registration delayed by the recovery reservation, including disposal or cancellation before registration and preventing later scheduling. + +For [#1624](https://github.com/Zoo-Code-Org/Zoo-Code/issues/1624), `recoverDeadDelegatedChild` repairs a nested delegated child whose descendant chain ends dead. Startup reconciliation and runtime re-delegation check provider-wide ownership while reserving registration through persistence; persisted delegated status alone is not evidence of a live session. Runtime re-delegation refreshes the chain with `TaskHistoryStore.refreshStrict`, so an unreadable or invalid descendant record aborts recovery instead of looking missing. If disposal cancels a delegation after the parent commit, rollback releases the parent through `recoverDelegationParent` and restores its pending action, without recreating a task on the disposed provider. + ## Shared-store concurrency model The same `pnpm lifecycle:model-check` command also runs a second bounded explorer over two `TaskHistoryStore` hosts. It imports the production `computeHistoryDelta` and `mergeHistoryDelta` functions, so its semantics match the store rather than assuming coherent caches or transactional pair writes: @@ -133,7 +137,8 @@ The task delegation checker currently enforces: 4. Every active or delegated linked child is the child its parent currently awaits. An interrupted prior child may retain lineage after re-delegation but cannot complete back into that parent. 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. +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. A delegated child may transition to `interrupted` only through dead-chain recovery after runtime liveness checks establish that neither it nor its descendants has a live owner. +8. Runtime owner loss does not mutate persisted status. A delegated chain with no remaining live owner has a reachable recovery transition within the bounded graph, after which its parent can re-delegate. The completion persistence checker additionally enforces: diff --git a/scripts/check-task-lifecycle.ts b/scripts/check-task-lifecycle.ts index 73e9078366..c51ea7c822 100644 --- a/scripts/check-task-lifecycle.ts +++ b/scripts/check-task-lifecycle.ts @@ -7,11 +7,14 @@ import { completeDelegatedChild, delegateTaskToChild, interruptDelegatedChild, + isDeadDelegationChain, + recoverDeadDelegatedChild, + recoverDelegationParent, } from "../src/core/task-persistence/taskLifecycle" const taskIds = ["parent", "child-a", "child-b"] as const type TaskId = (typeof taskIds)[number] -type ModelState = Record +type ModelState = Record & { liveTaskIds: TaskId[] } interface Transition { name: string @@ -25,7 +28,15 @@ interface TraceStep { const MAX_DEPTH = 12 const MAX_STATES = 10_000 -const expectedActions = ["delegate", "interrupt", "complete", "abandon"] as const +const expectedActions = [ + "delegate", + "owner-loss", + "interrupt", + "recover-active", + "recover", + "complete", + "abandon", +] as const const semanticLandmarks = { "interrupted-child-redelegation": (state: ModelState) => state.parent?.status === "delegated" && @@ -36,6 +47,14 @@ const semanticLandmarks = { state.parent.awaitingChildId === "child-a" && state["child-a"]?.status === "delegated" && state["child-a"].awaitingChildId === "child-b", + "delegated-owner-loss": (state: ModelState) => + state["child-a"]?.status === "delegated" && !state.liveTaskIds.includes("child-a"), + "dead-nested-chain-recovered": (state: ModelState) => + state.parent?.status === "delegated" && + state.parent.awaitingChildId === "child-a" && + state["child-a"]?.status === "interrupted" && + state["child-a"].awaitingChildId === undefined && + state["child-b"]?.status === "interrupted", } satisfies Record boolean> function task(id: TaskId, parentTaskId?: TaskId): HistoryItem { @@ -55,7 +74,7 @@ function task(id: TaskId, parentTaskId?: TaskId): HistoryItem { } function initialState(): ModelState { - return { parent: task("parent"), "child-a": undefined, "child-b": undefined } + return { parent: task("parent"), "child-a": undefined, "child-b": undefined, liveTaskIds: ["parent"] } } function replace(state: ModelState, ...updates: HistoryItem[]): ModelState { @@ -64,6 +83,10 @@ function replace(state: ModelState, ...updates: HistoryItem[]): ModelState { return next } +function withLiveTasks(state: ModelState, ...liveTaskIds: TaskId[]): ModelState { + return { ...state, liveTaskIds: Array.from(new Set(liveTaskIds)).sort() } +} + function transitions(state: ModelState): Transition[] { const result: Transition[] = [] for (const parentId of taskIds) { @@ -77,9 +100,20 @@ function transitions(state: ModelState): Transition[] { continue } const delegated = delegateTaskToChild(parent, childId, awaitedStatus) + const next = replace(state, delegated, task(childId, parentId)) result.push({ name: `delegate(${parentId}, ${childId})`, - next: replace(state, delegated, task(childId, parentId)), + next: withLiveTasks(next, ...state.liveTaskIds.filter((id) => id !== parentId), childId), + }) + } + } + + for (const taskId of state.liveTaskIds) { + const current = state[taskId] + if (current && current.parentTaskId && (current.status === "active" || current.status === "delegated")) { + result.push({ + name: `owner-loss(${taskId})`, + next: withLiveTasks(state, ...state.liveTaskIds.filter((id) => id !== taskId)), }) } } @@ -92,7 +126,30 @@ function transitions(state: ModelState): Transition[] { if (parent.status === "delegated" && parent.awaitingChildId === child.id && child.status === "active") { const interrupted = interruptDelegatedChild(parent, child) - result.push({ name: `interrupt(${childId})`, next: replace(state, interrupted) }) + result.push({ + name: `interrupt(${childId})`, + next: withLiveTasks(replace(state, interrupted), ...state.liveTaskIds.filter((id) => id !== childId)), + }) + if (!state.liveTaskIds.includes(childId) && !state.liveTaskIds.includes(parent.id as TaskId)) { + const ancestor = parent.parentTaskId ? state[parent.parentTaskId as TaskId] : undefined + result.push({ + name: `recover-active(${childId})`, + next: replace(state, interrupted, recoverDelegationParent(parent, ancestor)), + }) + } + } + + if ( + parent.status === "delegated" && + parent.awaitingChildId === child.id && + isDeadDelegationChain( + child, + (id) => state[id as TaskId], + (id) => state.liveTaskIds.includes(id as TaskId), + ) + ) { + const recovered = recoverDeadDelegatedChild(parent, child) + result.push({ name: `recover(${childId})`, next: replace(state, recovered) }) } if ( @@ -103,7 +160,11 @@ function transitions(state: ModelState): Transition[] { const completed = completeDelegatedChild(parent, child, `${childId} result`) result.push({ name: `complete(${childId})`, - next: replace(state, completed.parent, completed.child), + next: withLiveTasks( + replace(state, completed.parent, completed.child), + ...state.liveTaskIds.filter((id) => id !== childId), + child.parentTaskId as TaskId, + ), }) } @@ -111,12 +172,33 @@ function transitions(state: ModelState): Transition[] { const abandoned = abandonDelegatedChild(parent, child) result.push({ name: `abandon(${childId})`, - next: replace(state, abandoned.parent, abandoned.child), + next: withLiveTasks( + replace(state, abandoned.parent, abandoned.child), + ...state.liveTaskIds.filter((id) => id !== childId), + child.parentTaskId as TaskId, + ), }) } } + return result } +function deadDelegatedChildren(state: ModelState): TaskId[] { + return taskIds.filter((childId) => { + const child = state[childId] + if (!child?.parentTaskId) return false + const parent = state[child.parentTaskId as TaskId] + return ( + parent?.status === "delegated" && + parent.awaitingChildId === child.id && + isDeadDelegationChain( + child, + (id) => state[id as TaskId], + (id) => state.liveTaskIds.includes(id as TaskId), + ) + ) + }) +} function invariantViolations(state: ModelState): string[] { const violations: string[] = [] @@ -158,11 +240,16 @@ function invariantViolations(state: ModelState): string[] { cursor = state[cursor as TaskId]?.parentTaskId } } + for (const childId of deadDelegatedChildren(state)) { + if (!transitions(state).some((transition) => transition.name === `recover(${childId})`)) { + violations.push(`${childId}: dead delegated chain must be recoverable in the next transition`) + } + } return violations } function canonical(state: ModelState): string { - return JSON.stringify(taskIds.map((id) => state[id] ?? null)) + return JSON.stringify({ tasks: taskIds.map((id) => state[id] ?? null), liveTaskIds: state.liveTaskIds }) } function formatCounterexample(message: string, trace: TraceStep[]): string { @@ -274,6 +361,13 @@ function runRepresentativeScenarios(): void { const nestedCompletion = completeDelegatedChild(nestedParent, childB, "nested result") assert.equal(nestedCompletion.parent.status, "active") assert.equal(nestedCompletion.parent.completedByChildId, childB.id) + const interruptedNestedChild = interruptDelegatedChild(nestedParent, childB) + const recoveredNestedParent = recoverDeadDelegatedChild(delegated, { + ...nestedParent, + awaitingChildId: interruptedNestedChild.id, + }) + assert.equal(recoveredNestedParent.status, "interrupted") + assert.equal(recoveredNestedParent.awaitingChildId, undefined) const interruptedCompletion = completeDelegatedChild(delegated, interruptedA, "resumed result") assert.equal(interruptedCompletion.child.status, "completed") diff --git a/src/__tests__/ClineProvider.delegation.spec.ts b/src/__tests__/ClineProvider.delegation.spec.ts index 422c264e2c..05c229830f 100644 --- a/src/__tests__/ClineProvider.delegation.spec.ts +++ b/src/__tests__/ClineProvider.delegation.spec.ts @@ -21,6 +21,7 @@ function makeStoreStub( ) { return { invalidate: vi.fn().mockResolvedValue(undefined), + refreshStrict: vi.fn().mockResolvedValue(undefined), atomicReadAndUpdate: vi.fn(async (_taskId: string, updater: (h: HistoryItem) => HistoryItem) => { updater(parentHistoryItem) return [] @@ -90,6 +91,64 @@ describe("ClineProvider.delegateParentAndOpenChild()", () => { expect(taskHistoryStore.atomicReadAndUpdate).not.toHaveBeenCalled() }) + it("releases the committed parent when disposal cancels delegation after the commit", async () => { + const pendingAction = { + kind: "create_subtask" as const, + actionId: "create-action", + approvalText: "{}", + mode: "code", + message: "Do something", + todos: [], + } + let current: HistoryItem = { ...parentHistoryItem, status: "active", pendingAction } + const child = { taskId: "child-1", run: vi.fn().mockResolvedValue(undefined) } + const createTaskWithHistoryItem = vi.fn() + const deleteTaskWithId = vi.fn().mockResolvedValue(undefined) + const provider = { + _disposed: false, + taskScheduler: new TaskScheduler(), + emit: vi.fn(), + getCurrentTask: vi.fn(() => makeParentTask()), + removeClineFromStack: vi.fn().mockResolvedValue(undefined), + createTask: vi.fn().mockResolvedValue(child), + handleModeSwitch: vi.fn().mockResolvedValue(undefined), + deleteTaskWithId, + createTaskWithHistoryItem, + log: vi.fn(), + isViewLaunched: false, + } as unknown as ClineProvider + const parentTask = makeParentTask() + vi.mocked(provider.getCurrentTask).mockReturnValue(parentTask) + Reflect.set(provider, "taskHistoryStore", { + invalidate: vi.fn().mockResolvedValue(undefined), + get: vi.fn(() => current), + atomicReadAndUpdate: vi.fn(async (_taskId: string, updater: (item: HistoryItem) => HistoryItem) => { + current = updater(current) + // Disposal starts right after the parent commit lands. + Reflect.set(provider, "_disposed", true) + return [current] + }), + }) + + await expect( + ClineProvider.prototype.delegateParentAndOpenChild.call(provider, { + parentTaskId: "parent-1", + message: "Do something", + initialTodos: [], + mode: "code", + pendingActionId: "create-action", + }), + ).rejects.toThrow("disposed before child scheduling") + + expect(child.run).not.toHaveBeenCalled() + expect(deleteTaskWithId).toHaveBeenCalledWith("child-1", false) + // The parent no longer awaits the deleted child and keeps its unresolved action. + expect(current).toMatchObject({ status: "active", awaitingChildId: undefined, delegatedToId: undefined }) + expect(current.pendingAction).toEqual(pendingAction) + // A disposed provider must not recreate the parent task. + expect(createTaskWithHistoryItem).not.toHaveBeenCalled() + }) + it("clears a matching pending action when delegation commits", async () => { const pendingAction = { kind: "create_subtask" as const, @@ -519,11 +578,12 @@ describe("ClineProvider.delegateParentAndOpenChild()", () => { return [] }), }) + const parentTask = makeParentTask() const provider = { taskScheduler: new TaskScheduler(), emit: vi.fn(), - getCurrentTask: vi.fn(() => makeParentTask()), + getCurrentTask: vi.fn(() => parentTask), removeClineFromStack: vi.fn().mockResolvedValue(undefined), createTask: vi.fn().mockResolvedValue({ taskId: "child-2", start: vi.fn(), run: () => Promise.resolve() }), handleModeSwitch: vi.fn().mockResolvedValue(undefined), @@ -553,6 +613,555 @@ describe("ClineProvider.delegateParentAndOpenChild()", () => { expect(result.childIds).toContain("child-2") }) + it("recovers a dead nested chain and re-delegates the parent", async () => { + const oldChildId = "old-child" + const grandchildId = "grandchild" + const records = new Map([ + [ + "parent-1", + { + ...parentHistoryItem, + status: "delegated", + awaitingChildId: oldChildId, + delegatedToId: oldChildId, + childIds: [oldChildId], + }, + ], + [ + oldChildId, + { + ...parentHistoryItem, + id: oldChildId, + status: "delegated", + parentTaskId: "parent-1", + awaitingChildId: grandchildId, + delegatedToId: grandchildId, + }, + ], + [ + grandchildId, + { + ...parentHistoryItem, + id: grandchildId, + status: "interrupted", + parentTaskId: oldChildId, + }, + ], + ]) + const taskHistoryStore = makeStoreStub({ + get: vi.fn((id: string) => records.get(id)), + atomicReadAndUpdate: vi.fn(async (taskId: string, updater: (item: HistoryItem) => HistoryItem) => { + const updated = updater(records.get(taskId)!) + records.set(taskId, updated) + return [] + }), + }) + const parentTask = makeParentTask() + const provider = { + taskScheduler: new TaskScheduler(), + taskRegistry: { hasRunning: vi.fn().mockReturnValue(false) }, + emit: vi.fn(), + getCurrentTask: vi.fn(() => parentTask), + removeClineFromStack: vi.fn().mockResolvedValue(undefined), + createTask: vi.fn().mockResolvedValue({ taskId: "child-2", run: vi.fn().mockResolvedValue(undefined) }), + log: vi.fn(), + isViewLaunched: false, + recentTasksCache: undefined, + taskHistoryStore, + } as unknown as ClineProvider + Object.assign(provider, { + isTaskRunningInAnyProvider: ClineProvider.prototype["isTaskRunningInAnyProvider"], + refreshDelegationChain: ClineProvider.prototype["refreshDelegationChain"], + recoverDeadAwaitedChild: ClineProvider.prototype["recoverDeadAwaitedChild"], + }) + + await ClineProvider.prototype.delegateParentAndOpenChild.call(provider, { + parentTaskId: "parent-1", + message: "Continue", + initialTodos: [], + mode: "code", + }) + + expect(records.get(oldChildId)).toMatchObject({ + status: "interrupted", + awaitingChildId: undefined, + delegatedToId: undefined, + }) + expect(records.get("parent-1")).toMatchObject({ status: "delegated", awaitingChildId: "child-2" }) + }) + + it("does not recover a nested chain with a live task owner", async () => { + const oldChildId = "old-child" + const grandchildId = "grandchild" + const parent = { + ...parentHistoryItem, + status: "delegated" as const, + awaitingChildId: oldChildId, + delegatedToId: oldChildId, + } + const child = { + ...parentHistoryItem, + id: oldChildId, + status: "delegated" as const, + awaitingChildId: grandchildId, + delegatedToId: grandchildId, + } + const grandchild = { ...parentHistoryItem, id: grandchildId, status: "interrupted" as const } + const taskHistoryStore = makeStoreStub({ + get: vi.fn((id: string) => (id === "parent-1" ? parent : id === oldChildId ? child : grandchild)), + }) + const provider = { + taskRegistry: { hasRunning: vi.fn((id: string) => id === grandchildId) }, + getCurrentTask: vi.fn(() => makeParentTask()), + removeClineFromStack: vi.fn(), + createTask: vi.fn(), + taskHistoryStore, + } as unknown as ClineProvider + Object.assign(provider, { + isTaskRunningInAnyProvider: ClineProvider.prototype["isTaskRunningInAnyProvider"], + refreshDelegationChain: ClineProvider.prototype["refreshDelegationChain"], + recoverDeadAwaitedChild: ClineProvider.prototype["recoverDeadAwaitedChild"], + }) + await expect( + ClineProvider.prototype.delegateParentAndOpenChild.call(provider, { + parentTaskId: "parent-1", + message: "Continue", + initialTodos: [], + mode: "code", + }), + ).rejects.toThrow(`awaited child ${oldChildId} has status delegated`) + expect(taskHistoryStore.atomicReadAndUpdate).not.toHaveBeenCalled() + expect(provider.createTask).not.toHaveBeenCalled() + }) + + it("detects a live delegation owner in another provider instance", () => { + const provider = { + taskRegistry: { hasRunning: vi.fn().mockReturnValue(false) }, + } as unknown as ClineProvider + const otherProvider = { + taskRegistry: { hasRunning: vi.fn((taskId: string) => taskId === "live-child") }, + } as unknown as ClineProvider + const activeInstances = Reflect.get(ClineProvider, "activeInstances") + if (!(activeInstances instanceof Set)) throw new Error("ClineProvider active instance registry unavailable") + activeInstances.add(otherProvider) + + try { + expect(ClineProvider.prototype["isTaskRunningInAnyProvider"].call(provider, "live-child")).toBe(true) + } finally { + activeInstances.delete(otherProvider) + } + }) + + it("rejects when a dead delegation chain becomes live during its atomic recovery", async () => { + const oldChildId = "old-child" + const grandchildId = "grandchild" + const parent = { + ...parentHistoryItem, + status: "delegated" as const, + awaitingChildId: oldChildId, + delegatedToId: oldChildId, + } + const child = { + ...parentHistoryItem, + id: oldChildId, + status: "delegated" as const, + awaitingChildId: grandchildId, + delegatedToId: grandchildId, + } + const grandchild = { ...parentHistoryItem, id: grandchildId, status: "interrupted" as const } + let childIsLive = false + const taskHistoryStore = makeStoreStub({ + get: vi.fn((id: string) => (id === "parent-1" ? parent : id === oldChildId ? child : grandchild)), + atomicReadAndUpdate: vi.fn(async (_taskId: string, updater: (item: HistoryItem) => HistoryItem) => { + childIsLive = true + updater(child) + return [] + }), + }) + const provider = { + taskRegistry: { hasRunning: vi.fn((id: string) => id === oldChildId && childIsLive) }, + getCurrentTask: vi.fn(() => makeParentTask()), + removeClineFromStack: vi.fn(), + createTask: vi.fn(), + taskHistoryStore, + } as unknown as ClineProvider + Object.assign(provider, { + isTaskRunningInAnyProvider: ClineProvider.prototype["isTaskRunningInAnyProvider"], + refreshDelegationChain: ClineProvider.prototype["refreshDelegationChain"], + recoverDeadAwaitedChild: ClineProvider.prototype["recoverDeadAwaitedChild"], + }) + + await expect( + ClineProvider.prototype.delegateParentAndOpenChild.call(provider, { + parentTaskId: "parent-1", + message: "Continue", + initialTodos: [], + mode: "code", + }), + ).rejects.toThrow(`Delegation chain for child ${oldChildId} became live during recovery`) + expect(provider.createTask).not.toHaveBeenCalled() + }) + + it("delays task registration until runtime recovery persistence finishes", async () => { + let releasePersistence!: () => void + const persistenceBlocked = new Promise((resolve) => { + releasePersistence = resolve + }) + const parent = { ...parentHistoryItem, status: "delegated" as const, awaitingChildId: "old-child" } + const child = { + ...parentHistoryItem, + id: "old-child", + status: "delegated" as const, + awaitingChildId: "grandchild", + } + const grandchild = { ...parentHistoryItem, id: "grandchild", status: "interrupted" as const } + const taskHistoryStore = makeStoreStub({ + get: vi.fn((id: string) => (id === child.id ? child : grandchild)), + atomicReadAndUpdate: vi.fn(async (_id: string, updater: (item: HistoryItem) => HistoryItem) => { + updater(child) + await persistenceBlocked + return [] + }), + }) + const taskRegistry = { hasRunning: vi.fn().mockReturnValue(false), push: vi.fn() } + const provider = { + taskRegistry, + taskHistoryStore, + log: vi.fn(), + throwIfTaskRegistrationCancelled: ClineProvider.prototype["throwIfTaskRegistrationCancelled"], + performPreparationTasks: vi.fn(), + getState: vi.fn().mockResolvedValue({ mode: "code" }), + } as unknown as ClineProvider + Object.assign(provider, { + isTaskRunningInAnyProvider: ClineProvider.prototype["isTaskRunningInAnyProvider"], + refreshDelegationChain: ClineProvider.prototype["refreshDelegationChain"], + }) + + const recovery = ClineProvider.prototype["recoverDeadAwaitedChild"].call(provider, parent, child.id) + await vi.waitFor(() => expect(taskHistoryStore.atomicReadAndUpdate).toHaveBeenCalled()) + const registration = ClineProvider.prototype.addClineToStack.call(provider, { + taskId: child.id, + emit: vi.fn(), + } as never) + await Promise.resolve() + expect(taskRegistry.push).not.toHaveBeenCalled() + + releasePersistence() + await Promise.all([recovery, registration]) + expect(taskRegistry.push).toHaveBeenCalledOnce() + }) + + it("does not recover when a descendant is live in another provider", async () => { + const oldChildId = "old-child" + const grandchildId = "grandchild" + const parent = { + ...parentHistoryItem, + status: "delegated" as const, + awaitingChildId: oldChildId, + delegatedToId: oldChildId, + } + const child = { + ...parentHistoryItem, + id: oldChildId, + status: "delegated" as const, + awaitingChildId: grandchildId, + delegatedToId: grandchildId, + } + const grandchild = { ...parentHistoryItem, id: grandchildId, status: "active" as const } + const taskHistoryStore = makeStoreStub({ + get: vi.fn((id: string) => (id === oldChildId ? child : id === grandchildId ? grandchild : parent)), + }) + const provider = { + taskRegistry: { hasRunning: vi.fn().mockReturnValue(false) }, + taskHistoryStore, + } as unknown as ClineProvider + const otherProvider = { + taskRegistry: { hasRunning: vi.fn((id: string) => id === grandchildId) }, + } as unknown as ClineProvider + const activeInstances = Reflect.get(ClineProvider, "activeInstances") + if (!(activeInstances instanceof Set)) throw new Error("ClineProvider active instance registry unavailable") + activeInstances.add(otherProvider) + Object.assign(provider, { + isTaskRunningInAnyProvider: ClineProvider.prototype["isTaskRunningInAnyProvider"], + refreshDelegationChain: ClineProvider.prototype["refreshDelegationChain"], + }) + + try { + const result = await ClineProvider.prototype["recoverDeadAwaitedChild"].call(provider, parent, oldChildId) + expect(result).toBe(child) + expect(taskHistoryStore.atomicReadAndUpdate).not.toHaveBeenCalled() + } finally { + activeInstances.delete(otherProvider) + } + }) + + it("uses refreshed delegation records when deciding recovery", async () => { + const oldChildId = "old-child" + const grandchildId = "grandchild" + const parent = { + ...parentHistoryItem, + status: "delegated" as const, + awaitingChildId: oldChildId, + delegatedToId: oldChildId, + } + const refreshedChild = { + ...parentHistoryItem, + id: oldChildId, + status: "delegated" as const, + awaitingChildId: grandchildId, + delegatedToId: grandchildId, + } + const refreshedGrandchild = { ...parentHistoryItem, id: grandchildId, status: "interrupted" as const } + const cache = new Map([[oldChildId, { ...refreshedChild, status: "active" }]]) + const taskHistoryStore = makeStoreStub({ + get: vi.fn((id: string) => cache.get(id)), + atomicReadAndUpdate: vi.fn(async (id: string, updater: (item: HistoryItem) => HistoryItem) => { + cache.set(id, updater(cache.get(id)!)) + return [] + }), + }) + taskHistoryStore.refreshStrict.mockImplementation(async (id: string) => { + if (id === oldChildId) cache.set(id, refreshedChild) + if (id === grandchildId) cache.set(id, refreshedGrandchild) + }) + const provider = { + taskRegistry: { hasRunning: vi.fn().mockReturnValue(false) }, + taskHistoryStore, + log: vi.fn(), + } as unknown as ClineProvider + Object.assign(provider, { + isTaskRunningInAnyProvider: ClineProvider.prototype["isTaskRunningInAnyProvider"], + refreshDelegationChain: ClineProvider.prototype["refreshDelegationChain"], + }) + + await ClineProvider.prototype["recoverDeadAwaitedChild"].call(provider, parent, oldChildId) + + expect(taskHistoryStore.refreshStrict).toHaveBeenCalledWith(oldChildId) + expect(taskHistoryStore.refreshStrict).toHaveBeenCalledWith(grandchildId) + expect(cache.get(oldChildId)).toMatchObject({ status: "interrupted", awaitingChildId: undefined }) + }) + + it("fails closed when a descendant record cannot be read during recovery", async () => { + const oldChildId = "old-child" + const grandchildId = "grandchild" + const parent = { + ...parentHistoryItem, + status: "delegated" as const, + awaitingChildId: oldChildId, + delegatedToId: oldChildId, + } + const child = { + ...parentHistoryItem, + id: oldChildId, + status: "delegated" as const, + awaitingChildId: grandchildId, + delegatedToId: grandchildId, + } + const taskHistoryStore = makeStoreStub({ + get: vi.fn((id: string) => (id === oldChildId ? child : undefined)), + }) + taskHistoryStore.refreshStrict.mockImplementation(async (id: string) => { + if (id === grandchildId) throw Object.assign(new Error("EACCES: permission denied"), { code: "EACCES" }) + }) + const provider = { + taskRegistry: { hasRunning: vi.fn().mockReturnValue(false) }, + taskHistoryStore, + log: vi.fn(), + } as unknown as ClineProvider + Object.assign(provider, { + isTaskRunningInAnyProvider: ClineProvider.prototype["isTaskRunningInAnyProvider"], + refreshDelegationChain: ClineProvider.prototype["refreshDelegationChain"], + }) + + await expect( + ClineProvider.prototype["recoverDeadAwaitedChild"].call(provider, parent, oldChildId), + ).rejects.toThrow("EACCES") + expect(taskHistoryStore.atomicReadAndUpdate).not.toHaveBeenCalled() + }) + + it("does not recover when a missing descendant record is still live in a provider", async () => { + const oldChildId = "old-child" + const grandchildId = "grandchild" + const parent = { + ...parentHistoryItem, + status: "delegated" as const, + awaitingChildId: oldChildId, + delegatedToId: oldChildId, + } + const child = { + ...parentHistoryItem, + id: oldChildId, + status: "delegated" as const, + awaitingChildId: grandchildId, + delegatedToId: grandchildId, + } + const taskHistoryStore = makeStoreStub({ + // The grandchild record is absent, but its task still runs. + get: vi.fn((id: string) => (id === oldChildId ? child : undefined)), + }) + const provider = { + taskRegistry: { hasRunning: vi.fn((id: string) => id === grandchildId) }, + taskHistoryStore, + log: vi.fn(), + } as unknown as ClineProvider + Object.assign(provider, { + isTaskRunningInAnyProvider: ClineProvider.prototype["isTaskRunningInAnyProvider"], + refreshDelegationChain: ClineProvider.prototype["refreshDelegationChain"], + }) + + await expect( + ClineProvider.prototype["recoverDeadAwaitedChild"].call(provider, parent, oldChildId), + ).resolves.toMatchObject({ status: "delegated" }) + expect(taskHistoryStore.atomicReadAndUpdate).not.toHaveBeenCalled() + }) + + it("rejects a missing awaited-child record without creating a child", async () => { + const oldChildId = "missing-child" + const parent = { + ...parentHistoryItem, + status: "delegated" as const, + awaitingChildId: oldChildId, + delegatedToId: oldChildId, + } + const taskHistoryStore = makeStoreStub({ + get: vi.fn((id: string) => (id === "parent-1" ? parent : undefined)), + }) + const provider = { + getCurrentTask: vi.fn(() => makeParentTask()), + createTask: vi.fn(), + taskHistoryStore, + } as unknown as ClineProvider + + await expect( + ClineProvider.prototype.delegateParentAndOpenChild.call(provider, { + parentTaskId: "parent-1", + message: "Continue", + initialTodos: [], + mode: "code", + }), + ).rejects.toThrow(`awaited child ${oldChildId} has status missing`) + expect(provider.createTask).not.toHaveBeenCalled() + expect(taskHistoryStore.atomicReadAndUpdate).not.toHaveBeenCalled() + }) + + it("does not create a child when the provider is disposed during recovery", async () => { + const oldChildId = "old-child" + const grandchildId = "grandchild" + const parentHistory = { + ...parentHistoryItem, + status: "delegated" as const, + awaitingChildId: oldChildId, + delegatedToId: oldChildId, + } + const oldChild = { + ...parentHistoryItem, + id: oldChildId, + status: "delegated" as const, + awaitingChildId: grandchildId, + delegatedToId: grandchildId, + } + const grandchild = { ...parentHistoryItem, id: grandchildId, status: "interrupted" as const } + const records = new Map([ + ["parent-1", parentHistory], + [oldChildId, oldChild], + [grandchildId, grandchild], + ]) + const taskHistoryStore = makeStoreStub({ + get: vi.fn((id: string) => records.get(id)), + atomicReadAndUpdate: vi.fn(async (id: string, updater: (item: HistoryItem) => HistoryItem) => { + records.set(id, updater(records.get(id)!)) + return [] + }), + }) + const parentTask = makeParentTask() + const provider = { + _disposed: false, + taskRegistry: { hasRunning: vi.fn().mockReturnValue(false) }, + getCurrentTask: vi.fn(() => parentTask), + createTask: vi.fn(), + taskHistoryStore, + log: vi.fn(), + } as unknown as ClineProvider + taskHistoryStore.atomicReadAndUpdate.mockImplementation( + async (id: string, updater: (item: HistoryItem) => HistoryItem) => { + records.set(id, updater(records.get(id)!)) + Reflect.set(provider, "_disposed", true) + return [] + }, + ) + Object.assign(provider, { + isTaskRunningInAnyProvider: ClineProvider.prototype["isTaskRunningInAnyProvider"], + refreshDelegationChain: ClineProvider.prototype["refreshDelegationChain"], + recoverDeadAwaitedChild: ClineProvider.prototype["recoverDeadAwaitedChild"], + }) + + await expect( + ClineProvider.prototype.delegateParentAndOpenChild.call(provider, { + parentTaskId: "parent-1", + message: "Continue", + initialTodos: [], + mode: "code", + }), + ).rejects.toThrow("Provider was disposed during delegation") + expect(provider.createTask).not.toHaveBeenCalled() + }) + + it("does not create a child when the parent is cancelled during recovery", async () => { + const oldChildId = "old-child" + const grandchildId = "grandchild" + const parentHistory = { + ...parentHistoryItem, + status: "delegated" as const, + awaitingChildId: oldChildId, + delegatedToId: oldChildId, + } + const oldChild = { + ...parentHistoryItem, + id: oldChildId, + status: "delegated" as const, + awaitingChildId: grandchildId, + delegatedToId: grandchildId, + } + const grandchild = { ...parentHistoryItem, id: grandchildId, status: "interrupted" as const } + const records = new Map([ + ["parent-1", parentHistory], + [oldChildId, oldChild], + [grandchildId, grandchild], + ]) + const taskHistoryStore = makeStoreStub({ get: vi.fn((id: string) => records.get(id)) }) + const parentTask = makeParentTask() + const provider = { + _disposed: false, + taskRegistry: { hasRunning: vi.fn().mockReturnValue(false) }, + getCurrentTask: vi.fn(() => parentTask), + createTask: vi.fn(), + taskHistoryStore, + log: vi.fn(), + } as unknown as ClineProvider + taskHistoryStore.atomicReadAndUpdate.mockImplementation( + async (id: string, updater: (item: HistoryItem) => HistoryItem) => { + records.set(id, updater(records.get(id)!)) + parentTask.abort = true + return [] + }, + ) + Object.assign(provider, { + isTaskRunningInAnyProvider: ClineProvider.prototype["isTaskRunningInAnyProvider"], + refreshDelegationChain: ClineProvider.prototype["refreshDelegationChain"], + recoverDeadAwaitedChild: ClineProvider.prototype["recoverDeadAwaitedChild"], + }) + + await expect( + ClineProvider.prototype.delegateParentAndOpenChild.call(provider, { + parentTaskId: "parent-1", + message: "Continue", + initialTodos: [], + mode: "code", + }), + ).rejects.toThrow("Parent parent-1 was cancelled during delegation") + expect(provider.createTask).not.toHaveBeenCalled() + }) + it("rejects with 'Cannot re-delegate' when the existing awaited child is still active", async () => { const oldChildId = "old-child" const activeChild = { id: oldChildId, status: "active" } as unknown as HistoryItem @@ -604,7 +1213,7 @@ describe("ClineProvider.delegateParentAndOpenChild()", () => { initialTodos: [], mode: "code", }), - ).rejects.toThrow("Cannot re-delegate while the awaited child is not interrupted") + ).rejects.toThrow(`awaited child ${oldChildId} has status active`) // The authoritative preflight rejects before either provider mutates its stack. expect(child.run).not.toHaveBeenCalled() diff --git a/src/core/task-persistence/TaskHistoryStore.ts b/src/core/task-persistence/TaskHistoryStore.ts index 3d4cc47604..eaeeee81ff 100644 --- a/src/core/task-persistence/TaskHistoryStore.ts +++ b/src/core/task-persistence/TaskHistoryStore.ts @@ -9,12 +9,30 @@ import type { HistoryItem } from "@roo-code/types" import { GlobalFileNames } from "../../shared/globalFileNames" import { LOCK_STALE_MS, safeWriteJson } from "../../utils/safeWriteJson" import { getStorageBasePath } from "../../utils/storage" -import { assertValidTransition, type HistoryItemStatus } from "./taskLifecycle" +import { + assertValidTransition, + isDeadDelegationChain, + recoverDeadDelegatedChild, + recoverDelegationParent, + type HistoryItemStatus, +} from "./taskLifecycle" import { computeHistoryDelta, DeltaRejectedError, mergeHistoryDelta } from "./taskStoreConcurrency" export { assertValidTransition, type HistoryItemStatus } from "./taskLifecycle" export { DeltaRejectedError } from "./taskStoreConcurrency" +let taskOwnershipReservation: Promise = Promise.resolve() + +/** Serialize runtime task registration with dead-delegation recovery. */ +export function withTaskOwnershipReservation(fn: () => Promise): Promise { + const result = taskOwnershipReservation.then(fn, fn) + taskOwnershipReservation = result.then( + () => undefined, + () => undefined, + ) + return result +} + /** * Build a `safeWriteJson` merge callback that applies only `delta` to the * current disk state, preserving fields written by another process. @@ -47,7 +65,7 @@ interface DelegationRepairIntent { } target: { childStatus: "interrupted" - parentStatus: "active" + parentStatus: "active" | "interrupted" } } @@ -75,11 +93,14 @@ export interface TaskHistoryStoreOptions { * globalState during the transition period. */ onWrite?: (items: HistoryItem[]) => Promise + /** Return whether any provider currently owns a live runtime task. */ + isTaskOwned?: (taskId: string) => boolean } export class TaskHistoryStore { private readonly globalStoragePath: string private readonly onWrite?: (items: HistoryItem[]) => Promise + private readonly isTaskOwned: (taskId: string) => boolean private cache: Map = new Map() private taskFileMtimes: Map = new Map() private writeLock: Promise = Promise.resolve() @@ -100,6 +121,7 @@ export class TaskHistoryStore { constructor(globalStoragePath: string, options?: TaskHistoryStoreOptions) { this.globalStoragePath = globalStoragePath this.onWrite = options?.onWrite + this.isTaskOwned = options?.isTaskOwned ?? (() => false) this.initialized = new Promise((resolve) => { this.resolveInitialized = resolve }) @@ -122,15 +144,17 @@ export class TaskHistoryStore { // as an orphaned active child in the same startup pass. const persistedActiveIds = this.getPersistedActiveIds() - // 2. Complete any two-record repair interrupted after its intent was durable. - try { - await this.replayDelegationRepairIntent() - } catch (error) { - console.error("[TaskHistoryStore] Failed to replay delegation repair intent:", error) - } + await withTaskOwnershipReservation(async () => { + // 2. Complete any two-record repair interrupted after its intent was durable. + try { + await this.replayDelegationRepairIntent() + } catch (error) { + console.error("[TaskHistoryStore] Failed to replay delegation repair intent:", error) + } - // 3. Repair delegation inconsistencies left by a previous crash - await this.reconcileDelegationState(persistedActiveIds) + // 3. Repair delegation inconsistencies left by a previous crash + await this.reconcileDelegationState(persistedActiveIds) + }) // 4. Start fs.watch for cross-instance reactivity this.startWatcher() @@ -400,11 +424,12 @@ export class TaskHistoryStore { * - Parent `delegated` with no `awaitingChildId` → parent → `active` (invalid state) * - Parent `delegated`, child not found → parent → `active` (orphaned delegation) * - Parent `delegated`, child `completed` → parent → `active` (interrupted handoff) - * - Parent `delegated`, child `active` → child → `interrupted`, parent → `active` + * - Parent `delegated`, orphaned persisted `active` child → child → `interrupted`; + * parent → `interrupted` if its ancestor still awaits it, otherwise → `active` + * - Parent `delegated`, dead nested delegated chain → child → `interrupted` * - * A parent awaiting an `interrupted` or `delegated` child is left as-is — the child is - * resumable. An `active` child is treated as orphaned during startup recovery because - * no live task session exists to own it. + * A parent awaiting an `interrupted` child retains its delegation link. Runtime + * ownership prevents recovery even when the persisted chain appears dead. */ private async reconcileDelegationState(persistedActiveIds: ReadonlySet): Promise { return this.withLock(() => this.reconcileDelegationStateCore(persistedActiveIds)) @@ -431,8 +456,9 @@ export class TaskHistoryStore { // are visible when evaluating chained delegations. const byId = new Map(Array.from(this.cache.values()).map((i) => [i.id, i])) - for (const [, item] of byId) { - if (item.status !== "delegated") { + for (const [, snapshotItem] of byId) { + const item = this.cache.get(snapshotItem.id) + if (item?.status !== "delegated") { continue } @@ -465,14 +491,28 @@ export class TaskHistoryStore { `[TaskHistoryStore] Reconciled orphaned delegation: task ${item.id} → active (child ${item.awaitingChildId} not found)`, ) repairsInThisPass++ - } else if ((child.status ?? "active") === "active" && persistedActiveIds.has(child.id)) { - // An active child persisted across startup cannot have a live task session - // behind it. Mark it interrupted before releasing the parent's delegation - // link so the normal resume/re-delegate flow can take over. This is an - // administrative recovery, not a runtime delegation transition. - await this.repairActiveDelegation(item, child) + } else if ( + child.status === "delegated" && + isDeadDelegationChain(child, (id) => byId.get(id), this.isTaskOwned) + ) { + const recoveredChild = recoverDeadDelegatedChild(item, child) + await this.upsertCore(recoveredChild) + console.warn( + `[TaskHistoryStore] Reconciled dead nested delegation: child ${child.id} → interrupted, task ${item.id} remains delegated`, + ) + repairsInThisPass++ + } else if ( + (child.status ?? "active") === "active" && + persistedActiveIds.has(child.id) && + !this.isTaskOwned(child.id) && + !this.isTaskOwned(item.id) + ) { + // Recover only persisted sessions without a runtime owner. Preserve any + // ancestor's delegation while clearing this parent's dead child link. + const ancestor = item.parentTaskId ? byId.get(item.parentTaskId) : undefined + await this.repairActiveDelegation(item, child, ancestor) console.warn( - `[TaskHistoryStore] Reconciled orphaned active child: child ${child.id} → interrupted, task ${item.id} → active`, + `[TaskHistoryStore] Reconciled orphaned active child: child ${child.id} → interrupted, task ${item.id} → ${this.cache.get(item.id)?.status}`, ) repairsInThisPass++ } else if (child.status === "completed") { @@ -525,6 +565,9 @@ export class TaskHistoryStore { if (!intent) { return } + if (this.isTaskOwned(intent.childTaskId) || this.isTaskOwned(intent.parentTaskId)) { + return + } const child = this.cache.get(intent.childTaskId) const parent = this.cache.get(intent.parentTaskId) @@ -577,7 +620,12 @@ export class TaskHistoryStore { * Start and complete a guarded active-child repair while already holding the * store lock. The intent is durable before either task file is touched. */ - private async repairActiveDelegation(parent: HistoryItem, child: HistoryItem): Promise { + private async repairActiveDelegation( + parent: HistoryItem, + child: HistoryItem, + ancestor?: HistoryItem, + ): Promise { + const repairedParent = recoverDelegationParent(parent, ancestor) const intent: DelegationRepairIntent = { version: 1, operationId: crypto.randomUUID(), @@ -595,7 +643,7 @@ export class TaskHistoryStore { rootTaskId: child.rootTaskId, }, }, - target: { childStatus: "interrupted", parentStatus: "active" }, + target: { childStatus: "interrupted", parentStatus: repairedParent.status }, } await this.writeDelegationRepairIntent(intent) @@ -697,7 +745,7 @@ export class TaskHistoryStore { (expectedChild.rootTaskId === undefined || typeof expectedChild.rootTaskId === "string") && !!targetRecord && targetRecord.childStatus === "interrupted" && - targetRecord.parentStatus === "active" + (targetRecord.parentStatus === "active" || targetRecord.parentStatus === "interrupted") ) } @@ -770,6 +818,34 @@ export class TaskHistoryStore { }) } + /** + * Re-read a single task for an ownership decision. Unlike `invalidate()`, only a missing file + * removes the cache entry; an unreadable, malformed, or mismatched record throws and keeps the + * cached entry, so callers fail closed instead of treating the task as absent. + */ + async refreshStrict(taskId: string): Promise { + return this.withLock(async () => { + const filePath = await this.getTaskFilePath(taskId) + let raw: string + try { + raw = await fs.readFile(filePath, "utf8") + } catch (error) { + if (error && typeof error === "object" && "code" in error && error.code === "ENOENT") { + this.cache.delete(taskId) + this.taskFileMtimes.delete(taskId) + return + } + throw error + } + const item: unknown = JSON.parse(raw) + if (typeof item !== "object" || item === null || (item as { id?: unknown }).id !== taskId) { + throw new Error(`Invalid task history record: ${filePath}`) + } + this.cache.set(taskId, item as HistoryItem) + this.taskFileMtimes.delete(taskId) + }) + } + /** * Clear all in-memory cache entries; a subsequent `reconcile()` repopulates them from task files. */ @@ -792,40 +868,42 @@ export class TaskHistoryStore { return } - await this.withLock(async () => { - const tasksDir = await this.getTasksDir() + await withTaskOwnershipReservation(() => + this.withLock(async () => { + const tasksDir = await this.getTasksDir() - for (const item of taskHistoryEntries) { - if (!item.id) { - continue - } + for (const item of taskHistoryEntries) { + if (!item.id) { + continue + } - // Check if task directory exists on disk - const taskDir = path.join(tasksDir, item.id) + // Check if task directory exists on disk + const taskDir = path.join(tasksDir, item.id) - try { - await fs.access(taskDir) - } catch { - // Task directory doesn't exist; skip this entry as it's orphaned in globalState - continue - } + try { + await fs.access(taskDir) + } catch { + // Task directory doesn't exist; skip this entry as it's orphaned in globalState + continue + } - // Write history_item.json if it doesn't exist yet - const filePath = path.join(taskDir, GlobalFileNames.historyItem) - try { - await fs.access(filePath) - // File already exists, skip (don't overwrite existing per-task files) - } catch { - // File doesn't exist, write it - await safeWriteJson(filePath, item) - this.cache.set(item.id, item) + // Write history_item.json if it doesn't exist yet + const filePath = path.join(taskDir, GlobalFileNames.historyItem) + try { + await fs.access(filePath) + // File already exists, skip (don't overwrite existing per-task files) + } catch { + // File doesn't exist, write it + await safeWriteJson(filePath, item) + this.cache.set(item.id, item) + } } - } - // Repair any delegation inconsistencies introduced by the migrated entries. - // Run the lock-free core because migration already holds the store lock. - await this.reconcileDelegationStateCore(this.getPersistedActiveIds()) - }) + // Repair any delegation inconsistencies introduced by the migrated entries. + // Run the lock-free core because migration already holds the store lock. + await this.reconcileDelegationStateCore(this.getPersistedActiveIds()) + }), + ) } // ────────────────────────────── Private: Per-task file I/O ────────────────────────────── diff --git a/src/core/task-persistence/__tests__/TaskHistoryStore.reconciliation.spec.ts b/src/core/task-persistence/__tests__/TaskHistoryStore.reconciliation.spec.ts index e37fd1a25e..802032284c 100644 --- a/src/core/task-persistence/__tests__/TaskHistoryStore.reconciliation.spec.ts +++ b/src/core/task-persistence/__tests__/TaskHistoryStore.reconciliation.spec.ts @@ -8,6 +8,7 @@ import type { HistoryItem } from "@roo-code/types" import { GlobalFileNames } from "../../../shared/globalFileNames" import { TaskHistoryStore, assertValidTransition } from "../TaskHistoryStore" +import { delegateTaskToChild } from "../taskLifecycle" vi.mock("../../../utils/storage", () => ({ getStorageBasePath: vi.fn().mockImplementation((defaultPath: string) => defaultPath), @@ -81,6 +82,10 @@ describe("assertValidTransition", () => { expect(() => assertValidTransition("delegated", "active")).not.toThrow() }) + it("delegated → interrupted", () => { + expect(() => assertValidTransition("delegated", "interrupted")).not.toThrow() + }) + it("interrupted → completed", () => { expect(() => assertValidTransition("interrupted", "completed")).not.toThrow() }) @@ -644,29 +649,208 @@ describe("TaskHistoryStore reconcileDelegationState", () => { expect(store.get("parent-b")?.status).toBe("active") }) - it("repairs an orphaned link in a chained delegation without repairing its grandparent", async () => { - // C doesn't exist (orphaned). B is delegated waiting for C → repaired to active. - // A sees B as delegated in the persisted startup snapshot and remains delegated. - const parentA = makeItem({ id: "parent-a-chain", status: "delegated", awaitingChildId: "parent-b-chain" }) + it("recovers a delegated chain that terminates in a missing child", async () => { + const parentA = makeItem({ + id: "parent-a-chain", + status: "delegated", + awaitingChildId: "parent-b-chain", + delegatedToId: "parent-b-chain", + }) const parentB = makeItem({ id: "parent-b-chain", status: "delegated", awaitingChildId: "missing-child-chain", + delegatedToId: "missing-child-chain", + parentTaskId: parentA.id, }) await seedItems([parentA, parentB]) await store.initialize() - // B is repaired: its child (C) was missing - expect(store.get("parent-b-chain")?.status).toBe("active") - // A stays delegated: B was repaired from delegated to active and remains - // resumable rather than being mistaken for an active orphan from disk. - expect(store.get("parent-a-chain")?.status).toBe("delegated") - expect(store.get("parent-a-chain")?.awaitingChildId).toBe("parent-b-chain") - expect(store.get("parent-b-chain")?.status).toBe("active") - expect(store.get("parent-b-chain")?.awaitingChildId).toBeUndefined() + expect(store.get(parentA.id)).toMatchObject({ + status: "delegated", + awaitingChildId: parentB.id, + }) + expect(store.get(parentB.id)).toMatchObject({ + status: "interrupted", + parentTaskId: parentA.id, + awaitingChildId: undefined, + delegatedToId: undefined, + }) }) + it("recovers a delegated chain that terminates in an interrupted grandchild", async () => { + const parent = makeItem({ + id: "nested-parent", + status: "delegated", + awaitingChildId: "nested-child", + delegatedToId: "nested-child", + }) + const child = makeItem({ + id: "nested-child", + status: "delegated", + parentTaskId: parent.id, + awaitingChildId: "nested-grandchild", + delegatedToId: "nested-grandchild", + }) + const grandchild = makeItem({ + id: "nested-grandchild", + status: "interrupted", + parentTaskId: child.id, + }) + await seedItems([parent, child, grandchild]) + + await store.initialize() + + expect(store.get(parent.id)).toMatchObject({ status: "delegated", awaitingChildId: child.id }) + expect(store.get(child.id)).toMatchObject({ + status: "interrupted", + awaitingChildId: undefined, + delegatedToId: undefined, + }) + expect(store.get(grandchild.id)?.status).toBe("interrupted") + }) + + it.each([undefined, "child", "parent"])( + "recovers nested persisted-active history after a %s write failure", + async (failedWrite) => { + const grandparent = makeItem({ + id: "persisted-active-grandparent", + status: "delegated", + awaitingChildId: "persisted-active-parent", + delegatedToId: "persisted-active-parent", + childIds: ["persisted-active-parent"], + }) + const parent = makeItem({ + id: "persisted-active-parent", + status: "delegated", + parentTaskId: grandparent.id, + rootTaskId: grandparent.id, + awaitingChildId: "persisted-active-child", + delegatedToId: "persisted-active-child", + childIds: ["persisted-active-child"], + }) + const child = makeItem({ + id: "persisted-active-child", + status: "active", + parentTaskId: parent.id, + rootTaskId: grandparent.id, + }) + await seedItems([parent, grandparent, child]) + // Preserve this iteration order even if readdir returns a different order. + store["cache"].set(parent.id, parent) + const intentPath = path.join(tmpDir, "tasks", GlobalFileNames.delegationRepairIntent) + if (failedWrite) { + const failedTaskId = failedWrite === "child" ? child.id : parent.id + safeWriteJsonMock.mockImplementation(async (filePath, data) => { + if (filePath === path.join(tmpDir, "tasks", failedTaskId, GlobalFileNames.historyItem)) { + throw new Error("interrupted nested repair") + } + await writeJson(filePath, data) + }) + } + + await store.initialize() + + const recoveredGrandparent = store.get(grandparent.id)! + expect(recoveredGrandparent).toMatchObject({ + status: "delegated", + awaitingChildId: parent.id, + delegatedToId: parent.id, + }) + if (failedWrite) { + expect(JSON.parse(await fs.readFile(intentPath, "utf8"))).toMatchObject({ + target: { childStatus: "interrupted", parentStatus: "interrupted" }, + }) + } else { + expect(store.get(parent.id)?.status).toBe("interrupted") + expect(store.get(child.id)?.status).toBe("interrupted") + } + + store.dispose() + safeWriteJsonMock.mockImplementation(writeJson) + const subsequentStore = registerStore(new TaskHistoryStore(tmpDir)) + await subsequentStore.initialize() + const subsequentlyRecoveredGrandparent = subsequentStore.get(grandparent.id)! + expect(subsequentlyRecoveredGrandparent).toMatchObject({ + status: "delegated", + awaitingChildId: parent.id, + delegatedToId: parent.id, + }) + expect(subsequentStore.get(parent.id)).toMatchObject({ + status: "interrupted", + parentTaskId: grandparent.id, + }) + expect(subsequentStore.get(parent.id)?.awaitingChildId).toBeUndefined() + expect(subsequentStore.get(parent.id)?.delegatedToId).toBeUndefined() + expect(subsequentStore.get(child.id)?.status).toBe("interrupted") + await expect(fs.access(intentPath)).rejects.toThrow() + expect((await fs.readdir(path.join(tmpDir, "tasks"))).some((name) => name.includes("quarantine"))).toBe( + false, + ) + + const redelegated = delegateTaskToChild( + subsequentlyRecoveredGrandparent, + "replacement-child", + subsequentStore.get(parent.id)?.status, + ) + expect(redelegated).toMatchObject({ + status: "delegated", + awaitingChildId: "replacement-child", + delegatedToId: "replacement-child", + }) + }, + ) + + it("preserves a dead-looking delegated chain with a live runtime owner", async () => { + const parent = makeItem({ + id: "owned-parent", + status: "delegated", + awaitingChildId: "owned-child", + delegatedToId: "owned-child", + }) + const child = makeItem({ + id: "owned-child", + status: "delegated", + parentTaskId: parent.id, + awaitingChildId: "owned-grandchild", + delegatedToId: "owned-grandchild", + }) + const grandchild = makeItem({ id: "owned-grandchild", status: "interrupted", parentTaskId: child.id }) + await seedItems([parent, child, grandchild]) + store.dispose() + store = registerStore(new TaskHistoryStore(tmpDir, { isTaskOwned: (taskId) => taskId === grandchild.id })) + + await store.initialize() + + expect(store.get(child.id)).toMatchObject({ + status: "delegated", + awaitingChildId: grandchild.id, + }) + }) + + it.each(["owned-parent", "owned-child"])( + "preserves active-child history and repair intent while %s is live", + async (ownedId) => { + const parent = makeItem({ + id: "owned-parent", + status: "delegated", + awaitingChildId: "owned-child", + delegatedToId: "owned-child", + }) + const child = makeItem({ id: "owned-child", status: "active", parentTaskId: parent.id }) + await seedItems([parent, child]) + const intentPath = path.join(tmpDir, "tasks", GlobalFileNames.delegationRepairIntent) + await fs.writeFile(intentPath, JSON.stringify(makeRepairIntent(parent, child))) + store.dispose() + store = registerStore(new TaskHistoryStore(tmpDir, { isTaskOwned: (id) => id === ownedId })) + await store.initialize() + expect(store.get(parent.id)).toEqual(parent) + expect(store.get(child.id)).toEqual(child) + await expect(fs.access(intentPath)).resolves.toBeUndefined() + }, + ) + it("does not repair a grandparent when replay repairs the middle node", async () => { const grandparent = makeItem({ id: "grandparent-replay-chain", diff --git a/src/core/task-persistence/__tests__/TaskHistoryStore.spec.ts b/src/core/task-persistence/__tests__/TaskHistoryStore.spec.ts index 3e277ac867..6a7efde0bc 100644 --- a/src/core/task-persistence/__tests__/TaskHistoryStore.spec.ts +++ b/src/core/task-persistence/__tests__/TaskHistoryStore.spec.ts @@ -528,6 +528,49 @@ describe("TaskHistoryStore", () => { }) }) + describe("refreshStrict()", () => { + it.each([ + ["missing", undefined], + ["read-error", { code: "EISDIR" }], + ["malformed", { name: "SyntaxError" }], + ["other-task-id", { message: expect.stringContaining("Invalid task history record") }], + ] as const)("keeps the cached record unless the file is missing (%s)", async (scenario, error) => { + await store.initialize() + store.dispose() // Keep filesystem watcher reconciliation out of this explicit refresh test. + const owner = makeHistoryItem({ id: "owner", status: "delegated", awaitingChildId: "child" }) + await store.upsert(owner) + const filePath = path.join(tmpDir, "tasks", "owner", GlobalFileNames.historyItem) + if (scenario === "malformed") { + await fs.writeFile(filePath, "{") + } else if (scenario === "other-task-id") { + await fs.writeFile(filePath, JSON.stringify({ ...owner, id: "other-task" })) + } else { + await fs.unlink(filePath) + if (scenario === "read-error") await fs.mkdir(filePath) + } + + if (error) { + await expect(store.refreshStrict("owner")).rejects.toMatchObject(error) + expect(store.get("owner")).toEqual(owner) + } else { + await expect(store.refreshStrict("owner")).resolves.toBeUndefined() + expect(store.get("owner")).toBeUndefined() + } + }) + + it("re-reads a valid record from disk", async () => { + await store.initialize() + const item = makeHistoryItem({ id: "strict-task", tokensIn: 100 }) + await store.upsert(item) + const filePath = path.join(tmpDir, "tasks", "strict-task", GlobalFileNames.historyItem) + await fs.writeFile(filePath, JSON.stringify({ ...item, tokensIn: 999 })) + + await store.refreshStrict("strict-task") + + expect(store.get("strict-task")?.tokensIn).toBe(999) + }) + }) + describe("invalidateAll()", () => { it("waits for an in-flight write before clearing the cache", async () => { const onWrite = vi.fn().mockResolvedValue(undefined) diff --git a/src/core/task-persistence/__tests__/taskLifecycle.spec.ts b/src/core/task-persistence/__tests__/taskLifecycle.spec.ts index fe415f09f8..05efb70672 100644 --- a/src/core/task-persistence/__tests__/taskLifecycle.spec.ts +++ b/src/core/task-persistence/__tests__/taskLifecycle.spec.ts @@ -5,6 +5,9 @@ import { completeDelegatedChild, delegateTaskToChild, interruptDelegatedChild, + isDeadDelegationChain, + recoverDeadDelegatedChild, + recoverDelegationParent, } from "../taskLifecycle" function item(id: string, overrides: Partial = {}): HistoryItem { @@ -22,6 +25,21 @@ function item(id: string, overrides: Partial = {}): HistoryItem { } describe("task lifecycle transitions", () => { + it("preserves only an ancestor's current delegation when recovering a parent", () => { + const parent = delegateTaskToChild(item("middle", { parentTaskId: "root" }), "leaf") + const ancestor = delegateTaskToChild(item("root"), parent.id) + expect(recoverDelegationParent(parent, ancestor)).toMatchObject({ + status: "interrupted", + parentTaskId: "root", + awaitingChildId: undefined, + delegatedToId: undefined, + }) + expect(recoverDelegationParent(parent).status).toBe("active") + expect(recoverDelegationParent(parent, { ...ancestor, awaitingChildId: "replacement" }).status).toBe("active") + expect(recoverDelegationParent(parent, { ...ancestor, status: "active" }).status).toBe("active") + expect(() => recoverDelegationParent(item("not-delegated"))).toThrow("non-delegated parent") + }) + it("delegates an active parent and retains child history", () => { const parent = delegateTaskToChild(item("parent", { childIds: ["older"] }), "child") @@ -64,6 +82,71 @@ describe("task lifecycle transitions", () => { expect(interruptDelegatedChild(parent, child)).toMatchObject({ status: "interrupted", parentTaskId: "parent" }) }) + it("recognizes only delegated chains that terminate without a live owner", () => { + const interrupted = item("grandchild", { status: "interrupted" }) + const completed = item("completed-grandchild", { status: "completed" }) + const active = item("active-grandchild") + const child = item("child", { status: "delegated", awaitingChildId: interrupted.id }) + const tasks = new Map([interrupted, completed, active].map((task) => [task.id, task])) + + expect(isDeadDelegationChain(child, (id) => tasks.get(id))).toBe(true) + expect(isDeadDelegationChain({ ...child, awaitingChildId: completed.id }, (id) => tasks.get(id))).toBe(true) + expect(isDeadDelegationChain({ ...child, awaitingChildId: active.id }, (id) => tasks.get(id))).toBe(false) + expect(isDeadDelegationChain({ ...child, awaitingChildId: undefined }, (id) => tasks.get(id))).toBe(true) + expect( + isDeadDelegationChain( + child, + (id) => tasks.get(id), + (id) => id === child.id, + ), + ).toBe(false) + expect( + isDeadDelegationChain( + child, + (id) => tasks.get(id), + (id) => id === interrupted.id, + ), + ).toBe(false) + expect(isDeadDelegationChain({ ...child, status: "active" }, (id) => tasks.get(id))).toBe(false) + expect(isDeadDelegationChain({ ...child, awaitingChildId: "missing" }, (id) => tasks.get(id))).toBe(true) + // A missing record is not proof of death when its task still has a live owner. + expect( + isDeadDelegationChain( + { ...child, awaitingChildId: "missing" }, + (id) => tasks.get(id), + (id) => id === "missing", + ), + ).toBe(false) + expect( + isDeadDelegationChain({ ...child, awaitingChildId: child.id }, (id) => + id === child.id ? child : undefined, + ), + ).toBe(false) + }) + + it("recovers a dead delegated child without releasing its parent's ownership", () => { + const parent = item("parent", { status: "delegated", awaitingChildId: "child", delegatedToId: "child" }) + const child = item("child", { + status: "delegated", + parentTaskId: "parent", + awaitingChildId: "grandchild", + delegatedToId: "grandchild", + }) + + expect(recoverDeadDelegatedChild(parent, child)).toMatchObject({ + id: "child", + status: "interrupted", + parentTaskId: "parent", + awaitingChildId: undefined, + delegatedToId: undefined, + }) + expect(parent).toMatchObject({ status: "delegated", awaitingChildId: "child" }) + expect(() => recoverDeadDelegatedChild(parent, { ...child, status: "active" })).toThrow(/status active/) + expect(() => recoverDeadDelegatedChild({ ...parent, awaitingChildId: "other-child" }, child)).toThrow( + /not delegated to child/, + ) + }) + it("completes only the child the parent still awaits", () => { const parent = item("parent", { status: "delegated", awaitingChildId: "new-child", delegatedToId: "new-child" }) const staleChild = item("old-child", { status: "interrupted", parentTaskId: "parent" }) diff --git a/src/core/task-persistence/index.ts b/src/core/task-persistence/index.ts index 14adeedc68..d290969b11 100644 --- a/src/core/task-persistence/index.ts +++ b/src/core/task-persistence/index.ts @@ -13,13 +13,16 @@ export { } from "./taskMessages" export { taskMetadata } from "./taskMetadata" export { ensureMessageIdentifiers } from "./mergeMessageSnapshots" -export { TaskHistoryStore } from "./TaskHistoryStore" +export { TaskHistoryStore, withTaskOwnershipReservation } from "./TaskHistoryStore" export { abandonDelegatedChild, assertValidTransition, completeDelegatedChild, delegateTaskToChild, interruptDelegatedChild, + isDeadDelegationChain, + recoverDeadDelegatedChild, + recoverDelegationParent, LifecycleTransitionError, type HistoryItemStatus, VALID_TASK_STATUS_TRANSITIONS, diff --git a/src/core/task-persistence/taskLifecycle.ts b/src/core/task-persistence/taskLifecycle.ts index efd2e1148f..cb5b7102ef 100644 --- a/src/core/task-persistence/taskLifecycle.ts +++ b/src/core/task-persistence/taskLifecycle.ts @@ -5,7 +5,7 @@ export type HistoryItemStatus = NonNullable export const VALID_TASK_STATUS_TRANSITIONS: Readonly> = { active: ["delegated", "completed", "interrupted"], - delegated: ["active"], + delegated: ["active", "interrupted"], interrupted: ["completed"], completed: [], } @@ -62,6 +62,66 @@ export function interruptDelegatedChild(parent: HistoryItem, child: HistoryItem) return { ...child, status: "interrupted" } } +/** Release an orphaned active child's link without breaking an ancestor's delegation. */ +export function recoverDelegationParent( + parent: HistoryItem, + ancestor?: HistoryItem, +): HistoryItem & { status: "active" | "interrupted" } { + if (parent.status !== "delegated") { + throw new LifecycleTransitionError(`Cannot recover non-delegated parent ${parent.id}`) + } + const status = ancestor?.status === "delegated" && ancestor.awaitingChildId === parent.id ? "interrupted" : "active" + assertValidTransition(parent.status, status) + return { ...parent, status, awaitingChildId: undefined, delegatedToId: undefined } +} + +/** True when a delegated task has no live owner and its awaited chain ends dead. */ +export function isDeadDelegationChain( + child: HistoryItem, + getTask: (taskId: string) => HistoryItem | undefined, + isTaskLive: (taskId: string) => boolean = () => false, +): boolean { + if (child.status !== "delegated") return false + + const visited = new Set() + let current: HistoryItem | undefined = child + while (true) { + if (visited.has(current.id) || isTaskLive(current.id)) return false + visited.add(current.id) + if (current.status === "interrupted" || current.status === "completed") return true + if (current.status !== "delegated") return false + if (!current.awaitingChildId) return true + const awaitedId: string = current.awaitingChildId + current = getTask(awaitedId) + // A missing record is not proof of death: the descendant may still run elsewhere. + if (!current) return !isTaskLive(awaitedId) + } +} + +/** + * Recover a delegated child whose own execution chain has died. + * + * The caller must establish that neither this child nor any descendant in its + * active delegation chain has a live runtime owner. The parent deliberately + * keeps awaiting the now-interrupted child so existing resume, abandon, and + * re-delegation paths retain ownership semantics. + */ +export function recoverDeadDelegatedChild(parent: HistoryItem, child: HistoryItem): HistoryItem { + if (parent.status !== "delegated" || parent.awaitingChildId !== child.id) { + throw new LifecycleTransitionError(`Task ${parent.id} is not delegated to child ${child.id}`) + } + if (child.status !== "delegated") { + throw new LifecycleTransitionError(`Cannot recover child ${child.id} with status ${child.status}`) + } + assertValidTransition(child.status, "interrupted") + return { + ...child, + status: "interrupted", + awaitingChildId: undefined, + delegatedToId: undefined, + } +} + export function completeDelegatedChild( parent: HistoryItem, child: HistoryItem, diff --git a/src/core/webview/ClineProvider.ts b/src/core/webview/ClineProvider.ts index 185cfbb427..fef2ea95bf 100644 --- a/src/core/webview/ClineProvider.ts +++ b/src/core/webview/ClineProvider.ts @@ -122,9 +122,13 @@ import { saveApiMessages, saveTaskMessages, TaskHistoryStore, + withTaskOwnershipReservation, abandonDelegatedChild, completeDelegatedChild, delegateTaskToChild, + isDeadDelegationChain, + recoverDeadDelegatedChild, + recoverDelegationParent, interruptDelegatedChild, } from "../task-persistence" import { readTaskMessages } from "../task-persistence/taskMessages" @@ -196,6 +200,9 @@ export class ClineProvider public static readonly sideBarId = `${Package.name}.SidebarProvider` public static readonly tabPanelId = `${Package.name}.TabPanelProvider` private static activeInstances: Set = new Set() + private static isTaskRunningInAnyActiveProvider(taskId: string): boolean { + return Array.from(ClineProvider.activeInstances).some((provider) => provider.taskRegistry.hasRunning(taskId)) + } private disposables: vscode.Disposable[] = [] private webviewDisposables: vscode.Disposable[] = [] private pendingThemeFixtureProbes = new Map< @@ -343,7 +350,9 @@ export class ClineProvider void this.updateGlobalState("codebaseIndexModels", EMBEDDING_MODEL_PROFILES) // Initialize the authoritative per-task file-based history store. - this.taskHistoryStore = new TaskHistoryStore(this.contextProxy.globalStorageUri.fsPath) + this.taskHistoryStore = new TaskHistoryStore(this.contextProxy.globalStorageUri.fsPath, { + isTaskOwned: (taskId) => ClineProvider.isTaskRunningInAnyActiveProvider(taskId), + }) this.initializeTaskHistoryStore().catch((error) => { this.log(`Failed to initialize TaskHistoryStore: ${error}`) }) @@ -539,22 +548,54 @@ export class ClineProvider // When the task is completed, the top instance is removed, reactivating the // previous task. async addClineToStack(task: Task) { - // Add this cline instance into the stack that represents the order of - // all the called tasks. - this.taskRegistry.push(task) - task.emit(RooCodeEventName.TaskFocused) + await withTaskOwnershipReservation(async () => { + if (this._disposed || task.abort || task.abandoned) { + for (const cleanup of this.taskEventListeners.get(task) ?? []) cleanup() + this.taskEventListeners.delete(task) + await this.drainTaskDisposal(task) + throw new Error(`[addClineToStack] Task ${task.taskId} registration was cancelled`) + } + // Add this cline instance into the stack that represents the order of + // all the called tasks. + this.taskRegistry.push(task) + }) - // Perform special setup provider specific tasks. - await this.performPreparationTasks(task) + try { + this.throwIfTaskRegistrationCancelled(task) + task.emit(RooCodeEventName.TaskFocused) - // Ensure getState() resolves correctly. - const state = await this.getState() + // Perform special setup provider specific tasks. + await this.performPreparationTasks(task) + this.throwIfTaskRegistrationCancelled(task) - if (!state || typeof state.mode !== "string") { - throw new Error(t("common:errors.retrieve_current_mode")) + // Ensure getState() resolves correctly. + const state = await this.getState() + this.throwIfTaskRegistrationCancelled(task) + + if (!state || typeof state.mode !== "string") { + throw new Error(t("common:errors.retrieve_current_mode")) + } + } catch (error) { + await this.rollbackTaskRegistration(task) + throw error + } + } + + private throwIfTaskRegistrationCancelled(task: Task): void { + if (this._disposed || task.abort || task.abandoned) { + throw new Error(`[addClineToStack] Task ${task.taskId} registration was cancelled`) } } + private async rollbackTaskRegistration(task: Task): Promise { + if (this.taskRegistry.getById(task.taskId) === task) { + this.taskRegistry.remove(task.taskId) + } + for (const cleanup of this.taskEventListeners.get(task) ?? []) cleanup() + this.taskEventListeners.delete(task) + await this.drainTaskDisposal(task) + } + async performPreparationTasks(cline: Task) { // LMStudio: We need to force model loading in order to read its context // size; we do it now since we're starting a task with that model selected. @@ -3781,6 +3822,61 @@ export class ClineProvider return this.currentWorkspacePath || getWorkspacePath() } + private isTaskRunningInAnyProvider(taskId: string): boolean { + return ( + this.taskRegistry.hasRunning(taskId) || + Array.from(ClineProvider.activeInstances).some( + (provider) => provider !== this && provider.taskRegistry.hasRunning(taskId), + ) + ) + } + + private async refreshDelegationChain(taskId: string): Promise { + const visited = new Set() + let currentTaskId: string | undefined = taskId + while (currentTaskId && !visited.has(currentTaskId)) { + visited.add(currentTaskId) + // Strict: an unreadable descendant must abort recovery, not look missing (and so dead). + await this.taskHistoryStore.refreshStrict(currentTaskId) + const current = this.taskHistoryStore.get(currentTaskId) + currentTaskId = current?.status === "delegated" ? current.awaitingChildId : undefined + } + } + + private async recoverDeadAwaitedChild(parent: HistoryItem, childId: string): Promise { + return withTaskOwnershipReservation(async () => { + await this.refreshDelegationChain(childId) + const child = this.taskHistoryStore.get(childId) + if ( + !child || + !isDeadDelegationChain( + child, + (id) => this.taskHistoryStore.get(id), + (id) => this.isTaskRunningInAnyProvider(id), + ) + ) { + return child + } + await this.taskHistoryStore.atomicReadAndUpdate(childId, (currentChild) => { + if ( + !isDeadDelegationChain( + currentChild, + (id) => this.taskHistoryStore.get(id), + (id) => this.isTaskRunningInAnyProvider(id), + ) + ) { + throw new Error(`Delegation chain for child ${childId} became live during recovery`) + } + return recoverDeadDelegatedChild(parent, currentChild) + }) + const recovered = this.taskHistoryStore.get(childId) + this.log( + `[delegateParentAndOpenChild] Recovered dead delegated child ${childId} as interrupted before re-delegation`, + ) + return recovered + }) + } + /** * Delegate parent task and open child task. * @@ -3831,9 +3927,24 @@ export class ClineProvider const awaitedChildId = authoritativeParent.awaitingChildId if (!awaitedChildId) throw new Error("Cannot re-delegate a parent with no awaited child") await this.taskHistoryStore.invalidate(awaitedChildId) - if (this.taskHistoryStore.get(awaitedChildId)?.status !== "interrupted") { - throw new Error("Cannot re-delegate while the awaited child is not interrupted") + let awaitedChild = this.taskHistoryStore.get(awaitedChildId) + if (awaitedChild?.status === "delegated") { + awaitedChild = await this.recoverDeadAwaitedChild(authoritativeParent, awaitedChildId) } + if (awaitedChild?.status !== "interrupted") { + throw new Error( + `Cannot re-delegate while the awaited child ${awaitedChildId} has status ${awaitedChild?.status ?? "missing"}; expected interrupted`, + ) + } + } + if (this._disposed) { + throw new Error("[delegateParentAndOpenChild] Provider was disposed during delegation") + } + if (parent.abort || parent.abandoned) { + throw new Error(`[delegateParentAndOpenChild] Parent ${parent.taskId} was cancelled during delegation`) + } + if (this.getCurrentTask() !== parent) { + throw new Error(`[delegateParentAndOpenChild] Parent ${parent.taskId} is no longer current`) } if (pendingActionId) { const parentHistory = this.taskHistoryStore.get(parentTaskId) @@ -3913,6 +4024,16 @@ export class ClineProvider ) } + if (this._disposed) { + throw new Error("[delegateParentAndOpenChild] Provider was disposed during delegation") + } + if (parent.abort || parent.abandoned) { + throw new Error(`[delegateParentAndOpenChild] Parent ${parent.taskId} was cancelled during delegation`) + } + if (this.getCurrentTask() !== parent) { + throw new Error(`[delegateParentAndOpenChild] Parent ${parent.taskId} is no longer current`) + } + // 3) Enforce single-open invariant by closing/disposing the parent first // This ensures we never have >1 tasks open at any time during delegation. // Await abort completion to ensure clean disposal and prevent unhandled rejections. @@ -3927,6 +4048,10 @@ export class ClineProvider // Non-fatal: proceed with child creation even if parent cleanup had issues } + if (this._disposed) { + throw new Error("[delegateParentAndOpenChild] Provider was disposed during parent cleanup") + } + // 4) Bind the child directly to the delegating task's local provider // context. Delegation never mutates shared profile/global state. // Create child as sole active (parent reference preserved for lineage) @@ -3960,7 +4085,14 @@ export class ClineProvider // synchronously under the store lock) so a concurrent abandon or completion cannot // slip between the status snapshot and the write. An active child must never be // silently detached. + // Set once the parent record durably awaits this child, so rollback can undo that link + // even when disposal cancels delegation after the commit. + let parentCommitted = false + let pendingActionBeforeCommit: HistoryItem["pendingAction"] try { + if (this._disposed) { + throw new Error("[delegateParentAndOpenChild] Provider was disposed before delegation commit") + } await this.taskHistoryStore.atomicReadAndUpdate(parentTaskId, (historyItem) => { if (pendingActionId && historyItem.pendingAction?.actionId !== pendingActionId) { throw new Error( @@ -3971,12 +4103,17 @@ export class ClineProvider ? this.taskHistoryStore.get(historyItem.awaitingChildId)?.status : undefined const delegated = delegateTaskToChild(historyItem, child.taskId, awaitedChildStatus) + pendingActionBeforeCommit = historyItem.pendingAction return { ...delegated, pendingAction: delegated.pendingAction?.actionId === pendingActionId ? undefined : delegated.pendingAction, } }) + parentCommitted = true + if (this._disposed) { + throw new Error("[delegateParentAndOpenChild] Provider was disposed before child scheduling") + } this.recentTasksCache = undefined if (this.isViewLaunched) { const updatedItem = this.taskHistoryStore.get(parentTaskId) @@ -4012,15 +4149,41 @@ export class ClineProvider }`, ) } - try { - const { historyItem: parentHistory } = await this.getTaskWithId(parentTaskId) - await this.createTaskWithHistoryItem(parentHistory) - } catch (rollbackError) { - this.log( - `[delegateParentAndOpenChild] Failed to restore parent ${parentTaskId} during rollback: ${ - (rollbackError as Error)?.message ?? String(rollbackError) - }`, - ) + if (parentCommitted) { + // Undo the committed link even on a disposed provider: otherwise the persisted parent + // awaits a deleted child until a later reconciliation repairs it. + try { + await this.taskHistoryStore.atomicReadAndUpdate(parentTaskId, (current) => { + if (current.status !== "delegated" || current.awaitingChildId !== child.taskId) { + return current // A newer transition already owns the parent. + } + const ancestor = current.parentTaskId + ? this.taskHistoryStore.get(current.parentTaskId) + : undefined + return { + ...recoverDelegationParent(current, ancestor), + pendingAction: pendingActionBeforeCommit, + } + }) + } catch (repairError) { + this.log( + `[delegateParentAndOpenChild] Failed to release parent ${parentTaskId} from deleted child ${child.taskId}: ${ + (repairError as Error)?.message ?? String(repairError) + }`, + ) + } + } + if (!this._disposed) { + try { + const { historyItem: parentHistory } = await this.getTaskWithId(parentTaskId) + await this.createTaskWithHistoryItem(parentHistory) + } catch (rollbackError) { + this.log( + `[delegateParentAndOpenChild] Failed to restore parent ${parentTaskId} during rollback: ${ + (rollbackError as Error)?.message ?? String(rollbackError) + }`, + ) + } } throw err } diff --git a/src/core/webview/__tests__/ClineProvider.spec.ts b/src/core/webview/__tests__/ClineProvider.spec.ts index 72aa1bcb76..625330ff62 100644 --- a/src/core/webview/__tests__/ClineProvider.spec.ts +++ b/src/core/webview/__tests__/ClineProvider.spec.ts @@ -28,6 +28,7 @@ import { experimentDefault } from "../../../shared/experiments" import { setTtsEnabled } from "../../../utils/tts" import { ContextProxy } from "../../config/ContextProxy" import { Task, TaskOptions } from "../../task/Task" +import { withTaskOwnershipReservation } from "../../task-persistence/TaskHistoryStore" import { safeWriteJson } from "../../../utils/safeWriteJson" import { ClineProvider } from "../ClineProvider" @@ -646,8 +647,8 @@ describe("ClineProvider", () => { test("does not log when the active task is aborted", async () => { const task = new Task(defaultTaskOptions) Object.defineProperty(task, "taskId", { value: "aborted-task", writable: true }) - task.abort = true await provider.addClineToStack(task) + task.abort = true Object.defineProperty(mockWebviewView, "visible", { value: false, configurable: true }) visibilityCallback() expect(mockOutputChannel.appendLine).not.toHaveBeenCalled() @@ -1142,6 +1143,76 @@ describe("ClineProvider", () => { expect(disposeCalls).toHaveLength(1) }) + test.each(["dispose", "abort", "abandon"])("rejects queued task registration after %s", async (cancellation) => { + await provider.taskHistoryStore.initialized + let release!: () => void + const reservation = withTaskOwnershipReservation( + () => + new Promise((resolve) => { + release = resolve + }), + ) + await vi.waitFor(() => expect(release).toBeDefined()) + const task = new Task(defaultTaskOptions) + const cleanup = vi.fn() + provider["taskEventListeners"].set(task, [cleanup]) + vi.mocked(Task).mockImplementationOnce(function () { + return task + }) + const registrationSpy = vi.spyOn(provider, "addClineToStack") + const preparationSpy = vi.spyOn(provider, "performPreparationTasks") + const schedulingSpy = vi.spyOn(provider["taskScheduler"], "schedule") + const creation = provider.createTask("queued task") + const rejected = expect(creation).rejects.toThrow("registration was cancelled") + try { + await vi.waitFor(() => expect(registrationSpy).toHaveBeenCalledWith(task)) + if (cancellation === "dispose") await provider.dispose() + else if (cancellation === "abort") task.abort = true + else task.abandoned = true + } finally { + release() + } + await reservation + await rejected + expect(provider["taskRegistry"].hasRunning(task.taskId)).toBe(false) + expect(cleanup).toHaveBeenCalledOnce() + expect(provider["taskEventListeners"].has(task)).toBe(false) + expect(task.dispose).toHaveBeenCalledOnce() + expect(preparationSpy).not.toHaveBeenCalled() + expect(schedulingSpy).not.toHaveBeenCalled() + }) + + test.each(["preparation", "state"])("rolls back registration when cancellation occurs during %s", async (phase) => { + const task = new Task(defaultTaskOptions) + const cleanup = vi.fn() + provider["taskEventListeners"].set(task, [cleanup]) + let release!: () => void + const pending = new Promise((resolve) => { + release = resolve + }) + + if (phase === "preparation") { + vi.spyOn(provider, "performPreparationTasks").mockReturnValue(pending) + } else { + const getState = provider.getState.bind(provider) + vi.spyOn(provider, "getState").mockImplementation(async (options) => { + await pending + return getState(options) + }) + } + + const registration = provider.addClineToStack(task) + await vi.waitFor(() => expect(provider["taskRegistry"].getById(task.taskId)).toBe(task)) + task.abort = true + release() + + await expect(registration).rejects.toThrow("registration was cancelled") + expect(provider["taskRegistry"].getById(task.taskId)).toBeUndefined() + expect(cleanup).toHaveBeenCalledOnce() + expect(provider["taskEventListeners"].has(task)).toBe(false) + expect(task.dispose).toHaveBeenCalledOnce() + }) + test("dispose drains every task in abort-then-cleanup order", async () => { let resolveCurrentAbort!: () => void let resolveCurrentCleanup!: () => void