diff --git a/.changeset/fn-8799-recovery-notifications.md b/.changeset/fn-8799-recovery-notifications.md new file mode 100644 index 0000000000..ed6dd8c583 --- /dev/null +++ b/.changeset/fn-8799-recovery-notifications.md @@ -0,0 +1,7 @@ +--- +"@runfusion/fusion": patch +--- + +summary: Suppress task failure alerts while Fusion automatically recovers the task. +category: fix +dev: Wedge episodes now require recovery ownership to be absent or exhausted. diff --git a/docs/architecture.md b/docs/architecture.md index 24906f8eea..e7691493e9 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -10,6 +10,8 @@ This document describes the actual architecture of Fusion as implemented in this Actionable terminal task updates are classified into bounded reasons such as a named merge gate, retry exhaustion, or a completion blocker. The PostgreSQL-backed task row persists an active/resolved episode with an opaque identity, so `NotificationService` sends one `task-wedged` provider event and one dashboard system-mailbox message per active reason even across restarts. Repeated observations remain quiet until an authoritative non-wedge task update resolves the episode; changed and resolved-then-reentered reasons notify again. +A failed snapshot is not actionable while persisted automatic-recovery ownership remains: a scheduled recovery has both its retry counter and deadline, while transient merge recovery has an in-budget persisted retry counter. `NotificationService` re-reads the live task immediately before a wedge claim and again when a generic failure grace timer fires, so recovery that begins after a failed event cannot create a mailbox row or `task-wedged` provider event. Explicit operator-action parks and cleared/exhausted recovery markers remain terminal and claim exactly one episode. + 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. ## Planning dependency lifecycle lock diff --git a/packages/engine/src/errors/transient-merge-error-classifier.ts b/packages/engine/src/errors/transient-merge-error-classifier.ts index c1319343bf..dee90416fe 100644 --- a/packages/engine/src/errors/transient-merge-error-classifier.ts +++ b/packages/engine/src/errors/transient-merge-error-classifier.ts @@ -54,8 +54,12 @@ */ // Imports the import-free leaf, NOT `transient-error-detector.js` — that module pulls // `usage-limit-detector.js → logger.js`, the exact chain FN-5627 split this file out to avoid. +import type { Task } from "@fusion/core"; import { isTransientError } from "./transient-error-patterns.js"; +/** Shared persisted retry budget for automatic transient merge recovery. */ +export const MAX_AUTO_MERGE_TRANSIENT_RETRIES = 5; + /* FNXC:MergeReliability 2026-07-15-18:30: This classifier used to recognize only git/lease/spawn faults, while the inline retry gate in @@ -94,3 +98,17 @@ export function classifyTransientMergeError(error: string | null | undefined): s } return null; } + +/* +FNXC:TaskWedgeNotifications 2026-08-05-04:53: +A transient-error string is diagnostic evidence, not notification ownership by +itself. Fusion suppresses terminal alerts only after the merge writer persists +an in-budget retry marker; the same shared cap makes exhaustion operator-actionable. +*/ +export function hasTransientMergeRecoveryOwner(task: Pick): boolean { + return classifyTransientMergeError(task.error) !== null + && typeof task.mergeTransientRetryCount === "number" + && Number.isInteger(task.mergeTransientRetryCount) + && task.mergeTransientRetryCount >= 0 + && task.mergeTransientRetryCount < MAX_AUTO_MERGE_TRANSIENT_RETRIES; +} diff --git a/packages/engine/src/notification/__tests__/notification-service.test.ts b/packages/engine/src/notification/__tests__/notification-service.test.ts index e5265943a7..324e8055f6 100644 --- a/packages/engine/src/notification/__tests__/notification-service.test.ts +++ b/packages/engine/src/notification/__tests__/notification-service.test.ts @@ -98,6 +98,78 @@ describe("NotificationService deferred failure notifications", () => { return { store, service, sendNotification }; } + /* + FNXC:TaskWedgeNotifications 2026-08-05-04:53: + The reported sequence persists a failed snapshot while scheduler recovery owns + it. No delivery or durable wedge claim is allowed until the writer clears that + ownership at exhaustion, when the existing once-per-episode seam must alert. + */ + it("suppresses recovery-owned failed snapshots and alerts once after exhaustion", async () => { + const store = createStore(); + const sendMessageOnce = vi.fn(async () => ({ message: {} as any, inserted: true })); + const sendNotification = vi.fn(async () => ({ success: true, providerId: "mock" })); + let activeReason: string | undefined; + const claimTaskWedgeNotificationEpisode = vi.fn(async (_taskId: string, reasonKey: string | null) => { + if (reasonKey === null) { + activeReason = undefined; + return { claimed: false }; + } + if (activeReason === reasonKey) return { claimed: false }; + activeReason = reasonKey; + return { claimed: true, episodeId: `episode:${reasonKey}` }; + }); + Object.assign(store, { claimTaskWedgeNotificationEpisode }); + const service = new NotificationService(store as any, { + messageStore: { on: () => undefined, sendMessageOnce } as any, + failedNotificationGraceMs: 100, + }); + service.registerProvider({ getProviderId: () => "mock", isEventSupported: () => true, sendNotification }); + await service.start(); + + const recovering = task({ + id: "FN-recovering", + status: "failed", + error: "opaque executor failure", + recoveryRetryCount: 1, + nextRecoveryAt: "2026-08-05T05:00:00.000Z", + }); + store.setTask(recovering); + store.emit("task:updated", recovering); + await flushAsyncHandlers(); + + expect(sendMessageOnce).not.toHaveBeenCalled(); + expect(sendNotification).not.toHaveBeenCalled(); + expect(claimTaskWedgeNotificationEpisode).not.toHaveBeenCalled(); + expect(service.getPendingFailureCount()).toBe(0); + + const exhausted = task({ + ...recovering, + recoveryRetryCount: undefined, + nextRecoveryAt: undefined, + updatedAt: "2026-08-05T05:01:00.000Z", + }); + store.setTask(exhausted); + store.emit("task:updated", exhausted); + await vi.waitFor(() => expect(sendMessageOnce).toHaveBeenCalledTimes(1)); + expect(sendNotification).toHaveBeenCalledWith("task-wedged", expect.objectContaining({ taskId: "FN-recovering" })); + expect(claimTaskWedgeNotificationEpisode).toHaveBeenCalledTimes(1); + + store.emit("task:updated", exhausted); + await flushAsyncHandlers(); + expect(sendMessageOnce).toHaveBeenCalledTimes(1); + expect(sendNotification).toHaveBeenCalledTimes(1); + + await service.stop(); + const restarted = new NotificationService(store as any, { messageStore: { on: () => undefined, sendMessageOnce } as any }); + restarted.registerProvider({ getProviderId: () => "restarted", isEventSupported: () => true, sendNotification }); + await restarted.start(); + store.emit("task:updated", exhausted); + await flushAsyncHandlers(); + expect(sendMessageOnce).toHaveBeenCalledTimes(1); + expect(sendNotification).toHaveBeenCalledTimes(1); + await restarted.stop(); + }); + it("Failure that persists past grace dispatches exactly once", async () => { const { store, service, sendNotification } = await setup(); store.setTask(task({ id: "FN-1", status: "failed" })); @@ -145,11 +217,13 @@ describe("NotificationService deferred failure notifications", () => { id: "FN-5628", status: "failed", error: "Merge handoff refused (lease-handoff-failed): target-not-queued", + mergeTransientRetryCount: 1, })); store.emit("task:updated", task({ id: "FN-5628", status: "failed", error: "Merge handoff refused (lease-handoff-failed): target-not-queued", + mergeTransientRetryCount: 1, })); await vi.advanceTimersByTimeAsync(500); @@ -162,8 +236,8 @@ describe("NotificationService deferred failure notifications", () => { it("FN-5627: suppresses notification for transient same-SHA spurious-concurrent-advance failures", async () => { const { store, service, sendNotification } = await setup(); const transientError = "Integration branch main advanced concurrently (expected 694970b2f186fac31c1819d55ef30a2ad207b5c3, observed 694970b2f186fac31c1819d55ef30a2ad207b5c3) while applying b26f8fe1ee2d3dc36acf3571d42507b24bd8066b for FN-5626"; - store.setTask(task({ id: "FN-5626", status: "failed", error: transientError })); - store.emit("task:updated", task({ id: "FN-5626", status: "failed", error: transientError })); + store.setTask(task({ id: "FN-5626", status: "failed", error: transientError, mergeTransientRetryCount: 1 })); + store.emit("task:updated", task({ id: "FN-5626", status: "failed", error: transientError, mergeTransientRetryCount: 1 })); await vi.advanceTimersByTimeAsync(500); diff --git a/packages/engine/src/notification/__tests__/task-wedge-notification.test.ts b/packages/engine/src/notification/__tests__/task-wedge-notification.test.ts index 5d5d0270c3..b411785cea 100644 --- a/packages/engine/src/notification/__tests__/task-wedge-notification.test.ts +++ b/packages/engine/src/notification/__tests__/task-wedge-notification.test.ts @@ -1,7 +1,8 @@ import { describe, expect, it, vi } from "vitest"; import { WEDGE_RENOTIFY_COOLDOWN_MS, type NotificationProvider, type Settings, type Task } from "@fusion/core"; import { NotificationService } from "../notification-service.js"; -import { describeSelfHealingNoActionWedge, describeTaskWedge } from "../task-wedge-notification.js"; +import { MAX_AUTO_MERGE_TRANSIENT_RETRIES } from "../../errors/transient-merge-error-classifier.js"; +import { describeSelfHealingNoActionWedge, describeTaskRecoveryOwner, describeTaskWedge } from "../task-wedge-notification.js"; type Listener = (task: Task) => void; @@ -317,6 +318,29 @@ describe("task wedge notifications", () => { expect(describeTaskWedge(task({ status: "failed", paused: true, pausedReason }))).toMatchObject({ reasonKey }); }); + it.each([ + ["missing error and recovery metadata", { error: undefined }, false, null], + ["scheduled executor recovery", { error: "opaque failure", recoveryRetryCount: 1, nextRecoveryAt: "2026-08-05T05:00:00.000Z" }, true, null], + ["persisted transient merge retry", { error: "socket hang up", mergeTransientRetryCount: 1 }, true, null], + ["exhausted transient merge retry", { error: "socket hang up", mergeTransientRetryCount: MAX_AUTO_MERGE_TRANSIENT_RETRIES }, false, "terminal-failed"], + ["cleared recovery state", { error: "opaque failure", recoveryRetryCount: null, nextRecoveryAt: null }, false, "terminal-failed"], + ])("classifies %s without treating raw error text as automatic ownership", (_name, overrides, hasOwner, reasonKey) => { + const { task } = fixture(); + const failed = task(overrides as Partial); + expect(describeTaskRecoveryOwner(failed)).toEqual(hasOwner ? expect.anything() : null); + expect(describeTaskWedge(failed)).toEqual(reasonKey === null ? null : expect.objectContaining({ reasonKey })); + }); + + it("keeps explicit terminal pause reasons actionable despite stale recovery metadata", () => { + const { task } = fixture(); + expect(describeTaskWedge(task({ + paused: true, + pausedReason: "error-retry-exhausted", + recoveryRetryCount: 1, + nextRecoveryAt: "2026-08-05T05:00:00.000Z", + }))).toMatchObject({ reasonKey: "heartbeat-retry-exhausted" }); + }); + it("classifies an otherwise unknown persisted failure with a bounded fallback", () => { const { task } = fixture(); expect(describeTaskWedge(task({ error: "internal stack trace or opaque failure" }))).toMatchObject({ reasonKey: "terminal-failed" }); diff --git a/packages/engine/src/notification/notification-service.ts b/packages/engine/src/notification/notification-service.ts index 62a65f490a..1390f6118a 100644 --- a/packages/engine/src/notification/notification-service.ts +++ b/packages/engine/src/notification/notification-service.ts @@ -14,10 +14,9 @@ import type { LifecycleColumns, TaskMoveLanes, WorkflowIrResolverStore } from "@ import { DASHBOARD_USER_ID, 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 { classifyTransientMergeError } from "../errors/transient-merge-error-classifier.js"; import { NtfyNotificationProvider } from "./ntfy-provider.js"; import { WebhookNotificationProvider } from "./webhook-provider.js"; -import { describeTaskWedge, type TaskWedgeDescriptor } from "./task-wedge-notification.js"; +import { describeTaskRecoveryOwner, describeTaskWedge, type TaskWedgeDescriptor } from "./task-wedge-notification.js"; export interface NotificationServiceOptions { /** Project identifier for notification deep links */ @@ -372,14 +371,13 @@ export class NotificationService { private handleTaskUpdated = (task: Task, meta?: { lanes?: TaskMoveLanes }): void => { /* - FNXC:TaskWedgeNotifications 2026-07-22-20:00: - FN-5627 transient merge failures retain an active recovery owner despite - their temporary failed status. Classify them before claiming a durable wedge - episode: the generic terminal-failed fallback must not bypass its grace and - self-healing suppression, or turn a recoverable flap into an operator alert. + FNXC:TaskWedgeNotifications 2026-08-05-04:53: + A task update is a point-in-time snapshot. The pure classifier recognizes + only persisted recovery ownership, while `maybeNotifyTaskWedge` re-reads + live state before a durable claim so recovery cannot produce a false alert. */ - const transientFailure = task.status === "failed" ? classifyTransientMergeError(task.error) : null; - const wedge = transientFailure ? null : describeTaskWedge(task); + const recoveryOwner = task.status === "failed" ? describeTaskRecoveryOwner(task) : null; + const wedge = describeTaskWedge(task); /* FNXC:TaskWedgeNotifications 2026-07-22-14:30: A generic failed push may have been scheduled before a terminal error was @@ -387,7 +385,7 @@ export class NotificationService { only operator notification; dispatch-time suppression below covers races. */ if (wedge) this.cancelPendingFailureNotification(task.id, "classified-terminal-wedge"); - if (!transientFailure) void this.enqueueWedgeHandling(task.id, () => this.maybeNotifyTaskWedge(task, wedge)); + void this.enqueueWedgeHandling(task.id, () => this.maybeNotifyTaskWedge(task, wedge)); void this.maybeSuppressTransientFailedNotification(task, `status=${task.status ?? "undefined"}`); /* @@ -411,27 +409,12 @@ export class NotificationService { return; } - if (task.status === "failed" && !wedge) { - // FN-5627: Suppress notifications entirely for transient merge failure - // classes recognized by `classifyTransientMergeError`. These are - // recovered automatically by `SelfHealingManager.recoverTransientMergeFailures` - // and the per-tick auto-recovery in `project-engine.ts` fast-path; the - // task either lands cleanly on a retry or stays in in-review for the - // bounded recovery budget to handle. Without this guard, every flap - // cycle (typically every ~5 min when the merger keeps hitting the same - // transient class) fires another ntfy alarm even though the task is - // never genuinely stuck — producing user-facing alarm spam with no - // actionable information. - const transientClass = classifyTransientMergeError(task.error); - if (transientClass) { - this.failureNotificationSuppressedCount += 1; - schedulerLog.debug( - `[notify] ${task.id} transient merge failure (${transientClass}) — suppressed notification (self-heal in flight)`, - ); - return; - } + if (task.status === "failed" && recoveryOwner) { + this.failureNotificationSuppressedCount += 1; + schedulerLog.debug(`[notify] ${task.id} recovery-owned failure — suppressed notification`); + } else if (task.status === "failed" && !wedge) { if (this.failureNotificationMode === "all") { - this.maybeNotify(task.id, "failed", this.createTaskPayload(task, "failed")); + void this.maybeNotifyImmediateFailure(task); } else { this.scheduleFailureNotification(task); } @@ -543,7 +526,21 @@ export class NotificationService { } private async maybeNotifyTaskWedge(task: Task, suppliedDescriptor?: TaskWedgeDescriptor | null): Promise { - const descriptor = suppliedDescriptor ?? describeTaskWedge(task); + // Task events carry snapshots. Re-read before a durable claim so an immediate + // recovery update cannot turn a stale failed event into an operator alert. + const liveTask = this.store.getTask ? (await this.store.getTask(task.id)) ?? task : task; + const recoveryOwner = describeTaskRecoveryOwner(liveTask); + if (recoveryOwner) { + // 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; + } + const descriptor = suppliedDescriptor ?? describeTaskWedge(liveTask); + task = liveTask; let episode: string | undefined; if (!descriptor) { /* @@ -1021,6 +1018,19 @@ export class NotificationService { this.failureNotificationMode = settings.failureNotificationMode ?? "sticky-only"; } + private async maybeNotifyImmediateFailure(task: Task): Promise { + const liveTask = this.store.getTask ? (await this.store.getTask(task.id)) ?? task : task; + if ( + liveTask.status !== "failed" + || describeTaskRecoveryOwner(liveTask) + || describeTaskWedge(liveTask) + ) { + this.failureNotificationSuppressedCount += 1; + return; + } + this.maybeNotify(liveTask.id, "failed", this.createTaskPayload(liveTask, "failed")); + } + private scheduleFailureNotification(task: Task): void { if (this.pendingFailureNotifications.has(task.id)) { return; @@ -1094,17 +1104,15 @@ export class NotificationService { return; } - // FN-5627 defense-in-depth: even when a failure notification was scheduled - // (e.g., the failure happened slightly before the transient classifier - // suppression landed on a newer cycle), re-check at dispatch time before - // terminal-wedge classification. A transient failed task still has an - // automatic recovery owner and must not claim a wedge episode. - const transientClassAtDispatch = classifyTransientMergeError(task.error); - if (transientClassAtDispatch) { + /* + FNXC:TaskWedgeNotifications 2026-08-05-04:53: + A grace timer owns only delayed generic delivery, never the task lifecycle. + Re-check durable recovery ownership here so a retry scheduled after the + original failed event produces neither a mailbox row nor a provider dispatch. + */ + if (describeTaskRecoveryOwner(task)) { this.failureNotificationSuppressedCount += 1; - schedulerLog.debug( - `[notify] ${taskId} transient merge failure (${transientClassAtDispatch}) at dispatch time — suppressed notification (self-heal in flight)`, - ); + schedulerLog.debug(`[notify] ${taskId} recovery-owned failure at dispatch time — suppressed notification`); return; } diff --git a/packages/engine/src/notification/task-wedge-notification.ts b/packages/engine/src/notification/task-wedge-notification.ts index 7b8ff8b540..c4ec0684d0 100644 --- a/packages/engine/src/notification/task-wedge-notification.ts +++ b/packages/engine/src/notification/task-wedge-notification.ts @@ -1,4 +1,5 @@ import 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. */ export interface TaskWedgeDescriptor { @@ -8,6 +9,31 @@ export interface TaskWedgeDescriptor { gate?: string; } +/** Durable evidence that a failed snapshot remains assigned to a bounded automatic recovery path. */ +export interface TaskRecoveryOwner { + kind: "scheduled-recovery" | "transient-merge-retry"; +} + +/* +FNXC:TaskWedgeNotifications 2026-08-05-04:53: +Mailbox and push wedge alerts mean an operator must act. A failed snapshot with +both scheduler retry fields, or an in-budget transient-merge retry marker, is +still owned by Fusion and must not create an actionable notification episode. +*/ +export function describeTaskRecoveryOwner(task: Task): TaskRecoveryOwner | null { + if (hasTransientMergeRecoveryOwner(task)) return { kind: "transient-merge-retry" }; + if ( + typeof task.recoveryRetryCount === "number" + && Number.isInteger(task.recoveryRetryCount) + && task.recoveryRetryCount >= 0 + && typeof task.nextRecoveryAt === "string" + && Number.isFinite(Date.parse(task.nextRecoveryAt)) + ) { + return { kind: "scheduled-recovery" }; + } + return null; +} + /** * FNXC:TaskWedgeNotifications 2026-07-22-14:30: * Self-healing can deliberately decline a backward move without mutating task @@ -144,13 +170,12 @@ export function describeTaskWedge(task: Task): TaskWedgeDescriptor | null { return { reasonKey: "merge-blocked", reason: "Merge verification cannot progress without operator action.", action: "Fix the failing verification, then retry the task." }; } /* - FNXC:TaskWedgeNotifications 2026-07-22-20:00: - An opaque failure or exhausted merge-retry budget is terminal evidence when no - named writer classified it. NotificationService checks FN-5627 transient merge - failures before calling this classifier, preserving self-healing ownership. - Failures without either signal retain the generic grace path because recovery - may still own them. + FNXC:TaskWedgeNotifications 2026-08-05-04:53: + The generic fallback cannot convert a recovery-owned failed snapshot into an + operator alert. Explicit terminal pause and error contracts above remain + actionable; recovery exhaustion clears its durable marker and reaches this path. */ + if (describeTaskRecoveryOwner(task)) return null; if (!error && (task.mergeRetries ?? 0) < 3) return null; return { reasonKey: "terminal-failed", diff --git a/packages/engine/src/project-engine.ts b/packages/engine/src/project-engine.ts index b3ca6d60a1..85b5d73998 100644 --- a/packages/engine/src/project-engine.ts +++ b/packages/engine/src/project-engine.ts @@ -95,7 +95,7 @@ import { ResearchProviderRegistry } from "./research/provider-registry.js"; import { createRunAuditor, generateSyntheticRunId } from "./util/run-audit.js"; import { finalizeProvenAutoMergeTask } from "./merge/auto-merge-finalization.js"; import { isTransientError } from "./errors/transient-error-detector.js"; -import { classifyTransientMergeError } from "./errors/transient-merge-error-classifier.js"; +import { classifyTransientMergeError, MAX_AUTO_MERGE_TRANSIENT_RETRIES } from "./errors/transient-merge-error-classifier.js"; import { TunnelProcessManager } from "./remote-access/tunnel-process-manager.js"; import { deliverPostgresMigrationCompleteNoticeIfNeeded, @@ -566,7 +566,7 @@ export class ProjectEngine { * * Readable (not private) so tests derive the cap from this single source of truth rather * than hardcoding it — the FN-8004 bump broke two suites that had baked in the old `3`. */ - static readonly MAX_AUTO_MERGE_TRANSIENT_RETRIES = 5; + static readonly MAX_AUTO_MERGE_TRANSIENT_RETRIES = MAX_AUTO_MERGE_TRANSIENT_RETRIES; private static readonly MERGE_REQUEST_RETRY_EXHAUSTED_AGE_MS = 30 * 60 * 1000; /** Cap on outer in-review→in-progress bounces caused by deterministic * verification failures during auto-merge. After this many failed merges