FN-8958: fence orphaned merge-body writes
Prevent cancelled merge generations from writing stale task state. - Add a signal-aware merge write fence with orphan audit reporting. - Fence merge finalization, post-push metadata, and recovery-branch task logs. - Cover durable write callsites and cancellation behavior with tests and guidance. Files changed: .changeset/fn-8958-orphan-merge-write-fence.md | 7 + AGENTS.md | 1 + .../reliability/orphan-merge-body-write-fence.md | 68 ++ .../__tests__/_merge-durable-write-callsites.ts | 4 + .../merge-orphan-durable-write-inventory.json | 982 +++++++++++---------- .../merge-orphan-body-durable-writes.test.ts | 38 +- .../engine/src/__tests__/merge-write-fence.test.ts | 39 + .../engine/src/merge/auto-merge-finalization.ts | 9 + packages/engine/src/merge/merge-write-fence.ts | 92 ++ packages/engine/src/merge/merger-ai.ts | 173 ++-- 10 files changed, 875 insertions(+), 538 deletions(-) Fusion-Task-Id: FN-8958 Fusion-Task-Lineage: 5f398c44-4320-4f0c-be15-707184f66aa8 Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
7
.changeset/fn-8958-orphan-merge-write-fence.md
Normal file
7
.changeset/fn-8958-orphan-merge-write-fence.md
Normal file
@@ -0,0 +1,7 @@
|
||||
---
|
||||
"@runfusion/fusion": patch
|
||||
---
|
||||
|
||||
summary: Prevent canceled AI merge bodies from overwriting successor merge state.
|
||||
category: fix
|
||||
dev: Adds `merge-write-fence` with per-mutation ownership checks, optional squash-landing signals and ref-advance checkpoints. Aborts rethrow as `MergeAbortedError`; the injected `merge:orphan-write-fenced` audit emits once at first interaction with an emit-time suppression count.
|
||||
@@ -276,6 +276,7 @@ Scoped exception (FN-5819/FN-8823): while project auto-merge is On, shared-branc
|
||||
`testMode?: boolean` is now available in both project and global settings. If project `testMode === true` (or the resolved default provider is `"mock"` at any tier), every AI lane is forced to `mock/scripted`, overriding per-task and per-lane model selections. The dashboard exposes this via the Settings Modal "Enable test mode" toggle and a persistent "Test mode — no real AI calls" banner.
|
||||
|
||||
### Run Audit
|
||||
- FN-8958: `merge:orphan-write-fenced` is emitted once per orphan merge body at its fence's first interaction. Metadata is ids/counts/outcomes-only: `{ taskId, category, interaction, suppressedCount }`; `suppressedCount` is the emit-time count (`1` for `interaction:"suppressed"`, `0` for `interaction:"rejected"`), never a cumulative body total.
|
||||
|
||||
- Store-open provenance: every `TaskStore.init()` emits `store:open` with ids/paths-only metadata (`pid`, `ppid`, `execPath`, `entry`, `cwd`, `nodeVersion`). Purpose: attribute shared-DB mutations to the process that opened the store (the FN-7910 Ideas-evacuation writer was unidentifiable without it). Tests reading unfiltered `runAuditEvents` must filter out `store:open` rather than assert exact counts.
|
||||
- FN-8948: `mission:reconcile-pass` records a bounded automatic reconciliation result. Metadata contains optional mission ID, source enum, and scan/write/skip/conflict/failure counters only; never roadmap prose, titles, reasons, or secrets.
|
||||
|
||||
68
docs/solutions/reliability/orphan-merge-body-write-fence.md
Normal file
68
docs/solutions/reliability/orphan-merge-body-write-fence.md
Normal file
@@ -0,0 +1,68 @@
|
||||
---
|
||||
category: reliability
|
||||
module: "@fusion/engine"
|
||||
tags: [merge, cancellation, task-store, git]
|
||||
problem_type: race-condition
|
||||
applies_when: "A merge generation can outlive cancellation while a successor owns the same task."
|
||||
---
|
||||
|
||||
# Fence orphan merge-body writes
|
||||
|
||||
An abort signal is write authority for a claimed merge generation. An aborted signal means that body no longer owns the task row; the successor owns a fresh signal. Ownership must be read immediately before each mutation or irreversible action. A function-entry, closure, loop-head, or shared adjacent-writer check is unsound because an abort may arrive during the preceding await.
|
||||
|
||||
`merge-write-fence.ts` supplies one fence per merge body. Diagnostic and rebound writes use `fence.write`, which silently suppresses an orphaned call. Finalization, task completion, commit association, branch-group updates, group-PR synchronization, and cleanup use `assertOwned`, which throws `MergeAbortedError` so the body unwinds instead of leaving a partially finalized successor row. Every catch on the merge path rethrows that error; it is never classified as a transient failure, partial workspace land, rebound result, or group-PR sync failure.
|
||||
|
||||
## Boundary policy mapping
|
||||
|
||||
Each entry is a single action at its own boundary:
|
||||
|
||||
| Boundary/action | Bucket | Guard |
|
||||
| --- | --- | --- |
|
||||
| `runAiMerge` log entry | suppress | `fence.write("log")` |
|
||||
| `runAiMerge` agent log append | suppress | `fence.write("log")` |
|
||||
| `landWorkspaceTask` log entry | suppress | `fence.write("log")` |
|
||||
| `landWorkspaceTask` agent log append | suppress | `fence.write("log")` |
|
||||
| transient merge status update | suppress | `fence.write("lifecycle")` |
|
||||
| single-repo no-commits rebound error update | suppress | `fence.write("lifecycle")` |
|
||||
| single-repo no-commits rebound diagnostic | suppress | `fence.write("log")` |
|
||||
| single-repo no-commits rebound move | suppress | `fence.write("lifecycle")` |
|
||||
| single-repo missing-proof rebound error update | suppress | `fence.write("lifecycle")` |
|
||||
| single-repo missing-proof rebound diagnostic | suppress | `fence.write("log")` |
|
||||
| single-repo missing-proof rebound move | suppress | `fence.write("lifecycle")` |
|
||||
| single-repo overseer-veto rebound error update | suppress | `fence.write("lifecycle")` |
|
||||
| single-repo overseer-veto rebound diagnostic | suppress | `fence.write("log")` |
|
||||
| single-repo overseer-veto rebound move | suppress | `fence.write("lifecycle")` |
|
||||
| workspace missing-proof rebound error update | suppress | `fence.write("lifecycle")` |
|
||||
| workspace missing-proof rebound diagnostic | suppress | `fence.write("log")` |
|
||||
| workspace missing-proof rebound move | suppress | `fence.write("lifecycle")` |
|
||||
| auto-finalization hard-blocker update | throw | `fence.assertOwned()` |
|
||||
| workspace final `mergeDetails` update | throw | `fence.assertOwned()` |
|
||||
| workspace finalization hand-off | throw | `fence.assertOwned()` |
|
||||
| final merge-details/files update | throw | `fence.assertOwned()` |
|
||||
| commit association | throw | `fence.assertOwned()` |
|
||||
| cleanup merge-details update | throw | `fence.assertOwned()` |
|
||||
| branch deletion | throw | `fence.assertOwned()` |
|
||||
| worktree removal | throw | `fence.assertOwned()` |
|
||||
| worktree-clear update | throw | `fence.assertOwned()` |
|
||||
| branch-group member landing | throw | `fence.assertOwned()` |
|
||||
| managed group-PR sync | throw | `fence.assertOwned()` |
|
||||
| auto-finalization task update | throw | `fence.assertOwned()` |
|
||||
| auto-finalization column move | throw | `fence.assertOwned()` |
|
||||
| auto-finalization tail log | throw | `fence.assertOwned()` |
|
||||
| `task:merged` emit | throw | `fence.assertOwned()` |
|
||||
| squash case-B CAS ref advance | throw | `assertMergeGenerationOwned()` |
|
||||
| squash dirty-checkout CAS ref advance | throw | `assertMergeGenerationOwned()` |
|
||||
| squash case-A fast-forward merge | throw | `assertMergeGenerationOwned()` |
|
||||
| push-after-merge ref advance | throw | `assertMergeGenerationOwned()` |
|
||||
| `persistRepoLandedSha` | deliberately unfenced | Records an already-completed advance; suppressing it can recreate double-squash landing. |
|
||||
| append-only run audit | deliberately unfenced | Forensic data does not mutate task ownership; aborts are excluded from sync-failure classification. |
|
||||
|
||||
The squash helper receives an optional signal. Existing direct callers without a signal retain prior behavior. Its loop-head cancellation checks remain cheap early-outs, but each actual ref advance has its own immediate ownership check.
|
||||
|
||||
A non-abort failure in an already-orphaned rebound may have its next rebound write suppressed; no compensating mutation is made because the successor owns the row. Similarly, an abort at branch/worktree cleanup may leave cleanup for the successor or existing self-healing. Completion announcement remains last, preventing an abandoned body from announcing a merge.
|
||||
|
||||
## Observability
|
||||
|
||||
The fence gets an injected `recordRunAuditEvent` recorder; no ambient or module-global state is used. It emits one `merge:orphan-write-fenced` event on the fence's first interaction with `{ taskId, category, interaction, suppressedCount }`. The emit-time count is `1` for a suppression-first interaction and `0` for a rejection-first interaction. Later suppressions only increase the in-process counter: no end-of-body cumulative event is attempted because unwinding cannot reliably flush one. The rejection-first case is therefore unit-tested at the fence level.
|
||||
|
||||
This completes the progression from FN-8912's transient-status fence and FN-8923's call-site inventory. The regression suite drives real single-repository and workspace bodies so an orphan rejects with `MergeAbortedError` while successor-owned task state remains independent.
|
||||
@@ -67,6 +67,10 @@ const STORE_METHOD_CLASSIFICATION: Record<string, Omit<SurfaceClassification, "m
|
||||
addSteeringComment: { kind: "writer", reason: "persists or mutates TaskStore state" },
|
||||
addTaskComment: { kind: "writer", reason: "persists or mutates TaskStore state" },
|
||||
appendAgentLogBatch: { kind: "writer", reason: "persists or mutates TaskStore state" },
|
||||
// FNXC:MergeReliability 2026-08-11-21:39: Wedge pending markers persist task-scoped
|
||||
// notification state, so the durable-write ratchet must classify them as writers.
|
||||
clearTaskWedgeNotificationPending: { kind: "writer", reason: "persists or mutates TaskStore state" },
|
||||
markTaskWedgeNotificationPending: { kind: "writer", reason: "persists or mutates TaskStore state" },
|
||||
appendCurrentPlanEvidence: { kind: "writer", reason: "persists or mutates TaskStore state" },
|
||||
appendSpecDriftReport: { kind: "writer", reason: "persists or mutates TaskStore state" },
|
||||
appendSpecDriftReportWhilePlanningLocked: { kind: "writer", reason: "persists or mutates TaskStore state" },
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -162,12 +162,12 @@ describe("FN-8923 orphan merge-body durable writes", () => {
|
||||
expect(store.getTask).toHaveBeenCalledOnce();
|
||||
});
|
||||
|
||||
it("aborts from the real mid-review seam and records orphan finalization writes", async () => {
|
||||
it("aborts from the real mid-review seam before orphan finalization writes", async () => {
|
||||
const directory = createMergeRepo();
|
||||
const controller = new AbortController();
|
||||
const { store, records } = createRecordingStore(controller, { sharedGroup: true });
|
||||
let recordsAtAbort = -1;
|
||||
const result = await runAiMerge(store as never, directory, "FN-8923", { manual: true, signal: controller.signal }, {
|
||||
await expect(runAiMerge(store as never, directory, "FN-8923", { manual: true, signal: controller.signal }, {
|
||||
mergeAgent: async (cwd) => { git(cwd, "merge --squash fusion/FN-8923 && git commit -q -m squash"); },
|
||||
// This is the abort-at-boundary fixture: review completed, the squash exists, and the
|
||||
// production body is about to enter its landing/finalization stretch.
|
||||
@@ -176,23 +176,22 @@ describe("FN-8923 orphan merge-body durable writes", () => {
|
||||
recordsAtAbort = records.length;
|
||||
return "REVIEW_VERDICT: approve";
|
||||
},
|
||||
});
|
||||
})).rejects.toMatchObject({ name: "MergeAbortedError" });
|
||||
const orphanWrites = records.slice(recordsAtAbort);
|
||||
expect(recordsAtAbort).toBeGreaterThanOrEqual(0);
|
||||
expect(result.merged).toBe(true);
|
||||
// Positive writes are their own execution proof and demonstrate that this is a body that
|
||||
// outlived its abort rather than a pre-aborted fixture that stopped at an entry checkpoint.
|
||||
// FN-8923 follow-up: FN-8958 flips packages/engine/src/merge/merger-ai.ts::finalizeMerged::store.updateTask::#1.
|
||||
expect(orphanWrites.some(({ writer, args }) => writer === "updateTask" && (args[1] as Record<string, unknown>).mergeDetails)).toBe(true);
|
||||
expect(orphanWrites.some(({ writer, args }) => writer === "updateTask" && (args[1] as Record<string, unknown>).mergeDetails)).toBe(false);
|
||||
// FN-8923 follow-up: FN-8958 flips packages/engine/src/merge/auto-merge-finalization.ts::finalizeProvenAutoMergeTask>moved::store.moveTask::#1.
|
||||
expect(orphanWrites.some(({ writer }) => writer === "moveTask")).toBe(true);
|
||||
expect(orphanWrites.some(({ writer }) => writer === "moveTask")).toBe(false);
|
||||
// FN-8923 follow-up: FN-8958 flips packages/engine/src/merge/merger-ai.ts::finalizeTask::store.emit::#1.
|
||||
expect(orphanWrites.some(({ writer, args }) => writer === "emit" && args[0] === "task:merged")).toBe(true);
|
||||
expect(orphanWrites.some(({ writer, args }) => writer === "emit" && args[0] === "task:merged")).toBe(false);
|
||||
// FN-8923 follow-up: FN-8958 flips packages/engine/src/merge/merger-ai.ts::finalizeMerged::store.recordBranchGroupMemberLanded::#1.
|
||||
expect(orphanWrites.some(({ writer }) => writer === "recordBranchGroupMemberLanded")).toBe(true);
|
||||
expect(orphanWrites.some(({ writer }) => writer === "recordBranchGroupMemberLanded")).toBe(false);
|
||||
// FN-8923 follow-up: FN-8958 flips packages/engine/src/merge/merger-ai.ts::runAiMerge>log::store.logEntry::#1.
|
||||
expect(orphanWrites.some(({ writer }) => writer === "logEntry")).toBe(true);
|
||||
expect(orphanWrites.some(({ writer }) => writer === "appendAgentLog")).toBe(true);
|
||||
expect(orphanWrites.some(({ writer }) => writer === "logEntry")).toBe(false);
|
||||
expect(orphanWrites.some(({ writer }) => writer === "appendAgentLog")).toBe(false);
|
||||
});
|
||||
|
||||
it("attempts the real landOneRepo pre-land boundary", async () => {
|
||||
@@ -213,14 +212,13 @@ describe("FN-8923 orphan merge-body durable writes", () => {
|
||||
// `advanceIntegrationBranchRef` emits this durable audit only after its CAS ref write.
|
||||
onIntegrationRefAdvance: () => controller.abort("orphaned after integration ref advance"),
|
||||
});
|
||||
const result = await runAiMerge(store as never, directory, "FN-8923", { manual: true, signal: controller.signal }, {
|
||||
await expect(runAiMerge(store as never, directory, "FN-8923", { manual: true, signal: controller.signal }, {
|
||||
mergeAgent: async (cwd) => { git(cwd, "merge --squash fusion/FN-8923 && git commit -q -m squash"); },
|
||||
reviewAgent: async () => "REVIEW_VERDICT: approve",
|
||||
});
|
||||
})).rejects.toMatchObject({ name: "MergeAbortedError" });
|
||||
expect(controller.signal.aborted).toBe(true);
|
||||
expect(result.merged).toBe(true);
|
||||
// The post-boundary mergeDetails write proves the production body passed the ref-advance seam.
|
||||
expect(records.some(({ writer, args }) => writer === "updateTask" && (args[1] as Record<string, unknown>).mergeDetails)).toBe(true);
|
||||
expect(records.some(({ writer, args }) => writer === "updateTask" && (args[1] as Record<string, unknown>).mergeDetails)).toBe(false);
|
||||
});
|
||||
|
||||
it("aborts from the mergeDetails persistence seam before the private finalization tail", async () => {
|
||||
@@ -230,14 +228,13 @@ describe("FN-8923 orphan merge-body durable writes", () => {
|
||||
sharedGroup: true,
|
||||
onMergeDetailsPersist: () => controller.abort("orphaned after mergeDetails persistence"),
|
||||
});
|
||||
const result = await runAiMerge(store as never, directory, "FN-8923", { manual: true, signal: controller.signal }, {
|
||||
await expect(runAiMerge(store as never, directory, "FN-8923", { manual: true, signal: controller.signal }, {
|
||||
mergeAgent: async (cwd) => { git(cwd, "merge --squash fusion/FN-8923 && git commit -q -m squash"); },
|
||||
reviewAgent: async () => "REVIEW_VERDICT: approve",
|
||||
});
|
||||
})).rejects.toMatchObject({ name: "MergeAbortedError" });
|
||||
expect(controller.signal.aborted).toBe(true);
|
||||
expect(result.merged).toBe(true);
|
||||
// `task:merged` is emitted after mergeDetails, proving the private tail was reached through runAiMerge.
|
||||
expect(records.some(({ writer, args }) => writer === "emit" && args[0] === "task:merged")).toBe(true);
|
||||
expect(records.some(({ writer, args }) => writer === "emit" && args[0] === "task:merged")).toBe(false);
|
||||
});
|
||||
|
||||
it("aborts during a real second workspace-repository land and partitions successor writes", async () => {
|
||||
@@ -328,13 +325,12 @@ describe("FN-8923 orphan merge-body durable writes", () => {
|
||||
expect(successorIsHeldAtProductionSeam).toBe(true);
|
||||
releaseOrphan();
|
||||
// FNXC:MergeReliability 2026-08-11-00:37: This assertion proves the bodies overlap at a real production seam.
|
||||
const result = await body;
|
||||
await expect(body).rejects.toMatchObject({ name: "MergeAbortedError" });
|
||||
expect(successorIsHeldAtProductionSeam).toBe(true);
|
||||
releaseSuccessor();
|
||||
await successorBody;
|
||||
expect(mergeCount).toBe(2);
|
||||
expect(successorSignal).toBeDefined();
|
||||
expect(result.repos[0]).toMatchObject({ repo: "repo-a", status: "landed" });
|
||||
expect((task.workspaceWorktrees as Record<string, { landedSha?: string }>) ["repo-a"].landedSha).toEqual(expect.any(String));
|
||||
const postBoundary = records.slice(recordCountAtAbort);
|
||||
const orphanWrites = postBoundary.filter((record) => record.generation === "orphan");
|
||||
@@ -342,7 +338,7 @@ describe("FN-8923 orphan merge-body durable writes", () => {
|
||||
// FNXC:MergeReliability 2026-08-11-00:37: Separate facades prevent successor writes from satisfying orphan assertions.
|
||||
expect(successorWrites).toContainEqual(expect.objectContaining({ writer: "updateTask", args: ["FN-8923", { status: "merging" }] }));
|
||||
expect(orphanWrites).not.toContainEqual(expect.objectContaining({ writer: "updateTask", args: ["FN-8923", { status: "merging" }] }));
|
||||
expect(orphanWrites.some(({ writer }) => writer === "logEntry")).toBe(true);
|
||||
expect(orphanWrites.some(({ writer }) => writer === "logEntry")).toBe(false);
|
||||
});
|
||||
});
|
||||
});
|
||||
|
||||
39
packages/engine/src/__tests__/merge-write-fence.test.ts
Normal file
39
packages/engine/src/__tests__/merge-write-fence.test.ts
Normal file
@@ -0,0 +1,39 @@
|
||||
import { describe, expect, it, vi } from "vitest";
|
||||
import {
|
||||
assertMergeGenerationOwned,
|
||||
createMergeWriteFence,
|
||||
isMergeAbortedError,
|
||||
} from "../merge/merge-write-fence.js";
|
||||
|
||||
describe("merge write fence", () => {
|
||||
it("reads ownership at every individual boundary", async () => {
|
||||
const controller = new AbortController();
|
||||
const fence = createMergeWriteFence({ taskId: "FN-8958", signal: controller.signal });
|
||||
expect(fence.isOrphaned()).toBe(false);
|
||||
await fence.write("log", () => controller.abort());
|
||||
expect(fence.isOrphaned()).toBe(true);
|
||||
await expect(fence.write("log", () => { throw new Error("must not run"); })).resolves.toBeUndefined();
|
||||
expect(fence.suppressedCount).toBe(1);
|
||||
});
|
||||
|
||||
it("emits once at first suppression and keeps an in-process count", async () => {
|
||||
const record = vi.fn();
|
||||
const controller = new AbortController(); controller.abort();
|
||||
const fence = createMergeWriteFence({ taskId: "FN-8958", signal: controller.signal, recordAudit: record });
|
||||
await fence.write("log", () => undefined);
|
||||
await fence.write("lifecycle", () => undefined);
|
||||
expect(record).toHaveBeenCalledTimes(1);
|
||||
expect(record.mock.calls[0]).toEqual(["log", "suppressed", 1]);
|
||||
expect(fence.suppressedCount).toBe(2);
|
||||
});
|
||||
|
||||
it("emits rejection-first with count zero and has a common abort predicate", () => {
|
||||
const record = vi.fn(); const controller = new AbortController(); controller.abort();
|
||||
const fence = createMergeWriteFence({ taskId: "FN-8958", signal: controller.signal, recordAudit: record });
|
||||
let error: unknown;
|
||||
try { fence.assertOwned(); } catch (caught) { error = caught; }
|
||||
expect(isMergeAbortedError(error)).toBe(true);
|
||||
expect(record.mock.calls[0]).toEqual(["finalization", "rejected", 0]);
|
||||
expect(() => assertMergeGenerationOwned(controller.signal, "FN-8958")).toThrow(/AI merge aborted/);
|
||||
});
|
||||
});
|
||||
@@ -11,6 +11,7 @@ import {
|
||||
type TaskStore,
|
||||
} from "@fusion/core";
|
||||
import { createRunAuditor, generateSyntheticRunId, type DatabaseMutationType, type RunAuditor } from "../util/run-audit.js";
|
||||
import type { MergeWriteFence } from "./merge-write-fence.js";
|
||||
|
||||
/*
|
||||
FNXC:WorkflowMergeFinalization 2026-07-19-07:20 (U7 / R2/R3/KTD-1):
|
||||
@@ -82,6 +83,7 @@ export interface FinalizeProvenAutoMergeTaskOptions {
|
||||
auditPhase?: string;
|
||||
source: "direct-ai-merge" | "merge-confirmed-fast-path" | "self-healing" | "workflow-graph-merge-finalize";
|
||||
log?: (message: string) => void | Promise<void>;
|
||||
fence?: MergeWriteFence;
|
||||
}
|
||||
|
||||
export type WorkflowDoneMergeProofVerdict =
|
||||
@@ -232,6 +234,7 @@ export async function finalizeProvenAutoMergeTask({
|
||||
auditPhase,
|
||||
source,
|
||||
log,
|
||||
fence,
|
||||
}: FinalizeProvenAutoMergeTaskOptions): Promise<AutoMergeFinalizationResult> {
|
||||
const latest = await store.getTask(taskId).catch(() => null);
|
||||
if (!latest) {
|
||||
@@ -298,6 +301,9 @@ export async function finalizeProvenAutoMergeTask({
|
||||
error: undefined,
|
||||
});
|
||||
if (hardBlocker) {
|
||||
// FNXC:MergeReliability 2026-08-11-21:39: A blocker discovered before finalization
|
||||
// still writes task lifecycle state, so an orphan must reject rather than return a blocked result.
|
||||
fence?.assertOwned("finalization");
|
||||
await store.updateTask(taskId, {
|
||||
status: "failed",
|
||||
error: `Merge confirmed but finalization blocked: ${hardBlocker}`,
|
||||
@@ -333,6 +339,7 @@ export async function finalizeProvenAutoMergeTask({
|
||||
return { outcome: "blocked", task: latest, previousColumn: latest.column, reason: proofVerdict.reason };
|
||||
}
|
||||
|
||||
fence?.assertOwned("finalization");
|
||||
await store.updateTask(taskId, {
|
||||
paused: false,
|
||||
status: null,
|
||||
@@ -351,6 +358,7 @@ export async function finalizeProvenAutoMergeTask({
|
||||
}
|
||||
|
||||
try {
|
||||
fence?.assertOwned("finalization");
|
||||
const moved = await store.moveTask(taskId, completeColumn, shouldRecoveryRehome
|
||||
? { moveSource: "engine", recoveryRehome: true, preserveProgress: true }
|
||||
: { moveSource: "engine", preserveProgress: true });
|
||||
@@ -365,6 +373,7 @@ export async function finalizeProvenAutoMergeTask({
|
||||
auditAgentId,
|
||||
auditPhase,
|
||||
});
|
||||
fence?.assertOwned("finalization");
|
||||
await store.logEntry(
|
||||
taskId,
|
||||
`Auto-merge finalization repaired column mismatch: ${latest.column} → ${completeColumn} after proven merge; cleared stale status/blockers`,
|
||||
|
||||
92
packages/engine/src/merge/merge-write-fence.ts
Normal file
92
packages/engine/src/merge/merge-write-fence.ts
Normal file
@@ -0,0 +1,92 @@
|
||||
/*
|
||||
FNXC:MergeReliability 2026-08-11-21:12:
|
||||
An aborted per-claim signal means its merge body no longer owns the task row: a successor owns a
|
||||
fresh signal. Abort is asynchronous, so ownership is read immediately before each individual
|
||||
mutation or irreversible action; a closure, loop, or function-entry check cannot cover a later
|
||||
write, and one checkpoint never covers adjacent writers. Call sites choose their bucket: diagnostic
|
||||
writes suppress while finalization writes reject and unwind.
|
||||
|
||||
The recorder is injected so orphan observability has no ambient state. One best-effort audit row is
|
||||
emitted at the first fence interaction: suppression emits count 1 and rejection emits count 0.
|
||||
The running count is intentionally in-process only because unwinding cannot guarantee an end flush.
|
||||
*/
|
||||
|
||||
export type MergeWriteCategory = "log" | "lifecycle" | "finalization" | "audit";
|
||||
export type MergeFenceInteraction = "suppressed" | "rejected";
|
||||
|
||||
export interface OrphanFenceAuditRecorder {
|
||||
recordRunAuditEvent(event: {
|
||||
type: "merge:orphan-write-fenced";
|
||||
taskId: string;
|
||||
metadata: { taskId: string; category: MergeWriteCategory; interaction: MergeFenceInteraction; suppressedCount: number };
|
||||
}): Promise<unknown> | unknown;
|
||||
}
|
||||
|
||||
export interface MergeWriteFence {
|
||||
readonly taskId: string;
|
||||
readonly signal: AbortSignal | undefined;
|
||||
readonly suppressedCount: number;
|
||||
isOrphaned(): boolean;
|
||||
write<T>(category: MergeWriteCategory, write: () => Promise<T> | T): Promise<T | undefined>;
|
||||
assertOwned(category?: MergeWriteCategory): void;
|
||||
}
|
||||
|
||||
export function createMergeAbortedError(taskId: string): Error {
|
||||
const error = new Error(`AI merge aborted for ${taskId}`);
|
||||
error.name = "MergeAbortedError";
|
||||
return error;
|
||||
}
|
||||
|
||||
export function isMergeAbortedError(error: unknown): boolean {
|
||||
return error instanceof Error && error.name === "MergeAbortedError";
|
||||
}
|
||||
|
||||
export function assertMergeGenerationOwned(signal: AbortSignal | undefined, taskId: string): void {
|
||||
if (signal?.aborted === true) throw createMergeAbortedError(taskId);
|
||||
}
|
||||
|
||||
export function createMergeWriteFence(options: {
|
||||
taskId: string;
|
||||
signal?: AbortSignal;
|
||||
recordAudit?: OrphanFenceAuditRecorder | ((category: MergeWriteCategory, interaction: MergeFenceInteraction, suppressedCount: number) => Promise<unknown> | unknown);
|
||||
}): MergeWriteFence {
|
||||
let suppressedCount = 0;
|
||||
let auditEmitted = false;
|
||||
const emit = (category: MergeWriteCategory, interaction: MergeFenceInteraction): void => {
|
||||
if (auditEmitted) return;
|
||||
auditEmitted = true;
|
||||
try {
|
||||
const recorder = options.recordAudit;
|
||||
const result = typeof recorder === "function"
|
||||
? recorder(category, interaction, suppressedCount)
|
||||
: recorder?.recordRunAuditEvent({
|
||||
type: "merge:orphan-write-fenced",
|
||||
taskId: options.taskId,
|
||||
metadata: { taskId: options.taskId, category, interaction, suppressedCount },
|
||||
});
|
||||
void Promise.resolve(result).catch(() => undefined);
|
||||
} catch {
|
||||
// Observability must never replace the merge's original outcome.
|
||||
}
|
||||
};
|
||||
return {
|
||||
taskId: options.taskId,
|
||||
signal: options.signal,
|
||||
get suppressedCount() { return suppressedCount; },
|
||||
isOrphaned: () => options.signal?.aborted === true,
|
||||
async write<T>(category: MergeWriteCategory, write: () => Promise<T> | T): Promise<T | undefined> {
|
||||
if (options.signal?.aborted === true) {
|
||||
suppressedCount += 1;
|
||||
emit(category, "suppressed");
|
||||
return undefined;
|
||||
}
|
||||
return await write();
|
||||
},
|
||||
assertOwned(category: MergeWriteCategory = "finalization"): void {
|
||||
if (options.signal?.aborted === true) {
|
||||
emit(category, "rejected");
|
||||
throw createMergeAbortedError(options.taskId);
|
||||
}
|
||||
},
|
||||
};
|
||||
}
|
||||
@@ -66,6 +66,12 @@ import { selectUserCommentsForAgentContext } from "../agents/agent-user-comments
|
||||
import { resolveTaskWorkingBranch } from "../worktree/worktree-names.js";
|
||||
import { resolveIntegrationBranch } from "./integration-branch.js";
|
||||
import { advanceIntegrationBranchRef } from "./merger-ref-update-advance.js";
|
||||
import {
|
||||
assertMergeGenerationOwned,
|
||||
createMergeWriteFence,
|
||||
isMergeAbortedError,
|
||||
type MergeWriteFence,
|
||||
} from "./merge-write-fence.js";
|
||||
import { createResolvedAgentSession, resolveMergerSessionModel, resolveMergerThinkingLevel, resolveMergerFallbackThinkingLevel, resolveValidatorThinkingLevel } from "../agents/agent-session-helpers.js";
|
||||
import { promptWithFallback } from "../pi.js";
|
||||
import { AgentLogger } from "../agents/agent-logger.js";
|
||||
@@ -141,7 +147,8 @@ export function writeTransientMergeStatus(
|
||||
signal: AbortSignal | undefined,
|
||||
status: string | null,
|
||||
): Promise<unknown> {
|
||||
return signal?.aborted ? Promise.resolve(undefined) : store.updateTask(taskId, { status }).catch(() => undefined);
|
||||
const fence = createMergeWriteFence({ taskId, signal });
|
||||
return fence.write("lifecycle", () => store.updateTask(taskId, { status }).catch(() => undefined));
|
||||
}
|
||||
|
||||
async function git(args: string[], cwd: string, opts: { timeout?: number } = {}): Promise<string> {
|
||||
@@ -310,6 +317,7 @@ async function recoverApprovedPreexistingAiMergeWorktree(
|
||||
audit,
|
||||
resolveConflicts: stashResolveAgent,
|
||||
allowDirtyLocalCheckoutSync,
|
||||
signal,
|
||||
});
|
||||
if (land.outcome !== "advanced") return null;
|
||||
await log(`AI merge: recovered approved pre-existing clean-room commit ${short(selected.squashSha)} before pruning`);
|
||||
@@ -628,8 +636,9 @@ export async function landSquash(input: {
|
||||
* Resolved project settings default merger.allowDirtyLocalCheckoutSync to true for legacy operator UX, but this helper's parameter default intentionally remains false so direct/programmatic callers and tests fail closed unless they make the dirty-checkout sync policy explicit.
|
||||
*/
|
||||
allowDirtyLocalCheckoutSync?: boolean;
|
||||
signal?: AbortSignal;
|
||||
}): Promise<LandResult> {
|
||||
const { projectRootDir, mergeRoot, integrationBranch, tipSha, squashSha, taskId, audit, resolveConflicts, allowDirtyLocalCheckoutSync = false } = input;
|
||||
const { projectRootDir, mergeRoot, integrationBranch, tipSha, squashSha, taskId, audit, resolveConflicts, allowDirtyLocalCheckoutSync = false, signal } = input;
|
||||
const emit = (outcome: LocalSyncOutcome, extra: Record<string, unknown> = {}) =>
|
||||
audit.git({ type: "merge:ai-local-sync", target: integrationBranch, metadata: { taskId, outcome, squashSha, ...extra } }).catch(() => undefined);
|
||||
|
||||
@@ -637,6 +646,7 @@ export async function landSquash(input: {
|
||||
|
||||
// Case B — target not checked out here: bare CAS ref advance.
|
||||
if (currentBranch !== integrationBranch) {
|
||||
assertMergeGenerationOwned(signal, taskId);
|
||||
const adv = await advanceIntegrationBranchRef({
|
||||
rootDir: mergeRoot, projectRootDir, integrationBranch,
|
||||
newSha: squashSha, expectedCurrentSha: tipSha, taskId, audit,
|
||||
@@ -688,6 +698,7 @@ export async function landSquash(input: {
|
||||
// The dirty state couldn't be stashed (e.g. untracked/tracked collision or a
|
||||
// stash hook failure). Don't risk `merge --ff-only` aborting/clobbering:
|
||||
// advance the ref atomically and leave the user's working tree as-is.
|
||||
assertMergeGenerationOwned(signal, taskId);
|
||||
const adv = await advanceIntegrationBranchRef({
|
||||
rootDir: mergeRoot, projectRootDir, integrationBranch,
|
||||
newSha: squashSha, expectedCurrentSha: tipSha, taskId, audit,
|
||||
@@ -704,6 +715,7 @@ export async function landSquash(input: {
|
||||
}
|
||||
|
||||
// Fast-forward the checkout (and the branch ref) to the squash.
|
||||
assertMergeGenerationOwned(signal, taskId);
|
||||
if (!(await gitOk(["merge", "--ff-only", squashSha], projectRootDir))) {
|
||||
if (stashed) await gitOk(["stash", "pop"], projectRootDir); // restore the user's edits
|
||||
return { outcome: "concurrent", localSync: "skipped-other-branch" };
|
||||
@@ -1060,6 +1072,7 @@ export async function landOneRepo(
|
||||
projectRootDir: repoRootDir, mergeRoot, integrationBranch, tipSha, squashSha, taskId, audit,
|
||||
resolveConflicts: stashResolveAgent,
|
||||
allowDirtyLocalCheckoutSync: ctx.allowDirtyLocalCheckoutSync === true,
|
||||
signal,
|
||||
});
|
||||
if (landed.outcome === "concurrent") {
|
||||
if (advanceRetries < MAX_CONCURRENT_ADVANCE_RETRIES) {
|
||||
@@ -1320,18 +1333,27 @@ export async function runAiMerge(
|
||||
phase: "merge",
|
||||
});
|
||||
|
||||
const fence = createMergeWriteFence({
|
||||
taskId,
|
||||
signal: options.signal,
|
||||
recordAudit: (category, interaction, suppressedCount) => store.recordRunAuditEvent?.({
|
||||
taskId, agentId: "merger", runId: `merge-${taskId}`, domain: "git",
|
||||
mutationType: "merge:orphan-write-fenced", target: taskId,
|
||||
metadata: { taskId, category, interaction, suppressedCount },
|
||||
}),
|
||||
});
|
||||
// Surface progress on the task detail (status pill) + the task log stream.
|
||||
const log = async (message: string): Promise<void> => {
|
||||
await store.logEntry(taskId, message, "AiMerge").catch(() => undefined);
|
||||
await store.appendAgentLog(taskId, message, "status", undefined, "merger").catch(() => undefined);
|
||||
await fence.write("log", () => store.logEntry(taskId, message, "AiMerge").catch(() => undefined));
|
||||
await fence.write("log", () => store.appendAgentLog(taskId, message, "status", undefined, "merger").catch(() => undefined));
|
||||
};
|
||||
/*
|
||||
FNXC:MergeReliability 2026-08-09-22:35:
|
||||
`raceMergeWithAbort` rejects only the race; a body can outlive the bounded settle latch while a
|
||||
successor generation owns this task. Its per-claim signal remains aborted, so suppressing this
|
||||
status-only write prevents it from re-stamping `merging` (issue #3395) or clearing a successor's
|
||||
live stamp. Keep diagnostics unfenced and make this a no-op, not a throw, because finally paths
|
||||
must preserve the original failure.
|
||||
live stamp. Diagnostics use the same suppress-and-no-op policy rather than throwing, because
|
||||
finally paths must preserve the original failure.
|
||||
*/
|
||||
const setStatus = (status: string | null): Promise<unknown> =>
|
||||
writeTransientMergeStatus(store, taskId, options.signal, status);
|
||||
@@ -1366,7 +1388,7 @@ export async function runAiMerge(
|
||||
target: branch,
|
||||
metadata: { taskId, kind: alreadyMerged ? "already-merged" : "never-executed" },
|
||||
});
|
||||
return await finalizeTask(store, taskId, noOpResult(task, branch, alreadyMerged ? "already-merged" : "no-branch"), undefined, undefined, projectRootDir);
|
||||
return await finalizeTask(store, taskId, noOpResult(task, branch, alreadyMerged ? "already-merged" : "no-branch"), undefined, undefined, projectRootDir, fence);
|
||||
}
|
||||
|
||||
// The target branch must exist as a LOCAL ref to merge into it — surface a
|
||||
@@ -1416,9 +1438,13 @@ export async function runAiMerge(
|
||||
* FNXC:Lifecycle 2026-06-14-20:02:
|
||||
* FN-6461/FN-6455 requires the AI empty-merge lane to demote no-commits tasks whose skipped/incomplete steps outweigh done steps instead of finalizing the operational work as done.
|
||||
*/
|
||||
await store.updateTask(taskId, { error: reason });
|
||||
await fence.write("lifecycle", () => store.updateTask(taskId, { error: reason }));
|
||||
if (fence.isOrphaned()) return {
|
||||
task, branch, merged: false, noOp: false, ok: true, reason, error: reason,
|
||||
worktreeRemoved: false, branchDeleted: false,
|
||||
};
|
||||
const reboundColumn = await resolveFinalizeReboundColumn(store, taskId);
|
||||
await store.logEntry(
|
||||
await fence.write("log", () => store.logEntry(
|
||||
taskId,
|
||||
`Finalize blocked (no-commits incomplete-work guard): ${reason} — moving back to ${reboundColumn} with progress preserved`,
|
||||
JSON.stringify({
|
||||
@@ -1428,7 +1454,7 @@ export async function runAiMerge(
|
||||
integrationBranch,
|
||||
lane: "ai-empty-merge",
|
||||
}, null, 2),
|
||||
);
|
||||
));
|
||||
await audit.database({
|
||||
type: "task:no-commits-finalize-blocked-incomplete-steps" as Parameters<typeof audit.database>[0]["type"],
|
||||
target: taskId,
|
||||
@@ -1441,7 +1467,7 @@ export async function runAiMerge(
|
||||
lane: "ai-empty-merge",
|
||||
},
|
||||
});
|
||||
await store.moveTask(taskId, reboundColumn, { preserveProgress: true, moveSource: "engine" } as Parameters<TaskStore["moveTask"]>[2]);
|
||||
await fence.write("lifecycle", () => store.moveTask(taskId, reboundColumn, { preserveProgress: true, moveSource: "engine" } as Parameters<TaskStore["moveTask"]>[2]));
|
||||
return {
|
||||
task,
|
||||
branch,
|
||||
@@ -1476,13 +1502,17 @@ export async function runAiMerge(
|
||||
if (!landedProof) {
|
||||
const reason =
|
||||
"branch had no net changes vs main — work may have been reverted or lost; operator review required";
|
||||
await store.updateTask(taskId, { error: reason });
|
||||
await fence.write("lifecycle", () => store.updateTask(taskId, { error: reason }));
|
||||
if (fence.isOrphaned()) return {
|
||||
task, branch, merged: false, noOp: false, ok: true, reason, error: reason,
|
||||
worktreeRemoved: false, branchDeleted: false,
|
||||
};
|
||||
const reboundColumn = await resolveFinalizeReboundColumn(store, taskId);
|
||||
await store.logEntry(
|
||||
await fence.write("log", () => store.logEntry(
|
||||
taskId,
|
||||
`Finalize blocked (empty-merge no-landed-proof guard): ${reason} — moving back to ${reboundColumn} with progress preserved`,
|
||||
JSON.stringify({ branch, integrationBranch, lane: "ai-empty-merge", baseCommitSha: task.baseCommitSha }, null, 2),
|
||||
);
|
||||
));
|
||||
await audit.database({
|
||||
type: "task:empty-merge-finalize-blocked-no-landed-proof" as Parameters<typeof audit.database>[0]["type"],
|
||||
target: taskId,
|
||||
@@ -1495,7 +1525,7 @@ export async function runAiMerge(
|
||||
hadPriorNoOpProof: false,
|
||||
},
|
||||
});
|
||||
await store.moveTask(taskId, reboundColumn, { preserveProgress: true, moveSource: "engine" } as Parameters<TaskStore["moveTask"]>[2]);
|
||||
await fence.write("lifecycle", () => store.moveTask(taskId, reboundColumn, { preserveProgress: true, moveSource: "engine" } as Parameters<TaskStore["moveTask"]>[2]));
|
||||
return {
|
||||
task,
|
||||
branch,
|
||||
@@ -1542,9 +1572,13 @@ export async function runAiMerge(
|
||||
const executorVeto = evaluateNoOpFinalizeExecutorVeto({ mergeIsEmpty: true, task, memory: executorMemory, settings });
|
||||
if (executorVeto.veto) {
|
||||
const vetoReason = executorVeto.reason ?? "overseer failed-executor no-op-finalize veto";
|
||||
await store.updateTask(taskId, { error: vetoReason });
|
||||
await fence.write("lifecycle", () => store.updateTask(taskId, { error: vetoReason }));
|
||||
if (fence.isOrphaned()) return {
|
||||
task, branch, merged: false, noOp: false, ok: true, reason: vetoReason, error: vetoReason,
|
||||
worktreeRemoved: false, branchDeleted: false,
|
||||
};
|
||||
const reboundColumn = await resolveFinalizeReboundColumn(store, taskId);
|
||||
await store.logEntry(
|
||||
await fence.write("log", () => store.logEntry(
|
||||
taskId,
|
||||
`Finalize blocked (overseer failed-executor veto): ${vetoReason} — moving back to ${reboundColumn} with progress preserved`,
|
||||
JSON.stringify({
|
||||
@@ -1554,7 +1588,7 @@ export async function runAiMerge(
|
||||
integrationBranch,
|
||||
lane: "ai-empty-merge",
|
||||
}, null, 2),
|
||||
);
|
||||
));
|
||||
await audit.database({
|
||||
type: "overseer:no-op-finalize-vetoed-failed-executor" as Parameters<typeof audit.database>[0]["type"],
|
||||
target: taskId,
|
||||
@@ -1567,7 +1601,7 @@ export async function runAiMerge(
|
||||
lane: "ai-empty-merge",
|
||||
},
|
||||
});
|
||||
await store.moveTask(taskId, reboundColumn, { preserveProgress: true, moveSource: "engine" } as Parameters<TaskStore["moveTask"]>[2]);
|
||||
await fence.write("lifecycle", () => store.moveTask(taskId, reboundColumn, { preserveProgress: true, moveSource: "engine" } as Parameters<TaskStore["moveTask"]>[2]));
|
||||
return {
|
||||
task,
|
||||
branch,
|
||||
@@ -1582,13 +1616,13 @@ export async function runAiMerge(
|
||||
}
|
||||
|
||||
await log(`AI merge: ${branch} had no net changes vs ${integrationBranch} — finalizing as no-op`);
|
||||
const noOpFinalized = await finalizeMerged(store, projectRootDir, taskId, task, branch, integrationBranch, landResult.tipSha, audit, log, { empty: true }, mergeTarget, groupRouting, options.syncGroupPr);
|
||||
await runPushAfterMergeStep({ store, projectRootDir, taskId, settings, integrationBranch, audit, log, options, result: noOpFinalized });
|
||||
const noOpFinalized = await finalizeMerged(store, projectRootDir, taskId, task, branch, integrationBranch, landResult.tipSha, audit, log, { empty: true }, mergeTarget, groupRouting, options.syncGroupPr, fence);
|
||||
await runPushAfterMergeStep({ store, projectRootDir, taskId, settings, integrationBranch, audit, log, options, result: noOpFinalized, fence });
|
||||
return noOpFinalized;
|
||||
}
|
||||
|
||||
const finalized = await finalizeMerged(store, projectRootDir, taskId, task, branch, integrationBranch, landResult.squashSha, audit, log, { empty: false }, mergeTarget, groupRouting, options.syncGroupPr);
|
||||
await runPushAfterMergeStep({ store, projectRootDir, taskId, settings, integrationBranch, audit, log, options, result: finalized });
|
||||
const finalized = await finalizeMerged(store, projectRootDir, taskId, task, branch, integrationBranch, landResult.squashSha, audit, log, { empty: false }, mergeTarget, groupRouting, options.syncGroupPr, fence);
|
||||
await runPushAfterMergeStep({ store, projectRootDir, taskId, settings, integrationBranch, audit, log, options, result: finalized, fence });
|
||||
return finalized;
|
||||
}
|
||||
|
||||
@@ -1612,8 +1646,9 @@ async function runPushAfterMergeStep(input: {
|
||||
log: (message: string) => Promise<void>;
|
||||
options: MergerOptions;
|
||||
result: MergeResult;
|
||||
fence: MergeWriteFence;
|
||||
}): Promise<void> {
|
||||
const { store, projectRootDir, taskId, settings, integrationBranch, audit, log, options, result } = input;
|
||||
const { store, projectRootDir, taskId, settings, integrationBranch, audit, log, options, result, fence } = input;
|
||||
if (settings.pushAfterMerge !== true || settings.mergeStrategy === "pull-request") return;
|
||||
try {
|
||||
const pushOutcome = await pushAfterMergeToRemote({
|
||||
@@ -1627,6 +1662,7 @@ async function runPushAfterMergeStep(input: {
|
||||
signal: options.signal,
|
||||
onAgentText: options.onAgentText,
|
||||
onSession: options.onSession,
|
||||
fence,
|
||||
});
|
||||
result.pushedToRemote = pushOutcome.pushed;
|
||||
if (pushOutcome.error) result.pushError = pushOutcome.error;
|
||||
@@ -1653,9 +1689,9 @@ async function runPushAfterMergeStep(input: {
|
||||
const details = latest?.mergeDetails;
|
||||
if (details?.commitSha && details.commitSha !== pushOutcome.rebasedSha) {
|
||||
const { filesChanged, insertions, deletions } = await captureSingleCommitLandedMetadata(projectRootDir, pushOutcome.rebasedSha);
|
||||
await store.updateTask(taskId, {
|
||||
await fence.write("lifecycle", () => store.updateTask(taskId, {
|
||||
mergeDetails: { ...details, commitSha: pushOutcome.rebasedSha, filesChanged, insertions, deletions },
|
||||
});
|
||||
}));
|
||||
}
|
||||
} catch (refreshErr: unknown) {
|
||||
aiMergeLog.warn(`${taskId}: post-push mergeDetails refresh failed: ${getErrorMessage(refreshErr)}`);
|
||||
@@ -1663,11 +1699,11 @@ async function runPushAfterMergeStep(input: {
|
||||
}
|
||||
} else {
|
||||
aiMergeLog.warn(`${taskId}: push to remote failed: ${pushOutcome.error}`);
|
||||
await store.logEntry(
|
||||
await fence.write("log", () => store.logEntry(
|
||||
taskId,
|
||||
`Push to remote failed after merge — task finalized anyway; local ${integrationBranch} may diverge from ${pushOutcome.remote ?? "origin"}: ${pushOutcome.error}`,
|
||||
"PushToRemoteFailed",
|
||||
).catch(() => undefined);
|
||||
).catch(() => undefined));
|
||||
}
|
||||
} catch (err: unknown) {
|
||||
if (err instanceof Error && err.name === "MergeAbortedError") {
|
||||
@@ -1684,7 +1720,7 @@ async function runPushAfterMergeStep(input: {
|
||||
target: taskId,
|
||||
metadata: { integrationBranch, remote: settings.pushRemote ?? "origin", outcome: "aborted" },
|
||||
}).catch(() => undefined);
|
||||
await store.logEntry(taskId, message, "PushToRemoteFailed").catch(() => undefined);
|
||||
await fence.write("log", () => store.logEntry(taskId, message, "PushToRemoteFailed").catch(() => undefined));
|
||||
return;
|
||||
}
|
||||
const message = getErrorMessage(err);
|
||||
@@ -1696,11 +1732,11 @@ async function runPushAfterMergeStep(input: {
|
||||
target: taskId,
|
||||
metadata: { integrationBranch, remote: settings.pushRemote ?? "origin", outcome: "failed", stderrPreview: message.slice(0, 500) },
|
||||
}).catch(() => undefined);
|
||||
await store.logEntry(
|
||||
await fence.write("log", () => store.logEntry(
|
||||
taskId,
|
||||
`Push to remote threw after merge — task finalized anyway; local ${integrationBranch} may diverge from origin: ${message}`,
|
||||
"PushToRemoteFailed",
|
||||
).catch(() => undefined);
|
||||
).catch(() => undefined));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1873,9 +1909,18 @@ export async function landWorkspaceTask(
|
||||
taskId,
|
||||
phase: "merge",
|
||||
});
|
||||
const fence = createMergeWriteFence({
|
||||
taskId,
|
||||
signal: options.signal,
|
||||
recordAudit: (category, interaction, suppressedCount) => store.recordRunAuditEvent?.({
|
||||
taskId, agentId: "merger", runId: `merge-${taskId}`, domain: "git",
|
||||
mutationType: "merge:orphan-write-fenced", target: taskId,
|
||||
metadata: { taskId, category, interaction, suppressedCount },
|
||||
}),
|
||||
});
|
||||
const log = async (message: string): Promise<void> => {
|
||||
await store.logEntry(taskId, message, "AiMerge").catch(() => undefined);
|
||||
await store.appendAgentLog(taskId, message, "status", undefined, "merger").catch(() => undefined);
|
||||
await fence.write("log", () => store.logEntry(taskId, message, "AiMerge").catch(() => undefined));
|
||||
await fence.write("log", () => store.appendAgentLog(taskId, message, "status", undefined, "merger").catch(() => undefined));
|
||||
};
|
||||
/*
|
||||
FNXC:MergeReliability 2026-08-09-22:35:
|
||||
@@ -2046,6 +2091,7 @@ export async function landWorkspaceTask(
|
||||
repos.push({ repo: repoRel, repoRootDir, integrationBranch, branch: entry.branch, status: "empty" });
|
||||
}
|
||||
} catch (err: unknown) {
|
||||
if (isMergeAbortedError(err)) throw err;
|
||||
// A WorkspacePartialLandError from the persist-failure window above must PROPAGATE
|
||||
// (the engine parks/retries). The outer try/finally below resets status first (A3).
|
||||
if (err instanceof WorkspacePartialLandError) throw err;
|
||||
@@ -2113,22 +2159,23 @@ export async function landWorkspaceTask(
|
||||
if (hasRevertedEmptyRepo) {
|
||||
const reason =
|
||||
"branch had no net changes vs main — work may have been reverted or lost; operator review required";
|
||||
await store.updateTask(taskId, { error: reason });
|
||||
await fence.write("lifecycle", () => store.updateTask(taskId, { error: reason }));
|
||||
if (fence.isOrphaned()) return { taskId, repos, allLanded, finalized: false };
|
||||
const reboundColumn = await resolveFinalizeReboundColumn(store, taskId);
|
||||
await store.logEntry(
|
||||
await fence.write("log", () => store.logEntry(
|
||||
taskId,
|
||||
`Finalize blocked (empty-merge no-landed-proof guard, workspace): ${reason} — moving back to ${reboundColumn} with progress preserved`,
|
||||
JSON.stringify({ lane: "ai-empty-merge-workspace", repoCount: repos.length, landedCount, repos: repos.map((r) => r.repo) }, null, 2),
|
||||
).catch(() => undefined);
|
||||
).catch(() => undefined));
|
||||
await audit.database({
|
||||
type: "task:empty-merge-finalize-blocked-no-landed-proof" as Parameters<typeof audit.database>[0]["type"],
|
||||
target: taskId,
|
||||
metadata: { reason, lane: "ai-empty-merge-workspace", repoCount: repos.length, landedCount, hadPriorNoOpProof: false },
|
||||
}).catch(() => undefined);
|
||||
await store.moveTask(taskId, reboundColumn, { preserveProgress: true, moveSource: "engine" } as Parameters<TaskStore["moveTask"]>[2]);
|
||||
await fence.write("lifecycle", () => store.moveTask(taskId, reboundColumn, { preserveProgress: true, moveSource: "engine" } as Parameters<TaskStore["moveTask"]>[2]));
|
||||
return { taskId, repos, allLanded, finalized: false };
|
||||
}
|
||||
const finalized = await finalizeWorkspaceTask(store, taskId, task, repos);
|
||||
const finalized = await finalizeWorkspaceTask(store, taskId, task, repos, fence);
|
||||
return { taskId, repos, allLanded, finalized };
|
||||
}
|
||||
return { taskId, repos, allLanded, finalized: false };
|
||||
@@ -2183,6 +2230,7 @@ async function finalizeWorkspaceTask(
|
||||
taskId: string,
|
||||
task: Task,
|
||||
repos: WorkspaceRepoLandResult[],
|
||||
fence?: MergeWriteFence,
|
||||
): Promise<boolean> {
|
||||
const landed = repos.filter((r) => r.status === "landed" && r.landedSha);
|
||||
const workspaceLandedShas: Record<string, string> = {};
|
||||
@@ -2211,6 +2259,7 @@ async function finalizeWorkspaceTask(
|
||||
...(anyLanded ? { workspaceLandedShas } : {}),
|
||||
mergeConfirmed: anyLanded,
|
||||
};
|
||||
fence?.assertOwned("finalization");
|
||||
await store.updateTask(taskId, { mergeDetails });
|
||||
task.mergeDetails = mergeDetails;
|
||||
|
||||
@@ -2226,8 +2275,9 @@ async function finalizeWorkspaceTask(
|
||||
worktreeRemoved: false,
|
||||
branchDeleted: false,
|
||||
};
|
||||
await store.logEntry(taskId, `AI merge (workspace): all ${repos.length} sub-repo(s) landed — task → done`, "AiMerge").catch(() => undefined);
|
||||
await finalizeTask(store, taskId, result);
|
||||
await fence?.write("log", () => store.logEntry(taskId, `AI merge (workspace): all ${repos.length} sub-repo(s) landed — task → done`, "AiMerge").catch(() => undefined));
|
||||
fence?.assertOwned("finalization");
|
||||
await finalizeTask(store, taskId, result, undefined, undefined, undefined, fence);
|
||||
return true;
|
||||
}
|
||||
|
||||
@@ -2363,8 +2413,12 @@ export async function pushAfterMergeToRemote(input: {
|
||||
signal?: AbortSignal;
|
||||
onAgentText?: (delta: string) => void;
|
||||
onSession?: (session: { dispose: () => void }) => void;
|
||||
fence?: MergeWriteFence;
|
||||
}): Promise<{ pushed: boolean; remote?: string; targetBranch?: string; refAdvanced?: boolean; rebasedSha?: string; error?: string }> {
|
||||
const { store, projectRootDir, taskId, settings, integrationBranch, audit, log, signal } = input;
|
||||
// FNXC:MergeReliability 2026-08-11-22:17: Post-push recovery diagnostics can outlive
|
||||
// cancellation, so direct callers construct the same per-generation write fence.
|
||||
const fence = input.fence ?? createMergeWriteFence({ taskId, signal });
|
||||
|
||||
let remote: string;
|
||||
let targetBranch: string;
|
||||
@@ -2414,7 +2468,7 @@ export async function pushAfterMergeToRemote(input: {
|
||||
target: taskId,
|
||||
metadata: { taskId, remote, recoveryBranch, sha: localSha, outcome },
|
||||
}).catch(() => undefined);
|
||||
await store.logEntry(taskId, logMessage, logAction).catch(() => undefined);
|
||||
await fence.write("log", () => store.logEntry(taskId, logMessage, logAction).catch(() => undefined));
|
||||
};
|
||||
try {
|
||||
await git(["push", "--force", remote, `${localSha}:${recoveryRef}`], projectRootDir, { timeout: 120_000 });
|
||||
@@ -2500,6 +2554,7 @@ export async function pushAfterMergeToRemote(input: {
|
||||
if (!rebasedSha || rebasedSha === localSha) {
|
||||
return { pushed: true, remote, targetBranch };
|
||||
}
|
||||
assertMergeGenerationOwned(signal, taskId);
|
||||
const adv = await advanceIntegrationBranchRef({
|
||||
rootDir: canonicalPushRoot,
|
||||
projectRootDir,
|
||||
@@ -2568,6 +2623,7 @@ async function finalizeMerged(
|
||||
mergeTarget?: MergeTargetResolution,
|
||||
groupRouting?: BranchGroupMergeRouting | null,
|
||||
syncGroupPr?: SyncGroupPrFn,
|
||||
fence?: MergeWriteFence,
|
||||
): Promise<MergeResult> {
|
||||
/*
|
||||
FNXC:BranchGroupCompletion 2026-07-04-00:00:
|
||||
@@ -2601,10 +2657,12 @@ async function finalizeMerged(
|
||||
...mergeTargetPatch,
|
||||
};
|
||||
modifiedFiles = landedFiles.length > 0 ? landedFiles : undefined;
|
||||
fence?.assertOwned("finalization");
|
||||
await store.updateTask(taskId, { mergeDetails, modifiedFiles });
|
||||
task.mergeDetails = mergeDetails;
|
||||
task.modifiedFiles = modifiedFiles;
|
||||
if (task.lineageId && typeof (store as Partial<TaskStore>).upsertTaskCommitAssociation === "function") {
|
||||
fence?.assertOwned("finalization");
|
||||
await store.upsertTaskCommitAssociation({
|
||||
taskLineageId: task.lineageId,
|
||||
taskIdSnapshot: task.id,
|
||||
@@ -2619,6 +2677,7 @@ async function finalizeMerged(
|
||||
}
|
||||
} else if (mergeTargetPatch) {
|
||||
mergeDetails = { ...(task.mergeDetails ?? {}), ...mergeTargetPatch };
|
||||
fence?.assertOwned("finalization");
|
||||
await store.updateTask(taskId, { mergeDetails });
|
||||
task.mergeDetails = mergeDetails;
|
||||
}
|
||||
@@ -2626,6 +2685,7 @@ async function finalizeMerged(
|
||||
// NEVER delete the integration branch itself — a task whose branch name
|
||||
// coincides with the target (or merges into its own branch) must not have the
|
||||
// just-advanced integration ref force-deleted out from under it.
|
||||
fence?.assertOwned("finalization");
|
||||
if (branch !== integrationBranch && await gitOk(["branch", "-D", branch], projectRootDir)) {
|
||||
branchDeleted = true;
|
||||
await audit.git({ type: "branch:delete", target: branch, metadata: { taskId, force: true } }).catch(() => undefined);
|
||||
@@ -2633,7 +2693,9 @@ async function finalizeMerged(
|
||||
// Remove the task's own worktree if it still exists.
|
||||
let worktreeRemoved = false;
|
||||
if (task.worktree) {
|
||||
fence?.assertOwned("finalization");
|
||||
worktreeRemoved = await gitOk(["worktree", "remove", "--force", task.worktree], projectRootDir);
|
||||
fence?.assertOwned("finalization");
|
||||
await store.updateTask(taskId, { worktree: null }).catch(() => undefined);
|
||||
}
|
||||
|
||||
@@ -2655,27 +2717,27 @@ async function finalizeMerged(
|
||||
};
|
||||
await audit.git({ type: "merge:ai-landed", target: integrationBranch, metadata: { taskId, landedSha, empty: opts.empty } }).catch(() => undefined);
|
||||
await log(opts.empty ? `AI merge: finalized ${taskId} (no-op), finalizing task row` : `AI merge: landed ${short(landedSha)}, finalizing task row`);
|
||||
const finalized = await finalizeTask(store, taskId, result, audit, log, projectRootDir);
|
||||
await log(opts.empty ? `AI merge: finalized ${taskId} (no-op) → done` : `AI merge: landed ${short(landedSha)}, task → done`);
|
||||
|
||||
/*
|
||||
FNXC:BranchGroupCompletion 2026-07-04-00:00:
|
||||
FN-7532: mirror the legacy merger.ts executeMergeAttempt's shared-group landing bookkeeping so a
|
||||
member merged via the (now sole) runAiMerge path also updates the group row (worktreePath/status)
|
||||
and pushes the up-to-date checklist body onto any already-open managed group PR. Both are
|
||||
best-effort — a failure here must never fail an otherwise-successful merge.
|
||||
FNXC:MergeReliability 2026-08-11-21:39:
|
||||
Group bookkeeping is a finalization writer, so it must finish before the done-column move and
|
||||
`task:merged` announcement. An abort here rejects before external consumers see an announced
|
||||
merge whose managed-group state is still incomplete; each adjacent writer keeps its own fence.
|
||||
*/
|
||||
if (groupRouting) {
|
||||
try {
|
||||
fence?.assertOwned("finalization");
|
||||
await Promise.resolve((store as { recordBranchGroupMemberLanded?: TaskStore["recordBranchGroupMemberLanded"] }).recordBranchGroupMemberLanded?.(groupRouting.branchGroup.id, {
|
||||
worktreePath: task.worktree ?? null,
|
||||
status: "open",
|
||||
}));
|
||||
} catch {
|
||||
} catch (err) {
|
||||
if (isMergeAbortedError(err)) throw err;
|
||||
// best-effort persistence
|
||||
}
|
||||
if (syncGroupPr) {
|
||||
try {
|
||||
fence?.assertOwned("finalization");
|
||||
await syncGroupPrOnLanding({
|
||||
store,
|
||||
groupId: groupRouting.branchGroup.id,
|
||||
@@ -2683,6 +2745,7 @@ async function finalizeMerged(
|
||||
syncGroupPr,
|
||||
});
|
||||
} catch (err) {
|
||||
if (isMergeAbortedError(err)) throw err;
|
||||
try {
|
||||
store.recordRunAuditEvent?.({
|
||||
taskId,
|
||||
@@ -2700,6 +2763,9 @@ async function finalizeMerged(
|
||||
}
|
||||
}
|
||||
|
||||
fence?.assertOwned("finalization");
|
||||
const finalized = await finalizeTask(store, taskId, result, audit, log, projectRootDir, fence);
|
||||
await log(opts.empty ? `AI merge: finalized ${taskId} (no-op) → done` : `AI merge: landed ${short(landedSha)}, task → done`);
|
||||
return finalized;
|
||||
}
|
||||
|
||||
@@ -2711,6 +2777,7 @@ async function finalizeTask(
|
||||
audit?: RunAuditor,
|
||||
log?: (message: string) => Promise<void>,
|
||||
rootDir?: string,
|
||||
fence?: MergeWriteFence,
|
||||
): Promise<MergeResult> {
|
||||
const finalization = await finalizeProvenAutoMergeTask({
|
||||
store,
|
||||
@@ -2722,6 +2789,7 @@ async function finalizeTask(
|
||||
source: "direct-ai-merge",
|
||||
rootDir,
|
||||
log,
|
||||
fence,
|
||||
});
|
||||
if (finalization.outcome === "blocked") {
|
||||
throw new Error(`AI merge finalization blocked for ${taskId}: ${finalization.reason ?? "unknown"}`);
|
||||
@@ -2730,14 +2798,11 @@ async function finalizeTask(
|
||||
throw new Error(`AI merge finalization could not find task ${taskId}`);
|
||||
}
|
||||
result.task = finalization.task;
|
||||
fence?.assertOwned("finalization");
|
||||
store.emit("task:merged", result);
|
||||
return result;
|
||||
}
|
||||
|
||||
function throwIfAborted(signal: AbortSignal | undefined, taskId: string): void {
|
||||
if (signal?.aborted) {
|
||||
const err = new Error(`AI merge aborted for ${taskId}`);
|
||||
err.name = "MergeAbortedError";
|
||||
throw err;
|
||||
}
|
||||
assertMergeGenerationOwned(signal, taskId);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user