diff --git a/.changeset/fn-8691-wedge-notification-cooldown.md b/.changeset/fn-8691-wedge-notification-cooldown.md new file mode 100644 index 0000000000..3fba3dabe0 --- /dev/null +++ b/.changeset/fn-8691-wedge-notification-cooldown.md @@ -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. diff --git a/docs/architecture.md b/docs/architecture.md index a5b1cd3adf..afdd057624 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -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 diff --git a/packages/core/src/__tests__/postgres/store-wedge-resolution.pg.test.ts b/packages/core/src/__tests__/postgres/store-wedge-resolution.pg.test.ts index f220330cf1..2362f4113f 100644 --- a/packages/core/src/__tests__/postgres/store-wedge-resolution.pg.test.ts +++ b/packages/core/src/__tests__/postgres/store-wedge-resolution.pg.test.ts @@ -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(); diff --git a/packages/core/src/index.ts b/packages/core/src/index.ts index 8e4b4ccd0f..490ebabb20 100644 --- a/packages/core/src/index.ts +++ b/packages/core/src/index.ts @@ -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, diff --git a/packages/core/src/store.ts b/packages/core/src/store.ts index 392f809609..99b7e350e0 100644 --- a/packages/core/src/store.ts +++ b/packages/core/src/store.ts @@ -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 { 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; } diff --git a/packages/core/src/task-store/task-mutation-ops.ts b/packages/core/src/task-store/task-mutation-ops.ts index 92bf308c37..6df810a5dd 100644 --- a/packages/core/src/task-store/task-mutation-ops.ts +++ b/packages/core/src/task-store/task-mutation-ops.ts @@ -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, diff --git a/packages/core/src/types.ts b/packages/core/src/types.ts index 51ab69a803..49863e4d9e 100644 --- a/packages/core/src/types.ts +++ b/packages/core/src/types.ts @@ -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 { diff --git a/packages/core/src/types/task-core.ts b/packages/core/src/types/task-core.ts index 38746a0622..6fe68f711c 100644 --- a/packages/core/src/types/task-core.ts +++ b/packages/core/src/types/task-core.ts @@ -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; } export interface Task { 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 4c3c82fa5e..5d5d0270c3 100644 --- a/packages/engine/src/notification/__tests__/task-wedge-notification.test.ts +++ b/packages/engine/src/notification/__tests__/task-wedge-notification.test.ts @@ -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(); + let wedge: Task["wedgeNotification"]; + const stamps = new Map(); + 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 => ({ 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(); diff --git a/packages/engine/src/notification/notification-service.ts b/packages/engine/src/notification/notification-service.ts index e47d258e1b..8ea7b36778 100644 --- a/packages/engine/src/notification/notification-service.ts +++ b/packages/engine/src/notification/notification-service.ts @@ -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(); /* + 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>(); + /* 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(); + 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"