FN-8908: auto-recover terminal task failures
Recover generic terminal task failures through a bounded, durable retry budget before escalating them to operators. - add fenced task-store recovery claims, retries, budget resets, and audit events - defer terminal-failure notifications until recovery is exhausted while preserving a single escalation - expose operator retry budget reset and cover recovery lifecycle behavior Files changed: .../fn-8908-terminal-failure-auto-recovery.md | 7 + AGENTS.md | 1 + docs/agents.md | 2 + docs/architecture.md | 2 + packages/cli/src/commands/task.ts | 2 + packages/cli/src/extension.ts | 2 + ...terminal-failure-auto-recovery-store.pg.test.ts | 108 +++++++++ .../terminal-failure-auto-recovery.test.ts | 60 +++++ packages/core/src/index.gate.ts | 1 + packages/core/src/index.ts | 1 + packages/core/src/store.ts | 179 ++++++++++++++- .../core/src/task-store/archive-lifecycle-2.ts | 15 ++ packages/core/src/task-store/moves.ts | 48 +++- packages/core/src/task-store/persistence.ts | 18 +- packages/core/src/task-store/project-store-ops.ts | 4 +- .../src/task-store/workflow-task-create-ops.ts | 4 +- packages/core/src/tasks/index.ts | 1 + .../src/tasks/terminal-failure-auto-recovery.ts | 114 ++++++++++ packages/core/src/types/task/task-core.ts | 24 ++ .../src/routes/register-task-workflow-routes.ts | 2 + ...-healing-terminal-failure-auto-recovery.test.ts | 199 ++++++++++++++++ .../__tests__/notification-service.test.ts | 86 ++++++- .../__tests__/task-wedge-notification.test.ts | 2 +- .../src/notification/notification-service.ts | 103 +++++++-- .../src/notification/task-wedge-notification.ts | 30 ++- packages/engine/src/self-healing.ts | 251 ++++++++++++++++++++- packages/engine/src/util/run-audit.ts | 7 + 27 files changed, 1240 insertions(+), 33 deletions(-) Fusion-Task-Id: FN-8908 Fusion-Task-Lineage: 99e96b16-0306-41f1-87da-8623d69735f7 Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
7
.changeset/fn-8908-terminal-failure-auto-recovery.md
Normal file
7
.changeset/fn-8908-terminal-failure-auto-recovery.md
Normal file
@@ -0,0 +1,7 @@
|
||||
---
|
||||
"@runfusion/fusion": minor
|
||||
---
|
||||
|
||||
summary: Automatically retry generic terminal task failures before alerting operators.
|
||||
category: fix
|
||||
dev: Adds a durable recovery budget, fenced retry application, bounded escalation delivery, operator retry reset, and stale-mirror cleanup.
|
||||
@@ -295,6 +295,7 @@ Scoped exception (FN-5819/FN-8823): while project auto-merge is On, shared-branc
|
||||
- FN-7069: task-store open and self-healing housekeeping emit `task:reconcile-phantom-committed-reservation` when they prune orphaned child rows for a committed task-ID reservation that has no live/soft-deleted/archived task row and no task directory, while preserving the committed reservation so the ID is never reused.
|
||||
- FN-7074: task creation emits `task:reservation-commit-rolled-back` when a distributed reservation was committed atomically with a `tasks` row but a later create materialization step failed; metadata includes `reservationId`, `nodeId`, `reason: "failed-create"`, and `error`, and the reservation is moved to aborted so the sequence remains burned.
|
||||
- FN-6782/FN-6796: self-healing emits `task:auto-recover-paused-abort-park` when it clears a benign pause-abort operator park, requeueing safe `todo`/`in-progress` rows or preserving a clean auto-merge-eligible `in-review` row for review progression.
|
||||
- FN-8908: self-healing reserves `task:auto-recover-terminal-failure` and `task:auto-recover-terminal-failure-exhausted` for generic terminal-failure budget recovery. Metadata must remain ids/counts/outcomes-only and never include failure prose or the rotating `wedgeNotification.autoRecovery.applyToken`; that durable budget is the backoff source, and its apply fence—not the grace heuristic—authorizes the single clear/requeue transition.
|
||||
- FN-6793/FN-6797: self-healing emits `task:reconcile-in-review-unmet-dependencies` when it rebounds an `in-review` task whose declared dependencies are still unmet, and `task:reconcile-in-review-unmet-dependencies-no-action` when pause/user-pause, `autoMerge:false`, live execution/checkout proof, or a failed rebound mutation blocks that backward move.
|
||||
- Workspace (Phase D U1): self-healing emits `task:reconcile-workspace-partial-land` when it re-enqueues a partial/zero-landed workspace task's per-repo land (or parks it `failed` when a sub-repo's `fusion/<id>` branch is gone with no `landedSha`), and `task:reconcile-workspace-partial-land-no-action` when `autoMerge:false`, user-pause, or a live sub-repo worktree (workspace-aware liveness) blocks that backward move.
|
||||
- Workspace (Phase D U1): self-healing emits `task:reclaim-phantom-workspace-land-lease` when it clears a leaked `workspace-repo-land` lease whose owning task is terminal/dead and older than the FN-6736 staleness floor (a live merging owner is left untouched).
|
||||
|
||||
@@ -268,6 +268,8 @@ Separation of concerns:
|
||||
|
||||
When a task is terminally blocked (for example, by a merge gate, exhausted execution retries, or a completion blocker), Fusion posts a system message to the dashboard mailbox and sends a `task-wedged` notification through configured providers. Self-healing declines alert only when their proof shows no live session, no recent activity, and no intentional pause or auto-merge-off hold. Before delivery, Fusion revalidates the live row: progressing (including `reviewing`), paused, auto-merge-off, deleted, archived, and complete-lane rows do not alert. The message identifies the task, bounded reason/gate when known, and a recovery action. The active/resolved episode is persisted with the task, so it is sent once per active reason across service restarts; retrying or otherwise restoring progress clears the episode, so a later recurrence is visible again.
|
||||
|
||||
For an unclassified generic `terminal-failed` park, Fusion first uses a small durable automatic-recovery budget rather than immediately paging an operator. While retries are owed, all failure-alert channels are withheld. Budget exhaustion produces one confirmed terminal-failure escalation; turning automatic recovery off emits a reason-tagged drain alert without discarding remaining retries, so turning it back on resumes recovery. An operator Retry starts a fresh budget; success, archive, and deletion clear it automatically. Agent-initiated retry intentionally does not mint a fresh budget. A spent budget can also expire after its age bound only when a later foreign row write proves the episode moved on.
|
||||
|
||||
### CLI agent permission prompts and notifications
|
||||
|
||||
CLI-agent adapters keep their own autonomy posture and tool-permission handling separate from permanent-agent `permissionPolicy`. When an adapter reports a permission/input prompt (`PermissionRequest`, `Notification`, or a conservative approval-prompt heuristic), the CLI session moves to `waitingOnInput`; the dashboard shows the session banner, and Fusion dispatches the `cli-agent-awaiting-input` notification event through enabled ntfy/webhook providers. Repeated waiting events for the same CLI session are de-duplicated before provider delivery, while the in-app banner continues to reflect the live session state.
|
||||
|
||||
@@ -18,6 +18,8 @@ A failed snapshot is not actionable while persisted automatic-recovery ownership
|
||||
|
||||
Each task also stores `lastNotifiedAtByReason`, an independent timestamp map keyed by bounded reason. `WEDGE_RENOTIFY_COOLDOWN_MS` defaults to six hours: resolving an episode does not clear its reason's live stamp, so a scheduler/self-healing resolve→re-wedge flap sends neither a provider push nor a mailbox message until the window expires. A different reason notifies immediately, including X→Y→X while X remains within its own cooldown; expired or invalid entries are pruned during the atomic claim, and legacy rows without the map notify normally before initializing it. The no-durable-store fallback applies the same per-reason window in memory. Provider and mailbox delivery are independently best-effort after sharing this single claim decision, while run-audit metadata remains ids/counts/outcomes-only.
|
||||
|
||||
Generic `terminal-failed` parks are engine-owned before they become operator work. A durable `wedgeNotification.autoRecovery` budget supplies bounded attempts and backoff; its claim state and rotating apply fence ensure one observer is authorized to clear and requeue a park. The grace window detects an abandoned apply only—it is not a lock. An exhausted budget receives one reason-scoped escalation, confirmed durably at the shared dispatch seam; a write-once exhaustion marker prevents an earlier drain alert or cooldown suppression from satisfying that escalation. The budget resets on terminal success, archive, explicit operator Retry, soft-delete, or a sufficiently old foreign write. Disabling auto-recovery drains an owed alert without consuming the remaining retries; re-enabling resumes them.
|
||||
|
||||
## Planning dependency lifecycle lock
|
||||
|
||||
Dependency changes and planning finalization share one outer lifecycle lock keyed by the canonical project ID and task ID. In PostgreSQL mode this is a dedicated, single-connection session advisory lock: it is acquired before the normal task lock and released before the mutation/finalization Promise settles. The operational runtime pool is never borrowed for this purpose.
|
||||
|
||||
@@ -1524,6 +1524,8 @@ export async function runTaskRetry(id: string, projectName?: string) {
|
||||
const autoPauseClearPatch = buildAutoPauseClearPatch(task);
|
||||
const clearedDeadlockAutoPause = Object.keys(autoPauseClearPatch).length > 0;
|
||||
const retryLogSuffix = clearedDeadlockAutoPause ? ", cleared deadlock auto-pause" : "";
|
||||
// FNXC:TaskWedgeNotifications 2026-08-10-20:15: a human Retry proves intervention and mints a fresh bounded terminal-failure budget.
|
||||
await context.store.resetTerminalFailureAutoRecoveryBudget(id);
|
||||
|
||||
if (isMissingWorktreeSessionRetry) {
|
||||
await retryBoardCall(context, id, "move task", () => context.store.moveTask(id, retryHoldColumn as never, { preserveProgress: true }));
|
||||
|
||||
@@ -2419,6 +2419,8 @@ export default function kbExtension(pi: ExtensionAPI) {
|
||||
const autoPauseClearPatch = buildAutoPauseClearPatch(task);
|
||||
const clearedDeadlockAutoPause = Object.keys(autoPauseClearPatch).length > 0;
|
||||
const retryLogSuffix = clearedDeadlockAutoPause ? ", cleared deadlock auto-pause" : "";
|
||||
// FNXC:TaskWedgeNotifications 2026-08-10-20:15: an operator retry ends the prior terminal-failure episode and mints a fresh budget.
|
||||
await store.resetTerminalFailureAutoRecoveryBudget(params.id);
|
||||
|
||||
if (isMissingWorktreeSessionRetry) {
|
||||
await store.updateTask(params.id, {
|
||||
|
||||
@@ -0,0 +1,108 @@
|
||||
/*
|
||||
FNXC:TaskWedgeNotifications 2026-08-10-20:40:
|
||||
The terminal-failure apply must use the real PostgreSQL move transaction. A mock can prove the
|
||||
fence branch was selected while still missing the transaction boundary that keeps the failure clear,
|
||||
column move, and fence consumption atomic.
|
||||
*/
|
||||
import { afterAll, afterEach, beforeAll, beforeEach, expect, it } from "vitest";
|
||||
import "@fusion/core";
|
||||
import { pgDescribe, createSharedPgTaskStoreTestHarness } from "../../__test-utils__/pg-test-harness.js";
|
||||
|
||||
pgDescribe("terminal failure auto-recovery apply", () => {
|
||||
const harness = createSharedPgTaskStoreTestHarness({ prefix: "fusion_terminal_failure_apply" });
|
||||
|
||||
beforeAll(harness.beforeAll);
|
||||
afterAll(harness.afterAll);
|
||||
beforeEach(async () => { await harness.beforeEach(); });
|
||||
afterEach(async () => { await harness.afterEach(); });
|
||||
|
||||
it("consumes the matching fence in the same persisted move as the failure clear", async () => {
|
||||
const store = harness.store();
|
||||
const task = await store.createTask({ description: "fenced terminal failure" } as never);
|
||||
const token = "terminal-apply-token";
|
||||
await store.updateTask(task.id, {
|
||||
status: "failed",
|
||||
error: "opaque terminal failure",
|
||||
wedgeNotification: {
|
||||
reasonKey: "terminal-failed",
|
||||
episodeId: "episode",
|
||||
status: "resolved",
|
||||
transitionedAt: new Date().toISOString(),
|
||||
budgetRevision: 1,
|
||||
autoRecovery: {
|
||||
attempts: 1,
|
||||
lastAttemptAt: new Date().toISOString(),
|
||||
lastApplyStartedAt: new Date().toISOString(),
|
||||
applyToken: token,
|
||||
lastBudgetWriteAt: new Date().toISOString(),
|
||||
},
|
||||
},
|
||||
} as never);
|
||||
|
||||
const result = await store.applyTerminalFailureAutoRecoveryRetry(task.id, {
|
||||
applyToken: token,
|
||||
patch: { status: null, error: null, recoveryRetryCount: 1, nextRecoveryAt: new Date(Date.now() + 60_000).toISOString() },
|
||||
targetColumn: "todo",
|
||||
moveOptions: { preserveProgress: true, moveSource: "engine" },
|
||||
});
|
||||
|
||||
expect(result.outcome).toBe("applied");
|
||||
const applied = await store.getTask(task.id);
|
||||
expect(applied.column).toBe("todo");
|
||||
expect(applied.status).toBeUndefined();
|
||||
expect(applied.error).toBeUndefined();
|
||||
expect(applied.wedgeNotification?.autoRecovery?.retryAppliedAt).toBeTruthy();
|
||||
expect(applied.wedgeNotification?.autoRecovery?.applyToken).toBeUndefined();
|
||||
expect(applied.wedgeNotification?.autoRecovery?.lastApplyStartedAt).toBeUndefined();
|
||||
|
||||
const stale = await store.applyTerminalFailureAutoRecoveryRetry(task.id, {
|
||||
applyToken: token,
|
||||
patch: { status: null, error: null },
|
||||
targetColumn: "todo",
|
||||
moveOptions: { preserveProgress: true, moveSource: "engine" },
|
||||
});
|
||||
expect(stale.outcome).toBe("not-failed");
|
||||
});
|
||||
|
||||
it("does not spend a recovery attempt when the row moved on before the atomic claim", async () => {
|
||||
const store = harness.store();
|
||||
const task = await store.createTask({ description: "already recovered" } as never);
|
||||
|
||||
const claim = await store.claimTerminalFailureAutoRecoveryAttempt(task.id, {
|
||||
maxAttempts: 3,
|
||||
maxResumes: 1,
|
||||
minAttemptSpacingMs: 0,
|
||||
claimApplyGraceMs: 60_000,
|
||||
});
|
||||
|
||||
expect(claim).toEqual({ outcome: "already-claimed", attempt: 0 });
|
||||
expect((await store.getTask(task.id))?.wedgeNotification?.autoRecovery).toBeUndefined();
|
||||
});
|
||||
|
||||
it("clears the budget only after a backend delete wins its soft-delete claim", async () => {
|
||||
const store = harness.store();
|
||||
const task = await store.createTask({ description: "delete terminal failure budget" } as never);
|
||||
await store.updateTask(task.id, {
|
||||
wedgeNotification: {
|
||||
reasonKey: "terminal-failed",
|
||||
episodeId: "delete-episode",
|
||||
status: "active",
|
||||
transitionedAt: "2026-08-10T00:00:00.000Z",
|
||||
lastNotifiedAtByReason: { "terminal-failed": "2026-08-10T00:00:00.000Z", other: "2026-08-10T00:00:00.000Z" },
|
||||
budgetRevision: 4,
|
||||
autoRecovery: { attempts: 3, lastAttemptAt: "2026-08-10T00:00:00.000Z", escalationNotifiedAt: "2026-08-10T00:00:00.000Z" },
|
||||
},
|
||||
} as never);
|
||||
|
||||
const declined = await store.deleteTaskIf(task.id, async () => false);
|
||||
expect(declined.deleted).toBe(false);
|
||||
expect((await store.getTask(task.id))?.wedgeNotification?.autoRecovery?.attempts).toBe(3);
|
||||
|
||||
await store.deleteTaskIf(task.id, async () => true);
|
||||
const deleted = await store.getTask(task.id, { includeDeleted: true });
|
||||
expect(deleted?.wedgeNotification?.autoRecovery).toBeUndefined();
|
||||
expect(deleted?.wedgeNotification?.budgetRevision).toBe(5);
|
||||
expect(deleted?.wedgeNotification?.lastNotifiedAtByReason).toEqual({ other: "2026-08-10T00:00:00.000Z" });
|
||||
expect(deleted?.wedgeNotification?.status).toBe("resolved");
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,60 @@
|
||||
import { describe, expect, it } from "vitest";
|
||||
import type { Task } from "../types/task/task-core.js";
|
||||
import {
|
||||
BUDGET_WRITE_GUARD_MS,
|
||||
classifyTerminalFailureAutoRecovery,
|
||||
MAX_TERMINAL_FAILURE_AUTO_RETRIES,
|
||||
TERMINAL_FAILURE_BUDGET_MAX_AGE_MS,
|
||||
} from "../tasks/terminal-failure-auto-recovery.js";
|
||||
|
||||
function task(overrides: Partial<Task> = {}): Task {
|
||||
return {
|
||||
id: "FN-8908", description: "test", column: "todo", status: "failed",
|
||||
updatedAt: new Date(10_000).toISOString(), ...overrides,
|
||||
} as Task;
|
||||
}
|
||||
|
||||
const generic = {
|
||||
isGenericTerminalFailure: true, hasRecoveryOwner: false, isProgressing: false,
|
||||
inTerminalSuccessColumn: false, isArchivedOrDeleted: false, autoRecoveryEnabled: true,
|
||||
now: () => 20_000,
|
||||
};
|
||||
|
||||
describe("classifyTerminalFailureAutoRecovery", () => {
|
||||
it("retries a generic failed task from its durable budget", () => {
|
||||
expect(classifyTerminalFailureAutoRecovery(task(), generic)).toEqual({ action: "retry", attempt: 1 });
|
||||
expect(classifyTerminalFailureAutoRecovery(task({ wedgeNotification: {
|
||||
reasonKey: "terminal-failed", episodeId: "e", status: "active", transitionedAt: "",
|
||||
autoRecovery: { attempts: 1, lastAttemptAt: new Date(1_000).toISOString() },
|
||||
} }), generic)).toEqual({ action: "retry", attempt: 2 });
|
||||
});
|
||||
|
||||
it("escalates an exhausted budget before retry-path guards", () => {
|
||||
const exhausted = task({ paused: true, autoMerge: false, wedgeNotification: {
|
||||
reasonKey: "terminal-failed", episodeId: "e", status: "active", transitionedAt: "",
|
||||
autoRecovery: { attempts: MAX_TERMINAL_FAILURE_AUTO_RETRIES, lastAttemptAt: new Date(1_000).toISOString() },
|
||||
} });
|
||||
expect(classifyTerminalFailureAutoRecovery(exhausted, { ...generic, hasRecoveryOwner: true, isProgressing: true }))
|
||||
.toEqual({ action: "notify", reason: "budget-exhausted" });
|
||||
});
|
||||
|
||||
it("clears a budget before generic classification and never treats merely non-failed as progress", () => {
|
||||
const budget = { attempts: 1, lastAttemptAt: new Date(1_000).toISOString() };
|
||||
expect(classifyTerminalFailureAutoRecovery(task({ status: null, wedgeNotification: { reasonKey: "x", episodeId: "e", status: "active", transitionedAt: "", autoRecovery: budget } }), {
|
||||
...generic, isGenericTerminalFailure: false, inTerminalSuccessColumn: true,
|
||||
})).toEqual({ action: "reset-budget", reason: "terminal-success" });
|
||||
expect(classifyTerminalFailureAutoRecovery(task({ status: null, wedgeNotification: { reasonKey: "x", episodeId: "e", status: "active", transitionedAt: "", autoRecovery: budget } }), {
|
||||
...generic, isGenericTerminalFailure: false,
|
||||
})).toEqual({ action: "skip", reason: "not-generic-terminal-failure" });
|
||||
});
|
||||
|
||||
it("requires a foreign write after the budget watermark before stale cleanup", () => {
|
||||
const now = TERMINAL_FAILURE_BUDGET_MAX_AGE_MS + 10_000;
|
||||
const budget = { attempts: 3, lastAttemptAt: new Date(0).toISOString(), lastBudgetWriteAt: new Date(1_000).toISOString() };
|
||||
const base = { reasonKey: "terminal-failed", episodeId: "e", status: "active" as const, transitionedAt: "", autoRecovery: budget };
|
||||
expect(classifyTerminalFailureAutoRecovery(task({ updatedAt: new Date(1_000 + BUDGET_WRITE_GUARD_MS).toISOString(), wedgeNotification: base }), { ...generic, now: () => now }))
|
||||
.toEqual({ action: "notify", reason: "budget-exhausted" });
|
||||
expect(classifyTerminalFailureAutoRecovery(task({ updatedAt: new Date(1_001 + BUDGET_WRITE_GUARD_MS).toISOString(), wedgeNotification: base }), { ...generic, now: () => now }))
|
||||
.toEqual({ action: "reset-budget", reason: "budget-stale" });
|
||||
});
|
||||
});
|
||||
@@ -95,6 +95,7 @@ export { isActiveNearDuplicateColumn, isNearDuplicateCanonicalInactive } from ".
|
||||
export { resolveNearDuplicateCanonicalFlags } from "./duplicates/near-duplicate-canonical-flags.js";
|
||||
export type { NearDuplicateCanonicalState } from "./duplicates/near-duplicate-canonical.js";
|
||||
export * from "./tasks/frontend-ux-policy.js";
|
||||
export * from "./tasks/terminal-failure-auto-recovery.js";
|
||||
export * from "./tasks/original-description-policy.js";
|
||||
export * from "./planner/planning-plan-md.js";
|
||||
export * from "./tasks/file-scope-classification.js";
|
||||
|
||||
@@ -130,6 +130,7 @@ types/policy for severity-routed notes before they reach steering inject.
|
||||
export * from "./planner/overseer-advice.js";
|
||||
export * from "./planner/overseer-emission-guard.js";
|
||||
export * from "./tasks/frontend-ux-policy.js";
|
||||
export * from "./tasks/terminal-failure-auto-recovery.js";
|
||||
export * from "./tasks/original-description-policy.js";
|
||||
export * from "./planner/planning-plan-md.js";
|
||||
export * from "./tasks/file-scope-classification.js";
|
||||
|
||||
@@ -3,6 +3,7 @@ import type { TaskMoveLanes } from "./workflows/workflow-lifecycle-traits.js";
|
||||
import { TaskLaneCache } from "./task-lane-cache.js";
|
||||
import { randomUUID } from "node:crypto";
|
||||
import { WEDGE_RENOTIFY_COOLDOWN_MS } from "./types/task/task-core.js";
|
||||
import { clearTerminalFailureAutoRecoveryBudget } from "./tasks/terminal-failure-auto-recovery.js";
|
||||
import { join } from "node:path";
|
||||
import { and, desc, eq, isNull, ne, sql } from "drizzle-orm";
|
||||
import { createCurrentPlanEvidence, diffSpecLocks, isSpecLockActive, type CurrentPlanEvidence, type PlanEvidenceBindings, type SpecLock } from "./planner/spec-lock.js";
|
||||
@@ -111,7 +112,7 @@ import type { IntakeOwnershipExemption } from "./tasks/task-intake-owner-resolve
|
||||
// the single import source for all consumers (re-exports preserved below).
|
||||
import { TASK_JSONB_COLUMNS, type TaskRow, type TaskPersistSerializationContext, type TaskColumnDescriptor } from "./task-store/persistence.js";
|
||||
import { pgRowToTaskRow as pgRowToTaskRowExternal, rowToTask as rowToTaskExternal, rowToBranchGroup as rowToBranchGroupExternal, generateBranchGroupId as generateBranchGroupIdExternal, computeTimedExecutionMs as computeTimedExecutionMsExternal, archiveEntryToTask as archiveEntryToTaskExternal, summarizeAgentLog as summarizeAgentLogExternal, rowToTaskDocument as rowToTaskDocumentExternal, rowToArtifact as rowToArtifactExternal, rowToTaskDocumentRevision as rowToTaskDocumentRevisionExternal, rowToGoalCitation as rowToGoalCitationExternal } from "./task-store/serialization.js";
|
||||
import { moveTaskImpl, moveTaskIfImpl, handoffToReviewImpl, moveTaskInternalImpl, type MoveTaskIfResult } from "./task-store/moves.js";
|
||||
import { moveTaskImpl, moveTaskIfImpl, handoffToReviewImpl, moveTaskInternalImpl, TerminalFailureApplyRejected, type MoveTaskIfResult } from "./task-store/moves.js";
|
||||
import { recordGoalCitationsImpl, insertTaskWithFtsRecoveryImpl2, assertTaskIdAvailableImpl, atomicWriteTaskJsonImpl2, createTaskWithDistributedReservationImpl, toStoredWorkflowStepImpl, ensureWorkflowStepForTemplateImpl, resolveEnabledWorkflowStepsImpl, setTaskBranchGroupImpl, getTaskColumnsImpl, prepareWorkflowMovePolicyPreflightImpl, updateTaskCustomFieldsImpl, listWorkflowPromptOverridesForProjectImpl, listWorkflowWorkItemsForTaskImpl, listDueWorkflowWorkItemsImpl, rewriteBlockedByResidueDependentsForRemovalImpl, getAllDocumentsImpl, deleteWorkflowStepImpl, toWorkflowDefinitionImpl, materializeDefaultWorkflowStepsImpl, reconcileTaskCustomFieldsForSchemaImpl, getTaskMovedCountsByDayImpl, getGoalStoreImpl, upsertTaskCommitAssociationImpl } from "./task-store/workflow-task-create-ops.js";
|
||||
import { applyLegacyWorkflowStepOverridesImpl, archiveDbImpl, assertNoDependencyCycleImpl, atomicCreateTaskJsonImpl, buildActiveTaskDependencyLookupImpl, buildArchivedAgentLogFieldsImpl, buildTaskIdIntegrityFallbackReportImpl, createBranchGroupImpl, dbImpl, detectAndCacheTaskIdIntegrityReportImpl, findLiveDependentsImpl, findLiveLineageChildrenImpl, getLegacyWorkflowStepSnapshotImpl, getMalformedTaskMetadataReasonImpl, getMergeQueuedTaskIdsAsyncImpl, insertRunAuditEventRowImpl, insertTaskImpl, invokeTaskCreatedHookImpl, isTaskArchivedAsyncImpl, isTaskArchivedImpl, isTaskIdPresentInArchivedTasksTableAsyncImpl, isTaskIdPresentInArchivedTasksTableImpl, logTaskCreateConflictImpl, maybeResolveTombstonedTaskIdImpl, mergeTaskIdIntegrityReportsImpl, optionalGroupIdSetImpl, patchTaskRowInTransactionImpl, readConfigFastImpl, readConfigImpl, readPromptForArchiveImpl, readTaskFromDbImpl, reconcileDistributedTaskIdStateOnOpenImpl, recordActivityFromListenerImpl, recordDependencyCycleRejectedAuditImpl, refreshTaskIdIntegrityReportImpl, resolveLocalNodeIdForTaskAllocationImpl, runTaskFtsWriteWithRecoveryImpl, scanAndRecordCitationsImpl, taskIdExistsAnywhereImpl, throwSoftDeletedWriteBlockedImpl, toBuiltInWorkflowStepImpl, trackDeferredTaskCreatedWorkImpl, upsertTaskImpl, withConfigLockImpl, withTaskLockImpl, withWorktreeAllocationLockImpl } from "./task-store/task-id-integrity.js";
|
||||
import { claimNextToolFailureRetryImpl, createTaskVerificationRequestImpl, claimTaskVerificationRequestImpl, finishTaskVerificationRequestImpl, clearNearDuplicateReferencesToFailSoftImpl, clearWorkflowRunStepInstancesAsyncImpl, clearWorkflowRunStepInstancesImpl, computeMovedSettingsTargetWorkflowIdsImpl, ensureBranchGroupForSourceImpl, ensurePrEntityForSourceImpl, findRecentTasksByContentFingerprintImpl, getActiveMergingTaskImpl, getActivePrEntityBySourceImpl, getBranchGroupByBranchNameImpl, getBranchGroupBySourceImpl, getBranchGroupImpl, getBranchProgressByTaskImpl, getMutationsForRunImpl, getPrEntityByNumberImpl, getPrEntityImpl, getPrThreadStateImpl, getTasksByAssignedAgentImpl, getWorkflowPromptOverridesAsyncImpl, getWorkflowSettingValuesAsyncImpl, getWorkflowSettingValuesImpl, getWorkflowSettingsProjectIdImpl, getWorkflowWorkItemImpl, insertCompletionHandoffWorkflowWorkAuditImpl, listActivePrEntitiesImpl, listBranchGroupsImpl, listPrThreadStatesImpl, listTasksByBranchGroupImpl, listWorkflowSettingValuesForProjectImpl, loadWorkflowRunBranchesImpl, hasWorkflowRunStepInstancesForTaskImpl, loadWorkflowRunStepInstancesAsyncImpl, loadWorkflowRunStepInstancesImpl, markToolFailureRetryExhaustedAuditImpl, mergeCustomFieldPatchImpl, normalizeMergeRequestStateImpl, normalizeWorkflowWorkItemKindImpl, normalizeWorkflowWorkItemStateImpl, parseWorkflowPromptOverrideJsonImpl, recordPrThreadOutcomeImpl, resetAllStepsToPendingImpl, resetPromptCheckboxesImpl, resolveWorkflowMoveActorImpl, resolveWorkflowSettingDeclarationsImpl, saveWorkflowRunStepInstanceAsyncImpl, saveWorkflowRunStepInstanceImpl, transitionMergeRequestStateImpl, transitionWorkflowWorkItemSyncImpl, updateTaskImpl, updateWorkflowPromptOverridesImpl, upsertMergeRequestRecordImpl, workflowStateForMergeRequestStateImpl } from "./task-store/branch-and-pr-entities.js";
|
||||
@@ -306,6 +307,12 @@ export interface MoveTaskOptions {
|
||||
/** @internal Extracted to task-store/moves.ts */
|
||||
export interface MoveTaskInternalOptions {
|
||||
fromHandoff: boolean;
|
||||
/** @internal Fence-validated inside moveTaskInternal's transaction for terminal-failure recovery. */
|
||||
terminalFailureApply?: {
|
||||
applyToken: string;
|
||||
patch: Parameters<TaskStore["updateTask"]>[1];
|
||||
expectedColumn: ColumnId;
|
||||
};
|
||||
runContext?: Pick<RunMutationContext, "runId" | "agentId"> | { runId?: string; agentId?: string };
|
||||
ownerAgentId?: string | null;
|
||||
evidence?: HandoffToReviewOptions["evidence"];
|
||||
@@ -1741,11 +1748,179 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
|
||||
status: "active",
|
||||
transitionedAt: new Date(now).toISOString(),
|
||||
...(Object.keys(lastNotifiedAtByReason).length > 0 ? { lastNotifiedAtByReason } : {}),
|
||||
...(prior?.autoRecovery ? { autoRecovery: prior.autoRecovery } : {}),
|
||||
...(prior?.budgetRevision !== undefined ? { budgetRevision: prior.budgetRevision } : {}),
|
||||
},
|
||||
};
|
||||
});
|
||||
return result;
|
||||
}
|
||||
/*
|
||||
FNXC:TaskWedgeNotifications 2026-08-10-18:54:
|
||||
Generic terminal parks use a dedicated durable budget because transient recovery
|
||||
writers clear `recoveryRetryCount`. The apply token is only minted here; a later
|
||||
fenced move consumes it, so a clock-based abandonment check never authorizes a move.
|
||||
*/
|
||||
async claimTerminalFailureAutoRecoveryAttempt(
|
||||
taskId: string,
|
||||
options: { maxAttempts: number; maxResumes: number; minAttemptSpacingMs: number; claimApplyGraceMs: number },
|
||||
): Promise<{ outcome: "resume" | "claimed"; attempt: number; applyToken: string } | { outcome: "already-claimed"; attempt: number } | { outcome: "exhausted"; attempts: number }> {
|
||||
let result: { outcome: "resume" | "claimed"; attempt: number; applyToken: string } | { outcome: "already-claimed"; attempt: number } | { outcome: "exhausted"; attempts: number } = { outcome: "already-claimed", attempt: 0 };
|
||||
await this.updateTaskAtomic(taskId, (current) => {
|
||||
const now = Date.now();
|
||||
const prior = current.wedgeNotification?.autoRecovery;
|
||||
const attempts = typeof prior?.attempts === "number" && Number.isFinite(prior.attempts) && prior.attempts >= 0 ? Math.trunc(prior.attempts) : 0;
|
||||
// FNXC:TaskWedgeNotifications 2026-08-10-20:11: The sweep's pre-claim read can race
|
||||
// a real forward move or deletion. A claim is a park charge, so only a live failed row
|
||||
// may spend it; the fenced apply rechecks this again when it performs the lifecycle move.
|
||||
if (current.status !== "failed" || current.deletedAt != null) {
|
||||
result = { outcome: "already-claimed", attempt: attempts };
|
||||
return null;
|
||||
}
|
||||
const resumes = typeof prior?.resumeCount === "number" && Number.isFinite(prior.resumeCount) && prior.resumeCount >= 0 ? Math.trunc(prior.resumeCount) : 0;
|
||||
const lastAttempt = Date.parse(prior?.lastAttemptAt ?? "");
|
||||
const retryApplied = Date.parse(prior?.retryAppliedAt ?? "");
|
||||
const unapplied = !Number.isFinite(retryApplied) || !Number.isFinite(lastAttempt) || retryApplied < lastAttempt;
|
||||
const started = Date.parse(prior?.lastApplyStartedAt ?? prior?.lastAttemptAt ?? "");
|
||||
if (attempts > 0 && unapplied && Number.isFinite(started) && now - started < options.claimApplyGraceMs) {
|
||||
result = { outcome: "already-claimed", attempt: attempts };
|
||||
return null;
|
||||
}
|
||||
if (attempts > 0 && attempts <= options.maxAttempts && unapplied && resumes < options.maxResumes) {
|
||||
const applyToken = randomUUID();
|
||||
result = { outcome: "resume", attempt: attempts, applyToken };
|
||||
return { wedgeNotification: {
|
||||
...(current.wedgeNotification ?? { reasonKey: "terminal-failed", episodeId: randomUUID(), status: "resolved" as const, transitionedAt: new Date(now).toISOString() }),
|
||||
budgetRevision: (current.wedgeNotification?.budgetRevision ?? 0) + 1,
|
||||
autoRecovery: { ...prior!, resumeCount: resumes + 1, lastApplyStartedAt: new Date(now).toISOString(), applyToken, lastBudgetWriteAt: new Date(now).toISOString() },
|
||||
} };
|
||||
}
|
||||
if (attempts >= options.maxAttempts) {
|
||||
result = { outcome: "exhausted", attempts };
|
||||
if (Number.isFinite(Date.parse(prior?.exhaustedAt ?? ""))) return null;
|
||||
return { wedgeNotification: {
|
||||
...(current.wedgeNotification ?? { reasonKey: "terminal-failed", episodeId: randomUUID(), status: "resolved" as const, transitionedAt: new Date(now).toISOString() }),
|
||||
budgetRevision: (current.wedgeNotification?.budgetRevision ?? 0) + 1,
|
||||
autoRecovery: { ...prior!, exhaustedAt: new Date(now).toISOString(), lastBudgetWriteAt: new Date(now).toISOString() },
|
||||
} };
|
||||
}
|
||||
if (Number.isFinite(lastAttempt) && now - lastAttempt < options.minAttemptSpacingMs) {
|
||||
result = { outcome: "already-claimed", attempt: attempts };
|
||||
return null;
|
||||
}
|
||||
const applyToken = randomUUID();
|
||||
const startedAt = new Date(now).toISOString();
|
||||
result = { outcome: "claimed", attempt: attempts + 1, applyToken };
|
||||
return { wedgeNotification: {
|
||||
...(current.wedgeNotification ?? { reasonKey: "terminal-failed", episodeId: randomUUID(), status: "resolved" as const, transitionedAt: startedAt }),
|
||||
budgetRevision: (current.wedgeNotification?.budgetRevision ?? 0) + 1,
|
||||
autoRecovery: {
|
||||
...prior, attempts: attempts + 1, lastAttemptAt: startedAt, retryAppliedAt: undefined,
|
||||
resumeCount: 0, lastApplyStartedAt: startedAt, applyToken, lastBudgetWriteAt: startedAt,
|
||||
...(attempts === 0 ? { budgetStartedAt: startedAt } : {}),
|
||||
...(prior?.escalationReason === "auto-recovery-disabled" ? { escalationNotifiedAt: undefined, escalationReason: undefined } : {}),
|
||||
},
|
||||
} };
|
||||
});
|
||||
return result;
|
||||
}
|
||||
async markTerminalFailureAutoRecoveryBudgetExhausted(taskId: string, options: { maxAttempts: number }): Promise<"stamped" | "already-stamped" | "not-exhausted" | "no-budget"> {
|
||||
let result: "stamped" | "already-stamped" | "not-exhausted" | "no-budget" = "no-budget";
|
||||
await this.updateTaskAtomic(taskId, (current) => {
|
||||
const budget = current.wedgeNotification?.autoRecovery;
|
||||
if (!budget) { result = "no-budget"; return null; }
|
||||
const attempts = typeof budget.attempts === "number" && Number.isFinite(budget.attempts) && budget.attempts >= 0 ? Math.trunc(budget.attempts) : 0;
|
||||
if (attempts < options.maxAttempts) { result = "not-exhausted"; return null; }
|
||||
if (Number.isFinite(Date.parse(budget.exhaustedAt ?? ""))) { result = "already-stamped"; return null; }
|
||||
const now = new Date().toISOString(); result = "stamped";
|
||||
return { wedgeNotification: { ...current.wedgeNotification!, budgetRevision: (current.wedgeNotification!.budgetRevision ?? 0) + 1, autoRecovery: { ...budget, exhaustedAt: now, lastBudgetWriteAt: now } } };
|
||||
});
|
||||
return result;
|
||||
}
|
||||
async markTerminalFailureAutoRecoveryEscalationDelivered(
|
||||
taskId: string,
|
||||
input: { dispatchOutcome: "delivered" | "suppressed"; escalationReason: "budget-exhausted" | "auto-recovery-disabled" },
|
||||
): Promise<"stamped" | "already-stamped" | "not-stamped-stale-suppression" | "no-budget"> {
|
||||
let result: "stamped" | "already-stamped" | "not-stamped-stale-suppression" | "no-budget" = "no-budget";
|
||||
await this.updateTaskAtomic(taskId, (current) => {
|
||||
const wedge = current.wedgeNotification;
|
||||
const budget = wedge?.autoRecovery;
|
||||
if (!wedge || !budget) { result = "no-budget"; return null; }
|
||||
if (
|
||||
(budget.escalationNotifiedAt && budget.escalationReason === input.escalationReason)
|
||||
|| (budget.escalationReason === "budget-exhausted" && input.escalationReason === "auto-recovery-disabled")
|
||||
) { result = "already-stamped"; return null; }
|
||||
if (input.dispatchOutcome === "suppressed") {
|
||||
const notifiedAt = Date.parse(wedge.lastNotifiedAtByReason?.["terminal-failed"] ?? "");
|
||||
const floor = Date.parse(input.escalationReason === "budget-exhausted" ? budget.exhaustedAt ?? "" : budget.budgetStartedAt ?? "");
|
||||
if (!Number.isFinite(notifiedAt) || !Number.isFinite(floor) || notifiedAt < floor) {
|
||||
result = "not-stamped-stale-suppression";
|
||||
return null;
|
||||
}
|
||||
}
|
||||
const now = new Date().toISOString(); result = "stamped";
|
||||
return { wedgeNotification: {
|
||||
...wedge, budgetRevision: (wedge.budgetRevision ?? 0) + 1,
|
||||
autoRecovery: { ...budget, escalationNotifiedAt: now, escalationReason: input.escalationReason, lastBudgetWriteAt: now },
|
||||
} };
|
||||
});
|
||||
return result;
|
||||
}
|
||||
async resetTerminalFailureAutoRecoveryBudget(taskId: string): Promise<void> {
|
||||
await this.updateTaskAtomic(taskId, (current) => {
|
||||
const wedge = current.wedgeNotification;
|
||||
if (!wedge?.autoRecovery) return null;
|
||||
return { wedgeNotification: clearTerminalFailureAutoRecoveryBudget(wedge, new Date().toISOString()) };
|
||||
});
|
||||
}
|
||||
/*
|
||||
FNXC:TaskWedgeNotifications 2026-08-10-18:54:
|
||||
R38c is used because the current move implementation owns its transaction and cannot accept
|
||||
a companion task patch. Consume the fence before the internal move under the task lock: exactly
|
||||
one observer is authorized to move, while a crash leaves a cleared, non-failed row that no
|
||||
automatic recovery path re-moves. A wall-clock window is never a move authorization.
|
||||
*/
|
||||
async applyTerminalFailureAutoRecoveryRetry(
|
||||
taskId: string,
|
||||
input: { applyToken: string; patch: Parameters<TaskStore["updateTask"]>[1]; targetColumn: ColumnId; moveOptions: MoveTaskOptions },
|
||||
): Promise<{ outcome: "applied"; task: Task } | { outcome: "superseded" | "not-failed" | "deleted" | "no-budget" }> {
|
||||
return this.withTaskLock(taskId, async () => {
|
||||
const live = await this.readTaskForMove(taskId);
|
||||
const budget = live.wedgeNotification?.autoRecovery;
|
||||
if (!budget) return { outcome: "no-budget" };
|
||||
if (live.deletedAt != null) return { outcome: "deleted" };
|
||||
if (live.status !== "failed") return { outcome: "not-failed" };
|
||||
if (!budget.applyToken || budget.applyToken !== input.applyToken) return { outcome: "superseded" };
|
||||
/*
|
||||
FNXC:TaskWedgeNotifications 2026-08-10-20:40:
|
||||
The apply fence is validated only inside `moveTaskInternal`'s advisory-locked transaction,
|
||||
which consumes it while persisting the failure clear and column move. An in-process task lock
|
||||
is merely a local belt: separate engine processes must not be able to carry the same token
|
||||
into two moves, and a failed move must roll back all three changes together.
|
||||
*/
|
||||
const expectedColumn = live.column;
|
||||
try {
|
||||
const moved = await this.moveTaskInternal(
|
||||
taskId,
|
||||
input.targetColumn,
|
||||
input.moveOptions,
|
||||
{
|
||||
fromHandoff: false,
|
||||
terminalFailureApply: {
|
||||
applyToken: input.applyToken,
|
||||
patch: input.patch,
|
||||
expectedColumn,
|
||||
},
|
||||
},
|
||||
live,
|
||||
);
|
||||
return { outcome: "applied", task: moved };
|
||||
} catch (error) {
|
||||
if (error instanceof TerminalFailureApplyRejected) return { outcome: error.outcome };
|
||||
throw error;
|
||||
}
|
||||
});
|
||||
}
|
||||
async claimNextToolFailureRetry(taskId: string, expectedCursor: number, maxRetries: number): Promise<import("./task-store/branch-and-pr-entities.js").ToolFailureRetryClaim> {
|
||||
return claimNextToolFailureRetryImpl(this, taskId, expectedCursor, maxRetries);
|
||||
}
|
||||
@@ -2746,6 +2921,8 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
|
||||
public async recordRunAuditEventBackend( tx: DbTransaction, event: { domain: string; mutationType: string; target: string; taskId: string; agentId: string; runId: string; metadata: Record<string, unknown>; }, ): Promise<void> { return recordRunAuditEventBackendImpl(this, tx, event);
|
||||
}
|
||||
async deleteTask( id: string, options?: { removeDependencyReferences?: boolean; removeLineageReferences?: boolean; allowResurrection?: boolean; githubIssueAction?: GithubIssueAction; closureContext?: TaskDeleteClosureContext; auditContext?: TaskDeleteAuditContext; }, ): Promise<Task> {
|
||||
// FNXC:TaskWedgeNotifications 2026-08-10-20:30: The backend delete transaction clears
|
||||
// a terminal-failure budget only after it wins the soft-delete claim; never pre-clear here.
|
||||
return deleteTaskImpl(this, id, options);
|
||||
}
|
||||
async deleteTaskIf(
|
||||
|
||||
@@ -21,6 +21,7 @@ import {buildDeleteCallerAuditFields, buildDeleteClosureAuditFields, type TaskDe
|
||||
import {notifyOperatorOfNonOperatorDelete} from "../task-delete-notice.js";
|
||||
import "../builtin-traits.js";
|
||||
import {normalizeTaskPriority} from "../tasks/task-priority.js";
|
||||
import {clearTerminalFailureAutoRecoveryBudget} from "../tasks/terminal-failure-auto-recovery.js";
|
||||
import {generateTaskLineageId} from "../tasks/task-lineage.js";
|
||||
import {sanitizeFileScopeInPromptContent} from "../task-store/file-scope.js";
|
||||
import {__setTaskActivityLogLimitsForTesting} from "../task-store/comments.js";
|
||||
@@ -208,6 +209,20 @@ async function deleteTaskBackendWithClaimResultImpl(store: TaskStore, id: string
|
||||
if (!reloaded) throw new TaskNotFoundError(id);
|
||||
return { claimed: false, task: store.rowToTask(store.pgRowToTaskRow(reloaded)), lineageOutcome: { clearedChildIds: [] as string[], evidenceVersionByChild: new Map<string, number>(), evidenceUnavailableChildIds: [], evidenceInsertAttempts: 0 } };
|
||||
}
|
||||
/*
|
||||
FNXC:TaskWedgeNotifications 2026-08-10-20:30:
|
||||
A soft-deleted row is invisible to the recovery sweep. Clear its terminal-failure budget
|
||||
only after this transaction won the first-delete claim, so a declined conditional delete
|
||||
cannot mute a live card and every backend delete path shares the same atomic boundary.
|
||||
*/
|
||||
const deletedRow = await readTaskRowInTransaction(tx, id, { includeDeleted: true }, projectId);
|
||||
if (!deletedRow) throw new TaskNotFoundError(id);
|
||||
const deletedTask = store.rowToTask(store.pgRowToTaskRow(deletedRow));
|
||||
if (deletedTask.wedgeNotification?.autoRecovery) {
|
||||
await tx.update(schema.project.tasks)
|
||||
.set({ wedgeNotification: JSON.stringify(clearTerminalFailureAutoRecoveryBudget(deletedTask.wedgeNotification, deletedAt)) })
|
||||
.where(and(eq(schema.project.tasks.projectId, projectId), eq(schema.project.tasks.id, id)));
|
||||
}
|
||||
// Clear lineage references and approval only after locked candidates were revalidated.
|
||||
const lineageOutcome = context
|
||||
? await removeLineageReferences(tx, id, context.candidateIds, deletedAt, layer.projectId, context.promptByChildId, lineageEvidenceTargetVersionForTest(store))
|
||||
|
||||
@@ -9,6 +9,13 @@
|
||||
import {type TaskStore, type MoveTaskOptions, type MoveTaskInternalOptions, storeLog} from "../store.js";
|
||||
import * as schema from "../postgres/schema/index.js";
|
||||
import {TaskDeletedError, HandoffInvariantViolationError, TransitionRejectionError} from "./errors.js";
|
||||
|
||||
/** @internal A fenced terminal-failure apply lost ownership before its move transaction. */
|
||||
export class TerminalFailureApplyRejected extends Error {
|
||||
constructor(readonly outcome: "superseded" | "not-failed" | "deleted" | "no-budget") {
|
||||
super(`terminal failure apply ${outcome}`);
|
||||
}
|
||||
}
|
||||
import {and, eq, sql} from "drizzle-orm";
|
||||
import type {Task, Column, ColumnId, HandoffToReviewOptions} from "../types.js";
|
||||
import {VALID_TRANSITIONS, COLUMNS} from "../types.js";
|
||||
@@ -403,7 +410,7 @@ export async function handoffToReviewImpl(store: TaskStore, taskId: string, opts
|
||||
|
||||
export async function moveTaskInternalImpl(store: TaskStore, id: string, toColumn: ColumnId, options: MoveTaskOptions | undefined, internal: MoveTaskInternalOptions, currentTask?: Task,): Promise<Task> {
|
||||
const dir = store.taskDir(id);
|
||||
const task = currentTask ?? await store.readTaskForMove(id);
|
||||
let task = currentTask ?? await store.readTaskForMove(id);
|
||||
/*
|
||||
FNXC:TaskMovement 2026-06-22-18:20:
|
||||
Public moveTask calls without an explicit source keep the legacy emitted source of "engine", but they do not inherit workflow guard bypass. Engine, scheduler, handoff, and recovery call sites opt into bypass semantics with an explicit moveSource or skipMergeBlocker.
|
||||
@@ -487,7 +494,7 @@ export async function moveTaskInternalImpl(store: TaskStore, id: string, toColum
|
||||
*/
|
||||
const moveReviewColumns = workflowIr ? new Set(resolveReviewColumns(workflowIr)) : undefined;
|
||||
|
||||
if (task.column === toColumn) {
|
||||
if (task.column === toColumn && !internal.terminalFailureApply) {
|
||||
/*
|
||||
FNXC:WorkflowResolvedColumns 2026-07-31-17:20 (fleet — the same-column handoff target):
|
||||
A MOVE TARGET, resolved from the set already computed three lines above.
|
||||
@@ -604,6 +611,9 @@ export async function moveTaskInternalImpl(store: TaskStore, id: string, toColum
|
||||
}
|
||||
|
||||
const fromColumn = task.column;
|
||||
// A fenced recovery apply may already be in its rebound column; it still needs the transaction
|
||||
// to persist the failure clear and consume the fence, but must not become an invalid self-move.
|
||||
const terminalFailureApplySameColumn = internal.terminalFailureApply !== undefined && fromColumn === toColumn;
|
||||
|
||||
if (workflowIr) {
|
||||
// ── Flag-ON validation + sync guards (typed rejections, KTD-3/R13) ─────
|
||||
@@ -636,7 +646,7 @@ export async function moveTaskInternalImpl(store: TaskStore, id: string, toColum
|
||||
// re-home (recoveryRehome) skips this so a stranded card can reach its
|
||||
// new workflow's entry column from any current column.
|
||||
const allowed = resolveAllowedColumns(workflowIr, fromColumn);
|
||||
if (options?.recoveryRehome !== true && !allowed.includes(toColumn)) {
|
||||
if (!terminalFailureApplySameColumn && options?.recoveryRehome !== true && !allowed.includes(toColumn)) {
|
||||
throw new TransitionRejectionError(
|
||||
makeTransitionRejection(
|
||||
"guard-rejected",
|
||||
@@ -823,7 +833,7 @@ export async function moveTaskInternalImpl(store: TaskStore, id: string, toColum
|
||||
options?.recoveryRehome === true &&
|
||||
!(COLUMNS as readonly string[]).includes(toColumn) &&
|
||||
workflowHasColumn(await resolveTaskWorkflowIrForMove(store, id), toColumn);
|
||||
if (!isEvacuation && !isLegacyRecoveryRehome && !isWorkflowDeclaredRecoveryRehome) {
|
||||
if (!terminalFailureApplySameColumn && !isEvacuation && !isLegacyRecoveryRehome && !isWorkflowDeclaredRecoveryRehome) {
|
||||
/*
|
||||
FNXC:WorkflowColumns 2026-07-05-19:30:
|
||||
Workflow columns graduated to always-on (no experimental flag emitted), so this "flag-OFF"
|
||||
@@ -1128,6 +1138,36 @@ export async function moveTaskInternalImpl(store: TaskStore, id: string, toColum
|
||||
(the failure mode that made the R1 sentinel unbindable).
|
||||
*/
|
||||
await acquireTaskAdvisoryXactLock(tx, layer.projectId, id);
|
||||
if (internal.terminalFailureApply) {
|
||||
/*
|
||||
FNXC:TaskWedgeNotifications 2026-08-10-20:40:
|
||||
The apply fence is checked after the database advisory lock, inside the same transaction
|
||||
that writes the cleared task and its column move. `withTaskLock` serializes only one
|
||||
process; without this re-read, two engines can both carry a previously-valid token into
|
||||
separate move transactions and the loser can requeue a card after the winner has advanced it.
|
||||
*/
|
||||
const row = await readTaskRowInTransaction(tx, id, { includeDeleted: true }, layer.projectId);
|
||||
if (!row) throw new TerminalFailureApplyRejected("no-budget");
|
||||
const live = store.rowToTask(store.pgRowToTaskRow(row));
|
||||
const budget = live.wedgeNotification?.autoRecovery;
|
||||
if (!budget) throw new TerminalFailureApplyRejected("no-budget");
|
||||
if (live.deletedAt != null) throw new TerminalFailureApplyRejected("deleted");
|
||||
if (live.status !== "failed") throw new TerminalFailureApplyRejected("not-failed");
|
||||
if (live.column !== internal.terminalFailureApply.expectedColumn || budget.applyToken !== internal.terminalFailureApply.applyToken) {
|
||||
throw new TerminalFailureApplyRejected("superseded");
|
||||
}
|
||||
const now = new Date().toISOString();
|
||||
const { applyToken: _token, lastApplyStartedAt: _started, ...remainingBudget } = budget;
|
||||
task = {
|
||||
...live,
|
||||
...internal.terminalFailureApply.patch,
|
||||
wedgeNotification: {
|
||||
...live.wedgeNotification!,
|
||||
budgetRevision: (live.wedgeNotification!.budgetRevision ?? 0) + 1,
|
||||
autoRecovery: { ...remainingBudget, retryAppliedAt: now, lastBudgetWriteAt: now },
|
||||
},
|
||||
} as Task;
|
||||
}
|
||||
const capacityWorkflowId = await readTaskWorkflowSelectionInTransaction(tx, layer.projectId, id);
|
||||
const capacityPoolId = resolveCapacityPoolId(capacityWorkflowId);
|
||||
const capacityIr = await resolveWorkflowIrForSelectedWorkflowId(store, capacityWorkflowId);
|
||||
|
||||
@@ -193,15 +193,31 @@ export type TaskColumnDescriptor = {
|
||||
* FNXC:TaskStateReconciliation 2026-07-29-20:53:
|
||||
* A generic PostgreSQL task write may carry an active wedge snapshot that was read before the live API resolved that episode. Preserve the durable resolution for the same episode so changed-column persistence, task.json projection, and cache publication cannot reactivate an acknowledged operator notification; a genuinely new wedge must use a new episode ID.
|
||||
*/
|
||||
export function preserveResolvedTaskWedgeEpisode(existingRow: Pick<TaskRow, "wedgeNotification">, task: Task): void {
|
||||
export function preserveDurableTaskWedgeInvariants(existingRow: Pick<TaskRow, "wedgeNotification">, task: Task): void {
|
||||
const durable = fromJson<Task["wedgeNotification"]>(existingRow.wedgeNotification);
|
||||
const incoming = task.wedgeNotification;
|
||||
// Keep the legacy resolved-episode rule first and byte-for-byte equivalent in behavior.
|
||||
if (
|
||||
durable?.status === "resolved"
|
||||
&& incoming?.status === "active"
|
||||
&& durable.episodeId === incoming.episodeId
|
||||
) {
|
||||
task.wedgeNotification = durable;
|
||||
return;
|
||||
}
|
||||
/*
|
||||
FNXC:TaskWedgeNotifications 2026-08-10-18:54:
|
||||
Wedge JSON is persisted wholesale from task snapshots. A newer durable budget revision
|
||||
must win over an ordinary stale writer, including durable absence after a reset, without
|
||||
overwriting the caller's unrelated episode fields.
|
||||
*/
|
||||
const durableRevision = durable?.budgetRevision ?? 0;
|
||||
const incomingRevision = incoming?.budgetRevision ?? 0;
|
||||
if (durableRevision > incomingRevision) {
|
||||
task.wedgeNotification = {
|
||||
...(incoming ?? durable ?? {} as Task["wedgeNotification"]),
|
||||
...(durable ? { budgetRevision: durableRevision, autoRecovery: durable.autoRecovery } : {}),
|
||||
} as Task["wedgeNotification"];
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -33,7 +33,7 @@ import {CentralCore} from "../central/central-core.js";
|
||||
import {extractTaskIdTokens, normalizeTitleForTaskId} from "../tasks/task-title-id-drift.js";
|
||||
import {generateTaskLineageId} from "../tasks/task-lineage.js";
|
||||
import {sanitizeFileScopeInPromptContent} from "../task-store/file-scope.js";
|
||||
import {preserveResolvedTaskWedgeEpisode, type TaskRow} from "../task-store/persistence.js";
|
||||
import {preserveDurableTaskWedgeInvariants, type TaskRow} from "../task-store/persistence.js";
|
||||
import {__setTaskActivityLogLimitsForTesting} from "../task-store/comments.js";
|
||||
import {isWorkflowDefinitionIdPrimaryKeyCollision, nextWorkflowDefinitionIdAsyncImpl} from "../task-store/workflow-definitions.js";
|
||||
import {upsertTaskRowInTransaction, buildTaskInsertValues} from "./async/async-persistence.js";
|
||||
@@ -184,7 +184,7 @@ export async function atomicWriteTaskJsonWithAuditImpl(store: TaskStore, dir: st
|
||||
)) {
|
||||
throw new Error(`Planning dependency invalidation conflict for ${id}: dependencies changed before the lifecycle mutation committed`);
|
||||
}
|
||||
preserveResolvedTaskWedgeEpisode(existing, task);
|
||||
preserveDurableTaskWedgeInvariants(existing, task);
|
||||
const changedColumns = store.getChangedTaskColumns(existing, task);
|
||||
if (changedColumns.size > 0) {
|
||||
const context = store.createTaskPersistSerializationContext(task, existing);
|
||||
|
||||
@@ -34,7 +34,7 @@ import {normalizeTaskCommitAssociation} from "../tasks/task-lineage.js";
|
||||
import {__setTaskActivityLogLimitsForTesting} from "../task-store/comments.js";
|
||||
import {withTaskBranchContextInSourceMetadata} from "../task-store/branch-context.js";
|
||||
import {upsertTaskRowInTransaction, readTaskRowInTransaction, buildTaskInsertValues} from "../task-store/async/async-persistence.js";
|
||||
import {preserveResolvedTaskWedgeEpisode} from "../task-store/persistence.js";
|
||||
import {preserveDurableTaskWedgeInvariants} from "../task-store/persistence.js";
|
||||
import {isPlanReviewSatisfied} from "../planner/plan-approval.js";
|
||||
import {listDueWorkflowWorkItems as listDueWorkflowWorkItemsAsync, withTaskWorkflowSerialization} from "../task-store/async/async-workflow-workitems.js";
|
||||
import {getTaskMovedCountsByDay as getTaskMovedCountsByDayAsync} from "../task-store/async/async-audit.js";
|
||||
@@ -106,7 +106,7 @@ export async function atomicWriteTaskJsonImpl2(store: TaskStore, dir: string, ta
|
||||
return;
|
||||
}
|
||||
const existingRow = store.pgRowToTaskRow(pgRow);
|
||||
preserveResolvedTaskWedgeEpisode(existingRow, task);
|
||||
preserveDurableTaskWedgeInvariants(existingRow, task);
|
||||
const deletedAt = store.getSoftDeletedWriteConflict(id, task, existingRow);
|
||||
if (deletedAt) {
|
||||
store.throwSoftDeletedWriteBlocked(id, deletedAt, "atomicWriteTaskJson");
|
||||
|
||||
@@ -21,6 +21,7 @@ export * from "./step-parsers.js";
|
||||
export * from "./symbol-lock-lineage-approval.js";
|
||||
export * from "./symbol-lock-types.js";
|
||||
export * from "./task-age-staleness.js";
|
||||
export * from "./terminal-failure-auto-recovery.js";
|
||||
export * from "./task-creation-hooks.js";
|
||||
export * from "./task-fields.js";
|
||||
export * from "./task-helpers.js";
|
||||
|
||||
114
packages/core/src/tasks/terminal-failure-auto-recovery.ts
Normal file
114
packages/core/src/tasks/terminal-failure-auto-recovery.ts
Normal file
@@ -0,0 +1,114 @@
|
||||
import type { Task } from "../types/task/task-core.js";
|
||||
|
||||
/** The generic terminal-failure recovery budget is intentionally smaller than legacy transient recovery. */
|
||||
export const MAX_TERMINAL_FAILURE_AUTO_RETRIES = 3;
|
||||
/** A lost apply can be resumed once without spending a second park attempt. */
|
||||
export const MAX_TERMINAL_FAILURE_AUTO_RESUMES = 1;
|
||||
/** The claim window detects abandonment; it is not a cross-process lock. */
|
||||
export const TERMINAL_FAILURE_CLAIM_APPLY_GRACE_MS = 60_000;
|
||||
/** A budget is eligible for foreign-write staleness cleanup after one week. */
|
||||
export const TERMINAL_FAILURE_BUDGET_MAX_AGE_MS = 7 * 24 * 60 * 60 * 1_000;
|
||||
/** Avoid sub-millisecond ambiguity between a budget write and store-assigned updatedAt. */
|
||||
export const BUDGET_WRITE_GUARD_MS = 1_000;
|
||||
|
||||
export type TerminalFailureAutoRecoveryDecision =
|
||||
| { action: "retry"; attempt: number }
|
||||
| { action: "notify"; reason: "budget-exhausted" | "auto-recovery-disabled" }
|
||||
| { action: "reset-budget"; reason: "terminal-success" | "archived-or-deleted" | "budget-stale" }
|
||||
| { action: "skip"; reason: "not-generic-terminal-failure" | "no-action" | "escalation-already-delivered" | "auto-recovery-disabled-never-withheld" | "recovery-owned" | "paused-or-progressing" };
|
||||
|
||||
export interface TerminalFailureAutoRecoveryOptions {
|
||||
isGenericTerminalFailure: boolean;
|
||||
hasRecoveryOwner: boolean;
|
||||
isProgressing: boolean;
|
||||
inTerminalSuccessColumn: boolean;
|
||||
isArchivedOrDeleted: boolean;
|
||||
autoRecoveryEnabled: boolean;
|
||||
maxAttempts?: number;
|
||||
maxBudgetAgeMs?: number;
|
||||
now?: () => number;
|
||||
}
|
||||
|
||||
function nonNegativeInteger(value: unknown): number {
|
||||
return typeof value === "number" && Number.isFinite(value) && value >= 0 ? Math.trunc(value) : 0;
|
||||
}
|
||||
|
||||
/*
|
||||
FNXC:TaskWedgeNotifications 2026-08-10-20:30:
|
||||
Soft deletion is the last observable boundary for a task, so every delete entry point must
|
||||
share this reset shape with the public budget-reset primitive. The helper removes only the
|
||||
terminal-failure budget, its cooldown, and its active episode; other wedge reasons survive.
|
||||
*/
|
||||
export function clearTerminalFailureAutoRecoveryBudget(
|
||||
wedge: NonNullable<Task["wedgeNotification"]>,
|
||||
now: string,
|
||||
): NonNullable<Task["wedgeNotification"]> {
|
||||
const { autoRecovery: _budget, ...rest } = wedge;
|
||||
const lastNotifiedAtByReason = { ...(wedge.lastNotifiedAtByReason ?? {}) };
|
||||
delete lastNotifiedAtByReason["terminal-failed"];
|
||||
return {
|
||||
...rest,
|
||||
budgetRevision: (wedge.budgetRevision ?? 0) + 1,
|
||||
...(Object.keys(lastNotifiedAtByReason).length > 0 ? { lastNotifiedAtByReason } : { lastNotifiedAtByReason: undefined }),
|
||||
...(wedge.status === "active" && wedge.reasonKey === "terminal-failed" ? { status: "resolved" as const, transitionedAt: now } : {}),
|
||||
};
|
||||
}
|
||||
|
||||
function isBudgetStale(task: Task, now: number, maxAgeMs: number): boolean {
|
||||
const budget = task.wedgeNotification?.autoRecovery;
|
||||
if (!budget || nonNegativeInteger(budget.attempts) === 0) return false;
|
||||
const episodeAt = Date.parse(budget.escalationNotifiedAt ?? budget.lastAttemptAt ?? "");
|
||||
const taskUpdatedAt = Date.parse(task.updatedAt ?? "");
|
||||
const budgetWrittenAt = Date.parse(budget.lastBudgetWriteAt ?? "");
|
||||
return Number.isFinite(episodeAt)
|
||||
&& Number.isFinite(taskUpdatedAt)
|
||||
&& Number.isFinite(budgetWrittenAt)
|
||||
&& now - episodeAt > maxAgeMs
|
||||
&& taskUpdatedAt > budgetWrittenAt + BUDGET_WRITE_GUARD_MS;
|
||||
}
|
||||
|
||||
/*
|
||||
FNXC:TaskWedgeNotifications 2026-08-10-18:54:
|
||||
The generic terminal-failed park is engine-owned until its durable budget is spent.
|
||||
This pure classifier cannot duplicate engine error descriptors: callers supply the single
|
||||
`describeTaskWedge` determination. Reset precedes every failure gate so success, archival,
|
||||
and proven foreign-write staleness cannot leave a card permanently muted. The grace window is
|
||||
only an abandonment heuristic; the store apply fence, not a clock comparison, serializes moves.
|
||||
*/
|
||||
export function classifyTerminalFailureAutoRecovery(
|
||||
task: Task,
|
||||
options: TerminalFailureAutoRecoveryOptions,
|
||||
): TerminalFailureAutoRecoveryDecision {
|
||||
const budget = task.wedgeNotification?.autoRecovery;
|
||||
const now = options.now?.() ?? Date.now();
|
||||
const maxAttempts = options.maxAttempts ?? MAX_TERMINAL_FAILURE_AUTO_RETRIES;
|
||||
|
||||
if (budget) {
|
||||
if (options.inTerminalSuccessColumn) return { action: "reset-budget", reason: "terminal-success" };
|
||||
if (options.isArchivedOrDeleted) return { action: "reset-budget", reason: "archived-or-deleted" };
|
||||
if (isBudgetStale(task, now, options.maxBudgetAgeMs ?? TERMINAL_FAILURE_BUDGET_MAX_AGE_MS)) {
|
||||
return { action: "reset-budget", reason: "budget-stale" };
|
||||
}
|
||||
}
|
||||
|
||||
if (!options.isGenericTerminalFailure) return { action: "skip", reason: "not-generic-terminal-failure" };
|
||||
|
||||
const attempts = nonNegativeInteger(budget?.attempts);
|
||||
const owedReason = attempts >= maxAttempts
|
||||
? "budget-exhausted"
|
||||
: !options.autoRecoveryEnabled && budget
|
||||
? "auto-recovery-disabled"
|
||||
: undefined;
|
||||
if (owedReason) {
|
||||
if (budget?.escalationNotifiedAt && budget.escalationReason === owedReason) {
|
||||
return { action: "skip", reason: "escalation-already-delivered" };
|
||||
}
|
||||
return { action: "notify", reason: owedReason };
|
||||
}
|
||||
if (!options.autoRecoveryEnabled) return { action: "skip", reason: "auto-recovery-disabled-never-withheld" };
|
||||
if (options.hasRecoveryOwner) return { action: "skip", reason: "recovery-owned" };
|
||||
if (task.status !== "failed" || task.paused || task.userPaused || task.autoMerge === false || options.isProgressing) {
|
||||
return { action: "skip", reason: "paused-or-progressing" };
|
||||
}
|
||||
return { action: "retry", attempt: attempts + 1 };
|
||||
}
|
||||
@@ -574,6 +574,14 @@ engine's no-store fallback aligned across restarts.
|
||||
*/
|
||||
export const WEDGE_RENOTIFY_COOLDOWN_MS = 6 * 60 * 60 * 1_000;
|
||||
|
||||
/*
|
||||
FNXC:TaskWedgeNotifications 2026-08-10-18:54:
|
||||
`recoveryRetryCount` is a display mirror that executor, triage, and scheduler clear while
|
||||
creating terminal parks, so it cannot bound generic terminal-failure recovery. The durable
|
||||
budget instead lives in the existing wedge JSON and is revisioned because stale full-row writes
|
||||
can otherwise erase it. `applyToken` is a rotating fence consumed by the single clear/requeue
|
||||
transition; `lastApplyStartedAt` only detects abandonment and is never a lock.
|
||||
*/
|
||||
export interface TaskWedgeNotificationState {
|
||||
reasonKey: string;
|
||||
episodeId: string;
|
||||
@@ -586,6 +594,22 @@ export interface TaskWedgeNotificationState {
|
||||
would let X -> Y -> X re-notify X while its own cooldown is still active.
|
||||
*/
|
||||
lastNotifiedAtByReason?: Record<string, string>;
|
||||
/** Monotonic durable-state version used to reject stale whole-object writes. */
|
||||
budgetRevision?: number;
|
||||
/** Sweep-owned generic terminal-failure recovery state. */
|
||||
autoRecovery?: {
|
||||
attempts: number;
|
||||
lastAttemptAt: string;
|
||||
budgetStartedAt?: string;
|
||||
retryAppliedAt?: string;
|
||||
resumeCount?: number;
|
||||
lastApplyStartedAt?: string;
|
||||
applyToken?: string;
|
||||
lastBudgetWriteAt?: string;
|
||||
exhaustedAt?: string;
|
||||
escalationNotifiedAt?: string;
|
||||
escalationReason?: "budget-exhausted" | "auto-recovery-disabled";
|
||||
};
|
||||
}
|
||||
|
||||
export type TaskRecommendationCategory = "improvement" | "feature" | "bug" | "other";
|
||||
|
||||
@@ -3421,6 +3421,8 @@ export function registerTaskWorkflowRoutes(ctx: ApiRoutesContext, deps: TaskWork
|
||||
const autoPauseClearPatch = buildAutoPauseClearPatch(task);
|
||||
const clearedDeadlockAutoPause = Object.keys(autoPauseClearPatch).length > 0;
|
||||
const retryLogSuffix = clearedDeadlockAutoPause ? ", cleared deadlock auto-pause" : "";
|
||||
// FNXC:TaskWedgeNotifications 2026-08-10-20:15: dashboard Retry is explicit operator intervention, so it clears the spent generic-terminal budget.
|
||||
await scopedStore.resetTerminalFailureAutoRecoveryBudget(req.params.id);
|
||||
|
||||
if (isMissingWorktreeSessionRetry) {
|
||||
/*
|
||||
|
||||
@@ -0,0 +1,199 @@
|
||||
import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import type { Task, TaskStore } from "@fusion/core";
|
||||
import { createSharedPgTaskStoreTestHarness, pgDescribe } from "../../../core/src/__test-utils__/pg-test-harness.js";
|
||||
import { NotificationService } from "../notification/notification-service.js";
|
||||
import { SelfHealingManager } from "../self-healing.js";
|
||||
import { NtfyNotifier } from "../util/notifier.js";
|
||||
|
||||
function reviewFailure(): Task {
|
||||
return {
|
||||
id: "FN-8908-review",
|
||||
column: "in-review",
|
||||
status: "failed",
|
||||
error: "opaque terminal failure",
|
||||
updatedAt: "2026-08-01T00:00:00.000Z",
|
||||
columnMovedAt: "2026-08-01T00:00:00.000Z",
|
||||
wedgeNotification: undefined,
|
||||
} as Task;
|
||||
}
|
||||
|
||||
describe("SelfHealingManager terminal-failure auto recovery", () => {
|
||||
it("escalates a terminal failure discovered by the claim CAS", async () => {
|
||||
const task = {
|
||||
...reviewFailure(),
|
||||
id: "FN-8908-cas-exhausted",
|
||||
column: "todo",
|
||||
wedgeNotification: {
|
||||
budgetRevision: 2,
|
||||
autoRecovery: {
|
||||
attempts: 2,
|
||||
lastAttemptAt: "2026-08-01T00:00:00.000Z",
|
||||
retryAppliedAt: "2026-08-01T00:00:01.000Z",
|
||||
},
|
||||
},
|
||||
} as Task;
|
||||
const store = {
|
||||
getSettings: vi.fn().mockResolvedValue({
|
||||
globalPause: false,
|
||||
enginePaused: false,
|
||||
autoRecovery: { mode: "on" },
|
||||
maintenanceIntervalMs: 60_000,
|
||||
}),
|
||||
listTasks: vi.fn().mockResolvedValue([task]),
|
||||
getTask: vi.fn().mockResolvedValue(task),
|
||||
listWorkflowDefinitions: vi.fn().mockResolvedValue([]),
|
||||
claimTerminalFailureAutoRecoveryAttempt: vi.fn().mockResolvedValue({ outcome: "exhausted", attempts: 3 }),
|
||||
applyTerminalFailureAutoRecoveryRetry: vi.fn(),
|
||||
markTerminalFailureAutoRecoveryBudgetExhausted: vi.fn().mockResolvedValue("stamped"),
|
||||
markTerminalFailureAutoRecoveryEscalationDelivered: vi.fn().mockResolvedValue("stamped"),
|
||||
logEntry: vi.fn(),
|
||||
recordRunAuditEvent: vi.fn().mockResolvedValue(undefined),
|
||||
} as unknown as TaskStore;
|
||||
const notifyTaskWedge = vi.fn().mockResolvedValue("delivered");
|
||||
// FNXC:TaskWedgeNotifications 2026-08-10-20:01: NtfyNotifier is the production owner of
|
||||
// the active-service registration that the sweep reads, so this verifies the real bridge.
|
||||
new NtfyNotifier(store as never, {}, { notifyTaskWedge } as never);
|
||||
|
||||
await expect(new SelfHealingManager(store, { rootDir: process.cwd() }).autoRecoverTerminalFailures()).resolves.toBe(0);
|
||||
|
||||
expect(store.claimTerminalFailureAutoRecoveryAttempt).toHaveBeenCalledOnce();
|
||||
expect(store.applyTerminalFailureAutoRecoveryRetry).not.toHaveBeenCalled();
|
||||
expect(store.markTerminalFailureAutoRecoveryBudgetExhausted).toHaveBeenCalledWith(task.id, {
|
||||
maxAttempts: 3,
|
||||
});
|
||||
expect(notifyTaskWedge).toHaveBeenCalledWith(task, expect.objectContaining({ reasonKey: "terminal-failed" }), {
|
||||
source: "auto-recovery-escalation",
|
||||
});
|
||||
expect(store.markTerminalFailureAutoRecoveryEscalationDelivered).toHaveBeenCalledWith(task.id, {
|
||||
dispatchOutcome: "delivered",
|
||||
escalationReason: "budget-exhausted",
|
||||
});
|
||||
expect(store.recordRunAuditEvent).toHaveBeenCalledWith(expect.objectContaining({
|
||||
mutationType: "task:auto-recover-terminal-failure-exhausted",
|
||||
target: task.id,
|
||||
metadata: expect.objectContaining({
|
||||
taskId: task.id,
|
||||
escalationReason: "budget-exhausted",
|
||||
markedExhausted: true,
|
||||
outcome: "notified",
|
||||
}),
|
||||
}));
|
||||
expect(JSON.stringify((store.recordRunAuditEvent as ReturnType<typeof vi.fn>).mock.calls)).not.toContain("opaque terminal failure");
|
||||
});
|
||||
|
||||
it("requires backward-move triple proof before retrying an in-review terminal failure", async () => {
|
||||
const task = reviewFailure();
|
||||
const store = {
|
||||
getSettings: vi.fn().mockResolvedValue({
|
||||
globalPause: false,
|
||||
enginePaused: false,
|
||||
autoRecovery: { mode: "on" },
|
||||
maintenanceIntervalMs: 60_000,
|
||||
taskStuckTimeoutMs: 60_000,
|
||||
}),
|
||||
listTasks: vi.fn().mockResolvedValue([task]),
|
||||
getTask: vi.fn().mockResolvedValue(task),
|
||||
listWorkflowDefinitions: vi.fn().mockResolvedValue([]),
|
||||
recordRunAuditEvent: vi.fn().mockResolvedValue(undefined),
|
||||
claimTerminalFailureAutoRecoveryAttempt: vi.fn(),
|
||||
applyTerminalFailureAutoRecoveryRetry: vi.fn(),
|
||||
logEntry: vi.fn(),
|
||||
} as unknown as TaskStore;
|
||||
const manager = new SelfHealingManager(store, { rootDir: process.cwd() });
|
||||
const proof = vi.spyOn(manager as never, "evaluateBackwardMoveTripleProof" as never)
|
||||
.mockResolvedValue({
|
||||
ok: false,
|
||||
stalenessMs: 0,
|
||||
reason: "test-live-review-owner",
|
||||
metadata: {},
|
||||
} as never);
|
||||
|
||||
await expect(manager.autoRecoverTerminalFailures()).resolves.toBe(0);
|
||||
|
||||
expect(proof).toHaveBeenCalledWith(task, expect.objectContaining({
|
||||
stage: "auto-recover-terminal-failure",
|
||||
reason: "auto-recover-terminal-failure-review-candidate",
|
||||
}));
|
||||
expect(store.claimTerminalFailureAutoRecoveryAttempt).not.toHaveBeenCalled();
|
||||
expect(store.applyTerminalFailureAutoRecoveryRetry).not.toHaveBeenCalled();
|
||||
});
|
||||
});
|
||||
|
||||
/*
|
||||
FNXC:TaskWedgeNotifications 2026-08-10-20:30:
|
||||
This exercises the production TaskStore, NotificationService, and SelfHealingManager together.
|
||||
Mock-only branch tests cannot prove that durable claim, fenced apply, re-failure, and the
|
||||
service-first exhaustion dispatch share one persisted budget.
|
||||
*/
|
||||
pgDescribe("terminal-failure auto recovery production lifecycle", () => {
|
||||
const harness = createSharedPgTaskStoreTestHarness({ prefix: "fusion_terminal_failure_lifecycle" });
|
||||
|
||||
beforeAll(harness.beforeAll);
|
||||
afterAll(harness.afterAll);
|
||||
beforeEach(async () => { await harness.beforeEach(); });
|
||||
afterEach(async () => { await harness.afterEach(); });
|
||||
|
||||
it("withholds a generic failure, retries it durably, then escalates once", async () => {
|
||||
const store = harness.store();
|
||||
await store.updateSettings({ autoRecovery: { mode: "on" }, maintenanceIntervalMs: 60_000 } as never);
|
||||
const sendMessageOnce = vi.fn().mockResolvedValue({ message: {}, inserted: true });
|
||||
const service = new NotificationService(store, { messageStore: { on: () => undefined, sendMessageOnce } as never });
|
||||
const notifier = new NtfyNotifier(store, {}, service);
|
||||
await notifier.start();
|
||||
const notify = vi.spyOn(service, "notifyTaskWedge");
|
||||
const manager = new SelfHealingManager(store, { rootDir: process.cwd() });
|
||||
const task = await store.createTask({ description: "production generic terminal failure" } as never);
|
||||
|
||||
await store.updateTask(task.id, { status: "failed", error: "opaque terminal failure" } as never);
|
||||
await new Promise((resolve) => setImmediate(resolve));
|
||||
expect(sendMessageOnce).not.toHaveBeenCalled();
|
||||
|
||||
await manager.autoRecoverTerminalFailures();
|
||||
let current = await store.getTask(task.id);
|
||||
expect(current.status).toBeUndefined();
|
||||
expect(current.error).toBeUndefined();
|
||||
expect(current.wedgeNotification?.autoRecovery?.attempts).toBe(1);
|
||||
expect(current.wedgeNotification?.autoRecovery?.retryAppliedAt).toBeTruthy();
|
||||
expect(current.wedgeNotification?.autoRecovery?.applyToken).toBeUndefined();
|
||||
expect(current.nextRecoveryAt && Date.parse(current.nextRecoveryAt)).toBeGreaterThan(Date.now());
|
||||
expect(notify).not.toHaveBeenCalled();
|
||||
|
||||
// Retain a past display mirror across re-failure; only the durable budget owns attempts.
|
||||
for (const expectedAttempts of [2, 3]) {
|
||||
await store.updateTask(task.id, {
|
||||
status: "failed",
|
||||
error: "opaque terminal failure",
|
||||
nextRecoveryAt: new Date(Date.now() - 1).toISOString(),
|
||||
wedgeNotification: {
|
||||
...current.wedgeNotification!,
|
||||
autoRecovery: {
|
||||
...current.wedgeNotification!.autoRecovery!,
|
||||
lastAttemptAt: new Date(Date.now() - 3_600_000).toISOString(),
|
||||
retryAppliedAt: new Date(Date.now() - 3_599_000).toISOString(),
|
||||
},
|
||||
},
|
||||
} as never);
|
||||
await manager.autoRecoverTerminalFailures();
|
||||
current = await store.getTask(task.id);
|
||||
expect(current.wedgeNotification?.autoRecovery?.attempts).toBe(expectedAttempts);
|
||||
expect(current.status).toBeUndefined();
|
||||
}
|
||||
|
||||
await store.updateTask(task.id, { status: "failed", error: "opaque terminal failure" } as never);
|
||||
await new Promise((resolve) => setImmediate(resolve));
|
||||
await manager.autoRecoverTerminalFailures();
|
||||
current = await store.getTask(task.id);
|
||||
expect(current.wedgeNotification?.autoRecovery?.exhaustedAt).toBeTruthy();
|
||||
expect(current.wedgeNotification?.autoRecovery?.escalationNotifiedAt).toBeTruthy();
|
||||
expect(current.wedgeNotification?.autoRecovery?.escalationReason).toBe("budget-exhausted");
|
||||
// The task-update listener is the service-first ordering: it dispatches through the shared
|
||||
// private seam before the sweep sees the at-budget row, then the durable stamp suppresses it.
|
||||
expect(notify).not.toHaveBeenCalled();
|
||||
expect(sendMessageOnce).toHaveBeenCalledTimes(1);
|
||||
|
||||
await manager.autoRecoverTerminalFailures();
|
||||
expect(notify).not.toHaveBeenCalled();
|
||||
notifier.stop();
|
||||
manager.stop();
|
||||
});
|
||||
});
|
||||
@@ -75,6 +75,80 @@ function task(overrides: Partial<Task> = {}): Task {
|
||||
}
|
||||
|
||||
describe("NotificationService deferred failure notifications", () => {
|
||||
it("does not dispatch a stale source-tagged terminal escalation after the live budget advances", async () => {
|
||||
const store = createStore();
|
||||
const episode = vi.fn(async () => ({ claimed: true, episodeId: "should-not-claim" }));
|
||||
Object.assign(store, { claimTaskWedgeNotificationEpisode: episode });
|
||||
const service = new NotificationService(store as any);
|
||||
await service.start();
|
||||
const stale = task({
|
||||
id: "FN-stale-escalation",
|
||||
status: "failed",
|
||||
error: "opaque terminal failure",
|
||||
wedgeNotification: {
|
||||
reasonKey: "terminal-failed",
|
||||
episodeId: "stale-escalation",
|
||||
status: "active",
|
||||
transitionedAt: "2026-08-10T20:00:00.000Z",
|
||||
autoRecovery: { attempts: 1, lastAttemptAt: "2026-08-10T20:00:00.000Z" },
|
||||
},
|
||||
});
|
||||
store.setTask(stale);
|
||||
|
||||
await expect(service.notifyTaskWedge(stale, {
|
||||
reasonKey: "terminal-failed",
|
||||
reason: "The task entered a terminal failed state and needs operator intervention.",
|
||||
action: "Retry the task.",
|
||||
}, { source: "auto-recovery-escalation" })).resolves.toBe("unavailable");
|
||||
|
||||
expect(episode).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("stamps a durable suppressed exhaustion at the shared service seam", async () => {
|
||||
const store = createStore();
|
||||
const marker = vi.fn(async () => "already-stamped" as const);
|
||||
const stamp = vi.fn(async () => "stamped" as const);
|
||||
const episode = vi.fn(async () => ({ claimed: false }));
|
||||
Object.assign(store, {
|
||||
claimTaskWedgeNotificationEpisode: episode,
|
||||
markTerminalFailureAutoRecoveryBudgetExhausted: marker,
|
||||
markTerminalFailureAutoRecoveryEscalationDelivered: stamp,
|
||||
});
|
||||
const service = new NotificationService(store as any);
|
||||
await service.start();
|
||||
const exhausted = task({
|
||||
id: "FN-suppressed-exhaustion",
|
||||
status: "failed",
|
||||
error: "opaque terminal failure",
|
||||
wedgeNotification: {
|
||||
reasonKey: "terminal-failed",
|
||||
episodeId: "active",
|
||||
status: "active",
|
||||
transitionedAt: "2026-08-10T20:00:00.000Z",
|
||||
lastNotifiedAtByReason: { "terminal-failed": "2026-08-10T20:01:00.000Z" },
|
||||
autoRecovery: {
|
||||
attempts: 3,
|
||||
lastAttemptAt: "2026-08-10T20:00:00.000Z",
|
||||
budgetStartedAt: "2026-08-10T20:00:00.000Z",
|
||||
exhaustedAt: "2026-08-10T20:01:00.000Z",
|
||||
lastBudgetWriteAt: "2026-08-10T20:01:00.000Z",
|
||||
},
|
||||
},
|
||||
});
|
||||
store.setTask(exhausted);
|
||||
store.emit("task:updated", exhausted);
|
||||
await flushAsyncHandlers();
|
||||
|
||||
expect(marker).toHaveBeenCalledWith(exhausted.id, { maxAttempts: 3 });
|
||||
expect(episode).toHaveBeenCalledWith(exhausted.id, "terminal-failed");
|
||||
expect(stamp).toHaveBeenCalledWith(exhausted.id, {
|
||||
dispatchOutcome: "suppressed",
|
||||
escalationReason: "budget-exhausted",
|
||||
});
|
||||
expect(marker.mock.invocationCallOrder[0]).toBeLessThan(episode.mock.invocationCallOrder[0]);
|
||||
expect(episode.mock.invocationCallOrder[0]).toBeLessThan(stamp.mock.invocationCallOrder[0]);
|
||||
});
|
||||
|
||||
beforeEach(() => {
|
||||
vi.useFakeTimers();
|
||||
});
|
||||
@@ -147,6 +221,10 @@ describe("NotificationService deferred failure notifications", () => {
|
||||
recoveryRetryCount: undefined,
|
||||
nextRecoveryAt: undefined,
|
||||
updatedAt: "2026-08-05T05:01:00.000Z",
|
||||
wedgeNotification: {
|
||||
reasonKey: "terminal-failed", episodeId: "budget", status: "active", transitionedAt: "2026-08-05T05:00:00.000Z",
|
||||
autoRecovery: { attempts: 3, lastAttemptAt: "2026-08-05T05:00:00.000Z" },
|
||||
},
|
||||
});
|
||||
store.setTask(exhausted);
|
||||
store.emit("task:updated", exhausted);
|
||||
@@ -291,7 +369,7 @@ describe("NotificationService deferred failure notifications", () => {
|
||||
await service.stop();
|
||||
});
|
||||
|
||||
it("FN-5627: still dispatches notification for genuine concurrent-advance failures (different SHAs)", async () => {
|
||||
it("FN-5627: generic concurrent-advance failures are withheld for terminal auto-recovery", async () => {
|
||||
const { store, service, sendNotification } = await setup();
|
||||
const genuineError = "Integration branch main advanced concurrently (expected aaa1111aaa1111aaa1111aaa1111aaa1111aaaa, observed bbb2222bbb2222bbb2222bbb2222bbb2222bbbb) while applying ccc3333ccc3333ccc3333ccc3333ccc3333cccc for FN-genuine";
|
||||
store.setTask(task({ id: "FN-genuine", status: "failed", error: genuineError }));
|
||||
@@ -299,11 +377,7 @@ describe("NotificationService deferred failure notifications", () => {
|
||||
|
||||
await vi.advanceTimersByTimeAsync(500);
|
||||
|
||||
expect(sendNotification).toHaveBeenCalledTimes(1);
|
||||
expect(sendNotification).toHaveBeenCalledWith("task-wedged", expect.objectContaining({
|
||||
taskId: "FN-genuine",
|
||||
metadata: expect.objectContaining({ wedgeReason: "terminal-failed" }),
|
||||
}));
|
||||
expect(sendNotification).not.toHaveBeenCalled();
|
||||
await service.stop();
|
||||
});
|
||||
|
||||
|
||||
@@ -267,7 +267,7 @@ describe("task wedge notifications", () => {
|
||||
store.setLiveReadError(new TaskNotFoundError("FN-8501"));
|
||||
await service.start();
|
||||
|
||||
await expect(service.notifyTaskWedge(task(), describeTaskWedge(task())!)).resolves.toBeUndefined();
|
||||
await expect(service.notifyTaskWedge(task(), describeTaskWedge(task())!)).resolves.toBe("unavailable");
|
||||
store.emit(task());
|
||||
await flushWedgeHandling();
|
||||
|
||||
|
||||
@@ -11,12 +11,12 @@ import type {
|
||||
Task,
|
||||
} from "@fusion/core";
|
||||
import type { LifecycleColumns, TaskMoveLanes, WorkflowIrResolverStore } from "@fusion/core";
|
||||
import { DASHBOARD_USER_ID, isTaskNotFoundError, NotificationDispatcher, resolveProjectColumnsForRoles, resolveReviewColumns, resolveTaskLifecycleColumns, resolveWorkflowIrForTask, WEDGE_RENOTIFY_COOLDOWN_MS } from "@fusion/core";
|
||||
import { DASHBOARD_USER_ID, isTaskNotFoundError, MAX_TERMINAL_FAILURE_AUTO_RETRIES, NotificationDispatcher, resolveProjectColumnsForRoles, resolveReviewColumns, resolveTaskLifecycleColumns, resolveWorkflowIrForTask, WEDGE_RENOTIFY_COOLDOWN_MS } from "@fusion/core";
|
||||
import { DEFAULT_NTFY_EVENTS, buildNtfyClickUrl, formatTaskIdentifier } from "../util/notifier.js";
|
||||
import { schedulerLog } from "../logger.js";
|
||||
import { NtfyNotificationProvider } from "./ntfy-provider.js";
|
||||
import { WebhookNotificationProvider } from "./webhook-provider.js";
|
||||
import { describeTaskRecoveryOwner, describeTaskWedge, isTaskProgressing, type TaskWedgeDescriptor } from "./task-wedge-notification.js";
|
||||
import { classifyTerminalFailureAutoRecoveryForTask, describeTaskRecoveryOwner, describeTaskWedge, isTaskProgressing, shouldWithholdWedgeAlertForAutoRecovery, type TaskWedgeDescriptor } from "./task-wedge-notification.js";
|
||||
|
||||
export interface NotificationServiceOptions {
|
||||
/** Project identifier for notification deep links */
|
||||
@@ -46,6 +46,10 @@ interface NotificationServiceStore {
|
||||
getTask?(id: string): Promise<Task | undefined> | Task | undefined;
|
||||
/** Durable compare-and-set for restart-safe wedge delivery episodes. */
|
||||
claimTaskWedgeNotificationEpisode?(taskId: string, reasonKey: string | null): Promise<{ episodeId?: string; claimed: boolean }>;
|
||||
/** Optional so lightweight notification fakes retain their structural surface. */
|
||||
markTerminalFailureAutoRecoveryBudgetExhausted?(taskId: string, options: { maxAttempts: number }): Promise<"stamped" | "already-stamped" | "not-exhausted" | "no-budget">;
|
||||
/** Optional; self-healing's concrete store is the durable backstop when absent. */
|
||||
markTerminalFailureAutoRecoveryEscalationDelivered?(taskId: string, input: { dispatchOutcome: "delivered" | "suppressed"; escalationReason: "budget-exhausted" | "auto-recovery-disabled" }): Promise<"stamped" | "already-stamped" | "not-stamped-stale-suppression" | "no-budget">;
|
||||
/*
|
||||
FNXC:WorkflowResolvedColumns 2026-07-30-23:10 (fleet phase — notification lifecycle guards):
|
||||
The workflow-IR resolver surface, OPTIONAL. The real `TaskStore` implements all three; declaring them
|
||||
@@ -154,10 +158,10 @@ export class NotificationService {
|
||||
links never reject: `maybeNotifyTaskWedge` already owns its own error handling, and a rejected link
|
||||
would poison every later notification for that task.
|
||||
*/
|
||||
private readonly wedgeHandlingChains = new Map<string, Promise<void>>();
|
||||
private readonly wedgeHandlingChains = new Map<string, Promise<unknown>>();
|
||||
|
||||
/** Queues wedge handling for one task behind any handling already in flight for it. */
|
||||
private enqueueWedgeHandling(taskId: string, run: () => Promise<void>): Promise<void> {
|
||||
private enqueueWedgeHandling<T>(taskId: string, run: () => Promise<T>): Promise<T> {
|
||||
const previous = this.wedgeHandlingChains.get(taskId) ?? Promise.resolve();
|
||||
const next = previous.then(run, run);
|
||||
this.wedgeHandlingChains.set(taskId, next);
|
||||
@@ -369,6 +373,8 @@ export class NotificationService {
|
||||
}
|
||||
};
|
||||
|
||||
private autoRecoveryEnabled = true;
|
||||
|
||||
private handleTaskUpdated = (task: Task, meta?: { lanes?: TaskMoveLanes }): void => {
|
||||
/*
|
||||
FNXC:TaskWedgeNotifications 2026-08-09-06:30:
|
||||
@@ -378,6 +384,7 @@ export class NotificationService {
|
||||
*/
|
||||
const recoveryOwner = task.status === "failed" ? describeTaskRecoveryOwner(task) : null;
|
||||
const wedge = describeTaskWedge(task);
|
||||
const autoRecoveryWithheld = wedge != null && shouldWithholdWedgeAlertForAutoRecovery(task, { autoRecoveryEnabled: this.autoRecoveryEnabled });
|
||||
/*
|
||||
FNXC:TaskWedgeNotifications 2026-07-22-14:30:
|
||||
A generic failed push may have been scheduled before a terminal error was
|
||||
@@ -409,7 +416,7 @@ export class NotificationService {
|
||||
return;
|
||||
}
|
||||
|
||||
if (task.status === "failed" && recoveryOwner) {
|
||||
if (task.status === "failed" && (recoveryOwner || autoRecoveryWithheld)) {
|
||||
this.failureNotificationSuppressedCount += 1;
|
||||
schedulerLog.debug(`[notify] ${task.id} recovery-owned failure — suppressed notification`);
|
||||
} else if (task.status === "failed" && !wedge) {
|
||||
@@ -530,11 +537,20 @@ export class NotificationService {
|
||||
resolved-then-reparked reason to begin a new actionable episode.
|
||||
*/
|
||||
/** Delivers a self-healing no-action escalation through the durable wedge episode seam. */
|
||||
async notifyTaskWedge(task: Task, descriptor: TaskWedgeDescriptor): Promise<void> {
|
||||
await this.enqueueWedgeHandling(task.id, () => this.maybeNotifyTaskWedge(task, descriptor));
|
||||
async notifyTaskWedge(
|
||||
task: Task,
|
||||
descriptor: TaskWedgeDescriptor,
|
||||
options?: { source?: "auto-recovery-escalation" },
|
||||
): Promise<"delivered" | "suppressed" | "unavailable"> {
|
||||
return this.enqueueWedgeHandling(task.id, () => this.maybeNotifyTaskWedge(task, descriptor, options));
|
||||
}
|
||||
|
||||
private async maybeNotifyTaskWedge(task: Task, suppliedDescriptor?: TaskWedgeDescriptor | null): Promise<void> {
|
||||
private async maybeNotifyTaskWedge(
|
||||
task: Task,
|
||||
suppliedDescriptor?: TaskWedgeDescriptor | null,
|
||||
options?: { source?: "auto-recovery-escalation" },
|
||||
): Promise<"delivered" | "suppressed" | "unavailable"> {
|
||||
if (options?.source !== "auto-recovery-escalation" && shouldWithholdWedgeAlertForAutoRecovery(task, { autoRecoveryEnabled: this.autoRecoveryEnabled })) return "unavailable";
|
||||
/*
|
||||
FNXC:TaskWedgeNotifications 2026-08-10-04:35:
|
||||
`getTask` throws TaskNotFoundError for soft-deleted rows instead of returning
|
||||
@@ -547,17 +563,34 @@ export class NotificationService {
|
||||
} catch (error) {
|
||||
const detail = isTaskNotFoundError(error) ? "not found" : error instanceof Error ? error.message : String(error);
|
||||
schedulerLog.warn(`[notify] ${task.id} wedge live-read failed (${detail}) — suppressed notification`);
|
||||
return;
|
||||
return "unavailable";
|
||||
}
|
||||
const classification = classifyTerminalFailureAutoRecoveryForTask(liveTask, { autoRecoveryEnabled: this.autoRecoveryEnabled });
|
||||
const requestedAutoRecoveryEscalation = options?.source === "auto-recovery-escalation";
|
||||
const ownEscalation = requestedAutoRecoveryEscalation
|
||||
|| (suppliedDescriptor === undefined && classification.action === "notify");
|
||||
/*
|
||||
FNXC:TaskWedgeNotifications 2026-08-10-20:32:
|
||||
The sweep can hold a snapshot while a manual retry, reset, or another engine
|
||||
advances the live row. A source-tagged terminal-failed dispatch is valid only
|
||||
while the live budget still owes the same escalation; otherwise it must not
|
||||
turn a stale snapshot into an alert or delivery stamp. Non-terminal supplied
|
||||
descriptors retain their existing behavior and never become budget stamps.
|
||||
*/
|
||||
if (
|
||||
requestedAutoRecoveryEscalation
|
||||
&& suppliedDescriptor?.reasonKey === "terminal-failed"
|
||||
&& classification.action !== "notify"
|
||||
) return "unavailable";
|
||||
const recoveryOwner = describeTaskRecoveryOwner(liveTask);
|
||||
if (recoveryOwner) {
|
||||
if (recoveryOwner && !ownEscalation) {
|
||||
// Recovery ownership is not a wedge episode. Resolve only an episode we
|
||||
// can prove active, avoiding a write/claim for a never-notified snapshot.
|
||||
if (this.activeWedgeReasons.has(task.id) || liveTask.wedgeNotification?.status === "active") {
|
||||
this.activeWedgeReasons.delete(task.id);
|
||||
await this.store.claimTaskWedgeNotificationEpisode?.(task.id, null);
|
||||
}
|
||||
return;
|
||||
return "unavailable";
|
||||
}
|
||||
/*
|
||||
FNXC:TaskWedgeNotifications 2026-08-10-04:35:
|
||||
@@ -621,15 +654,43 @@ export class NotificationService {
|
||||
this.activeWedgeReasons.delete(task.id);
|
||||
await this.store.claimTaskWedgeNotificationEpisode?.(task.id, null);
|
||||
}
|
||||
return;
|
||||
return "unavailable";
|
||||
}
|
||||
const isAutoRecoveryEscalationDispatch = ownEscalation && descriptor.reasonKey === "terminal-failed";
|
||||
if (isAutoRecoveryEscalationDispatch && classification.action === "notify" && classification.reason === "budget-exhausted") {
|
||||
try {
|
||||
await this.store.markTerminalFailureAutoRecoveryBudgetExhausted?.(task.id, { maxAttempts: MAX_TERMINAL_FAILURE_AUTO_RETRIES });
|
||||
} catch (error) {
|
||||
schedulerLog.debug(`[notify] ${task.id} could not mark terminal recovery exhaustion: ${error instanceof Error ? error.message : String(error)}`);
|
||||
}
|
||||
}
|
||||
if (this.store.claimTaskWedgeNotificationEpisode) {
|
||||
const claim = await this.store.claimTaskWedgeNotificationEpisode(task.id, descriptor.reasonKey);
|
||||
if (!claim.claimed || !claim.episodeId) return;
|
||||
if (!claim.claimed || !claim.episodeId) {
|
||||
/*
|
||||
FNXC:TaskWedgeNotifications 2026-08-10-20:40:
|
||||
A durable episode-CAS decline proves an equivalent terminal-failed dispatch is already
|
||||
inside the current cooldown. Stamp this escalation at the shared seam, subject to the
|
||||
store's budget/owed-marker floor validation, so a service-first suppression cannot wait
|
||||
for a later sweep and reopen into a repeat alert. The fallback episode is intentionally
|
||||
excluded: its process-local cooldown cannot prove a durable budget delivery.
|
||||
*/
|
||||
if (isAutoRecoveryEscalationDispatch && classification.action === "notify") {
|
||||
try {
|
||||
await this.store.markTerminalFailureAutoRecoveryEscalationDelivered?.(task.id, {
|
||||
dispatchOutcome: "suppressed",
|
||||
escalationReason: classification.reason,
|
||||
});
|
||||
} catch (error) {
|
||||
schedulerLog.debug(`[notify] ${task.id} could not stamp suppressed terminal recovery delivery: ${error instanceof Error ? error.message : String(error)}`);
|
||||
}
|
||||
}
|
||||
return "suppressed";
|
||||
}
|
||||
episode = claim.episodeId;
|
||||
} else {
|
||||
episode = this.claimFallbackWedgeNotificationEpisode(task.id, descriptor.reasonKey, task.updatedAt);
|
||||
if (!episode) return;
|
||||
if (!episode) return "suppressed";
|
||||
}
|
||||
const link = buildNtfyClickUrl({ dashboardHost: this.dashboardHost, projectId: this.options.projectId, taskId: task.id });
|
||||
const content = [
|
||||
@@ -653,6 +714,17 @@ export class NotificationService {
|
||||
} catch (error) {
|
||||
schedulerLog.log(`[notify] ${task.id} wedge mailbox message failed: ${error instanceof Error ? error.message : String(error)}`);
|
||||
}
|
||||
if (isAutoRecoveryEscalationDispatch && classification.action === "notify") {
|
||||
try {
|
||||
await this.store.markTerminalFailureAutoRecoveryEscalationDelivered?.(task.id, {
|
||||
dispatchOutcome: "delivered",
|
||||
escalationReason: classification.reason,
|
||||
});
|
||||
} catch (error) {
|
||||
schedulerLog.debug(`[notify] ${task.id} could not stamp terminal recovery delivery: ${error instanceof Error ? error.message : String(error)}`);
|
||||
}
|
||||
}
|
||||
return "delivered";
|
||||
}
|
||||
|
||||
private claimFallbackWedgeNotificationEpisode(taskId: string, reasonKey: string, updatedAt: string): string | undefined {
|
||||
@@ -1044,6 +1116,9 @@ export class NotificationService {
|
||||
}
|
||||
|
||||
private refreshFailureNotificationSettings(settings: Settings): void {
|
||||
// Fail open when no maintenance pass can own the retry budget; classic wedge notification remains available.
|
||||
this.autoRecoveryEnabled = settings.autoRecovery?.mode !== "off"
|
||||
&& (settings.maintenanceIntervalMs === undefined || settings.maintenanceIntervalMs > 0);
|
||||
this.failureNotificationDelayMs =
|
||||
typeof settings.failureNotificationDelayMs === "number" && settings.failureNotificationDelayMs >= 0
|
||||
? settings.failureNotificationDelayMs
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import type { Task } from "@fusion/core";
|
||||
import { classifyTerminalFailureAutoRecovery, type TerminalFailureAutoRecoveryDecision, type Task } from "@fusion/core";
|
||||
import { hasTransientMergeRecoveryOwner } from "../errors/transient-merge-error-classifier.js";
|
||||
|
||||
/** A bounded, operator-safe description of a task that cannot make progress. */
|
||||
@@ -141,6 +141,34 @@ Terminal task updates are the shared delivery seam for merger, executor, heartbe
|
||||
and self-healing writers. Classify only states that have no scheduled owner; raw
|
||||
error output is never used as an idempotency key or forwarded into audit metadata.
|
||||
*/
|
||||
/*
|
||||
FNXC:TaskWedgeNotifications 2026-08-10-18:54:
|
||||
Core owns only the durable budget rules. The generic failure classification stays here so
|
||||
specific terminal descriptors cannot drift from notification withholding or self-healing.
|
||||
A past display mirror is not a live recovery owner for this adapter.
|
||||
*/
|
||||
export function classifyTerminalFailureAutoRecoveryForTask(
|
||||
task: Task,
|
||||
options: { autoRecoveryEnabled: boolean; inTerminalSuccessColumn?: boolean; isArchivedOrDeleted?: boolean; now?: number },
|
||||
): TerminalFailureAutoRecoveryDecision {
|
||||
const now = options.now ?? Date.now();
|
||||
const nextRecoveryAt = Date.parse(task.nextRecoveryAt ?? "");
|
||||
return classifyTerminalFailureAutoRecovery(task, {
|
||||
isGenericTerminalFailure: describeTaskWedge(task)?.reasonKey === "terminal-failed",
|
||||
hasRecoveryOwner: describeTaskRecoveryOwner(task) !== null && Number.isFinite(nextRecoveryAt) && nextRecoveryAt > now,
|
||||
isProgressing: isTaskProgressing(task),
|
||||
inTerminalSuccessColumn: options.inTerminalSuccessColumn === true,
|
||||
isArchivedOrDeleted: options.isArchivedOrDeleted === true || task.deletedAt != null,
|
||||
autoRecoveryEnabled: options.autoRecoveryEnabled,
|
||||
now: () => now,
|
||||
});
|
||||
}
|
||||
|
||||
export function shouldWithholdWedgeAlertForAutoRecovery(task: Task, options: { autoRecoveryEnabled: boolean }): boolean {
|
||||
const decision = classifyTerminalFailureAutoRecoveryForTask(task, options);
|
||||
return decision.action === "retry" || (decision.action === "skip" && decision.reason === "escalation-already-delivered");
|
||||
}
|
||||
|
||||
export function describeTaskWedge(task: Task): TaskWedgeDescriptor | null {
|
||||
if (isTaskProgressing(task)) return null;
|
||||
const error = task.error ?? "";
|
||||
|
||||
@@ -119,7 +119,13 @@ import { runSurfacingSweep, hours, type SurfacingCycle } from "./surfacing-sweep
|
||||
self-healing-git-evidence.ts. Imported back here because call sites remain. */
|
||||
import { SelfHealingGitEvidence, execAsync, shellQuote } from "./self-healing-git-evidence.js";
|
||||
import { evaluateParkedAgentTaskLink, PARKED_AGENT_LINK_FRESH_RUN_MS } from "./agents/task-agent-sync.js";
|
||||
import { describeSelfHealingNoActionWedge } from "./notification/task-wedge-notification.js";
|
||||
import { classifyTerminalFailureAutoRecoveryForTask, describeSelfHealingNoActionWedge, describeTaskWedge } from "./notification/task-wedge-notification.js";
|
||||
import {
|
||||
MAX_TERMINAL_FAILURE_AUTO_RESUMES,
|
||||
MAX_TERMINAL_FAILURE_AUTO_RETRIES,
|
||||
TERMINAL_FAILURE_CLAIM_APPLY_GRACE_MS,
|
||||
} from "@fusion/core";
|
||||
import { BASE_DELAY_MS, computeRecoveryDecision, MAX_DELAY_MS } from "./healing/recovery-policy.js";
|
||||
|
||||
export {
|
||||
COMPLETED_BLOCKED_PAUSE_REASON,
|
||||
@@ -879,6 +885,8 @@ export class SelfHealingManager extends SelfHealingGitEvidence {
|
||||
// ── Maintenance timer ───────────────────────────────────────────────
|
||||
private maintenanceInterval: ReturnType<typeof setInterval> | null = null;
|
||||
private maintenanceRunning = false;
|
||||
/** In-process belt only; the durable apply token remains the cross-process fence. */
|
||||
private autoRecoverTerminalFailuresInFlight = false;
|
||||
|
||||
// ── Event listener cleanup ──────────────────────────────────────────
|
||||
private settingsListener: ((data: { settings: Settings; previous: Settings }) => void) | null = null;
|
||||
@@ -1760,6 +1768,7 @@ export class SelfHealingManager extends SelfHealingGitEvidence {
|
||||
{ name: "recover-orphan-only-scope-violations", fn: () => this.recoverOrphanOnlyScopeViolations().then(() => undefined) },
|
||||
{ name: "recover-stuck-merge-deadlocks", fn: () => this.recoverStuckMergeDeadlocks().then(() => undefined) },
|
||||
{ name: "misclassified-failures", fn: () => this.recoverMisclassifiedFailures().then(() => undefined) },
|
||||
{ name: "auto-recover-terminal-failures", fn: () => this.autoRecoverTerminalFailures().then(() => undefined) },
|
||||
{ name: "partial-progress-no-task-done", fn: () => this.recoverPartialProgressNoTaskDoneFailures().then(() => undefined) },
|
||||
{ name: "orphaned-executions", fn: () => this.recoverOrphanedExecutions().then(() => undefined) },
|
||||
{ name: "approved-triage", fn: () => this.recoverApprovedTriageTasks().then(() => undefined) },
|
||||
@@ -2867,6 +2876,7 @@ export class SelfHealingManager extends SelfHealingGitEvidence {
|
||||
{ name: "recover-misclassified-failures", fn: () => this.recoverMisclassifiedFailures() },
|
||||
{ name: "recover-no-progress-no-task-done", fn: () => this.recoverNoProgressNoTaskDoneFailures() },
|
||||
{ name: "recover-paused-abort-failures", fn: () => this.recoverPausedAbortFailures() },
|
||||
{ name: "auto-recover-terminal-failures", fn: () => this.autoRecoverTerminalFailures() },
|
||||
{ name: "recover-partial-progress-no-task-done", fn: () => this.recoverPartialProgressNoTaskDoneFailures() },
|
||||
{ name: "recover-orphaned-executions", fn: () => this.recoverOrphanedExecutions() },
|
||||
{ name: "recover-approved-triage", fn: () => this.recoverApprovedTriageTasks() },
|
||||
@@ -12519,6 +12529,245 @@ const movedTask = await this.store.moveTask(task.id, completeLane);
|
||||
*
|
||||
* @returns Number of tasks recovered
|
||||
*/
|
||||
/*
|
||||
FNXC:TaskWedgeNotifications 2026-08-10-20:15:
|
||||
Generic terminal parks are retried from their durable wedge budget, not from the
|
||||
transient display mirror. The claim token is consumed by TaskStore's fenced apply
|
||||
transition; this sweep intentionally never moves a card itself.
|
||||
*/
|
||||
async autoRecoverTerminalFailures(): Promise<number> {
|
||||
if (this.autoRecoverTerminalFailuresInFlight) return 0;
|
||||
const settings = await this.store.getSettings();
|
||||
if (settings.globalPause || settings.enginePaused) return 0;
|
||||
this.autoRecoverTerminalFailuresInFlight = true;
|
||||
try {
|
||||
// A disabled maintenance interval drains already-withheld escalations but must not claim retries.
|
||||
const enabled = settings.autoRecovery?.mode !== "off"
|
||||
&& (settings.maintenanceIntervalMs === undefined || settings.maintenanceIntervalMs > 0);
|
||||
const tasks = await this.store.listTasks({ slim: true });
|
||||
const completeColumns = await resolveProjectColumnsForRoles(this.store, ["complete"]);
|
||||
const archivedColumns = await resolveProjectColumnsForRoles(this.store, ["archived"]);
|
||||
const reviewColumnsByWorkflow = new Map<string, Awaited<ReturnType<typeof resolveWorkflowIrForTask>>>();
|
||||
let changed = 0;
|
||||
type TerminalFailureAuditOutcome = "retried" | "resumed-claim" | "already-claimed" | "apply-superseded" | "apply-aborted-not-failed" | "budget-reset-on-success" | "budget-reset-stale" | "cleared-stale-recovery-mirror" | "notified" | "notify-suppressed" | "notify-suppressed-stale" | "notify-unavailable" | "escalation-already-delivered";
|
||||
/*
|
||||
FNXC:TaskWedgeNotifications 2026-08-10-20:32:
|
||||
A terminal-failure recovery is otherwise invisible after it clears status.
|
||||
Emit a bounded audit record for every material recovery or escalation result;
|
||||
retain only ids, counters, columns, and fixed outcomes so opaque task errors
|
||||
and the rotating apply token never enter durable telemetry.
|
||||
*/
|
||||
const auditTerminalFailure = async (
|
||||
task: Task,
|
||||
outcome: TerminalFailureAuditOutcome,
|
||||
input: { attempt?: number; escalationReason?: "budget-exhausted" | "auto-recovery-disabled"; markedExhausted?: boolean } = {},
|
||||
): Promise<void> => {
|
||||
try {
|
||||
await createRunAuditor(this.store, {
|
||||
runId: generateSyntheticRunId("auto-recover-terminal-failure", task.id),
|
||||
agentId: "self-healing",
|
||||
taskId: task.id,
|
||||
taskLineageId: task.lineageId,
|
||||
phase: "auto-recover-terminal-failure",
|
||||
}).database({
|
||||
type: input.escalationReason ? "task:auto-recover-terminal-failure-exhausted" : "task:auto-recover-terminal-failure",
|
||||
target: task.id,
|
||||
metadata: {
|
||||
taskId: task.id,
|
||||
attempt: input.attempt,
|
||||
maxAttempts: MAX_TERMINAL_FAILURE_AUTO_RETRIES,
|
||||
column: task.column,
|
||||
escalationOwed: input.escalationReason !== undefined,
|
||||
escalationReason: input.escalationReason,
|
||||
markedExhausted: input.markedExhausted === true,
|
||||
outcome,
|
||||
},
|
||||
});
|
||||
} catch (error) {
|
||||
log.debug(`autoRecoverTerminalFailures: audit write failed for ${task.id}: ${error instanceof Error ? error.message : String(error)}`);
|
||||
}
|
||||
};
|
||||
const escalateTerminalFailure = async (
|
||||
task: Task,
|
||||
reason: "budget-exhausted" | "auto-recovery-disabled",
|
||||
): Promise<{ outcome: Extract<TerminalFailureAuditOutcome, "notified" | "notify-suppressed" | "notify-suppressed-stale" | "notify-unavailable">; markedExhausted: boolean }> => {
|
||||
let markedExhausted = false;
|
||||
if (reason === "budget-exhausted") {
|
||||
// FNXC:TaskWedgeNotifications 2026-08-10-20:01: A claim CAS can discover that a
|
||||
// concurrent writer spent the final attempt after the classifier selected retry. That
|
||||
// race still owes the same marker-first escalation as the ordinary notify branch; never
|
||||
// drop it merely because the CAS, rather than the classifier, found exhaustion.
|
||||
try {
|
||||
markedExhausted = await this.store.markTerminalFailureAutoRecoveryBudgetExhausted(task.id, { maxAttempts: MAX_TERMINAL_FAILURE_AUTO_RETRIES }) === "stamped";
|
||||
} catch (error) {
|
||||
log.debug(`autoRecoverTerminalFailures: could not mark exhaustion for ${task.id}: ${error instanceof Error ? error.message : String(error)}`);
|
||||
}
|
||||
}
|
||||
const descriptor = describeTaskWedge(task);
|
||||
const service = getActiveNotificationService();
|
||||
let outcome: "delivered" | "suppressed" | "unavailable" = "unavailable";
|
||||
if (descriptor && service) {
|
||||
try {
|
||||
outcome = await service.notifyTaskWedge(task, descriptor, { source: "auto-recovery-escalation" });
|
||||
} catch (error) {
|
||||
log.debug(`autoRecoverTerminalFailures: terminal recovery notification failed for ${task.id}: ${error instanceof Error ? error.message : String(error)}`);
|
||||
}
|
||||
}
|
||||
if ((outcome === "delivered" || outcome === "suppressed") && service) {
|
||||
try {
|
||||
const stamped = await this.store.markTerminalFailureAutoRecoveryEscalationDelivered(task.id, {
|
||||
dispatchOutcome: outcome,
|
||||
escalationReason: reason,
|
||||
});
|
||||
const auditOutcome = outcome === "delivered"
|
||||
? "notified"
|
||||
: stamped === "not-stamped-stale-suppression"
|
||||
? "notify-suppressed-stale"
|
||||
: "notify-suppressed";
|
||||
await auditTerminalFailure(task, auditOutcome, { escalationReason: reason, markedExhausted });
|
||||
return { outcome: auditOutcome, markedExhausted };
|
||||
} catch (error) {
|
||||
// FNXC:TaskWedgeNotifications 2026-08-10-21:05: A failed durable
|
||||
// confirmation leaves the escalation owed for the next sweep.
|
||||
log.debug(`autoRecoverTerminalFailures: could not stamp terminal recovery delivery for ${task.id}: ${error instanceof Error ? error.message : String(error)}`);
|
||||
}
|
||||
}
|
||||
await auditTerminalFailure(task, "notify-unavailable", { escalationReason: reason, markedExhausted });
|
||||
return { outcome: "notify-unavailable", markedExhausted };
|
||||
};
|
||||
for (const snapshot of tasks) {
|
||||
if (snapshot.status !== "failed" && !snapshot.wedgeNotification?.autoRecovery) continue;
|
||||
let task = await this.store.getTask(snapshot.id);
|
||||
if (!task) continue;
|
||||
/*
|
||||
FNXC:TaskWedgeNotifications 2026-08-10-19:22:
|
||||
The retry display mirror is not recovery ownership after its window expires. The wedge
|
||||
descriptor intentionally recognizes any parseable mirror, including a past one, so clear
|
||||
stale mirrors before classification or a re-failed card is silently classified non-generic.
|
||||
This write preserves the durable budget and its apply fence; only a budget-owned watermark
|
||||
and revision may advance with the mirror clear.
|
||||
*/
|
||||
const hasRecoveryMirror = typeof task.recoveryRetryCount === "number";
|
||||
const mirrorIsLive = typeof task.nextRecoveryAt === "string" && Date.parse(task.nextRecoveryAt) > Date.now();
|
||||
// FNXC:TaskWedgeNotifications 2026-08-10-21:05: Disabled recovery drains
|
||||
// owed alerts and resets budgets only; it must not mutate a display mirror.
|
||||
if (enabled && task.status === "failed" && hasRecoveryMirror && !mirrorIsLive) {
|
||||
await this.store.updateTaskAtomic(task.id, (current) => {
|
||||
const currentHasMirror = typeof current.recoveryRetryCount === "number";
|
||||
const currentMirrorIsLive = typeof current.nextRecoveryAt === "string" && Date.parse(current.nextRecoveryAt) > Date.now();
|
||||
if (current.status !== "failed" || !currentHasMirror || currentMirrorIsLive) return null;
|
||||
const budget = current.wedgeNotification?.autoRecovery;
|
||||
if (!budget) return { recoveryRetryCount: null, nextRecoveryAt: null };
|
||||
const now = new Date().toISOString();
|
||||
return {
|
||||
recoveryRetryCount: null,
|
||||
nextRecoveryAt: null,
|
||||
wedgeNotification: {
|
||||
...current.wedgeNotification!,
|
||||
budgetRevision: (current.wedgeNotification!.budgetRevision ?? 0) + 1,
|
||||
autoRecovery: { ...budget, lastBudgetWriteAt: now },
|
||||
},
|
||||
};
|
||||
});
|
||||
await auditTerminalFailure(task, "cleared-stale-recovery-mirror");
|
||||
task = await this.store.getTask(snapshot.id);
|
||||
if (!task) continue;
|
||||
}
|
||||
const decision = classifyTerminalFailureAutoRecoveryForTask(task, {
|
||||
autoRecoveryEnabled: enabled,
|
||||
inTerminalSuccessColumn: completeColumns.has(task.column),
|
||||
isArchivedOrDeleted: archivedColumns.has(task.column) || task.deletedAt != null,
|
||||
});
|
||||
if (decision.action === "reset-budget") {
|
||||
await this.store.resetTerminalFailureAutoRecoveryBudget(task.id);
|
||||
await auditTerminalFailure(task, decision.reason === "budget-stale" ? "budget-reset-stale" : "budget-reset-on-success");
|
||||
changed += 1;
|
||||
continue;
|
||||
}
|
||||
if (decision.action === "notify") {
|
||||
await escalateTerminalFailure(task, decision.reason);
|
||||
continue;
|
||||
}
|
||||
if (decision.action === "skip") {
|
||||
if (decision.reason === "escalation-already-delivered") {
|
||||
await auditTerminalFailure(task, "escalation-already-delivered");
|
||||
}
|
||||
continue;
|
||||
}
|
||||
if (!enabled) continue;
|
||||
if (this.options.hasLiveSessionSurface?.(task.id) || executingTaskLock.has(task.id) || this.options.getExecutingTaskIds?.().has(task.id)) continue;
|
||||
/*
|
||||
FNXC:TaskWedgeNotifications 2026-08-10-19:53:
|
||||
Requeuing a failed review-lane card is a backward lifecycle move, so it needs the
|
||||
same triple proof as every other review recovery before spending an auto-recovery
|
||||
attempt. The fenced apply prevents duplicate moves; this proof independently proves
|
||||
that no live session, usable worktree, or recent activity still owns the card.
|
||||
*/
|
||||
if ((await this.resolveReviewColumnsFor(task.id, reviewColumnsByWorkflow)).has(task.column)) {
|
||||
const proof = await this.evaluateBackwardMoveTripleProof(task, {
|
||||
stage: "auto-recover-terminal-failure",
|
||||
graceMs: settings.taskStuckTimeoutMs ?? STALE_ACTIVE_BRANCH_EXECUTION_GRACE_MS,
|
||||
stalenessAnchor: task.columnMovedAt ?? task.updatedAt,
|
||||
reason: "auto-recover-terminal-failure-review-candidate",
|
||||
});
|
||||
if (!proof.ok) {
|
||||
await this.emitBackwardMoveNoAction(
|
||||
task,
|
||||
"auto-recover-terminal-failure",
|
||||
"task:auto-recover-terminal-failure",
|
||||
proof,
|
||||
);
|
||||
continue;
|
||||
}
|
||||
}
|
||||
const rawAttempts = task.wedgeNotification?.autoRecovery?.attempts;
|
||||
const attempts = typeof rawAttempts === "number" && Number.isFinite(rawAttempts) && rawAttempts >= 0
|
||||
? Math.trunc(rawAttempts)
|
||||
: 0;
|
||||
const spacing = attempts <= 0 ? 0 : 0.9 * Math.min(BASE_DELAY_MS * 2 ** (attempts - 1), MAX_DELAY_MS);
|
||||
const claim = await this.store.claimTerminalFailureAutoRecoveryAttempt(task.id, {
|
||||
maxAttempts: MAX_TERMINAL_FAILURE_AUTO_RETRIES,
|
||||
maxResumes: MAX_TERMINAL_FAILURE_AUTO_RESUMES,
|
||||
minAttemptSpacingMs: spacing,
|
||||
claimApplyGraceMs: TERMINAL_FAILURE_CLAIM_APPLY_GRACE_MS,
|
||||
});
|
||||
if (claim.outcome === "exhausted") {
|
||||
await escalateTerminalFailure(task, "budget-exhausted");
|
||||
continue;
|
||||
}
|
||||
if (claim.outcome === "already-claimed") {
|
||||
await auditTerminalFailure(task, "already-claimed", { attempt: claim.attempt });
|
||||
continue;
|
||||
}
|
||||
const delayMs = computeRecoveryDecision({ recoveryRetryCount: claim.attempt - 1 }).delayMs;
|
||||
const applied = await this.store.applyTerminalFailureAutoRecoveryRetry(task.id, {
|
||||
applyToken: claim.applyToken,
|
||||
patch: {
|
||||
status: null,
|
||||
error: null,
|
||||
paused: false,
|
||||
recoveryRetryCount: claim.attempt,
|
||||
nextRecoveryAt: new Date(Date.now() + delayMs).toISOString(),
|
||||
},
|
||||
targetColumn: await resolveReboundTargetForTask(this.store, task.id),
|
||||
moveOptions: { preserveProgress: true, moveSource: "engine" },
|
||||
});
|
||||
if (applied.outcome === "applied") {
|
||||
await this.store.logEntry(task.id, `Auto-recovered generic terminal failure (attempt ${claim.attempt}/${MAX_TERMINAL_FAILURE_AUTO_RETRIES})`);
|
||||
await auditTerminalFailure(task, claim.outcome === "resume" ? "resumed-claim" : "retried", { attempt: claim.attempt });
|
||||
changed += 1;
|
||||
} else if (applied.outcome === "superseded" || applied.outcome === "no-budget") {
|
||||
await auditTerminalFailure(task, "apply-superseded", { attempt: claim.attempt });
|
||||
} else {
|
||||
await auditTerminalFailure(task, "apply-aborted-not-failed", { attempt: claim.attempt });
|
||||
}
|
||||
}
|
||||
return changed;
|
||||
} finally {
|
||||
this.autoRecoverTerminalFailuresInFlight = false;
|
||||
}
|
||||
}
|
||||
|
||||
async recoverMisclassifiedFailures(): Promise<number> {
|
||||
try {
|
||||
/*
|
||||
|
||||
@@ -517,6 +517,13 @@ export type DatabaseMutationType =
|
||||
| "mergeQueue:stale-lease-on-column-exit"
|
||||
| "mergeQueue:auto-cleanup-stale-row"
|
||||
| "task:auto-recover-already-merged"
|
||||
/*
|
||||
FNXC:TaskWedgeNotifications 2026-08-10-19:12:
|
||||
Generic terminal recovery records only durable identifiers and bounded outcomes.
|
||||
The apply token is a fencing capability, so audit rows must never persist it or task error prose.
|
||||
*/
|
||||
| "task:auto-recover-terminal-failure"
|
||||
| "task:auto-recover-terminal-failure-exhausted"
|
||||
| "task:auto-recover-finalize-already-on-main"
|
||||
/** Metadata: { taskId, previousColumn, targetColumn, commitSha, status, blockedBy, overlapBlockedBy, reason } */
|
||||
| "task:auto-merge-finalize-column-mismatch-reconciled"
|
||||
|
||||
Reference in New Issue
Block a user