FN-8837: fix pull-request merge retry recovery
Make pull-request merge failures retry safely and report actionable terminal states. - Persist exponential retry backoff and schedule durable retry wakeups. - Classify policy, transient, and non-retryable GitHub merge failures correctly. - Clear retry state on successful merges and cover notifications and recovery paths. Files changed: .changeset/fn-8837-pr-merge-retries.md | 7 + docs/architecture.md | 1 + .../task-update-awaiting-approval-reason.test.ts | 9 + packages/core/src/store.ts | 2 +- packages/core/src/task-store/task-update.ts | 6 +- packages/core/src/types/task/task-core.ts | 7 +- .../src/__tests__/merge-error-recovery.test.ts | 440 ++++++++++++++++++++- .../engine/src/__tests__/ntfy-provider.test.ts | 14 + .../engine/src/__tests__/webhook-provider.test.ts | 15 + .../__tests__/notification-service.test.ts | 48 +++ .../src/notification/notification-service.ts | 15 +- packages/engine/src/notification/ntfy-provider.ts | 15 +- .../engine/src/notification/webhook-provider.ts | 9 +- packages/engine/src/project-engine.ts | 240 +++++++++-- 14 files changed, 781 insertions(+), 47 deletions(-) Fusion-Task-Id: FN-8837 Fusion-Task-Lineage: a39b3491-d5cb-4509-b450-3e5671e0ba2d Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
7
.changeset/fn-8837-pr-merge-retries.md
Normal file
7
.changeset/fn-8837-pr-merge-retries.md
Normal file
@@ -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.
|
||||
@@ -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.
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -1426,7 +1426,7 @@ 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; 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<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; awaitingApprovalReason?: import("./types.js").Task["awaitingApprovalReason"] | 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; 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<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);
|
||||
}
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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<string, unknown>;
|
||||
updatedAt: string;
|
||||
log: Array<{ action?: string }>;
|
||||
};
|
||||
@@ -108,6 +110,7 @@ type MockTaskStore = {
|
||||
recordRunAuditEvent: ReturnType<typeof vi.fn>;
|
||||
on: ReturnType<typeof vi.fn>;
|
||||
off: ReturnType<typeof vi.fn>;
|
||||
emit: ReturnType<typeof vi.fn>;
|
||||
};
|
||||
|
||||
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<MockTask>) => {
|
||||
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<string, unknown>) => {
|
||||
if (patch.customFields !== undefined) {
|
||||
const validation = validateCustomFieldPatch([], patch.customFields as Record<string, unknown>);
|
||||
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<MockTask>) => {
|
||||
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<MockTask>) => {
|
||||
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<Settings, "autoMerge" | "maxAutoMergeRetries">) => Promise<number>;
|
||||
};
|
||||
|
||||
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<MockTask>) => {
|
||||
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<Settings, "autoMerge" | "maxAutoMergeRetries">) => Promise<number>;
|
||||
};
|
||||
|
||||
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";
|
||||
|
||||
@@ -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");
|
||||
|
||||
@@ -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"],
|
||||
|
||||
@@ -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<string>();
|
||||
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" }));
|
||||
|
||||
@@ -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)}`,
|
||||
|
||||
@@ -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": {
|
||||
|
||||
@@ -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":
|
||||
|
||||
@@ -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<unknown> | null = null;
|
||||
private mergeAbortController: AbortController | null = null;
|
||||
private mergeRetryTimer: ReturnType<typeof setTimeout> | null = null;
|
||||
/** One durable PR-retry wake per task; admission can safely re-request it after restart/races. */
|
||||
private readonly prMergeRetryTimers = new Map<string, ReturnType<typeof setTimeout>>();
|
||||
private autostashSweepTimer: ReturnType<typeof setTimeout> | null = null;
|
||||
private mergeActiveReconcileTimer: ReturnType<typeof setInterval> | 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<MergeResult>((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<Settings, "autoMerge" | "maxAutoMergeRetries">): Promise<number> {
|
||||
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<boolean> {
|
||||
if (!taskOnErr || (!isTransientError(errorMsg) && classifyTransientMergeError(errorMsg) === null)) {
|
||||
if (!taskOnErr || (!structuredTransient && !isTransientError(errorMsg) && classifyTransientMergeError(errorMsg) === null)) {
|
||||
return false;
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user