FN-8691: suppress repeated task wedge notifications

Prevent resolved-and-rewedged tasks from flooding operator notification channels.

- Persist per-reason notification timestamps with a six-hour cooldown.
- Apply cooldown claims consistently in durable storage and in-memory fallback paths.
- Cover reason transitions, legacy state, invalid timestamps, and cooldown expiry.

Files changed:
 .changeset/fn-8691-wedge-notification-cooldown.md  |   7 ++
 docs/architecture.md                               |   4 +-
 .../postgres/store-wedge-resolution.pg.test.ts     |  67 +++++++++++-
 packages/core/src/index.ts                         |   1 +
 packages/core/src/store.ts                         |  21 +++-
 packages/core/src/task-store/task-mutation-ops.ts  |   6 ++
 packages/core/src/types.ts                         |   2 +
 packages/core/src/types/task-core.ts               |  15 +++
 .../__tests__/task-wedge-notification.test.ts      | 117 ++++++++++++++++++++-
 .../src/notification/notification-service.ts       |  30 +++++-
 10 files changed, 261 insertions(+), 9 deletions(-)

Fusion-Task-Id: FN-8691

Fusion-Task-Lineage: b7460508-f9b0-488b-828c-a51e9477304d

Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
gsxdsm
2026-08-01 08:58:49 -07:00
parent 79d2a73a10
commit ac8ce148d5
10 changed files with 261 additions and 9 deletions

View File

@@ -0,0 +1,7 @@
---
"@runfusion/fusion": patch
---
summary: Prevent repeat task-wedge alerts from flooding operator inboxes.
category: fix
dev: Adds a six-hour durable per-reason cooldown that survives resolve/re-wedge flaps.

View File

@@ -8,7 +8,9 @@ This document describes the actual architecture of Fusion as implemented in this
## Terminal task wedge notifications
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. Provider and mailbox delivery are independently best-effort, while run-audit metadata remains ids/counts/outcomes-only.
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.
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.
## 1) Overview

View File

