FN-8505: notify operators of terminal task wedges
Deliver durable, actionable notifications when terminal task recovery wedges. - Persist and deduplicate terminal wedge notification episodes across task updates and service restarts. - Classify terminal failure and self-healing escalation states, then deliver actionable ntfy and mailbox alerts. - Align PostgreSQL baseline and upgrade migration registration for the durable wedge field. Files changed: .changeset/fn-8505-task-wedge-notifications.md | 7 + docs/agents.md | 4 + docs/architecture.md | 4 + .../core/src/postgres/migrations/0000_initial.sql | 2 + .../migrations/0033_fn-8505_wedge_notification.sql | 5 + packages/core/src/postgres/schema-applier.ts | 29 +++- packages/core/src/postgres/schema/project.ts | 1 + packages/core/src/store.ts | 22 ++- packages/core/src/task-store/persistence.ts | 2 + packages/core/src/task-store/serialization.ts | 1 + packages/core/src/task-store/task-row-mappers.ts | 2 +- packages/core/src/task-store/task-update.ts | 5 + packages/core/src/types.ts | 14 ++ packages/core/src/types/workflow-steps.ts | 2 + .../src/__tests__/notification-service.test.ts | 20 +++ packages/engine/src/__tests__/notifier.test.ts | 26 +++- packages/engine/src/__tests__/self-healing.test.ts | 9 +- .../__tests__/notification-service.test.ts | 34 ++++- .../__tests__/task-wedge-notification.test.ts | 134 +++++++++++++++++ .../src/notification/notification-service.ts | 108 +++++++++++++- packages/engine/src/notification/ntfy-provider.ts | 11 ++ .../src/notification/task-wedge-notification.ts | 160 +++++++++++++++++++++ packages/engine/src/notifier.ts | 3 + packages/engine/src/self-healing.ts | 17 +++ 24 files changed, 610 insertions(+), 12 deletions(-) Fusion-Task-Id: FN-8505 Fusion-Task-Lineage: eee85220-18ba-475d-9d01-dc96e2b923e6 Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
7
.changeset/fn-8505-task-wedge-notifications.md
Normal file
7
.changeset/fn-8505-task-wedge-notifications.md
Normal file
@@ -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.
|
||||
@@ -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.
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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;
|
||||
@@ -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 };
|
||||
});
|
||||
}
|
||||
|
||||
@@ -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"),
|
||||
|
||||
@@ -1240,10 +1240,30 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
|
||||
}
|
||||
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<string, unknown>; 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<string, unknown> | 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<string, unknown>; 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<string, unknown> | 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<Task> {
|
||||
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<import("./task-store/branch-and-pr-entities.js").ToolFailureRetryClaim> {
|
||||
return claimNextToolFailureRetryImpl(this, taskId, expectedCursor, maxRetries);
|
||||
}
|
||||
|
||||
@@ -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),
|
||||
|
||||
@@ -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<Task["wedgeNotification"]>(row.wedgeNotification) ?? undefined,
|
||||
userPaused: row.userPaused ? true : undefined,
|
||||
baseBranch: row.baseBranch || undefined,
|
||||
executionStartBranch: row.executionStartBranch || undefined,
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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. */
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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 });
|
||||
|
||||
@@ -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:'),
|
||||
})
|
||||
);
|
||||
});
|
||||
|
||||
@@ -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" });
|
||||
|
||||
@@ -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();
|
||||
});
|
||||
|
||||
|
||||
@@ -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<Listener>();
|
||||
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> = {}): 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();
|
||||
});
|
||||
});
|
||||
@@ -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> | Settings;
|
||||
getTask?(id: string): Promise<Task | undefined> | Task | undefined;
|
||||
/** Durable compare-and-set for restart-safe wedge delivery episodes. */
|
||||
claimTaskWedgeNotificationEpisode?(taskId: string, reasonKey: string | null): Promise<{ episodeId?: string; claimed: boolean }>;
|
||||
on<K extends keyof NotificationServiceStoreEvents>(
|
||||
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<string, string>();
|
||||
|
||||
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<void> {
|
||||
await this.maybeNotifyTaskWedge(task, descriptor);
|
||||
}
|
||||
|
||||
private async maybeNotifyTaskWedge(task: Task, suppliedDescriptor?: TaskWedgeDescriptor | null): Promise<void> {
|
||||
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;
|
||||
|
||||
@@ -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<SupportedNtfyEvent>([
|
||||
"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}`
|
||||
|
||||
160
packages/engine/src/notification/task-wedge-notification.ts
Normal file
160
packages/engine/src/notification/task-wedge-notification.ts
Normal file
@@ -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<string, Omit<TaskWedgeDescriptor, "reasonKey">> = {
|
||||
"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<string, unknown> | 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<string, TaskWedgeDescriptor> = {
|
||||
"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.",
|
||||
};
|
||||
}
|
||||
@@ -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",
|
||||
|
||||
@@ -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)`);
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user