diff --git a/.changeset/fn-8505-task-wedge-notifications.md b/.changeset/fn-8505-task-wedge-notifications.md new file mode 100644 index 0000000000..4d09fb228f --- /dev/null +++ b/.changeset/fn-8505-task-wedge-notifications.md @@ -0,0 +1,7 @@ +--- +"@runfusion/fusion": patch +--- + +summary: Notify operators when a task is terminally blocked or exhausts automated recovery. +category: fix +dev: Adds deduped task-wedge provider and dashboard mailbox delivery. diff --git a/docs/agents.md b/docs/agents.md index c59d7e6c3e..573c2d836c 100644 --- a/docs/agents.md +++ b/docs/agents.md @@ -248,6 +248,10 @@ Separation of concerns: - `permissionPolicy` determines how sensitive runtime actions are gated (`allow`, `block`, `require-approval`) once the capability path is in play. `require-approval` creates an approval request with the permanent or ephemeral actor identity, pauses the associated task safely, and resumes through the existing approval lifecycle. - Dashboard persona presets (`packages/dashboard/app/components/agent-presets/`) are UI templates for identity/behavior and are **not** the source of truth for permission-policy enforcement. +### Task wedge operator notifications + +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. 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. + ### 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. diff --git a/docs/architecture.md b/docs/architecture.md index a5346aae26..4aa2b4a907 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -6,6 +6,10 @@ 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. + ## 1) Overview Fusion is an AI-orchestrated task board. It takes tasks through a structured lifecycle (`planning → todo → in-progress → in-review → done → archived`) and automates planning, execution, review, merge, and operational recovery. diff --git a/packages/core/src/postgres/migrations/0000_initial.sql b/packages/core/src/postgres/migrations/0000_initial.sql index 085033b4bb..db13b4dd4b 100644 --- a/packages/core/src/postgres/migrations/0000_initial.sql +++ b/packages/core/src/postgres/migrations/0000_initial.sql @@ -55,6 +55,8 @@ CREATE TABLE IF NOT EXISTS project.tasks ( paused integer DEFAULT 0, user_paused integer DEFAULT 0, paused_reason text, + -- FNXC:TaskWedgeNotifications 2026-10-19-00:00: Fresh PostgreSQL baselines must include the durable terminal-wedge episode field that upgrade migration 0033 adds to existing databases. + wedge_notification text, base_branch text, branch text, auto_merge integer, diff --git a/packages/core/src/postgres/migrations/0033_fn-8505_wedge_notification.sql b/packages/core/src/postgres/migrations/0033_fn-8505_wedge_notification.sql new file mode 100644 index 0000000000..9520c9b3a4 --- /dev/null +++ b/packages/core/src/postgres/migrations/0033_fn-8505_wedge_notification.sql @@ -0,0 +1,5 @@ +-- FNXC:TaskWedgeNotifications 2026-07-22-14:00: +-- Persist the active/resolved terminal-wedge episode on the task row. The opaque +-- episode id is used by push and mailbox idempotency; human-readable error output +-- never enters this durable dedupe state. +ALTER TABLE project.tasks ADD COLUMN IF NOT EXISTS wedge_notification text; diff --git a/packages/core/src/postgres/schema-applier.ts b/packages/core/src/postgres/schema-applier.ts index b8257d6c10..51dee47efb 100644 --- a/packages/core/src/postgres/schema-applier.ts +++ b/packages/core/src/postgres/schema-applier.ts @@ -44,8 +44,13 @@ continuations at workflow column boundaries. FNXC:LegacyAdoption 2026-07-21-17:30: SCHEMA_BASELINE_VERSION advances to 0032 for fusion_runtime SELECT + SECURITY DEFINER write access to the legacy-adoption drained marker. + +FNXC:TaskWedgeNotifications 2026-10-19-00:00: +Advance the PostgreSQL schema ceiling for the durable wedge episode column. The +forward migration must run before TaskStore writes the new field on fresh and +upgraded databases. */ -export const SCHEMA_BASELINE_VERSION = "0032"; +export const SCHEMA_BASELINE_VERSION = "0033"; /** FNXC:SymbolLock 2026-07-31-10:00: upgrades need durable task declarations before admission resolves symbols. */ export const TASK_DECLARED_SYMBOLS_VERSION = "0028"; const INITIAL_SCHEMA_VERSION = "0000"; @@ -139,6 +144,8 @@ export const SQLITE_MIGRATION_RUNTIME_READ_VERSION = "0030"; export const WORKFLOW_TASK_CONTINUATIONS_VERSION = "0031"; /** FNXC:LegacyAdoption 2026-07-21-17:30: runtime role needs drained-marker read + restricted write. */ export const LEGACY_ADOPTION_DRAINED_MARKER_RUNTIME_GRANTS_VERSION = "0032"; +/** FNXC:TaskWedgeNotifications 2026-10-19-00:00: manually register the durable wedge episode migration for PostgreSQL upgrades. */ +export const TASK_WEDGE_NOTIFICATION_VERSION = "0033"; /** SECURITY DEFINER helper that only inserts LEGACY_ADOPTION_DRAINED_MARKER. */ export const LEGACY_ADOPTION_DRAINED_MARKER_FUNCTION = "fusion_mark_legacy_adoption_drained"; @@ -339,6 +346,10 @@ const LEGACY_ADOPTION_DRAINED_MARKER_RUNTIME_GRANTS_PATH = join( MIGRATIONS_DIR, "0032_legacy_adoption_drained_marker_runtime_grants.sql", ); +const TASK_WEDGE_NOTIFICATION_MIGRATION_PATH = join( + MIGRATIONS_DIR, + "0033_fn-8505_wedge_notification.sql", +); /** * Ensure the migration bookkeeping table exists. Lives in the public schema so @@ -441,6 +452,7 @@ export async function applySchemaBaseline( const legacyAdoptionDrainedMarkerRuntimeGrantsAlreadyApplied = applied.includes( LEGACY_ADOPTION_DRAINED_MARKER_RUNTIME_GRANTS_VERSION, ); + const taskWedgeNotificationAlreadyApplied = applied.includes(TASK_WEDGE_NOTIFICATION_VERSION); assertBinaryNotOlderThanDatabase(applied); let schemaChanged = false; @@ -912,6 +924,21 @@ export async function applySchemaBaseline( schemaChanged = true; } + /* + FNXC:TaskWedgeNotifications 2026-10-19-00:00: + PostgreSQL migrations are explicitly registered rather than discovered. + Apply the wedge episode column before persistence writes it, including on + databases that already recorded the prior schema ceiling. + */ + if (!taskWedgeNotificationAlreadyApplied) { + const migrationSql = await readFile(TASK_WEDGE_NOTIFICATION_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(migrationSql)); + await tx.execute( + sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${TASK_WEDGE_NOTIFICATION_VERSION}) ON CONFLICT (version) DO NOTHING`, + ); + schemaChanged = true; + } + return { applied: schemaChanged, pluginHooksRun: pluginHooks.length }; }); } diff --git a/packages/core/src/postgres/schema/project.ts b/packages/core/src/postgres/schema/project.ts index 854703759f..0e8203a8bc 100644 --- a/packages/core/src/postgres/schema/project.ts +++ b/packages/core/src/postgres/schema/project.ts @@ -83,6 +83,7 @@ export const tasks = projectSchema.table("tasks", { paused: integer("paused").default(0), userPaused: integer("user_paused").default(0), pausedReason: text("paused_reason"), + wedgeNotification: text("wedge_notification"), baseBranch: text("base_branch"), branch: text("branch"), autoMerge: integer("auto_merge"), diff --git a/packages/core/src/store.ts b/packages/core/src/store.ts index d7a37cb363..9a6f3ff210 100644 --- a/packages/core/src/store.ts +++ b/packages/core/src/store.ts @@ -1240,10 +1240,30 @@ export class TaskStore extends EventEmitter { } async updateTask( id: string, - updates: { title?: string; description?: string; priority?: TaskPriority | null; prompt?: string; worktree?: string | null; workspaceWorktrees?: import("./types.js").Task["workspaceWorktrees"]; status?: string | null; dependencies?: string[]; steps?: import("./types.js").TaskStep[]; customFields?: Record; currentStep?: number; blockedBy?: string | null; overlapBlockedBy?: string | null; assignedAgentId?: string | null; pausedByAgentId?: string | null; pausedReason?: string | null; tokenBudgetSoftAlertedAt?: string | null; worktrunkFallbackAlertedAt?: string | null; worktrunkFailure?: import("./types.js").Task["worktrunkFailure"] | null; tokenBudgetHardAlertedAt?: string | null; tokenBudgetOverride?: import("./types.js").TaskTokenBudgetOverride | null; dispatchStormCount?: number | null; lastDispatchAt?: string | null; assigneeUserId?: string | null; scopeOverride?: boolean | null; scopeOverrideReason?: string | null; scopeAutoWiden?: string[] | null; nodeId?: string | null; effectiveNodeId?: string | null; effectiveNodeSource?: string | null; checkedOutBy?: string | null; checkedOutAt?: string | null; checkoutNodeId?: string | null; checkoutRunId?: string | null; checkoutLeaseRenewedAt?: string | null; checkoutLeaseEpoch?: number | null; paused?: boolean; baseBranch?: string | null; autoMerge?: boolean | null; branch?: string | null; executionStartBranch?: string | null; baseCommitSha?: string | null; size?: "S" | "M" | "L"; reviewLevel?: number; executionMode?: import("./types.js").ExecutionMode | null; mergeRetries?: number; workflowStepRetries?: number; stuckKillCount?: number | null; resumeLimboCount?: number | null; executeRequeueLoopCount?: number | null; graphResumeRetryCount?: number | null; consecutiveToolFailureRetryCount?: number | null; executorEscalationAttempted?: boolean | null; toolFailureDetectorLogCursor?: number | null; toolFailureRetryExhaustedAuditEmitted?: boolean | null; resumeLimboTipSha?: string | null; resumeLimboStepSignature?: string | null; executeRequeueLoopSignature?: string | null; postReviewFixCount?: number | null; planReviewReplanCount?: number | null; recoveryRetryCount?: number | null; taskDoneRetryCount?: number | null; bulkCompletionRefusalAt?: string | null; workflowIrPin?: string | null; workflowIrPinNodeId?: string | null; workflowIrPinColumnId?: string | null; legacyAdoptedAt?: string | null; worktreeSessionRetryCount?: number | null; completionHandoffLimboRecoveryCount?: number | null; verificationFailureCount?: number | null; mergeConflictBounceCount?: number | null; mergeAuditBounceCount?: number | null; mergeTransientRetryCount?: number | null; branchConflictRecoveryCount?: number | null; reviewerContextRetryCount?: number | null; reviewerFallbackRetryCount?: number | null; nextRecoveryAt?: string | null; enabledWorkflowSteps?: string[]; noCommitsExpected?: boolean | null; modelProvider?: string | null; modelId?: string | null; validatorModelProvider?: string | null; validatorModelId?: string | null; planningModelProvider?: string | null; planningModelId?: string | null; mergerModelProvider?: string | null; mergerModelId?: string | null; thinkingLevel?: string | null; validatorThinkingLevel?: string | null; planningThinkingLevel?: string | null; mergerThinkingLevel?: string | null; error?: string | null; summary?: string | null; sessionFile?: string | null; firstExecutionAt?: string | null; cumulativeActiveMs?: number | null; cumulativePlanningMs?: number | null; planningStartedAt?: string | null; executionStartedAt?: string | null; executionCompletedAt?: string | null; review?: import("./types.js").TaskReview | null; reviewState?: import("./types.js").TaskReviewState | null; workflowStepResults?: import("./types.js").WorkflowStepResult[] | null; mergeDetails?: import("./types.js").MergeDetails | null; sourceIssue?: import("./types.js").TaskSourceIssue | null; sourceMetadataPatch?: Record | null; githubTracking?: import("./types.js").TaskGithubTracking | null; tokenUsage?: import("./types.js").TaskTokenUsage | null; modifiedFiles?: string[] | null; declaredSymbols?: string[] | null | undefined; missionId?: string | null; sliceId?: string | null; workflowTransitionNotification?: import("./types.js").WorkflowTransitionNotificationMarker | undefined; plannerOversightLevel?: string | null; sessionAdvisorEnabled?: boolean | null; approvedPlanFingerprint?: string | null }, runContext?: RunMutationContext, + updates: { title?: string; description?: string; priority?: TaskPriority | null; prompt?: string; worktree?: string | null; workspaceWorktrees?: import("./types.js").Task["workspaceWorktrees"]; status?: string | null; dependencies?: string[]; steps?: import("./types.js").TaskStep[]; customFields?: Record; currentStep?: number; blockedBy?: string | null; overlapBlockedBy?: string | null; assignedAgentId?: string | null; pausedByAgentId?: string | null; pausedReason?: string | null; wedgeNotification?: import("./types.js").TaskWedgeNotificationState | null; tokenBudgetSoftAlertedAt?: string | null; worktrunkFallbackAlertedAt?: string | null; worktrunkFailure?: import("./types.js").Task["worktrunkFailure"] | null; tokenBudgetHardAlertedAt?: string | null; tokenBudgetOverride?: import("./types.js").TaskTokenBudgetOverride | null; dispatchStormCount?: number | null; lastDispatchAt?: string | null; assigneeUserId?: string | null; scopeOverride?: boolean | null; scopeOverrideReason?: string | null; scopeAutoWiden?: string[] | null; nodeId?: string | null; effectiveNodeId?: string | null; effectiveNodeSource?: string | null; checkedOutBy?: string | null; checkedOutAt?: string | null; checkoutNodeId?: string | null; checkoutRunId?: string | null; checkoutLeaseRenewedAt?: string | null; checkoutLeaseEpoch?: number | null; paused?: boolean; baseBranch?: string | null; autoMerge?: boolean | null; branch?: string | null; executionStartBranch?: string | null; baseCommitSha?: string | null; size?: "S" | "M" | "L"; reviewLevel?: number; executionMode?: import("./types.js").ExecutionMode | null; mergeRetries?: number; workflowStepRetries?: number; stuckKillCount?: number | null; resumeLimboCount?: number | null; executeRequeueLoopCount?: number | null; graphResumeRetryCount?: number | null; consecutiveToolFailureRetryCount?: number | null; executorEscalationAttempted?: boolean | null; toolFailureDetectorLogCursor?: number | null; toolFailureRetryExhaustedAuditEmitted?: boolean | null; resumeLimboTipSha?: string | null; resumeLimboStepSignature?: string | null; executeRequeueLoopSignature?: string | null; postReviewFixCount?: number | null; planReviewReplanCount?: number | null; recoveryRetryCount?: number | null; taskDoneRetryCount?: number | null; bulkCompletionRefusalAt?: string | null; workflowIrPin?: string | null; workflowIrPinNodeId?: string | null; workflowIrPinColumnId?: string | null; legacyAdoptedAt?: string | null; worktreeSessionRetryCount?: number | null; completionHandoffLimboRecoveryCount?: number | null; verificationFailureCount?: number | null; mergeConflictBounceCount?: number | null; mergeAuditBounceCount?: number | null; mergeTransientRetryCount?: number | null; branchConflictRecoveryCount?: number | null; reviewerContextRetryCount?: number | null; reviewerFallbackRetryCount?: number | null; nextRecoveryAt?: string | null; enabledWorkflowSteps?: string[]; noCommitsExpected?: boolean | null; modelProvider?: string | null; modelId?: string | null; validatorModelProvider?: string | null; validatorModelId?: string | null; planningModelProvider?: string | null; planningModelId?: string | null; mergerModelProvider?: string | null; mergerModelId?: string | null; thinkingLevel?: string | null; validatorThinkingLevel?: string | null; planningThinkingLevel?: string | null; mergerThinkingLevel?: string | null; error?: string | null; summary?: string | null; sessionFile?: string | null; firstExecutionAt?: string | null; cumulativeActiveMs?: number | null; cumulativePlanningMs?: number | null; planningStartedAt?: string | null; executionStartedAt?: string | null; executionCompletedAt?: string | null; review?: import("./types.js").TaskReview | null; reviewState?: import("./types.js").TaskReviewState | null; workflowStepResults?: import("./types.js").WorkflowStepResult[] | null; mergeDetails?: import("./types.js").MergeDetails | null; sourceIssue?: import("./types.js").TaskSourceIssue | null; sourceMetadataPatch?: Record | null; githubTracking?: import("./types.js").TaskGithubTracking | null; tokenUsage?: import("./types.js").TaskTokenUsage | null; modifiedFiles?: string[] | null; declaredSymbols?: string[] | null | undefined; missionId?: string | null; sliceId?: string | null; workflowTransitionNotification?: import("./types.js").WorkflowTransitionNotificationMarker | undefined; plannerOversightLevel?: string | null; sessionAdvisorEnabled?: boolean | null; approvedPlanFingerprint?: string | null }, runContext?: RunMutationContext, ): Promise { return updateTaskImpl(this, id, updates, runContext); } + /** + * Atomically opens or resolves the durable task-wedge episode. The task lock is + * shared by SQLite and PostgreSQL persistence writes, so concurrent engine + * observers cannot allocate two episodes for one active normalized reason. + */ + async claimTaskWedgeNotificationEpisode(taskId: string, reasonKey: string | null): Promise<{ episodeId?: string; claimed: boolean }> { + let result: { episodeId?: string; claimed: boolean } = { claimed: false }; + await this.updateTaskAtomic(taskId, (current) => { + const prior = current.wedgeNotification; + if (reasonKey === null) { + if (!prior || prior.status === "resolved") return null; + return { wedgeNotification: { ...prior, status: "resolved", transitionedAt: new Date().toISOString() } }; + } + if (prior?.status === "active" && prior.reasonKey === reasonKey) return null; + const episodeId = randomUUID(); + result = { episodeId, claimed: true }; + return { wedgeNotification: { reasonKey, episodeId, status: "active", transitionedAt: new Date().toISOString() } }; + }); + return result; + } async claimNextToolFailureRetry(taskId: string, expectedCursor: number, maxRetries: number): Promise { return claimNextToolFailureRetryImpl(this, taskId, expectedCursor, maxRetries); } diff --git a/packages/core/src/task-store/persistence.ts b/packages/core/src/task-store/persistence.ts index 02310a7220..de7e289ae3 100644 --- a/packages/core/src/task-store/persistence.ts +++ b/packages/core/src/task-store/persistence.ts @@ -28,6 +28,7 @@ export interface TaskRow { overlapBlockedBy: string | null; paused: number | null; pausedReason: string | null; + wedgeNotification: string | null; userPaused: number | null; baseBranch: string | null; executionStartBranch: string | null; @@ -228,6 +229,7 @@ export const TASK_COLUMN_DESCRIPTORS: TaskColumnDescriptor[] = [ defineTaskColumn("overlapBlockedBy", (task) => task.overlapBlockedBy ?? null), defineTaskColumn("paused", (task) => task.paused ? 1 : 0), defineTaskColumn("pausedReason", (task) => task.pausedReason ?? null), + defineTaskColumn("wedgeNotification", (task) => toJsonNullable(task.wedgeNotification)), defineTaskColumn("userPaused", (task) => task.userPaused ? 1 : 0), defineTaskColumn("baseBranch", (task) => task.baseBranch ?? null), defineTaskColumn("branch", (task) => task.branch ?? null), diff --git a/packages/core/src/task-store/serialization.ts b/packages/core/src/task-store/serialization.ts index 1eab34f70d..12fb7ab228 100644 --- a/packages/core/src/task-store/serialization.ts +++ b/packages/core/src/task-store/serialization.ts @@ -78,6 +78,7 @@ export function rowToTask(row: TaskRow): Task { overlapBlockedBy: row.overlapBlockedBy || undefined, paused: row.paused ? true : undefined, pausedReason: row.pausedReason || undefined, + wedgeNotification: fromJson(row.wedgeNotification) ?? undefined, userPaused: row.userPaused ? true : undefined, baseBranch: row.baseBranch || undefined, executionStartBranch: row.executionStartBranch || undefined, diff --git a/packages/core/src/task-store/task-row-mappers.ts b/packages/core/src/task-store/task-row-mappers.ts index cede070382..ffd462f35d 100644 --- a/packages/core/src/task-store/task-row-mappers.ts +++ b/packages/core/src/task-store/task-row-mappers.ts @@ -33,7 +33,7 @@ export function getTaskSelectClauseImpl2(store: TaskStore, slim: boolean, tableA const prefix = tableAlias ? `${tableAlias}.` : ""; return [ "id", "lineageId", "title", "description", "priority", "\"column\"", "status", "size", "reviewLevel", "currentStep", - "worktree", "blockedBy", "overlapBlockedBy", "paused", "pausedReason", "userPaused", "baseBranch", "branch", "autoMerge", "autoMergeProvenance", "executionStartBranch", "baseCommitSha", + "worktree", "blockedBy", "overlapBlockedBy", "paused", "pausedReason", "wedgeNotification", "userPaused", "baseBranch", "branch", "autoMerge", "autoMergeProvenance", "executionStartBranch", "baseCommitSha", "modelPresetId", "modelProvider", "modelId", "validatorModelProvider", "validatorModelId", "planningModelProvider", "planningModelId", "mergerModelProvider", "mergerModelId", diff --git a/packages/core/src/task-store/task-update.ts b/packages/core/src/task-store/task-update.ts index ca266fc0db..83a733bc2f 100644 --- a/packages/core/src/task-store/task-update.ts +++ b/packages/core/src/task-store/task-update.ts @@ -206,6 +206,11 @@ export async function updateTaskUnlockedImpl(store: TaskStore, id: string, updat } else if (updates.pausedReason !== undefined) { task.pausedReason = updates.pausedReason; } + if (updates.wedgeNotification === null) { + task.wedgeNotification = undefined; + } else if (updates.wedgeNotification !== undefined) { + task.wedgeNotification = updates.wedgeNotification; + } if (updates.tokenBudgetSoftAlertedAt === null) { task.tokenBudgetSoftAlertedAt = undefined; } else if (updates.tokenBudgetSoftAlertedAt !== undefined) { diff --git a/packages/core/src/types.ts b/packages/core/src/types.ts index a1db3b827d..77afb38bee 100644 --- a/packages/core/src/types.ts +++ b/packages/core/src/types.ts @@ -1146,6 +1146,13 @@ export interface ExecutorOverseerSignalMemory { observedAt: number; } +export interface TaskWedgeNotificationState { + reasonKey: string; + episodeId: string; + status: "active" | "resolved"; + transitionedAt: string; +} + export interface Task { id: string; /** Immutable lineage identity used for durable commit/task attribution. */ @@ -1224,6 +1231,13 @@ export interface Task { userPaused?: boolean; /** Optional machine-readable reason for automated pauses (for example dispatch-storm). */ pausedReason?: string; + /* + FNXC:TaskWedgeNotifications 2026-07-22-14:00: + A terminal wedge is an episode, not a permanently suppressed error. Persist its + normalized reason and opaque episode id so delivery remains quiet across restart + until an authoritative lifecycle update resolves the episode. + */ + wedgeNotification?: TaskWedgeNotificationState; /** ISO timestamp set when the task first crossed the soft token budget cap. */ tokenBudgetSoftAlertedAt?: string; /** ISO timestamp marking first one-shot alert when worktrunk failed and fell back to native backend. */ diff --git a/packages/core/src/types/workflow-steps.ts b/packages/core/src/types/workflow-steps.ts index 936137d1bf..72930d10cc 100644 --- a/packages/core/src/types/workflow-steps.ts +++ b/packages/core/src/types/workflow-steps.ts @@ -93,6 +93,7 @@ export type NtfyNotificationEvent = | "in-review" | "merged" | "failed" + | "task-wedged" | "awaiting-approval" | "awaiting-user-review" | "planning-awaiting-input" @@ -115,6 +116,7 @@ export const NOTIFICATION_EVENTS = [ "in-review", "merged", "failed", + "task-wedged", "awaiting-approval", "awaiting-user-review", "planning-awaiting-input", diff --git a/packages/engine/src/__tests__/notification-service.test.ts b/packages/engine/src/__tests__/notification-service.test.ts index b37d2d2595..94291a566c 100644 --- a/packages/engine/src/__tests__/notification-service.test.ts +++ b/packages/engine/src/__tests__/notification-service.test.ts @@ -1240,6 +1240,26 @@ describe("NotificationService", () => { }); }); + it("does not claim a wedge episode for a transient failed merge", async () => { + const store = createStore({ ntfyEnabled: true, ntfyTopic: "topic" }); + const claimTaskWedgeNotificationEpisode = vi.fn(); + (store as any).claimTaskWedgeNotificationEpisode = claimTaskWedgeNotificationEpisode; + const sendNotification = vi.fn(async () => ({ success: true, providerId: "mock" })); + const service = new NotificationService(store as any); + service.registerProvider({ getProviderId: () => "mock", isEventSupported: () => true, sendNotification }); + await service.start(); + + const failed = task({ id: "FN-5627", status: "failed", column: "in-review", error: "spawn git ENOENT" }); + store.setTask(failed); + store.emit("task:updated", failed); + await Promise.resolve(); + await Promise.resolve(); + + expect(claimTaskWedgeNotificationEpisode).not.toHaveBeenCalled(); + expect(sendNotification).not.toHaveBeenCalledWith("task-wedged", expect.anything()); + await service.stop(); + }); + it("suppresses transient failed notification after Auto-recovered status clear", async () => { vi.useFakeTimers(); const store = createStore({ ntfyEnabled: true, ntfyTopic: "topic", failureNotificationMode: "sticky-only", failureNotificationDelayMs: 50 }); diff --git a/packages/engine/src/__tests__/notifier.test.ts b/packages/engine/src/__tests__/notifier.test.ts index 9a36143e51..888f2cf87d 100644 --- a/packages/engine/src/__tests__/notifier.test.ts +++ b/packages/engine/src/__tests__/notifier.test.ts @@ -27,6 +27,22 @@ describe("Ntfy notifier helpers", () => { expect(DEFAULT_NTFY_EVENTS).toContain("message:agent-to-agent"); expect(DEFAULT_NTFY_EVENTS).toContain("message:room"); expect(DEFAULT_NTFY_EVENTS).toContain("workflow-notify"); + expect(DEFAULT_NTFY_EVENTS).toContain("task-wedged"); + }); + + it("allows task-wedged notifications through the real Ntfy default filter", async () => { + const provider = new NtfyNotificationProvider(); + await provider.initialize({ topic: "operator-alerts" }); + expect(provider.isEventSupported("task-wedged")).toBe(true); + expect(provider.isEventSupported("task-wedged")).toBe(true); + await provider.shutdown(); + }); + + it("respects an explicit Ntfy event filter for task-wedged notifications", async () => { + const provider = new NtfyNotificationProvider(); + await provider.initialize({ topic: "operator-alerts", events: ["failed"] }); + expect(provider.isEventSupported("task-wedged")).toBe(false); + await provider.shutdown(); }); it("checks awaiting-input event enablement", () => { @@ -714,11 +730,15 @@ describe("NtfyNotifier", () => { expect(fetchMock).not.toHaveBeenCalled(); }); - it("sends high priority notification when task fails", async () => { + it("sends a high-priority wedge notification when a task terminally fails", async () => { notifier = new NtfyNotifier(store); await notifier.start(); const failedTask = createTask("FN-001", "Test Task", "failed"); + // FNXC:TaskWedgeNotifications 2026-07-22-20:00: A wedge must originate + // from a census-classified terminal writer, not a bare failed status that + // could still be in transient merge recovery. + failedTask.error = "merge verification failed: check:changeset-format"; store.triggerTaskUpdated(failedTask); await flushAsyncWork(); @@ -729,10 +749,10 @@ describe("NtfyNotifier", () => { expect.objectContaining({ method: "POST", headers: expect.objectContaining({ - "Title": "Task FN-001 failed", + "Title": "Task FN-001 needs operator action", "Priority": "high", }), - body: 'Task "Test Task" has failed and needs attention', + body: expect.stringContaining('Task "Test Task" is wedged:'), }) ); }); diff --git a/packages/engine/src/__tests__/self-healing.test.ts b/packages/engine/src/__tests__/self-healing.test.ts index 856f77ccea..88b2a38ddd 100644 --- a/packages/engine/src/__tests__/self-healing.test.ts +++ b/packages/engine/src/__tests__/self-healing.test.ts @@ -6320,7 +6320,10 @@ describe("SelfHealingManager", () => { await managerWithRecovery.recoverAlreadyMergedReviewTasks(); await vi.advanceTimersByTimeAsync(60); - expect(sendNotification).toHaveBeenCalledWith("failed", expect.objectContaining({ taskId: "FN-1" })); + expect(sendNotification).toHaveBeenCalledWith("task-wedged", expect.objectContaining({ + taskId: "FN-1", + metadata: expect.objectContaining({ wedgeReason: "terminal-failed" }), + })); await notificationService.stop(); managerWithRecovery.stop(); @@ -11252,6 +11255,10 @@ describe("FN-5335 triple-proof no-action unit coverage", () => { it("emits reclaim-pr-conflict no-action when triple proof fails", async () => { const store = createMockStore({ getSettings: vi.fn().mockResolvedValue({ autoMerge: true, globalPause: false, enginePaused: false, taskStuckTimeoutMs: 1_000 } as any), + // FNXC:TaskWedgeNotifications 2026-07-22-19:15: The recovery path inventories + // active worktrees before triple-proof; keep this fixture deterministic so it + // exercises the intended ownerless no-action seam. + listTasks: vi.fn().mockResolvedValue([]), getTask: vi.fn().mockResolvedValue({ id: "FN-PR", column: "in-review", paused: false, status: null, worktree: "/tmp/wt-pr", branch: "fusion/fn-pr", prInfo: { number: 1, mergeable: "conflicting" }, updatedAt: new Date(Date.now() - 10_000).toISOString() }), }); const manager = new SelfHealingManager(store, { rootDir: "/tmp/test-project" }); diff --git a/packages/engine/src/notification/__tests__/notification-service.test.ts b/packages/engine/src/notification/__tests__/notification-service.test.ts index 391b3cf2fe..550f14b806 100644 --- a/packages/engine/src/notification/__tests__/notification-service.test.ts +++ b/packages/engine/src/notification/__tests__/notification-service.test.ts @@ -98,6 +98,35 @@ describe("NotificationService deferred failure notifications", () => { await service.stop(); }); + it("delivers a new wedge episode when an opaque terminal failure gains a specific cause", async () => { + const { store, service, sendNotification } = await setup(); + const genericFailure = task({ id: "FN-wedge", status: "failed", error: "unexpected failure" }); + store.setTask(genericFailure); + store.emit("task:updated", genericFailure); + + const wedge = task({ + id: "FN-wedge", + status: "failed", + column: "in-review", + error: "merge verification failed: check:changeset-format", + }); + store.setTask(wedge); + store.emit("task:updated", wedge); + await vi.advanceTimersByTimeAsync(100); + + expect(sendNotification).toHaveBeenCalledTimes(2); + expect(sendNotification).toHaveBeenCalledWith("task-wedged", expect.objectContaining({ + taskId: "FN-wedge", + metadata: expect.objectContaining({ wedgeReason: "terminal-failed" }), + })); + expect(sendNotification).toHaveBeenCalledWith("task-wedged", expect.objectContaining({ + taskId: "FN-wedge", + metadata: expect.objectContaining({ wedgeReason: "merge-blocked:changeset-format" }), + })); + expect(service.getPendingFailureCount()).toBe(0); + await service.stop(); + }); + it("FN-5627: suppresses notification for transient lease-handoff-target-not-queued failures", async () => { const { store, service, sendNotification } = await setup(); store.setTask(task({ @@ -140,7 +169,10 @@ describe("NotificationService deferred failure notifications", () => { await vi.advanceTimersByTimeAsync(500); expect(sendNotification).toHaveBeenCalledTimes(1); - expect(sendNotification).toHaveBeenCalledWith("failed", expect.objectContaining({ taskId: "FN-genuine" })); + expect(sendNotification).toHaveBeenCalledWith("task-wedged", expect.objectContaining({ + taskId: "FN-genuine", + metadata: expect.objectContaining({ wedgeReason: "terminal-failed" }), + })); await service.stop(); }); diff --git a/packages/engine/src/notification/__tests__/task-wedge-notification.test.ts b/packages/engine/src/notification/__tests__/task-wedge-notification.test.ts new file mode 100644 index 0000000000..8d3c0ef6bf --- /dev/null +++ b/packages/engine/src/notification/__tests__/task-wedge-notification.test.ts @@ -0,0 +1,134 @@ +import { describe, expect, it, vi } from "vitest"; +import type { NotificationProvider, Settings, Task } from "@fusion/core"; +import { NotificationService } from "../notification-service.js"; +import { describeSelfHealingNoActionWedge, describeTaskWedge } from "../task-wedge-notification.js"; + +type Listener = (task: Task) => void; +function fixture() { + const listeners = new Set(); + let wedge: Task["wedgeNotification"]; + 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)), + 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 }; + wedge = { reasonKey, episodeId: `${taskId}-${reasonKey}-${Date.now()}`, status: "active", transitionedAt: new Date().toISOString() }; + return { claimed: true, episodeId: wedge.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, failedNotificationGraceMs: 60_000 }); + const sendNotification = vi.fn(async () => ({ success: true, providerId: "test" })); + const provider: NotificationProvider = { getProviderId: () => "test", isEventSupported: () => true, sendNotification }; + service.registerProvider(provider); + const task = (overrides: Partial = {}): Task => ({ id: "FN-8501", title: "Fix changeset", description: "", column: "in-review", status: "failed", error: "merge verification failed: check:changeset-format", dependencies: [], steps: [], currentStep: 0, log: [], createdAt: "2026-07-22T12:00:00.000Z", updatedAt: "2026-07-22T12:00:00.000Z", ...overrides } as Task); + return { store, service, sendMessageOnce, sendNotification, task, getWedge: () => wedge }; +} + +describe("task wedge notifications", () => { + it("sends one actionable push and mailbox message per active terminal episode", async () => { + const { store, service, sendMessageOnce, sendNotification, task } = fixture(); + await service.start(); + store.emit(task()); + store.emit(task({ updatedAt: "2026-07-22T12:01:00.000Z" })); + await vi.waitFor(() => expect(sendMessageOnce).toHaveBeenCalledTimes(1)); + expect(sendNotification).toHaveBeenCalledWith("task-wedged", expect.objectContaining({ taskId: "FN-8501", metadata: expect.objectContaining({ gate: "check:changeset-format" }) })); + const firstMessage = sendMessageOnce.mock.calls[0]?.[0] as { content: string } | undefined; + expect(firstMessage?.content).toContain("Fix changeset"); + expect(firstMessage?.content).toContain("Recommended action"); + store.emit(task({ status: "queued", error: undefined, column: "todo", updatedAt: "2026-07-22T12:02:00.000Z" })); + store.emit(task({ updatedAt: "2026-07-22T12:03:00.000Z" })); + await vi.waitFor(() => expect(sendMessageOnce).toHaveBeenCalledTimes(2)); + await service.stop(); + }); + + it("does not re-deliver an unchanged durable episode after service restart", async () => { + const { store, service, sendNotification, sendMessageOnce, task } = fixture(); + await service.start(); + store.emit(task()); + await vi.waitFor(() => expect(sendMessageOnce).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({ updatedAt: "2026-07-22T12:01:00.000Z" })); + await new Promise((resolve) => setTimeout(resolve, 10)); + expect(sendMessageOnce).toHaveBeenCalledTimes(1); + expect(sendNotification).toHaveBeenCalledTimes(1); + await restarted.stop(); + }); + + it("delivers one durable episode for an ownerless self-healing no-action escalation", async () => { + const { service, sendMessageOnce, sendNotification, task } = fixture(); + await service.start(); + const ownerless = task({ status: "in-review", error: undefined, paused: false, userPaused: false }); + const descriptor = describeSelfHealingNoActionWedge(ownerless, "reconcile-in-review-unmet-dependencies", { + taskActive: false, + hasExecutingTaskLock: false, + livePaths: [], + }); + expect(descriptor).toMatchObject({ reasonKey: "self-healing-no-action:reconcile-in-review-unmet-dependencies" }); + await service.notifyTaskWedge(ownerless, descriptor!); + await service.notifyTaskWedge(ownerless, descriptor!); + await vi.waitFor(() => expect(sendMessageOnce).toHaveBeenCalledTimes(1)); + expect(sendNotification).toHaveBeenCalledWith("task-wedged", expect.objectContaining({ taskId: "FN-8501" })); + expect(describeSelfHealingNoActionWedge(ownerless, "reconcile-in-review-unmet-dependencies", { taskActive: true })).toBeNull(); + await service.stop(); + }); + + it("does not resolve an ownerless no-action episode on an incidental in-review update", async () => { + const { store, service, sendMessageOnce, task, getWedge } = fixture(); + await service.start(); + const ownerless = task({ status: "in-review", error: undefined, paused: false, userPaused: false }); + const descriptor = describeSelfHealingNoActionWedge(ownerless, "reconcile-in-review-unmet-dependencies", { + taskActive: false, + hasExecutingTaskLock: false, + livePaths: [], + }); + await service.notifyTaskWedge(ownerless, descriptor!); + await vi.waitFor(() => expect(sendMessageOnce).toHaveBeenCalledTimes(1)); + + store.emit({ ...ownerless, title: "Unrelated update", wedgeNotification: getWedge() }); + await Promise.resolve(); + await Promise.resolve(); + expect(getWedge()?.status).toBe("active"); + await service.notifyTaskWedge(ownerless, descriptor!); + expect(sendMessageOnce).toHaveBeenCalledTimes(1); + await service.stop(); + }); + + it.each([ + ["branch-cross-contamination", "branch-cross-contamination"], + ["branch-conflict-tripwire", "branch-conflict-tripwire"], + ["branch-conflict-recovery-exhausted", "branch-conflict-recovery-exhausted"], + ["branch-conflict-unrecoverable", "branch-conflict-unrecoverable"], + ["stuck-loop-exhausted-manual-intervention-required", "stuck-loop-exhausted"], + ["non-retryable-provider-error", "non-retryable-provider-error"], + ["in-review-stall-deadlock", "in-review-stall-deadlock"], + ])("classifies automated terminal pause %s as %s", (pausedReason, reasonKey) => { + const { task } = fixture(); + expect(describeTaskWedge(task({ status: "failed", paused: true, pausedReason }))).toMatchObject({ reasonKey }); + }); + + 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" }); + }); + + it("changes the active reason into a new episode without raw error keys", async () => { + const { store, service, sendMessageOnce, task } = fixture(); + await service.start(); + store.emit(task({ error: "EXECUTION_DISPATCH_LOOP_EXHAUSTED: details" })); + store.emit(task({ error: "Tool failure retries exhausted", updatedAt: "2026-07-22T12:01:00.000Z", column: "in-progress" })); + await vi.waitFor(() => expect(sendMessageOnce).toHaveBeenCalledTimes(2)); + expect(sendMessageOnce.mock.calls.map((call) => call[1])).not.toContain(expect.stringContaining("details")); + await service.stop(); + }); +}); diff --git a/packages/engine/src/notification/notification-service.ts b/packages/engine/src/notification/notification-service.ts index 09c832d806..7e72acb819 100644 --- a/packages/engine/src/notification/notification-service.ts +++ b/packages/engine/src/notification/notification-service.ts @@ -16,6 +16,7 @@ import { schedulerLog } from "../logger.js"; import { classifyTransientMergeError } from "../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"; export interface NotificationServiceOptions { /** Project identifier for notification deep links */ @@ -43,6 +44,8 @@ interface NotificationServiceStoreEvents { interface NotificationServiceStore { getSettings(): Promise | Settings; getTask?(id: string): Promise | Task | undefined; + /** Durable compare-and-set for restart-safe wedge delivery episodes. */ + claimTaskWedgeNotificationEpisode?(taskId: string, reasonKey: string | null): Promise<{ episodeId?: string; claimed: boolean }>; on( event: K, listener: (...args: NotificationServiceStoreEvents[K]) => void, @@ -99,6 +102,8 @@ export class NotificationService { private failureNotificationSuppressedCount = 0; private failureNotificationDelayMs = 60_000; private failureNotificationMode: "sticky-only" | "all" | "terminal-only" = "sticky-only"; + /** Compatibility fallback for lightweight test stores without the durable TaskStore CAS. */ + private readonly activeWedgeReasons = new Map(); constructor( private readonly store: NotificationServiceStore, @@ -247,6 +252,23 @@ export class NotificationService { }; private handleTaskUpdated = (task: Task): 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. + */ + const transientFailure = task.status === "failed" ? classifyTransientMergeError(task.error) : null; + const wedge = transientFailure ? null : describeTaskWedge(task); + /* + FNXC:TaskWedgeNotifications 2026-07-22-14:30: + A generic failed push may have been scheduled before a terminal error was + classified. Cancel it at wedge entry so the immediate episode alert is the + only operator notification; dispatch-time suppression below covers races. + */ + if (wedge) this.cancelPendingFailureNotification(task.id, "classified-terminal-wedge"); + if (!transientFailure) void this.maybeNotifyTaskWedge(task, wedge); void this.maybeSuppressTransientFailedNotification(task, `status=${task.status ?? "undefined"}`); /* @@ -270,7 +292,7 @@ export class NotificationService { return; } - if (task.status === "failed") { + 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` @@ -390,6 +412,75 @@ export class NotificationService { } } + /* + FNXC:TaskWedgeNotifications 2026-07-22-12:00: + Terminal task updates are the shared seam for merger, executor, heartbeat, and + self-healing parks. Clear only on a non-wedge lifecycle update, allowing a + 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 { + await this.maybeNotifyTaskWedge(task, descriptor); + } + + private async maybeNotifyTaskWedge(task: Task, suppliedDescriptor?: TaskWedgeDescriptor | null): Promise { + const descriptor = suppliedDescriptor ?? describeTaskWedge(task); + let episode: string | undefined; + if (!descriptor) { + /* + FNXC:TaskWedgeNotifications 2026-07-22-14:45: + A self-healing no-action escalation may remain in `in-review` without a + status/error mutation. Incidental task updates are not resolution evidence; + resolve only when the lifecycle has visibly resumed in an active or terminal + column, or when it carries a non-failed workflow status. + */ + const isActiveSelfHealingNoAction = task.wedgeNotification?.status === "active" + && task.wedgeNotification.reasonKey.startsWith("self-healing-no-action:"); + // FNXC:TaskWedgeNotifications 2026-07-22-15:00: A no-action task normally + // stays in review, so arbitrary in-review/status writes are not resolution + // evidence. Only an active owner state or real lifecycle advance can close it. + const hasProgressed = task.column === "todo" || task.column === "in-progress" || task.column === "done" || task.column === "archived" + || (!isActiveSelfHealingNoAction && typeof task.status === "string" && task.status !== "failed") + || (isActiveSelfHealingNoAction && ["queued", "planning", "in-progress", "merging", "merging-pr", "merged", "done"].includes(task.status ?? "")); + if (hasProgressed) { + this.activeWedgeReasons.delete(task.id); + await this.store.claimTaskWedgeNotificationEpisode?.(task.id, null); + } + return; + } + if (this.store.claimTaskWedgeNotificationEpisode) { + const claim = await this.store.claimTaskWedgeNotificationEpisode(task.id, descriptor.reasonKey); + 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}`; + } + const link = buildNtfyClickUrl({ dashboardHost: this.dashboardHost, projectId: this.options.projectId, taskId: task.id }); + const content = [ + `**${formatTaskIdentifier(task)} needs operator action**`, "", descriptor.reason, + ...(descriptor.gate ? [`Gate: \`${descriptor.gate}\`.`] : []), + `Recommended action: ${descriptor.action}`, + ...(link ? ["", `[Open ${task.id}](${link})`] : []), + ].join("\n"); + const payload: NotificationPayload = { + taskId: task.id, taskTitle: task.title, taskDescription: task.description, event: "task-wedged", + // Bounded descriptor text is operator-facing provider content, never audit metadata. + metadata: { wedgeReason: descriptor.reasonKey, reason: descriptor.reason, action: descriptor.action, ...(descriptor.gate ? { gate: descriptor.gate } : {}), notificationDedupeKey: `task-wedge:${episode}` }, + }; + // Push and mailbox delivery are independently best-effort. + void this.dispatch("task-wedged", payload); + try { + await this.options.messageStore?.sendMessageOnce?.({ + fromId: "system", fromType: "system", toId: DASHBOARD_USER_ID, toType: "user", type: "system", + content, metadata: { taskId: task.id, kind: "task-wedge", wedgeReason: descriptor.reasonKey, ...(descriptor.gate ? { gate: descriptor.gate } : {}) }, + }, `task-wedge:${episode}`); + } catch (error) { + schedulerLog.log(`[notify] ${task.id} wedge mailbox message failed: ${error instanceof Error ? error.message : String(error)}`); + } + } + private isTriageDuplicateDecision(task: Task): boolean { return task.paused === true && task.pausedReason === "duplicate-decision-required" @@ -839,9 +930,9 @@ export class NotificationService { // 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. Self- - // healing may have flipped the error to a transient class via FN-5627 - // auto-recovery, in which case ntfy stays silent. + // 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) { this.failureNotificationSuppressedCount += 1; @@ -851,6 +942,15 @@ export class NotificationService { return; } + // A previously generic failure can become a terminal wedge while its grace + // timer is pending. The wedge path owns delivery and its durable episode + // idempotency; never let this delayed generic event become a second alert. + if (describeTaskWedge(task)) { + this.failureNotificationSuppressedCount += 1; + schedulerLog.log(`[notify] ${taskId} classified terminal wedge at dispatch time — suppressed generic failed notification`); + return; + } + const isTerminal = task.paused === true || task.column === "in-review"; if (this.failureNotificationMode === "terminal-only" && !isTerminal) { this.failureNotificationSuppressedCount += 1; diff --git a/packages/engine/src/notification/ntfy-provider.ts b/packages/engine/src/notification/ntfy-provider.ts index de0152caf3..f3587f0281 100644 --- a/packages/engine/src/notification/ntfy-provider.ts +++ b/packages/engine/src/notification/ntfy-provider.ts @@ -34,6 +34,7 @@ type SupportedNtfyEvent = | "in-review" | "merged" | "failed" + | "task-wedged" | "awaiting-approval" | "awaiting-user-review" | "planning-awaiting-input" @@ -50,6 +51,7 @@ const SUPPORTED_EVENTS = new Set([ "in-review", "merged", "failed", + "task-wedged", "awaiting-approval", "awaiting-user-review", "planning-awaiting-input", @@ -214,6 +216,15 @@ export class NtfyNotificationProvider implements NotificationProvider { message: `Task "${identifier}" has failed and needs attention`, priority: "high", }, + "task-wedged": { + title: `Task ${taskId} needs operator action`, + message: [ + `Task "${identifier}" is wedged${typeof payload.metadata?.reason === "string" ? `: ${payload.metadata.reason}` : " and needs operator action"}`, + typeof payload.metadata?.gate === "string" ? `Gate: ${payload.metadata.gate}` : null, + typeof payload.metadata?.action === "string" ? `Recommended action: ${payload.metadata.action}` : null, + ].filter((part): part is string => part !== null).join("\n"), + priority: "high", + }, "awaiting-approval": { title: payload.metadata?.awaitingApprovalReason === "plan-review-replan-cap" ? `Plan Review cap reached for ${taskId}` diff --git a/packages/engine/src/notification/task-wedge-notification.ts b/packages/engine/src/notification/task-wedge-notification.ts new file mode 100644 index 0000000000..7b8ff8b540 --- /dev/null +++ b/packages/engine/src/notification/task-wedge-notification.ts @@ -0,0 +1,160 @@ +import type { Task } from "@fusion/core"; + +/** A bounded, operator-safe description of a task that cannot make progress. */ +export interface TaskWedgeDescriptor { + reasonKey: string; + reason: string; + action: string; + gate?: string; +} + +/** + * FNXC:TaskWedgeNotifications 2026-07-22-14:30: + * Self-healing can deliberately decline a backward move without mutating task + * status. These bounded stage keys make that ownerless escalation visible through + * the same durable episode seam as failed and paused terminal parks. + */ +const SELF_HEALING_NO_ACTIONS: Record> = { + "reclaim-pr-conflict": { + reason: "Self-healing could not reclaim a stalled pull-request conflict.", + action: "Inspect the branch conflict and retry or reset the task to todo.", + }, + "reclaim-self-owned-branch-conflict": { + reason: "Self-healing could not recover a stalled branch conflict.", + action: "Inspect the branch conflict and retry or reset the task to todo.", + }, + "reconcile-in-review-unmet-dependencies": { + reason: "Unmet dependencies are preventing review from progressing.", + action: "Resolve the dependencies, then retry or reset the task to todo.", + }, + "reconcile-dependency-blocking-lease": { + reason: "A dependency-blocking task lease could not be safely reclaimed.", + action: "Inspect the task owner and lease, then retry or reset the task to todo.", + }, + "auto-rebound-paused-scope-decay": { + reason: "A paused task with stale scope could not be safely resumed.", + action: "Inspect the task scope and retry or reset the task to todo.", + }, + "stuck-merge-deadlock": { + reason: "A merge deadlock needs operator intervention.", + action: "Inspect merge ownership and retry or reset the task to todo.", + }, + "missing-worktree-merge-active": { + reason: "An active merge has an unusable worktree and could not be recovered.", + action: "Repair the worktree or reset the task to todo and retry.", + }, + "missing-worktree-review": { + reason: "Review has an unusable worktree and could not be recovered.", + action: "Repair the worktree or reset the task to todo and retry.", + }, + "finalize-no-op-review": { + reason: "A no-op review task could not be safely finalized.", + action: "Inspect merge state and retry or reset the task to todo.", + }, + "stale-incomplete-review": { + reason: "Incomplete review work could not be safely resumed.", + action: "Inspect incomplete steps and retry or reset the task to todo.", + }, + "ghost-review": { + reason: "A review task has no recoverable workflow owner.", + action: "Inspect workflow state and retry or reset the task to todo.", + }, + "no-progress-no-task-done": { + reason: "Execution stopped without progress and could not be safely requeued.", + action: "Inspect the task worktree and retry or reset the task to todo.", + }, + "partial-progress-no-task-done": { + reason: "Partial execution progress could not be safely resumed.", + action: "Inspect partial work and retry or reset the task to todo.", + }, +}; + +/** Returns an actionable descriptor only for an ownerless self-healing escalation. */ +export function describeSelfHealingNoActionWedge(task: Task, stage: string, metadata: Record | undefined): TaskWedgeDescriptor | null { + const description = SELF_HEALING_NO_ACTIONS[stage]; + if (!description || task.userPaused || task.paused || task.autoMerge === false) return null; + // Test and legacy proof producers may omit metadata; absent ownership evidence + // remains ownerless rather than turning best-effort notification into a park failure. + const proof = metadata ?? {}; + // These proof signals mean a live executor, checkout, or queued merge owns the task. + if (proof.taskActive === true || proof.hasExecutingTaskLock === true || proof.mergePending === true) return null; + if (Array.isArray(proof.livePaths) && proof.livePaths.length > 0) return null; + return { reasonKey: `self-healing-no-action:${stage}`, ...description }; +} + +/* +FNXC:TaskWedgeNotifications 2026-07-22-12:00: +Terminal task updates are the shared delivery seam for merger, executor, heartbeat, +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. +*/ +export function describeTaskWedge(task: Task): TaskWedgeDescriptor | null { + const error = task.error ?? ""; + if (task.pausedReason === "completed-blocked") { + return { reasonKey: "completion-blocked", reason: "Completed work is blocked from advancing to review.", action: "Clear the blocker or reset the task to todo." }; + } + if (task.pausedReason === "error-retry-exhausted") { + return { reasonKey: "heartbeat-retry-exhausted", reason: "The assigned agent exhausted its heartbeat recovery budget.", action: "Repair the agent configuration, then retry the task." }; + } + if (task.pausedReason === "error-unrecoverable") { + return { reasonKey: "heartbeat-error-unrecoverable", reason: "The assigned agent needs operator repair before it can resume.", action: "Repair credentials, access, or configuration, then retry the task." }; + } + /* + FNXC:TaskWedgeNotifications 2026-07-22-19:00: + Branch and remediation safety parks deliberately stop automatic recovery. They + are actionable terminal writers rather than user-controlled approval pauses. + Keep their reason keys stable so a changed safety failure opens a new episode. + */ + const pausedDescriptors: Record = { + "branch-cross-contamination": { reasonKey: "branch-cross-contamination", reason: "Branch contamination recovery requires operator intervention.", action: "Inspect the branch history, repair the contamination, then retry or reset to todo." }, + "branch-conflict-tripwire": { reasonKey: "branch-conflict-tripwire", reason: "A repeated branch conflict stopped automatic recovery.", action: "Resolve the branch conflict, then retry or reset to todo." }, + "branch-conflict-recovery-exhausted": { reasonKey: "branch-conflict-recovery-exhausted", reason: "Branch conflict recovery retries were exhausted.", action: "Resolve the conflict, then retry or reset to todo." }, + "branch-conflict-unrecoverable": { reasonKey: "branch-conflict-unrecoverable", reason: "An unrecoverable branch conflict needs operator intervention.", action: "Resolve the conflict or reset the task to todo and retry." }, + "stuck-loop-exhausted-manual-intervention-required": { reasonKey: "stuck-loop-exhausted", reason: "Self-healing exhausted its stalled-task recovery loop.", action: "Inspect the task state, then retry or reset to todo." }, + "non-retryable-provider-error": { reasonKey: "non-retryable-provider-error", reason: "A non-retryable provider error stopped the task.", action: "Repair provider access or configuration, then retry the task." }, + "in-review-stall-deadlock": { reasonKey: "in-review-stall-deadlock", reason: "Review stalled in a deadlock that needs operator intervention.", action: "Inspect review ownership and retry or reset to todo." }, + }; + if (task.pausedReason && pausedDescriptors[task.pausedReason]) return pausedDescriptors[task.pausedReason]; + if (task.status !== "failed") return null; + if (error.startsWith("EXECUTION_DISPATCH_LOOP_EXHAUSTED")) { + return { reasonKey: "execution-dispatch-loop-exhausted", reason: "Execution re-queued without progress until its retry budget was exhausted.", action: "Retry, decompose, or rescope the task." }; + } + if (error.includes("tool failure") || error.includes("Tool failure")) { + return { reasonKey: "tool-failure-retry-exhausted", reason: "Execution tool-failure retries were exhausted.", action: "Inspect the failing tool and retry the task." }; + } + if (/escalat(?:ion|ed).*exhaust|exhaust.*escalat/i.test(error)) { + return { reasonKey: "execution-escalation-exhausted", reason: "The alternate execution escalation was exhausted.", action: "Inspect the failure and retry or rescope the task." }; + } + if (/review.*(?:retry|rework|revision).*(?:exhaust|limit)|(?:exhaust|limit).*review/i.test(error)) { + return { reasonKey: "review-retry-exhausted", reason: "The configured review recovery budget was exhausted.", action: "Review the feedback, fix the task, or use the approved review recovery action." }; + } + const gate = error.match(/check:([\w:-]+)/)?.[1]; + if (gate) { + return { + reasonKey: `merge-blocked:${gate}`, + reason: `Merge is blocked by the check:${gate} verification gate.`, + action: gate === "changeset-format" ? "Fix the changeset format, then retry the task." : "Fix the failing gate or reset the task to todo and retry.", + gate: `check:${gate}`, + }; + } + if (error.startsWith("BLOCKED:")) { + return { reasonKey: "execution-blocked", reason: "Execution was parked because an external blocker requires operator action.", action: "Resolve the dependency or blocker, then retry the task." }; + } + if (/merge.*(?:blocked|verification|gate)|(?:verification|gate).*(?:failed|exhaust)/i.test(error)) { + 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. + */ + if (!error && (task.mergeRetries ?? 0) < 3) return null; + return { + reasonKey: "terminal-failed", + reason: "The task entered a terminal failed state and needs operator intervention.", + action: "Inspect the task error, fix the underlying issue, then retry or reset to todo.", + }; +} diff --git a/packages/engine/src/notifier.ts b/packages/engine/src/notifier.ts index c6059f141c..3053ddccda 100644 --- a/packages/engine/src/notifier.ts +++ b/packages/engine/src/notifier.ts @@ -26,7 +26,10 @@ const DEFAULT_NTFY_RETRY_DELAY_MS = 500; export const DEFAULT_NTFY_EVENTS: readonly NtfyNotificationEvent[] = [ "in-review", "merged", + // FNXC:TaskWedgeNotifications 2026-07-22-19:00: Wedge episodes are operator + // escalation events, so normal Ntfy defaults must not silently filter them. "failed", + "task-wedged", "awaiting-approval", "awaiting-user-review", "planning-awaiting-input", diff --git a/packages/engine/src/self-healing.ts b/packages/engine/src/self-healing.ts index 37e994b124..8596a46db5 100644 --- a/packages/engine/src/self-healing.ts +++ b/packages/engine/src/self-healing.ts @@ -96,6 +96,7 @@ import type { GhostBugDecision } from "./triage-preflight.js"; import { DependencyBlockedTodoReporter } from "./dependency-blocked-todo-reporter.js"; import { filterPathsByIgnoreList, getUnmetSchedulingDependencies, isCoordinationOnlyTask, pathsOverlap, shouldHoldActiveFileScopeLease } from "./scheduler.js"; import { evaluateParkedAgentTaskLink, PARKED_AGENT_LINK_FRESH_RUN_MS } from "./task-agent-sync.js"; +import { describeSelfHealingNoActionWedge } from "./notification/task-wedge-notification.js"; export { COMPLETED_BLOCKED_PAUSE_REASON, @@ -1074,6 +1075,22 @@ export class SelfHealingManager { const message = error instanceof Error ? error.message : String(error); log.warn(`[${stage}] ${task.id}: no-action audit emission failed: ${message}`); } + + /* + FNXC:TaskWedgeNotifications 2026-07-22-14:30: + No-action reconciles do not always write `failed` or `paused`, so task-updated + classification cannot see them. Deliver only the bounded ownerless stages + through NotificationService; its durable CAS suppresses repeated sweeps. + */ + const descriptor = describeSelfHealingNoActionWedge(task, stage, proof.metadata); + if (descriptor) { + try { + await getActiveNotificationService()?.notifyTaskWedge(task, descriptor); + } catch (error) { + const message = error instanceof Error ? error.message : String(error); + log.warn(`[${stage}] ${task.id}: wedge notification failed: ${message}`); + } + } log.log(`[${stage}] ${task.id}: triple-proof not satisfied — no action (operator-decides)`); }