diff --git a/.changeset/fn-8837-pr-merge-retries.md b/.changeset/fn-8837-pr-merge-retries.md new file mode 100644 index 0000000000..7eb855e84f --- /dev/null +++ b/.changeset/fn-8837-pr-merge-retries.md @@ -0,0 +1,7 @@ +--- +"@runfusion/fusion": patch +--- + +summary: Make pull-request merge retries honest and pause branch-policy blocks for operator action. +category: fix +dev: Enforces persisted PR retry backoff and resumes policy holds through manual merge. diff --git a/docs/architecture.md b/docs/architecture.md index 60c218e7f7..cae08c4848 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -2277,6 +2277,7 @@ This section preserves the detailed lifecycle/self-healing contracts that were f - **Worktrunk-managed lifecycles**: when `worktrunk.enabled`, self-healing defers prune/idle/worktree-cap sweeps to the worktrunk backend; branch-level stale/ conflict reclaim stays native. Orphan `fusion/*` branches are operator-managed via standard git tooling (no auto-rescue task filing). - **Post-finalize verification no-op (FN-4944)**: when auto-merge receives a delayed `VerificationError` after a task is already `done` with `mergeDetails.mergeConfirmed === true` (already-on-main fast-path), it must log one `[verification] ... no action` diagnostic and must not bounce the task back to `in-progress` / `merging-fix`. Defense-in-depth now re-checks the done+mergeConfirmed condition immediately before each verification-failure status write site, and emits `task:post-finalize-verification-no-op` database audit events with failure metadata for forensics. - **Transient auto-merge retry classification (FN-5697)**: non-conflict auto-merge errors now run through `isTransientError(...)` before terminal parking. Transient provider/network failures (for example `This operation was aborted`, `socket hang up`, and `server_error` payloads) are retried with bounded exponential backoff (`5s/10s/20s`) and `status=null` for both direct and pull-request merge strategies; once `MAX_AUTO_MERGE_TRANSIENT_RETRIES` is exhausted, tasks are parked `in-review/failed` with explicit transient-exhaustion logs. +- **Pull-request merge accounting (FN-8837)**: retryable non-transient PR failures increment `mergeRetries` one attempt at a time; their atomic task-update timestamp derives the 5s/10s/20s not-before deadline enforced by merge admission, keeping workflow-defined `customFields` free of engine metadata. Structured branch-policy blocks park as `awaiting-approval` without consuming either retry budget; an operator's manual merge action clears that hold and re-enters the normal serialized PR merge path. Other non-retryable GitHub errors park as failed while retaining their truthful retry count. - **Merge-seam abort provenance (FN-6568/FN-6735/FN-7749)**: workflow graph merge-node failures must not be classified as pause/resume aborts merely because the merge seam hard-canceled an in-flight session. `TaskExecutor` tracks paused-abort provenance separately (`global-pause`, `merge-seam`, `hard-cancel`); genuine user/global pauses still preserve FN-6478/FN-5147 parking, while non-paused merge-seam graph failures (`merge`, `requestMerge`, built-in merge-region node ids, `merge-manual-hold`, and `merge-retry`) route back into the bounded auto-merge retry path instead of being parked `status:"failed"` with `mergeRetries=NULL`. Benign pause/resume aborts at these seams are retryable when the task is already `in-review`, has no durable failure/status, has not confirmed a merge, remains auto-merge eligible (or is a shared-branch local integration), and has merge retries remaining. When auto-merge is off (`settings.autoMerge:false` or an explicit task-level `autoMerge:false`), a benign hard-cancel at a non-terminal merge-region/manual-hold node is instead preserved cleanly in `in-review` for human Merge & Close; stale pause-abort status/error of this exact shape is cleared in place and never requeued. Conflict/contamination/foreign-work/retry-exhaustion values, pre-existing failures, global/user pauses, shared-branch member integrations, and post-confirmation partial landings retain their existing terminal/retry/finalize behavior. - **Worktree pool exclusivity (FN-4954)**: `WorktreePool.acquire(taskId)` / `release(path, taskId?)` track a `leased` map so every pooled path is either idle or leased, never both. Cross-task double-lease detection throws `PoolDoubleLeaseError` and emits `worktree:pool-double-lease-detected`; merger Step 8 now detaches HEAD and clears `task.worktree` / `task.branch` before releasing paths back to the pool. - **Stale registration recovery (FN-5056)**: `NativeWorktreeBackend.create` and `executor.tryCreateWorktree` detect `missing but already registered worktree` failures, run `git worktree prune` (plus `remove --force` / `add -f` fallbacks) before retrying, and emit `worktree:stale-registration-{detected,recovered,recovery-failed}` audit events. diff --git a/packages/core/src/__tests__/task-update-awaiting-approval-reason.test.ts b/packages/core/src/__tests__/task-update-awaiting-approval-reason.test.ts index d1b9305ce2..7c3cf26ddc 100644 --- a/packages/core/src/__tests__/task-update-awaiting-approval-reason.test.ts +++ b/packages/core/src/__tests__/task-update-awaiting-approval-reason.test.ts @@ -67,6 +67,15 @@ describe("awaitingApprovalReason survives updateTask", () => { expect(row.awaitingApprovalReason).toBe("plan-review-replan-cap"); }); + it("persists the structured PR policy-block reason alongside its human hold", async () => { + const { store, row } = harness({ status: null }); + + await run(store, { status: "awaiting-approval", awaitingApprovalReason: "merge-blocked-by-policy" }); + + expect(row.status).toBe("awaiting-approval"); + expect(row.awaitingApprovalReason).toBe("merge-blocked-by-policy"); + }); + it("clears the reason on an explicit null (the manual plan gate's stale-reason guard)", async () => { const { store, row } = harness({ status: "needs-replan", diff --git a/packages/core/src/store.ts b/packages/core/src/store.ts index 69fc2ff8e5..bfb5a8dd81 100644 --- a/packages/core/src/store.ts +++ b/packages/core/src/store.ts @@ -1426,7 +1426,7 @@ 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; 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; credentialInstanceId?: string | null; modelId?: string | null; validatorModelProvider?: string | null; validatorCredentialInstanceId?: string | null; validatorModelId?: string | null; planningModelProvider?: string | null; planningCredentialInstanceId?: string | null; planningModelId?: string | null; mergerModelProvider?: string | null; mergerCredentialInstanceId?: string | null; mergerModelId?: string | null; thinkingLevel?: string | null; validatorThinkingLevel?: string | null; planningThinkingLevel?: string | null; mergerThinkingLevel?: string | null; error?: string | null; summary?: string | null; recommendations?: import("./types.js").TaskRecommendation[]; 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; awaitingApprovalReason?: import("./types.js").Task["awaitingApprovalReason"] | 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; credentialInstanceId?: string | null; modelId?: string | null; validatorModelProvider?: string | null; validatorCredentialInstanceId?: string | null; validatorModelId?: string | null; planningModelProvider?: string | null; planningCredentialInstanceId?: string | null; planningModelId?: string | null; mergerModelProvider?: string | null; mergerCredentialInstanceId?: string | null; mergerModelId?: string | null; thinkingLevel?: string | null; validatorThinkingLevel?: string | null; planningThinkingLevel?: string | null; mergerThinkingLevel?: string | null; error?: string | null; summary?: string | null; recommendations?: import("./types.js").TaskRecommendation[]; 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); } diff --git a/packages/core/src/task-store/task-update.ts b/packages/core/src/task-store/task-update.ts index 4fe14ffcb6..58599e9b61 100644 --- a/packages/core/src/task-store/task-update.ts +++ b/packages/core/src/task-store/task-update.ts @@ -333,7 +333,11 @@ export async function updateTaskUnlockedImpl(store: TaskStore, id: string, updat but never applied by this field-by-field merge, so EVERY writer silently lost it — the executor's Plan Review replan-cap park (`plan-review-replan-cap`) and the triage manual gate's explicit null-clear both no-oped, and FN-8647's non-converging Plan Review loop - surfaced on the board as a generic "needs approval" with no explanation. Merge it like the + surfaced on the board as a generic "needs approval" with no explanation. + + FNXC:PullRequestMerge 2026-08-09-05:07: + The PR merge queue also writes `merge-blocked-by-policy`; persist it through this same + nullable contract so its notification and manual-resume lifecycle are durable. Merge it like the other nullable fields (null clears), and auto-clear the stored reason whenever a status write moves the task OFF `awaiting-approval` without the caller addressing the reason, so an approved/replanned card can never carry a stale escalation reason into its next park. diff --git a/packages/core/src/types/task/task-core.ts b/packages/core/src/types/task/task-core.ts index 06a1a1cea6..1da48dce21 100644 --- a/packages/core/src/types/task/task-core.ts +++ b/packages/core/src/types/task/task-core.ts @@ -1068,9 +1068,14 @@ export interface Task { * `plan-review-replan-cap` when automatic REVISE replans hit PLAN_REVIEW_GATE_REPLAN_CAP. * Dashboard badge/detail banner/notifications must surface that reason so operators know * approval is required because Plan Review did not converge — not a generic require-all gate. + * + * FNXC:PullRequestMerge 2026-08-09-05:07: + * The PR merge queue stamps `merge-blocked-by-policy` for branch-protection holds. + * Its notification must ask for policy remediation and a manual merge retry, never + * mislabel a completed implementation as a plan awaiting approval. * Undefined means either no hold or a routine manual plan-approval hold. */ - awaitingApprovalReason?: "release-authorization" | "plan-review-replan-cap"; + awaitingApprovalReason?: "release-authorization" | "plan-review-replan-cap" | "merge-blocked-by-policy"; /* * FNXC:PlanApproval 2026-07-04-22:41: * FN-7569 — records the computePlanApprovalFingerprint (packages/core/src/plan-approval.ts) diff --git a/packages/engine/src/__tests__/merge-error-recovery.test.ts b/packages/engine/src/__tests__/merge-error-recovery.test.ts index b4b48ed923..e16e9c1159 100644 --- a/packages/engine/src/__tests__/merge-error-recovery.test.ts +++ b/packages/engine/src/__tests__/merge-error-recovery.test.ts @@ -1,5 +1,5 @@ import { beforeEach, describe, expect, it, vi, type MockInstance } from "vitest"; -import type { Settings, Task } from "@fusion/core"; +import { validateCustomFieldPatch, type Settings, type Task } from "@fusion/core"; const testState = vi.hoisted(() => { class MockVerificationError extends Error { @@ -87,10 +87,12 @@ type MockTask = { verificationFailureCount?: number; mergeConflictBounceCount?: number; mergeTransientRetryCount?: number; + awaitingApprovalReason?: "merge-blocked-by-policy" | null; branch?: string; worktree?: string; sourceType?: string; sourceParentTaskId?: string; + customFields?: Record; updatedAt: string; log: Array<{ action?: string }>; }; @@ -108,6 +110,7 @@ type MockTaskStore = { recordRunAuditEvent: ReturnType; on: ReturnType; off: ReturnType; + emit: ReturnType; }; const TASK_ID = "FN-2084"; @@ -169,6 +172,7 @@ function makeStore({ recordRunAuditEvent: vi.fn(async () => undefined), on: vi.fn(), off: vi.fn(), + emit: vi.fn(), }; } @@ -632,6 +636,25 @@ describe("ProjectEngine merge error recovery", () => { vi.useRealTimers(); }); + it("uses the transient budget for structured GitHub transport outcomes", async () => { + vi.useFakeTimers(); + const setTimeoutSpy = vi.spyOn(globalThis, "setTimeout"); + const structuredRateLimit = Object.assign(new Error("GitHub rate limiting response"), { code: "rate-limited" }); + const store = makeStore(); + const processPullRequestMerge = vi.fn(async () => { throw structuredRateLimit; }); + const engine = createEngine(store, { getMergeStrategy: () => "pull-request", processPullRequestMerge }); + + await runMergeCycle(engine); + + expect(store.updateTask).toHaveBeenCalledWith(TASK_ID, { + mergeTransientRetryCount: 1, + status: null, + }); + expect(store.updateTask).not.toHaveBeenCalledWith(TASK_ID, expect.objectContaining({ mergeRetries: expect.any(Number) })); + expect(setTimeoutSpy).toHaveBeenCalledWith(expect.any(Function), 5_000); + vi.useRealTimers(); + }); + it("logs when non-direct merge strategy recovery update fails", async () => { const store = makeStore({ updateTask: vi.fn(async () => { @@ -650,17 +673,422 @@ describe("ProjectEngine merge error recovery", () => { await expect(runMergeCycle(engine)).resolves.toBeUndefined(); expect(processPullRequestMerge).toHaveBeenCalledTimes(1); - expect(store.updateTask).toHaveBeenCalledWith(TASK_ID, { - status: "failed", - mergeRetries: 3, - error: "PR API timeout", - }); + expect(store.updateTask).toHaveBeenCalledWith(TASK_ID, expect.objectContaining({ + status: null, + mergeRetries: 1, + error: null, + })); expect(hasErrorLog(errorSpy, `failed to update ${TASK_ID} after merge strategy error`)).toBe( true, ); expect(hasErrorLog(errorSpy, "persist failed")).toBe(true); }); + it("accounts pull-request retryable failures one at a time with durable backoff", async () => { + vi.useFakeTimers(); + const store = makeStore(); + const processPullRequestMerge = vi.fn(async () => { throw new Error("unexpected GitHub response"); }); + const engine = createEngine(store, { getMergeStrategy: () => "pull-request", processPullRequestMerge }); + + await runMergeCycle(engine); + + expect(store.updateTask).toHaveBeenCalledWith(TASK_ID, expect.objectContaining({ + mergeRetries: 1, + status: null, + })); + expect(store.updateTask).not.toHaveBeenCalledWith(TASK_ID, expect.objectContaining({ mergeRetries: 3, status: "failed" })); + expect(vi.getTimerCount()).toBeGreaterThanOrEqual(1); + vi.useRealTimers(); + }); + + it("does not let its retry log move the durable PR backoff past the scheduled timer", async () => { + vi.useFakeTimers(); + vi.setSystemTime(new Date("2026-08-09T03:11:00.000Z")); + const task = makeTask({ updatedAt: "2026-08-09T03:10:00.000Z" }); + const store = makeStore({ tasks: [task] }); + store.getTask.mockImplementation(async () => task); + store.updateTask.mockImplementation(async (_taskId: string, patch: Partial) => { + Object.assign(task, patch, { updatedAt: new Date().toISOString() }); + }); + // Task logs also update `updatedAt` in the production store. Model a later + // timestamp to prove the retry patch, rather than its log, owns the anchor. + store.logEntry.mockImplementation(async () => { + task.updatedAt = new Date(Date.now() + 1).toISOString(); + }); + const processPullRequestMerge = vi + .fn<(...args: unknown[]) => Promise<"merged" | "waiting" | "skipped">>() + .mockRejectedValueOnce(new Error("unexpected GitHub response")) + .mockResolvedValueOnce("merged"); + const engine = createEngine(store, { getMergeStrategy: () => "pull-request", processPullRequestMerge }); + (engine as unknown as { started: boolean }).started = true; + + await runMergeCycle(engine); + await vi.advanceTimersByTimeAsync(5_000); + + expect(processPullRequestMerge).toHaveBeenCalledTimes(2); + vi.useRealTimers(); + }); + + it("uses the persisted retry count for the 10s ladder and truthful final boundary", async () => { + vi.useFakeTimers(); + const setTimeoutSpy = vi.spyOn(globalThis, "setTimeout"); + const retrying = makeTask({ mergeRetries: 1, updatedAt: new Date(Date.now() - 20_000).toISOString() }); + const retryingStore = makeStore({ tasks: [retrying, retrying] }); + const retryingEngine = createEngine(retryingStore, { + getMergeStrategy: () => "pull-request", + processPullRequestMerge: vi.fn(async () => { throw new Error("unexpected GitHub response"); }), + }); + + await runMergeCycle(retryingEngine); + expect(retryingStore.updateTask).toHaveBeenCalledWith(TASK_ID, { + mergeRetries: 2, + status: null, + error: null, + }); + expect(setTimeoutSpy).toHaveBeenCalledWith(expect.any(Function), 10_000); + await retryingEngine.stop(); + + const twentySecondAttempt = makeTask({ mergeRetries: 2, updatedAt: new Date(Date.now() - 30_000).toISOString() }); + const twentySecondStore = makeStore({ tasks: [twentySecondAttempt, twentySecondAttempt], settings: { maxAutoMergeRetries: 4 } }); + const twentySecondEngine = createEngine(twentySecondStore, { + getMergeStrategy: () => "pull-request", + processPullRequestMerge: vi.fn(async () => { throw new Error("unexpected GitHub response"); }), + }); + await runMergeCycle(twentySecondEngine); + expect(twentySecondStore.updateTask).toHaveBeenCalledWith(TASK_ID, { + mergeRetries: 3, + status: null, + error: null, + }); + expect(setTimeoutSpy).toHaveBeenCalledWith(expect.any(Function), 20_000); + await twentySecondEngine.stop(); + + const finalAttempt = makeTask({ mergeRetries: 2, updatedAt: new Date(Date.now() - 30_000).toISOString() }); + const finalStore = makeStore({ tasks: [finalAttempt, finalAttempt], settings: { maxAutoMergeRetries: 3 } }); + const finalEngine = createEngine(finalStore, { + getMergeStrategy: () => "pull-request", + processPullRequestMerge: vi.fn(async () => { throw new Error("unexpected GitHub response"); }), + }); + + await runMergeCycle(finalEngine); + expect(finalStore.updateTask).toHaveBeenCalledWith(TASK_ID, { + status: "failed", + mergeRetries: 3, + error: "unexpected GitHub response", + }); + expect(finalStore.logEntry).toHaveBeenCalledWith( + TASK_ID, + expect.stringContaining("3/3 actual failures"), + "MergeRetriesExhausted", + ); + await finalEngine.stop(); + vi.useRealTimers(); + }); + + it("reschedules an early PR retry rejected by drain admission", async () => { + vi.useFakeTimers(); + vi.setSystemTime(new Date("2026-08-09T04:05:00.000Z")); + const task = makeTask({ mergeRetries: 1, updatedAt: new Date().toISOString() }); + const store = makeStore({ tasks: [task] }); + const processPullRequestMerge = vi.fn(async () => "merged" as const); + const engine = createEngine(store, { getMergeStrategy: () => "pull-request", processPullRequestMerge }); + (engine as unknown as { started: boolean }).started = true; + + // Model a duplicate/restart enqueue which arrives before the persisted not-before. + await runMergeCycle(engine); + expect(processPullRequestMerge).not.toHaveBeenCalled(); + expect(vi.getTimerCount()).toBeGreaterThanOrEqual(1); + await vi.advanceTimersByTimeAsync(5_000); + expect(processPullRequestMerge).toHaveBeenCalledTimes(1); + await engine.stop(); + vi.useRealTimers(); + }); + + it("parks exhausted pull-request transient retries without consuming the normal retry budget", async () => { + const atCap = ProjectEngine.MAX_AUTO_MERGE_TRANSIENT_RETRIES; + const task = makeTask({ + mergeRetries: 1, + mergeTransientRetryCount: atCap, + updatedAt: new Date(Date.now() - 6_000).toISOString(), + }); + const store = makeStore({ tasks: [task, task] }); + const processPullRequestMerge = vi.fn(async () => { throw new Error("socket hang up"); }); + const engine = createEngine(store, { getMergeStrategy: () => "pull-request", processPullRequestMerge }); + + await runMergeCycle(engine); + + expect(store.updateTask).toHaveBeenCalledWith(TASK_ID, { + status: "failed", + error: "socket hang up", + }); + expect(store.updateTask).not.toHaveBeenCalledWith(TASK_ID, expect.objectContaining({ mergeRetries: expect.any(Number) })); + expect(store.logEntry).toHaveBeenCalledWith( + TASK_ID, + expect.stringContaining("transient retries exhausted"), + "MergeTransientRetryExhausted", + ); + }); + + it("parks exhausted structured GitHub transport retries without consuming mergeRetries", async () => { + const task = makeTask({ + mergeRetries: 1, + mergeTransientRetryCount: ProjectEngine.MAX_AUTO_MERGE_TRANSIENT_RETRIES, + updatedAt: new Date(Date.now() - 6_000).toISOString(), + }); + const store = makeStore({ tasks: [task, task] }); + const structuredTimeout = Object.assign(new Error("GitHub request timed out"), { code: "timeout" }); + const processPullRequestMerge = vi.fn(async () => { throw structuredTimeout; }); + const engine = createEngine(store, { getMergeStrategy: () => "pull-request", processPullRequestMerge }); + + await runMergeCycle(engine); + + expect(store.updateTask).toHaveBeenCalledWith(TASK_ID, { + status: "failed", + error: "GitHub request timed out", + }); + expect(store.updateTask).not.toHaveBeenCalledWith(TASK_ID, expect.objectContaining({ mergeRetries: expect.any(Number) })); + expect(store.logEntry).toHaveBeenCalledWith( + TASK_ID, + expect.stringContaining("transient retries exhausted"), + "MergeTransientRetryExhausted", + ); + }); + + it("keeps PR retry metadata outside workflow custom fields", async () => { + const customFieldPatch = { __fusionPrMergeRetryNotBefore: "2026-08-09T02:40:00.000Z" }; + expect(validateCustomFieldPatch([], customFieldPatch)).toMatchObject({ + ok: false, + rejection: { code: "no-fields-defined" }, + }); + + const updateTask = vi.fn(async (_taskId: string, patch: Record) => { + if (patch.customFields !== undefined) { + const validation = validateCustomFieldPatch([], patch.customFields as Record); + if (!validation.ok) throw new Error(validation.rejection.detail); + } + }); + const store = makeStore({ updateTask }); + const processPullRequestMerge = vi.fn(async () => { throw new Error("unexpected GitHub response"); }); + const engine = createEngine(store, { getMergeStrategy: () => "pull-request", processPullRequestMerge }); + + await runMergeCycle(engine); + + expect(updateTask).toHaveBeenCalledWith(TASK_ID, { + mergeRetries: 1, + status: null, + error: null, + }); + }); + + it("parks structured pull-request policy blocks without consuming retries", async () => { + const store = makeStore(); + const policyError = Object.assign(new Error("Pull request is blocked by branch protection."), { code: "merge-blocked-by-policy" }); + const processPullRequestMerge = vi.fn(async () => { throw policyError; }); + const engine = createEngine(store, { getMergeStrategy: () => "pull-request", processPullRequestMerge }); + + await runMergeCycle(engine); + + expect(store.updateTask).toHaveBeenCalledWith(TASK_ID, expect.objectContaining({ + status: "awaiting-approval", + error: policyError.message, + awaitingApprovalReason: "merge-blocked-by-policy", + })); + expect(store.updateTask).not.toHaveBeenCalledWith(TASK_ID, expect.objectContaining({ mergeRetries: expect.any(Number) })); + }); + + it("parks non-retryable structured pull-request failures honestly", async () => { + const priorAttempt = makeTask({ mergeRetries: 2, updatedAt: new Date(Date.now() - 30_000).toISOString() }); + const store = makeStore({ tasks: [priorAttempt, priorAttempt] }); + const permissionError = Object.assign(new Error("GitHub denied access to this resource."), { code: "permission" }); + const processPullRequestMerge = vi.fn(async () => { throw permissionError; }); + const engine = createEngine(store, { getMergeStrategy: () => "pull-request", processPullRequestMerge }); + + await runMergeCycle(engine); + + expect(store.updateTask).toHaveBeenCalledWith(TASK_ID, expect.objectContaining({ + status: "failed", + error: permissionError.message, + })); + expect(store.updateTask).not.toHaveBeenCalledWith(TASK_ID, expect.objectContaining({ mergeRetries: 3 })); + }); + + it("keeps a policy hold parked when the public auto-enqueue surface is invoked", async () => { + const parked = makeTask({ + status: "awaiting-approval", + error: "Pull request is blocked by branch protection.", + }); + const store = makeStore({ tasks: [parked] }); + const processPullRequestMerge = vi.fn(async () => "merged" as const); + const engine = createEngine(store, { getMergeStrategy: () => "pull-request", processPullRequestMerge }); + await engine.start(); + + expect(engine.enqueueMerge(TASK_ID)).toBe(true); + await vi.waitFor(() => expect(store.getTask).toHaveBeenCalled()); + expect(processPullRequestMerge).not.toHaveBeenCalled(); + expect(store.updateTask).not.toHaveBeenCalled(); + await engine.stop(); + }); + + it("resumes a policy hold through manual onMerge without changing retry counters", async () => { + const parked = makeTask({ + status: "awaiting-approval", + error: "Pull request is blocked by branch protection.", + mergeRetries: 2, + mergeTransientRetryCount: 1, + }); + const store = makeStore({ tasks: [parked, parked, parked] }); + const processPullRequestMerge = vi.fn(async () => "merged" as const); + const engine = createEngine(store, { getMergeStrategy: () => "pull-request", processPullRequestMerge }); + await engine.start(); + await engine.onMerge(TASK_ID); + + expect(store.updateTask).toHaveBeenCalledWith(TASK_ID, { + status: null, + error: null, + awaitingApprovalReason: null, + }); + expect(processPullRequestMerge).toHaveBeenCalledTimes(1); + // The resume itself preserves both budgets; successful completion then closes + // that retry episode and clears the durable retry/backoff state. + expect(store.updateTask).toHaveBeenCalledWith(TASK_ID, { + mergeRetries: 0, + mergeTransientRetryCount: 0, + }); + await engine.stop(); + }); + + it("resumes a policy hold through the interpreter merge requester", async () => { + const parked = makeTask({ + status: "awaiting-approval", + error: "Pull request is blocked by branch protection.", + mergeRetries: 2, + mergeTransientRetryCount: 1, + }); + const store = makeStore({ tasks: [parked, parked, parked] }); + const processPullRequestMerge = vi.fn(async () => "merged" as const); + const engine = createEngine(store, { getMergeStrategy: () => "pull-request", processPullRequestMerge }); + await engine.start(); + + await engine.requestInterpreterMerge(TASK_ID); + + expect(store.updateTask).toHaveBeenCalledWith(TASK_ID, { + status: null, + error: null, + awaitingApprovalReason: null, + }); + expect(processPullRequestMerge).toHaveBeenCalledTimes(1); + await engine.stop(); + }); + + it("cancels a pending PR retry wake when an operator merges during backoff", async () => { + vi.useFakeTimers(); + vi.setSystemTime(new Date("2026-08-09T04:18:00.000Z")); + const task = makeTask(); + const store = makeStore({ tasks: [task] }); + store.getTask.mockImplementation(async () => task); + store.updateTask.mockImplementation(async (_taskId: string, patch: Partial) => { + Object.assign(task, patch, { updatedAt: new Date().toISOString() }); + }); + const processPullRequestMerge = vi + .fn<(...args: unknown[]) => Promise<"merged" | "waiting" | "skipped">>() + .mockRejectedValueOnce(new Error("unexpected GitHub response")) + .mockResolvedValueOnce("merged"); + const engine = createEngine(store, { getMergeStrategy: () => "pull-request", processPullRequestMerge }); + (engine as unknown as { started: boolean }).started = true; + + await runMergeCycle(engine); + expect(processPullRequestMerge).toHaveBeenCalledTimes(1); + await engine.onMerge(TASK_ID); + expect(processPullRequestMerge).toHaveBeenCalledTimes(2); + + await vi.advanceTimersByTimeAsync(5_000); + expect(processPullRequestMerge).toHaveBeenCalledTimes(2); + await engine.stop(); + vi.useRealTimers(); + }); + + it("blocks the real periodic sweep until the durable PR retry backoff elapses", async () => { + vi.useFakeTimers(); + vi.setSystemTime(new Date("2026-08-09T02:39:00.000Z")); + const task = makeTask({ updatedAt: new Date().toISOString() }); + const store = makeStore({ tasks: [task] }); + store.getTask.mockImplementation(async () => task); + store.updateTask.mockImplementation(async (_taskId: string, patch: Partial) => { + Object.assign(task, patch, { updatedAt: new Date().toISOString() }); + }); + const processPullRequestMerge = vi + .fn<(...args: unknown[]) => Promise<"merged" | "waiting" | "skipped">>() + .mockRejectedValueOnce(new Error("unexpected GitHub response")) + .mockResolvedValueOnce("merged"); + const engine = createEngine(store, { getMergeStrategy: () => "pull-request", processPullRequestMerge }); + (engine as unknown as { started: boolean }).started = true; + const privateEngine = engine as unknown as { + enqueueEligibleInReviewTasks: (tasks: Task[], settings: Pick) => Promise; + }; + + await runMergeCycle(engine); + expect(task.mergeRetries).toBe(1); + expect(processPullRequestMerge).toHaveBeenCalledTimes(1); + + // FNXC:AutoMergeRetries 2026-08-09-03:23: This production periodic-sweep + // dispatcher, rather than a predicate unit test, must honor the retry anchor. + await expect(privateEngine.enqueueEligibleInReviewTasks([task as Task], { + autoMerge: true, + maxAutoMergeRetries: 3, + })).resolves.toBe(0); + expect(processPullRequestMerge).toHaveBeenCalledTimes(1); + + await vi.advanceTimersByTimeAsync(5_000); + expect(processPullRequestMerge).toHaveBeenCalledTimes(2); + vi.useRealTimers(); + }); + + it("keeps a re-blocked policy hold out of sweeps and preserves both retry counters", async () => { + const task = makeTask({ + status: "awaiting-approval", + error: "Pull request is blocked by branch protection.", + mergeRetries: 2, + mergeTransientRetryCount: 1, + }); + const store = makeStore({ tasks: [task] }); + store.getTask.mockImplementation(async () => task); + store.updateTask.mockImplementation(async (_taskId: string, patch: Partial) => { + Object.assign(task, patch, { updatedAt: new Date().toISOString() }); + }); + const policyError = Object.assign(new Error("Pull request is blocked by branch protection."), { + code: "merge-blocked-by-policy", + }); + const processPullRequestMerge = vi.fn(async () => { throw policyError; }); + const engine = createEngine(store, { getMergeStrategy: () => "pull-request", processPullRequestMerge }); + await engine.start(); + const privateEngine = engine as unknown as { + enqueueEligibleInReviewTasks: (tasks: Task[], settings: Pick) => Promise; + }; + + await expect(privateEngine.enqueueEligibleInReviewTasks([task as Task], { + autoMerge: true, + maxAutoMergeRetries: 3, + })).resolves.toBe(0); + expect(processPullRequestMerge).not.toHaveBeenCalled(); + + await engine.onMerge(TASK_ID); + expect(processPullRequestMerge).toHaveBeenCalledTimes(1); + expect(task.status).toBe("awaiting-approval"); + expect(task.mergeRetries).toBe(2); + expect(task.mergeTransientRetryCount).toBe(1); + expect(task.error).toBe(policyError.message); + await engine.stop(); + }); + + it("treats absent or malformed pull-request retry anchors as elapsed", () => { + const store = makeStore(); + const engine = createEngine(store); + const privateEngine = engine as unknown as { canMergeTask: (task: MockTask, cap: number, review?: boolean, enforcePrBackoff?: boolean) => boolean }; + expect(privateEngine.canMergeTask(makeTask(), 3, undefined, true)).toBe(true); + expect(privateEngine.canMergeTask(makeTask({ mergeRetries: 1, updatedAt: "not-a-date" }), 3, undefined, true)).toBe(true); + expect(privateEngine.canMergeTask(makeTask({ status: "awaiting-approval" }), 3, undefined, true)).toBe(false); + }); + it("treats post-finalize verification failures as a no-op diagnostic", async () => { const verificationError = new Error("Deterministic test verification failed: assertion mismatch in workspace"); verificationError.name = "VerificationError"; diff --git a/packages/engine/src/__tests__/ntfy-provider.test.ts b/packages/engine/src/__tests__/ntfy-provider.test.ts index 588fd846a2..25e15c7f25 100644 --- a/packages/engine/src/__tests__/ntfy-provider.test.ts +++ b/packages/engine/src/__tests__/ntfy-provider.test.ts @@ -99,6 +99,20 @@ describe("NtfyNotificationProvider", () => { ); }); + it("describes a pull-request policy hold without calling it plan approval", async () => { + await provider.sendNotification("awaiting-approval", { + taskId: "FN-1", + taskTitle: "T", + event: "awaiting-approval", + metadata: { awaitingApprovalReason: "merge-blocked-by-policy" }, + }); + + expect(mocks.sendNtfyNotificationWithResult).toHaveBeenCalledWith(expect.objectContaining({ + title: "Pull-request policy block for FN-1", + message: expect.stringContaining("Resolve the policy requirement, then retry the merge."), + })); + }); + it("resolveParticipantLabel prefers names and falls back to ids", () => { expect(resolveParticipantLabel({ fromName: "Triage Bot", fromId: "agent-1" }, "from")).toBe("Triage Bot"); expect(resolveParticipantLabel({ toId: "agent-2" }, "to")).toBe("agent-2"); diff --git a/packages/engine/src/__tests__/webhook-provider.test.ts b/packages/engine/src/__tests__/webhook-provider.test.ts index 3cd82f6fae..92906bb7ae 100644 --- a/packages/engine/src/__tests__/webhook-provider.test.ts +++ b/packages/engine/src/__tests__/webhook-provider.test.ts @@ -156,6 +156,21 @@ describe("WebhookNotificationProvider", () => { ).resolves.toEqual({ success: false, providerId: "webhook", error: "Not initialized" }); }); + it("describes a pull-request policy hold without calling it plan approval", async () => { + fetchMock.mockResolvedValue({ ok: true, status: 200, statusText: "OK" }); + await provider.initialize({ webhookUrl: "https://example.com/hook", webhookFormat: "slack" }); + + await provider.sendNotification("awaiting-approval", { + taskId: "FN-1", + taskTitle: "My Task", + event: "awaiting-approval", + metadata: { awaitingApprovalReason: "merge-blocked-by-policy" }, + }); + + const [, requestInit] = fetchMock.mock.calls[0] as [string, RequestInit]; + expect(JSON.parse(String(requestInit.body)).text).toContain("Resolve the policy requirement, then retry the merge."); + }); + it.each([ ["in-review", "ready for review"], ["merged", "has been merged to main"], diff --git a/packages/engine/src/notification/__tests__/notification-service.test.ts b/packages/engine/src/notification/__tests__/notification-service.test.ts index 324e8055f6..c75fc2194a 100644 --- a/packages/engine/src/notification/__tests__/notification-service.test.ts +++ b/packages/engine/src/notification/__tests__/notification-service.test.ts @@ -170,6 +170,54 @@ describe("NotificationService deferred failure notifications", () => { await restarted.stop(); }); + it("deduplicates durable awaiting-approval messages across policy re-parks and restart", async () => { + const store = createStore(); + const messageKeys = new Set(); + let insertedCount = 0; + const sendMessageOnce = vi.fn(async (_input: unknown, idempotencyKey: string) => { + const inserted = !messageKeys.has(idempotencyKey); + messageKeys.add(idempotencyKey); + if (inserted) insertedCount += 1; + return { message: {} as any, inserted }; + }); + const service = new NotificationService(store as any, { + messageStore: { on: () => undefined, sendMessageOnce } as any, + }); + await service.start(); + const policyHold = task({ + id: "FN-policy-hold", + column: "in-review", + status: "awaiting-approval", + error: "Pull request is blocked by branch protection.", + awaitingApprovalReason: "merge-blocked-by-policy", + }); + + store.emit("task:updated", policyHold); + store.emit("task:updated", policyHold); + await vi.waitFor(() => expect(sendMessageOnce).toHaveBeenCalledTimes(2)); + expect(insertedCount).toBe(1); + expect(sendMessageOnce).toHaveBeenLastCalledWith( + expect.objectContaining({ + content: expect.stringContaining("pull-request merge is blocked"), + metadata: expect.objectContaining({ + taskId: "FN-policy-hold", + awaitingApprovalReason: "merge-blocked-by-policy", + }), + }), + "merge-policy-block:FN-policy-hold", + ); + + await service.stop(); + const restarted = new NotificationService(store as any, { + messageStore: { on: () => undefined, sendMessageOnce } as any, + }); + await restarted.start(); + store.emit("task:updated", policyHold); + await vi.waitFor(() => expect(sendMessageOnce).toHaveBeenCalledTimes(3)); + expect(insertedCount).toBe(1); + await restarted.stop(); + }); + it("Failure that persists past grace dispatches exactly once", async () => { const { store, service, sendNotification } = await setup(); store.setTask(task({ id: "FN-1", status: "failed" })); diff --git a/packages/engine/src/notification/notification-service.ts b/packages/engine/src/notification/notification-service.ts index 1390f6118a..59e61cb768 100644 --- a/packages/engine/src/notification/notification-service.ts +++ b/packages/engine/src/notification/notification-service.ts @@ -482,17 +482,26 @@ export class NotificationService { } const identifier = formatTaskIdentifier(task); const reason = task.awaitingApprovalReason ?? "manual"; + /* + FNXC:PullRequestMerge 2026-08-09-05:07: + A policy-blocked PR uses the durable awaiting-approval mailbox transport, + but it is not a plan gate. Keep its operator instruction explicit and use a + distinct once-key so an earlier plan-approval notice cannot hide the block. + */ + const isMergePolicyBlock = reason === "merge-blocked-by-policy"; const reasonLine = reason === "plan-review-replan-cap" ? "Plan Review exhausted its automatic revision attempts and escalated this plan for a human decision." - : "The generated plan is ready and needs your approval before execution begins."; + : isMergePolicyBlock + ? "The pull request is blocked by repository policy. Resolve the policy requirement, then retry the merge." + : "The generated plan is ready and needs your approval before execution begins."; const link = buildNtfyClickUrl({ dashboardHost: this.dashboardHost, projectId: this.options.projectId, taskId: task.id, }); const content = [ - `**${identifier} needs plan approval**`, + isMergePolicyBlock ? `**${identifier} pull-request merge is blocked**` : `**${identifier} needs plan approval**`, "", reasonLine, ...(link ? ["", `[Open ${task.id}](${link})`] : []), @@ -506,7 +515,7 @@ export class NotificationService { content, metadata: { taskId: task.id, awaitingApprovalReason: reason }, }; - await messageStore.sendMessageOnce(input, `plan-approval:${task.id}`); + await messageStore.sendMessageOnce(input, `${isMergePolicyBlock ? "merge-policy-block" : "plan-approval"}:${task.id}`); } catch (error) { schedulerLog.log( `[notify] ${task.id} awaiting-approval mailbox message failed: ${error instanceof Error ? error.message : String(error)}`, diff --git a/packages/engine/src/notification/ntfy-provider.ts b/packages/engine/src/notification/ntfy-provider.ts index 671dd008e6..668886e423 100644 --- a/packages/engine/src/notification/ntfy-provider.ts +++ b/packages/engine/src/notification/ntfy-provider.ts @@ -228,15 +228,20 @@ export class NtfyNotificationProvider implements NotificationProvider { "awaiting-approval": { title: payload.metadata?.awaitingApprovalReason === "plan-review-replan-cap" ? `Plan Review cap reached for ${taskId}` - : `Plan needs approval for ${taskId}`, + : payload.metadata?.awaitingApprovalReason === "merge-blocked-by-policy" + ? `Pull-request policy block for ${taskId}` + : `Plan needs approval for ${taskId}`, /* - FNXC:PlanReviewReplan 2026-07-15-11:09: - Replan-cap escalations must say Plan Review failed to converge so the push is - actionable, not a generic "needs approval" ping. + FNXC:PullRequestMerge 2026-08-09-05:07: + Policy holds reuse awaiting-approval's delivery channel but require a + merge-specific instruction; calling them plan approvals sends operators + to the wrong remediation surface. */ message: payload.metadata?.awaitingApprovalReason === "plan-review-replan-cap" ? `Task "${identifier}" needs approval because Plan Review requested revisions repeatedly without converging. Approve the current plan or reject to regenerate.` - : `Task "${identifier}" needs your approval before implementation can start`, + : payload.metadata?.awaitingApprovalReason === "merge-blocked-by-policy" + ? `Task "${identifier}" has a pull request blocked by repository policy. Resolve the policy requirement, then retry the merge.` + : `Task "${identifier}" needs your approval before implementation can start`, priority: "high", }, "awaiting-user-review": { diff --git a/packages/engine/src/notification/webhook-provider.ts b/packages/engine/src/notification/webhook-provider.ts index 7e104eb985..ee4c7956ea 100644 --- a/packages/engine/src/notification/webhook-provider.ts +++ b/packages/engine/src/notification/webhook-provider.ts @@ -166,12 +166,15 @@ export class WebhookNotificationProvider implements NotificationProvider { return `Task "${identifier}" has failed and needs attention`; case "awaiting-approval": /* - FNXC:PlanReviewReplan 2026-07-15-11:09: - Mirror ntfy: replan-cap holds must state that Plan Review did not converge. + FNXC:PullRequestMerge 2026-08-09-05:07: + Mirror the mailbox and ntfy wording: a PR policy hold is actionable by + resolving repository policy and retrying merge, not approving a plan. */ return payload.metadata?.awaitingApprovalReason === "plan-review-replan-cap" ? `Task "${identifier}" needs approval because Plan Review requested revisions repeatedly without converging. Approve the current plan or reject to regenerate.` - : `Task "${identifier}" needs your approval before implementation can start`; + : payload.metadata?.awaitingApprovalReason === "merge-blocked-by-policy" + ? `Task "${identifier}" has a pull request blocked by repository policy. Resolve the policy requirement, then retry the merge.` + : `Task "${identifier}" needs your approval before implementation can start`; case "awaiting-user-review": return `Task "${identifier}" needs human review before it can proceed`; case "planning-awaiting-input": diff --git a/packages/engine/src/project-engine.ts b/packages/engine/src/project-engine.ts index 728ba5dac6..90e050c99e 100644 --- a/packages/engine/src/project-engine.ts +++ b/packages/engine/src/project-engine.ts @@ -46,6 +46,7 @@ import { resolveWipTargetForTask, resolveReboundTargetForTask, REVIEW_ELIGIBLE_SENTINEL_COLUMN, clearMergeConfirmedTransientStatus, + classifyGhError, } from "@fusion/core"; import { assemblePlannerOverseerRuntimeSnapshot } from "./overseer/planner-overseer-runtime-snapshot.js"; import { resolveIntegrationBranch } from "./merge/integration-branch.js"; @@ -136,6 +137,32 @@ const execFileAsync = promisify(execFile); */ const MERGE_HANDOFF_GRACE_MS = 300; +const PR_MERGE_RETRY_BACKOFF_BASE_MS = 5_000; + +/** + * Derive the PR retry not-before instant from the atomic task update timestamp. + * `customFields` is user/workflow-owned and validates unknown keys, so it cannot + * safely carry engine lifecycle metadata. + */ +function getPrMergeRetryNotBefore(task: { mergeRetries?: number | null; updatedAt?: string | null }): number | null { + const retries = task.mergeRetries ?? 0; + const updatedAt = task.updatedAt ? Date.parse(task.updatedAt) : NaN; + if (!Number.isInteger(retries) || retries <= 0 || !Number.isFinite(updatedAt)) return null; + const delayMs = PR_MERGE_RETRY_BACKOFF_BASE_MS * Math.pow(2, retries - 1); + const notBefore = updatedAt + delayMs; + return Number.isFinite(notBefore) ? notBefore : null; +} + +/** + * FNXC:PullRequestMerge 2026-08-09-05:07: + * GitHub's structured network, timeout, and rate-limit outcomes are transport + * failures even when their human-readable message has no legacy transient token. + * They must use the fenced transient budget, not the PR retry budget. + */ +function isStructuredTransientGhOutcome(code: unknown): boolean { + return code === "network" || code === "timeout" || code === "rate-limited"; +} + /* FNXC:MergerUnification 2026-06-21-19:05: Master-plan U0 made `runAiMerge` the SOLE merge path; `merger.mode` is now inert @@ -489,6 +516,8 @@ export class ProjectEngine { private mergeBodyInFlight: Promise | null = null; private mergeAbortController: AbortController | null = null; private mergeRetryTimer: ReturnType | null = null; + /** One durable PR-retry wake per task; admission can safely re-request it after restart/races. */ + private readonly prMergeRetryTimers = new Map>(); private autostashSweepTimer: ReturnType | null = null; private mergeActiveReconcileTimer: ReturnType | null = null; @@ -1341,6 +1370,8 @@ export class ProjectEngine { clearTimeout(this.mergeRetryTimer); this.mergeRetryTimer = null; } + for (const timer of this.prMergeRetryTimers.values()) clearTimeout(timer); + this.prMergeRetryTimers.clear(); if (this.autostashSweepTimer) { clearTimeout(this.autostashSweepTimer); this.autostashSweepTimer = null; @@ -2342,6 +2373,26 @@ export class ProjectEngine { throw new Error(`Merge request for ${taskId} aborted`); } + const store = this.runtime.getTaskStore(); + // FNXC:PullRequestMerge 2026-08-09-04:18: A manual merge supersedes a pending + // PR backoff wakeup. Cancel it before admitting the manual attempt so a later + // stale callback cannot repeat a merge that the operator already completed. + this.clearPrMergeRetryTimer(taskId); + const existing = await store.getTask(taskId); + if (existing?.status === "awaiting-approval") { + /* + FNXC:PullRequestMerge 2026-08-09-02:39: + A branch-policy hold is operator-resumable only through this manual merge + entry point. Clear its durable wait marker before enqueueing so the normal + single-flight pump performs one fresh PR merge without consuming either retry budget. + */ + await store.updateTask(taskId, { + status: null, + error: null, + awaitingApprovalReason: null, + }); + } + return new Promise((resolve, reject) => { let settled = false; let abort: () => void = () => undefined; @@ -2754,7 +2805,7 @@ export class ProjectEngine { log?: Array<{ action?: string }>; updatedAt?: string | null; mergeDetails?: { mergeConfirmed?: boolean } | null; - }, maxAutoMergeRetries: number, isReviewColumn?: boolean): boolean { + }, maxAutoMergeRetries: number, isReviewColumn?: boolean, enforcePrRetryBackoff = false): boolean { // Merge-confirmed tasks use the fast-path finalizer, which applies blocker // checks after clearing transient status/error state. Once that path parks // a blocked task as failed, skip future auto-merge retries. @@ -2765,7 +2816,18 @@ export class ProjectEngine { // Terminal failure: don't let the cooldown sweep re-attempt a merge that // already gave up (verification cap, conflict-bounce cap, or non-conflict // error). The task is parked for human/follow-up intervention. - if (task.status === "failed") return false; + if (task.status === "failed" || task.status === "awaiting-approval" || task.status === "awaiting-user-review") return false; + /* + FNXC:AutoMergeRetries 2026-08-09-03:02: + Retry backoff must be enforced at this shared admission point, not only by + its timer: periodic sweeps, duplicate enqueue calls, and a restarted engine + all reach canMergeTask. The retry update's timestamp is the durable anchor; + absent, invalid, or elapsed timestamps fail open for legacy rows. + */ + if (enforcePrRetryBackoff) { + const notBefore = getPrMergeRetryNotBefore(task); + if (notBefore !== null && notBefore > Date.now()) return false; + } return ( (task.mergeRetries ?? 0) < maxAutoMergeRetries || this.hasAutoHealableVerificationBufferFailure(task, maxAutoMergeRetries, isReviewColumn) || @@ -2873,6 +2935,31 @@ export class ProjectEngine { }); } + /* + FNXC:PullRequestMerge 2026-08-09-04:05: + PR retry backoff is persisted through the retry update timestamp, but its wakeup + is process-local. Keep one shutdown-safe timer per task and allow an admission + rejection to restore that wakeup, so a restart or duplicate enqueue cannot drop + a real retry until the cooldown sweep happens to notice it. + */ + private clearPrMergeRetryTimer(taskId: string): void { + const timer = this.prMergeRetryTimers.get(taskId); + if (!timer) return; + clearTimeout(timer); + this.prMergeRetryTimers.delete(taskId); + } + + private schedulePrMergeRetry(taskId: string, notBefore: number): void { + if (this.shuttingDown || this.prMergeRetryTimers.has(taskId)) return; + const delayMs = Math.max(0, notBefore - Date.now()); + const timer = setTimeout(() => { + this.prMergeRetryTimers.delete(taskId); + if (!this.shuttingDown) this.internalEnqueueMerge(taskId); + }, delayMs); + timer.unref?.(); + this.prMergeRetryTimers.set(taskId, timer); + } + private internalEnqueueMerge(taskId: string): boolean { if (this.shuttingDown || !this.started) return false; if (this.capacityDeferredMergeTaskIds.has(taskId)) return false; @@ -3004,6 +3091,7 @@ export class ProjectEngine { private async enqueueEligibleInReviewTasks(tasks: readonly Task[], settings: Pick): Promise { const maxAutoMergeRetries = resolveMaxAutoMergeRetries(settings); + const enforcePrRetryBackoff = (this.options.getMergeStrategy?.(settings as Settings) ?? "direct") === "pull-request"; // FNXC:PostgresCutover 2026-07-10: allowInReviewMergeProcessing awaits the // async getBranchGroup read on the PG branch, so eligibility resolves per // task before the sync priority sort. @@ -3031,6 +3119,7 @@ export class ProjectEngine { t as any, maxAutoMergeRetries, reviewLane === undefined ? undefined : t.column === reviewLane, + enforcePrRetryBackoff, ); }) as Task[]; const allowFlags = await Promise.all(candidates.map((t) => this.allowInReviewMergeProcessing(t, settings, this.runtime.getTaskStore()))); @@ -3470,11 +3559,21 @@ export class ProjectEngine { "not in review", which would disable auto-heal outright. */ const mergeLoopReviewLane = (await resolveTaskLifecycleColumns(store, taskId).catch(() => undefined))?.review; + const pullRequestMerge = (this.options.getMergeStrategy?.(settings) ?? "direct") === "pull-request"; if (!this.canMergeTask( task as any, maxAutoMergeRetries, mergeLoopReviewLane === undefined ? undefined : task.column === mergeLoopReviewLane, + pullRequestMerge, )) { + // A queued retry can be rejected after an engine restart or a racing + // task update. Reinstall the single-flight wake instead of dropping it. + if (pullRequestMerge) { + const notBefore = getPrMergeRetryNotBefore(task); + if (notBefore !== null && notBefore > Date.now()) { + this.schedulePrMergeRetry(taskId, notBefore); + } + } continue; } @@ -4133,6 +4232,19 @@ export class ProjectEngine { mergeTargetBranch: mergedTask.mergeDetails?.mergeTargetBranch, } as MergeResult); } + /* + FNXC:PullRequestMerge 2026-08-09-03:32: + A successful PR merge ends the retry episode. Reset both independent + counters so persisted completed work never carries stale retry exhaustion + or a derived backoff anchor into a later recovery/finalization read. + */ + this.clearPrMergeRetryTimer(taskId); + if (mergedTask && ((mergedTask.mergeRetries ?? 0) > 0 || (mergedTask.mergeTransientRetryCount ?? 0) > 0)) { + await store.updateTask(taskId, { + mergeRetries: 0, + mergeTransientRetryCount: 0, + }); + } await attemptBranchGroupPromotion(mergedTask, this.mergeAbortController?.signal); } else if (result === "waiting") { runtimeLog.log(`${hasManualResolver ? "Manual" : "Auto"}-merge PR waiting: ${taskId}`); @@ -4294,7 +4406,7 @@ export class ProjectEngine { // Reset retries on success const latestTask = await store.getTask(taskId).catch(() => null); - if (latestTask?.mergeRetries && latestTask.mergeRetries > 0) { + if (latestTask && (latestTask.mergeRetries ?? 0) > 0) { await store.updateTask(taskId, { mergeRetries: 0 }); } // FNXC:Workspace 2026-06-22-05:10 (Phase C review B4): clear the in-memory busy @@ -4527,8 +4639,10 @@ export class ProjectEngine { ); }); - // If this was a manual merge, reject the promise and skip auto-retry logic - if (hasManualResolver) { + // A manual policy-resume attempt must re-park through the same durable + // handoff path; other manual merge failures still reject their caller. + const isPolicyBlock = (err as { code?: unknown })?.code === "merge-blocked-by-policy"; + if (hasManualResolver && !isPolicyBlock) { this.rejectMergeResolvers(taskId, err instanceof Error ? err : new Error(errorMsg)); continue; } @@ -4942,16 +5056,55 @@ export class ProjectEngine { } } } else { - // Non-direct merge strategy (e.g. pull-request) errored — park as - // failed so the cooldown sweep stops re-attempting silently. + /* + FNXC:AutoMergeRetries 2026-08-09-03:02: + PR failures have four mutually-exclusive dispositions: policy holds wait for + an operator, transient ownership uses its separate counter, non-retryable gh + outcomes park honestly, and only retryable failures consume mergeRetries. + The atomic retry update timestamp is the durable backoff anchor because sweeps + and restarts can outrun a timer; workflow custom fields reject engine metadata. + */ try { - if (await this.maybeRetryTransientMerge(store, taskId, taskOnErr, errorMsg)) { + const classified = classifyGhError(err); + const structuredCode = (err as { code?: unknown })?.code; + const isStructuredNonRetryable = structuredCode === "merge-conflict" + || structuredCode === "validation" + || structuredCode === "permission" + || structuredCode === "not-found" + || structuredCode === "not-installed"; + const diagnosis = isPolicyBlock + ? { ...classified, code: "merge-blocked-by-policy" as const, retryable: false } + : isStructuredNonRetryable + ? { ...classified, code: structuredCode, retryable: false } + : classified; + if (diagnosis.code === "merge-blocked-by-policy") { + this.clearPrMergeRetryTimer(taskId); + await store.updateTask(taskId, { + status: "awaiting-approval", + error: diagnosis.message, + awaitingApprovalReason: "merge-blocked-by-policy", + }); + await store.logEntry(taskId, `Pull-request merge blocked by policy; awaiting operator resume: ${diagnosis.message}`, "MergePolicyBlocked"); continue; } - if (this.isTransientMergeRetryExhausted(taskOnErr, errorMsg)) { + const structuredTransient = isStructuredTransientGhOutcome(diagnosis.code); + if (await this.maybeRetryTransientMerge(store, taskId, taskOnErr, errorMsg, structuredTransient)) { + continue; + } + /* + FNXC:PullRequestMerge 2026-08-09-05:07: + A transient failure has a separately owned retry budget. Structured + GitHub transport results join the legacy text classifier here, so + neither form can fall through and consume mergeRetries. + */ + if (this.isTransientMergeRetryExhausted(taskOnErr, errorMsg, structuredTransient)) { + this.clearPrMergeRetryTimer(taskId); + // FNXC:PullRequestMerge 2026-08-09-04:18: The shadow merge-request + // contract owns transient exhaustion independently of task status. + // Keep its retrying -> exhausted transition intact; only non-shadow + // tasks park failed, and neither path spends mergeRetries. const settings = await store.getSettings().catch(() => null); - const useMergeRequestContract = settings?.mergeRequestContractShadowEnabled === true; - if (useMergeRequestContract) { + if (settings?.mergeRequestContractShadowEnabled === true) { const record = await store.getMergeRequestRecordAsync(taskId); if (record && record.state !== "exhausted" && record.state !== "cancelled" && record.state !== "succeeded") { if (record.state === "running") { @@ -4968,24 +5121,52 @@ export class ProjectEngine { }); } } - await store.logEntry( - taskId, - `Auto-merge transient retries exhausted (${ProjectEngine.MAX_AUTO_MERGE_TRANSIENT_RETRIES}/${ProjectEngine.MAX_AUTO_MERGE_TRANSIENT_RETRIES}); marked merge request exhausted without column rebound: ${errorMsg}`, - "MergeTransientRetryExhausted", - ); + await store.logEntry(taskId, `Pull-request transient retries exhausted (${ProjectEngine.MAX_AUTO_MERGE_TRANSIENT_RETRIES}/${ProjectEngine.MAX_AUTO_MERGE_TRANSIENT_RETRIES}); marked merge request exhausted without consuming merge retries: ${errorMsg}`, "MergeTransientRetryExhausted"); continue; } - await store.logEntry( - taskId, - `Auto-merge transient retries exhausted (${ProjectEngine.MAX_AUTO_MERGE_TRANSIENT_RETRIES}/${ProjectEngine.MAX_AUTO_MERGE_TRANSIENT_RETRIES}); parking task as failed: ${errorMsg}`, - "MergeTransientRetryExhausted", - ); + await store.updateTask(taskId, { + status: "failed", + error: errorMsg, + }); + await store.logEntry(taskId, `Pull-request transient retries exhausted (${ProjectEngine.MAX_AUTO_MERGE_TRANSIENT_RETRIES}/${ProjectEngine.MAX_AUTO_MERGE_TRANSIENT_RETRIES}); task parked without consuming merge retries: ${errorMsg}`, "MergeTransientRetryExhausted"); + continue; } + if (!diagnosis.retryable) { + this.clearPrMergeRetryTimer(taskId); + await store.updateTask(taskId, { + status: "failed", + error: diagnosis.message, + }); + await store.logEntry(taskId, `Pull-request merge failed without retry (${diagnosis.code}): ${diagnosis.message}`, "MergeNonRetryableFailure"); + continue; + } + const currentRetries = taskOnErr?.mergeRetries ?? 0; + const nextRetries = currentRetries + 1; + if (nextRetries >= maxAutoMergeRetriesOnErr) { + this.clearPrMergeRetryTimer(taskId); + await store.updateTask(taskId, { + status: "failed", + mergeRetries: nextRetries, + error: errorMsg, + }); + await store.logEntry(taskId, `Pull-request merge retries exhausted after ${nextRetries}/${maxAutoMergeRetriesOnErr} actual failures: ${errorMsg}`, "MergeRetriesExhausted"); + continue; + } + const delayMs = PR_MERGE_RETRY_BACKOFF_BASE_MS * Math.pow(2, currentRetries); + /* + FNXC:AutoMergeRetries 2026-08-09-03:11: + The retry log must precede the atomic retry patch. `updatedAt` on + that patch is the durable not-before anchor; logging afterwards can + advance it past the timer deadline and make the queue reject its own + scheduled retry as still early. + */ + await store.logEntry(taskId, `Pull-request merge retry ${nextRetries}/${maxAutoMergeRetriesOnErr} scheduled in ${delayMs / 1000}s: ${errorMsg}`, "MergeRetry"); await store.updateTask(taskId, { - status: "failed", - mergeRetries: maxAutoMergeRetriesOnErr, - error: errorMsg, + mergeRetries: nextRetries, + status: null, + error: null, }); + this.schedulePrMergeRetry(taskId, Date.now() + delayMs); } catch (recoveryErr) { runtimeLog.error( `Auto-merge: failed to update ${taskId} after merge strategy error: ${recoveryErr instanceof Error ? recoveryErr.message : String(recoveryErr)}`, @@ -5032,8 +5213,12 @@ export class ProjectEngine { } } - private isTransientMergeRetryExhausted(task: Task | null, errorMsg: string): boolean { - if (!task || (!isTransientError(errorMsg) && classifyTransientMergeError(errorMsg) === null)) { + private isTransientMergeRetryExhausted( + task: Task | null, + errorMsg: string, + structuredTransient = false, + ): boolean { + if (!task || (!structuredTransient && !isTransientError(errorMsg) && classifyTransientMergeError(errorMsg) === null)) { return false; } const current = task.mergeTransientRetryCount ?? 0; @@ -5045,8 +5230,9 @@ export class ProjectEngine { taskId: string, taskOnErr: Task | null, errorMsg: string, + structuredTransient = false, ): Promise { - if (!taskOnErr || (!isTransientError(errorMsg) && classifyTransientMergeError(errorMsg) === null)) { + if (!taskOnErr || (!structuredTransient && !isTransientError(errorMsg) && classifyTransientMergeError(errorMsg) === null)) { return false; }