FN-8799: suppress alerts during automatic task recovery

Prevent terminal task alerts while persisted automatic recovery owns a failed task.

- Classify scheduled and transient merge recovery ownership before wedge notification claims.
- Re-read live task state before claims and deferred failure delivery to avoid recovery races.
- Add recovery suppression coverage, architecture documentation, and a patch changeset.

Files changed:
 .changeset/fn-8799-recovery-notifications.md       |  7 ++
 docs/architecture.md                               |  2 +
 .../src/errors/transient-merge-error-classifier.ts | 18 +++++
 .../__tests__/notification-service.test.ts         | 78 ++++++++++++++++++-
 .../__tests__/task-wedge-notification.test.ts      | 26 ++++++-
 .../src/notification/notification-service.ts       | 90 ++++++++++++----------
 .../src/notification/task-wedge-notification.ts    | 37 +++++++--
 packages/engine/src/project-engine.ts              |  4 +-
 8 files changed, 210 insertions(+), 52 deletions(-)

Fusion-Task-Id: FN-8799

Fusion-Task-Lineage: f8296e81-6b84-4697-8b20-63823f09d584

Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
gsxdsm
2026-08-04 22:15:25 -07:00
parent 3071018593
commit 26219b8d34
8 changed files with 210 additions and 52 deletions

View File

@@ -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.

View File

@@ -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

View File

@@ -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<Task, "error" | "mergeTransientRetryCount">): boolean {
return classifyTransientMergeError(task.error) !== null
&& typeof task.mergeTransientRetryCount === "number"
&& Number.isInteger(task.mergeTransientRetryCount)
&& task.mergeTransientRetryCount >= 0
&& task.mergeTransientRetryCount < MAX_AUTO_MERGE_TRANSIENT_RETRIES;
}

View File

@@ -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);

View File

@@ -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<Task>);
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" });

View File

@@ -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<void> {
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<void> {
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;
}

View File

@@ -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",

View File

@@ -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