@@ -9,6 +9,7 @@
* A deletion that commits after the compare-and-set but before projection must retain deleted-task semantics instead of returning a stale successful resolution.
*/
import { afterAll, afterEach, beforeAll, beforeEach, expect, it, vi } from "vitest";
import { WEDGE_RENOTIFY_COOLDOWN_MS } from "../../types/task-core.js";
import {
createSharedPgTaskStoreTestHarness,
pgDescribe,
@@ -26,9 +27,73 @@ pgTest("TaskStore wedge episode resolution (PostgreSQL)", () => {
beforeAll(h.beforeAll);
beforeEach(h.beforeEach);
afterEach(h.afterEach);
afterEach(() => vi.restoreAllMocks());
afterEach(() => {
vi.restoreAllMocks();
vi.useRealTimers();
});
afterAll(h.afterAll);
/*
FNXC:TaskWedgeNotifications 2026-08-01-15:35:
Cooldown timestamps are keyed by normalized reason rather than the current episode.
This covers the reported resolve/re-wedge flap and proves X -> Y -> X retains X's
own suppression window while a distinct operator action remains immediately visible.
*/
it("claims wedge notifications once per reason cooldown window", async () => {
let now = Date.parse("2026-08-01T12:00:00.000Z");
vi.spyOn(Date, "now").mockImplementation(() => now);
const store = h.store();
const task = await h.createTestTask();
const first = await store.claimTaskWedgeNotificationEpisode(task.id, "execution-blocked");
expect(first.claimed).toBe(true);
await store.resolveTaskWedgeNotificationEpisode(task.id, first.episodeId!);
const repeated = await store.claimTaskWedgeNotificationEpisode(task.id, "execution-blocked");
expect(repeated).toEqual({ claimed: false });
const suppressed = await store.getTask(task.id);
expect(suppressed.wedgeNotification).toMatchObject({
reasonKey: "execution-blocked",
status: "active",
lastNotifiedAtByReason: { "execution-blocked": "2026-08-01T12:00:00.000Z" },
});
await store.resolveTaskWedgeNotificationEpisode(task.id, suppressed.wedgeNotification!.episodeId);
const different = await store.claimTaskWedgeNotificationEpisode(task.id, "merge-blocked:check:lint");
expect(different.claimed).toBe(true);
await store.resolveTaskWedgeNotificationEpisode(task.id, different.episodeId!);
expect(await store.claimTaskWedgeNotificationEpisode(task.id, "execution-blocked")).toEqual({ claimed: false });
const secondSuppressed = await store.getTask(task.id);
await store.resolveTaskWedgeNotificationEpisode(task.id, secondSuppressed.wedgeNotification!.episodeId);
now += WEDGE_RENOTIFY_COOLDOWN_MS;
const afterCooldown = await store.claimTaskWedgeNotificationEpisode(task.id, "execution-blocked");
expect(afterCooldown.claimed).toBe(true);
const legacy = await h.createTestTask();
await store.updateTask(legacy.id, {
wedgeNotification: {
reasonKey: "legacy",
episodeId: "legacy-episode",
status: "resolved",
transitionedAt: "2026-08-01T00:00:00.000Z",
},
});
expect((await store.claimTaskWedgeNotificationEpisode(legacy.id, "execution-blocked")).claimed).toBe(true);
const invalid = await h.createTestTask();
await store.updateTask(invalid.id, {
wedgeNotification: {
reasonKey: "invalid",
episodeId: "invalid-episode",
status: "resolved",
transitionedAt: "2026-08-01T00:00:00.000Z",
lastNotifiedAtByReason: { "execution-blocked": "not-a-timestamp" },
},
});
expect((await store.claimTaskWedgeNotificationEpisode(invalid.id, "execution-blocked")).claimed).toBe(true);
});
it("resolves the exact active episode", async () => {
const store = h.store();
const task = await h.createTestTask();

View File

@@ -39,6 +39,7 @@ export type {
MissionLineageSnapshot,
} from "./symbol-lock-lineage-approval.js";
export { PLANNER_AGENT_ROLE, AGENT_VALID_TRANSITIONS, DUPLICATE_OF_METADATA_KEY, REPORT_ATTACHMENT_SOURCE, assertNotWorkspaceTaskMerge, isWorkspaceTask, WorkspaceTaskMergeError } from "./types.js";
export { WEDGE_RENOTIFY_COOLDOWN_MS } from "./types.js";
export {
resolveEntryPointBranchAssignment,
sanitizeBranchSegment,

View File

@@ -2,6 +2,7 @@ import { EventEmitter } from "node:events";
import type { TaskMoveLanes } from "./workflow-lifecycle-traits.js";
import { TaskLaneCache } from "./task-lane-cache.js";
import { randomUUID } from "node:crypto";
import { WEDGE_RENOTIFY_COOLDOWN_MS } from "./types/task-core.js";
import { join } from "node:path";
import { and, eq, isNull, ne, sql } from "drizzle-orm";
import * as schema from "./postgres/schema/index.js";
@@ -1394,9 +1395,25 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
return { wedgeNotification: { ...prior, status: "resolved", transitionedAt: new Date().toISOString() } };
}
if (prior?.status === "active" && prior.reasonKey === reasonKey) return null;
const now = Date.now();
const lastNotifiedAtByReason = Object.fromEntries(Object.entries(prior?.lastNotifiedAtByReason ?? {}).filter(([, timestamp]) => {
const notifiedAt = Date.parse(timestamp);
return Number.isFinite(notifiedAt) && notifiedAt <= now && now - notifiedAt < WEDGE_RENOTIFY_COOLDOWN_MS;
}));
const episodeId = randomUUID();
result = { episodeId, claimed: true };
return { wedgeNotification: { reasonKey, episodeId, status: "active", transitionedAt: new Date().toISOString() } };
const suppressed = reasonKey in lastNotifiedAtByReason;
if (!suppressed) lastNotifiedAtByReason[reasonKey] = new Date(now).toISOString();
result = suppressed ? { claimed: false } : { episodeId, claimed: true };
return {
wedgeNotification: {
reasonKey,
episodeId,
status: "active",
transitionedAt: new Date(now).toISOString(),
...(Object.keys(lastNotifiedAtByReason).length > 0 ? { lastNotifiedAtByReason } : {}),
},
};
});
return result;
}

View File

@@ -374,6 +374,12 @@ export async function updateTaskAtomicImpl(store: TaskStore, id: string, updater
});
}
/*
FNXC:TaskWedgeNotifications 2026-08-01-15:35:
Resolution changes only the active episode status. The PostgreSQL compare-and-set
merges that field into the existing JSON, preserving per-reason cooldown stamps so
resolving X or notifying Y cannot reopen X's live spam window.
*/
export async function resolveTaskWedgeNotificationEpisodeImpl(
store: TaskStore,
id: string,

View File

@@ -544,6 +544,7 @@ import {
CheckoutConflictError,
WorkspaceTaskMergeError,
DUPLICATE_OF_METADATA_KEY,
WEDGE_RENOTIFY_COOLDOWN_MS,
} from "./types/task-core.js";
export {
assertNotWorkspaceTaskMerge,
@@ -551,6 +552,7 @@ export {
CheckoutConflictError,
WorkspaceTaskMergeError,
DUPLICATE_OF_METADATA_KEY,
WEDGE_RENOTIFY_COOLDOWN_MS,
};
import type {

View File

@@ -566,11 +566,26 @@ export interface ExecutorOverseerSignalMemory {
observedAt: number;
}
/*
FNXC:TaskWedgeNotifications 2026-08-01-15:35:
A task that resolves and re-wedges with the same actionable reason must notify once
per six-hour window. A durable shared constant keeps the PostgreSQL claim and the
engine's no-store fallback aligned across restarts.
*/
export const WEDGE_RENOTIFY_COOLDOWN_MS = 6 * 60 * 60 * 1_000;
export interface TaskWedgeNotificationState {
reasonKey: string;
episodeId: string;
status: "active" | "resolved";
transitionedAt: string;
/*
FNXC:TaskWedgeNotifications 2026-08-01-15:35:
A resolve/re-wedge flap used to mint an episode id that defeated mailbox dedupe.
Keep notification times independently per reason: one global last-reason pair
would let X -> Y -> X re-notify X while its own cooldown is still active.
*/
lastNotifiedAtByReason?: Record<string, string>;
}
export interface Task {

View File

@@ -1,5 +1,5 @@
import { describe, expect, it, vi } from "vitest";
import type { NotificationProvider, Settings, Task } from "@fusion/core";
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";
@@ -50,7 +50,122 @@ function fixture(workflowIr?: unknown) {
return { store, service, sendMessageOnce, sendNotification, task, getWedge: () => wedge };
}
/* Creates a restart-safe claim fake so NotificationService tests exercise delivery policy, not storage implementation. */
function cooldownFixture({ durable = true }: { durable?: boolean } = {}) {
const listeners = new Set<Listener>();
let wedge: Task["wedgeNotification"];
const stamps = new Map<string, number>();
const store = {
getSettings: async () => ({ ntfyEnabled: true, ntfyTopic: "test" }) as Settings,
on: (event: string, listener: Listener) => { if (event === "task:updated") listeners.add(listener); },
off: () => undefined,
emit: (task: Task) => listeners.forEach((listener) => listener(task)),
...(durable ? {
claimTaskWedgeNotificationEpisode: async (taskId: string, reasonKey: string | null) => {
if (reasonKey === null) {
if (wedge?.status === "active") wedge = { ...wedge, status: "resolved" };
return { claimed: false };
}
if (wedge?.status === "active" && wedge.reasonKey === reasonKey) return { claimed: false };
const now = Date.now();
for (const [key, notifiedAt] of stamps) {
if (now - notifiedAt >= WEDGE_RENOTIFY_COOLDOWN_MS) stamps.delete(key);
}
const episodeId = `${taskId}-${reasonKey}-${now}`;
wedge = { reasonKey, episodeId, status: "active", transitionedAt: new Date(now).toISOString() };
if (stamps.has(reasonKey)) return { claimed: false };
stamps.set(reasonKey, now);
return { claimed: true, episodeId };
},
} : {}),
};
const sendMessageOnce = vi.fn(async (_input: unknown, _key: string) => ({ message: {} as any, inserted: true }));
const service = new NotificationService(store as any, { messageStore: { on: () => undefined, sendMessageOnce } as any });
const sendNotification = vi.fn(async () => ({ success: true, providerId: "test" }));
service.registerProvider({ getProviderId: () => "test", isEventSupported: () => true, sendNotification });
const task = (overrides: Partial<Task> = {}): Task => ({ id: "FN-8691", title: "Blocked task", description: "", column: "in-review", status: "failed", error: "BLOCKED: dependency unavailable", dependencies: [], steps: [], currentStep: 0, log: [], createdAt: "2026-08-01T12:00:00.000Z", updatedAt: new Date().toISOString(), ...overrides } as Task);
return { store, service, sendMessageOnce, sendNotification, task };
}
async function flushWedgeHandling() {
await Promise.resolve();
await Promise.resolve();
await Promise.resolve();
}
describe("task wedge notifications", () => {
/*
FNXC:TaskWedgeNotifications 2026-08-01-15:35:
A BLOCKED task can resolve and re-wedge as scheduler/self-healing touch it. A
fresh episode id previously defeated sendMessageOnce; both delivery lanes now
notify once per reason cooldown and re-open only after that window expires.
*/
it("suppresses five resolve/re-wedge execution-blocked flaps until the cooldown expires", async () => {
vi.useFakeTimers();
vi.setSystemTime(new Date("2026-08-01T12:00:00.000Z"));
const { store, service, sendMessageOnce, sendNotification, task } = cooldownFixture();
await service.start();
for (let index = 0; index < 5; index += 1) {
store.emit(task({ updatedAt: new Date().toISOString() }));
await flushWedgeHandling();
store.emit(task({ status: "in-progress", error: undefined, column: "in-progress", updatedAt: new Date().toISOString() }));
await flushWedgeHandling();
}
expect(sendMessageOnce).toHaveBeenCalledTimes(1);
expect(sendNotification).toHaveBeenCalledTimes(1);
vi.advanceTimersByTime(WEDGE_RENOTIFY_COOLDOWN_MS);
store.emit(task({ updatedAt: new Date().toISOString() }));
await flushWedgeHandling();
expect(sendMessageOnce).toHaveBeenCalledTimes(2);
expect(sendNotification).toHaveBeenCalledTimes(2);
await service.stop();
vi.useRealTimers();
});
it("retains X cooldown across Y, preserves distinct gates, and covers self-healing entry", async () => {
vi.useFakeTimers();
vi.setSystemTime(new Date("2026-08-01T12:00:00.000Z"));
const { store, service, sendMessageOnce, task } = cooldownFixture();
await service.start();
const executionBlocked = describeTaskWedge(task())!;
const lint = describeTaskWedge(task({ error: "merge verification failed: check:lint" }))!;
const changeset = describeTaskWedge(task({ error: "merge verification failed: check:changeset-format" }))!;
await service.notifyTaskWedge(task(), executionBlocked);
await service.notifyTaskWedge(task(), lint);
await service.notifyTaskWedge(task(), executionBlocked);
await service.notifyTaskWedge(task(), changeset);
expect(sendMessageOnce).toHaveBeenCalledTimes(3);
const noAction = describeSelfHealingNoActionWedge(task({ status: "in-review", error: undefined }), "reconcile-in-review-unmet-dependencies", { taskActive: false });
await service.notifyTaskWedge(task({ status: "in-review", error: undefined }), noAction!);
store.emit(task({ status: "in-progress", error: undefined, column: "in-progress" }));
await flushWedgeHandling();
await service.notifyTaskWedge(task({ status: "in-review", error: undefined }), noAction!);
expect(sendMessageOnce).toHaveBeenCalledTimes(4);
await service.stop();
vi.useRealTimers();
});
it("gives the in-memory fallback the same resolve/re-wedge cooldown", async () => {
vi.useFakeTimers();
vi.setSystemTime(new Date("2026-08-01T12:00:00.000Z"));
const { store, service, sendMessageOnce, sendNotification, task } = cooldownFixture({ durable: false });
await service.start();
for (let index = 0; index < 5; index += 1) {
store.emit(task());
await flushWedgeHandling();
store.emit(task({ status: "in-progress", error: undefined, column: "in-progress" }));
await flushWedgeHandling();
}
expect(sendMessageOnce).toHaveBeenCalledTimes(1);
expect(sendNotification).toHaveBeenCalledTimes(1);
await service.stop();
vi.useRealTimers();
});
it("sends one actionable push and mailbox message per active terminal episode", async () => {
const { store, service, sendMessageOnce, sendNotification, task } = fixture();
await service.start();

View File

@@ -11,7 +11,7 @@ import type {
Task,
} from "@fusion/core";
import type { LifecycleColumns, TaskMoveLanes, WorkflowIrResolverStore } from "@fusion/core";
import { DASHBOARD_USER_ID, NotificationDispatcher, resolveProjectColumnsForRoles, resolveReviewColumns, resolveTaskLifecycleColumns, resolveWorkflowIrForTask } from "@fusion/core";
import { DASHBOARD_USER_ID, NotificationDispatcher, resolveProjectColumnsForRoles, resolveReviewColumns, resolveTaskLifecycleColumns, resolveWorkflowIrForTask, WEDGE_RENOTIFY_COOLDOWN_MS } from "@fusion/core";
import { DEFAULT_NTFY_EVENTS, buildNtfyClickUrl, formatTaskIdentifier } from "../notifier.js";
import { schedulerLog } from "../logger.js";
import { classifyTransientMergeError } from "../transient-merge-error-classifier.js";
@@ -127,6 +127,13 @@ export class NotificationService {
/** Compatibility fallback for lightweight test stores without the durable TaskStore CAS. */
private readonly activeWedgeReasons = new Map<string, string>();
/*
FNXC:TaskWedgeNotifications 2026-08-01-15:35:
Lightweight stores lack the durable compare-and-set, but must retain the same
per-task, per-reason cooldown. Clearing an active episode deliberately does not
clear these timestamps: X -> Y -> X must not re-notify X inside its window.
*/
private readonly fallbackWedgeNotificationTimestamps = new Map<string, Map<string, number>>();
/*
FNXC:TaskWedgeNotifications 2026-07-31-21:10:
PER-TASK SERIALISATION OF WEDGE HANDLING — the blocker two earlier fleet passes recorded and
declined to take on.
@@ -591,9 +598,8 @@ export class NotificationService {
if (!claim.claimed || !claim.episodeId) return;
episode = claim.episodeId;
} else {
if (this.activeWedgeReasons.get(task.id) === descriptor.reasonKey) return;
this.activeWedgeReasons.set(task.id, descriptor.reasonKey);
episode = `${task.id}:${descriptor.reasonKey}:${task.updatedAt}`;
episode = this.claimFallbackWedgeNotificationEpisode(task.id, descriptor.reasonKey, task.updatedAt);
if (!episode) return;
}
const link = buildNtfyClickUrl({ dashboardHost: this.dashboardHost, projectId: this.options.projectId, taskId: task.id });
const content = [
@@ -619,6 +625,22 @@ export class NotificationService {
}
}
private claimFallbackWedgeNotificationEpisode(taskId: string, reasonKey: string, updatedAt: string): string | undefined {
if (this.activeWedgeReasons.get(taskId) === reasonKey) return undefined;
const now = Date.now();
const timestamps = this.fallbackWedgeNotificationTimestamps.get(taskId) ?? new Map<string, number>();
for (const [key, notifiedAt] of timestamps) {
if (!Number.isFinite(notifiedAt) || notifiedAt > now || now - notifiedAt >= WEDGE_RENOTIFY_COOLDOWN_MS) timestamps.delete(key);
}
this.activeWedgeReasons.set(taskId, reasonKey);
this.fallbackWedgeNotificationTimestamps.set(taskId, timestamps);
if (timestamps.has(reasonKey)) return undefined;
timestamps.set(reasonKey, now);
return `${taskId}:${reasonKey}:${updatedAt}`;
}
private isTriageDuplicateDecision(task: Task): boolean {
return task.paused === true
&& task.pausedReason === "duplicate-decision-required"