FN-8838: refresh automated PR heads before GitHub mutations

Refresh isolated automated PR heads against their target immediately before PR creation or merge.

- Rebase and lease-publish task and group heads at every automated GitHub boundary.
- Fail closed on concurrent head updates and reconcile retained refresh worktrees.
- Add lifecycle coverage and document the configurable integration remote.

Files changed:
 .changeset/fn-8838-refresh-pr-heads.md             |   7 +
 docs/settings-reference.md                         |   6 +
 .../task-lifecycle-refresh.integration.test.ts     | 438 ++++++++++++++++
 .../src/commands/__tests__/task-lifecycle.test.ts  | 309 ++++++++++-
 packages/cli/src/commands/daemon.ts                |   4 +-
 packages/cli/src/commands/dashboard.ts             |   4 +-
 packages/cli/src/commands/serve.ts                 |   4 +-
 packages/cli/src/commands/task-lifecycle.ts        | 581 +++++++++++++++++++--
 .../src/__tests__/group-merge-coordinator.test.ts  |  43 ++
 .../engine/src/merge/group-merge-coordinator.ts    |  34 +-
 packages/engine/src/merge/pr-nodes.ts              |  35 +-
 packages/engine/src/project-engine.ts              |  16 +-
 12 files changed, 1425 insertions(+), 56 deletions(-)

Fusion-Task-Id: FN-8838
Fusion-Task-Lineage: a9b9e800-7d73-441f-bc1b-2488d244e0b1
Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
gsxdsm
2026-08-08 22:35:59 -07:00
parent 1e5d80dd4d
commit b3504f01a4
12 changed files with 1425 additions and 56 deletions

View File

@@ -0,0 +1,7 @@
---
"@runfusion/fusion": patch
---
summary: Refresh automated pull-request heads before creating or merging them.
category: fix
dev: Automated task, group, promotion, and workflow PR paths use verified checkout refreshes and leased rewrites.

View File

@@ -512,6 +512,12 @@ Security-sensitive file-browser escape hatches are project-only. `allowAbsoluteF
| `autoRecovery.maxRetries` | `number` | `3` | Retry budget for dispatcher decisions. When `retryCount >= maxRetries`, dispatcher forces `pause` with rationale `retry-budget-exhausted`. | | `autoRecovery.maxRetries` | `number` | `3` | Retry budget for dispatcher decisions. When `retryCount >= maxRetries`, dispatcher forces `pause` with rationale `retry-budget-exhausted`. |
| `reliabilityStatsResetAt` | `string` (ISO-8601) | `undefined` | Optional reliability baseline cursor used by `/api/health/reliability`; events older than this timestamp are excluded from reliability aggregates but retained in storage. | | `reliabilityStatsResetAt` | `string` (ISO-8601) | `undefined` | Optional reliability baseline cursor used by `/api/health/reliability`; events older than this timestamp are excluded from reliability aggregates but retained in storage. |
### Automated pull-request freshness
With `mergeStrategy: "pull-request"`, Fusion refreshes an automated task, managed-group, promotion-group, or workflow PR head immediately before it creates a PR and immediately before it requests a merge. It resolves the task's effective integration target and configured/tracked integration remote, fetches that target, then rebases the exact head checkout onto the remote target and distinct local target when needed.
Fusion never performs this refresh in an arbitrary project-root checkout. It first proves an existing worktree owns `refs/heads/<head>`; otherwise it materializes a remote-only head at an explicit local ref after verifying the observed remote OID, then creates a short-lived managed worktree attached to that exact ref. A head checked out at the project root is refused, so user checkout state is never changed. Temporary refresh paths have durable reservations: a cleanup failure quarantines the exact managed path, prevents GitHub mutation, and requires a later bounded inactive-worktree reconciliation to remove it before that path can be reused. Rewritten published heads use `--force-with-lease` against the observed remote OID; a rejected lease restores the locally rebased head through a compare-and-swap ref update. Lease loss, fetch/rebase conflicts, missing heads, or a failed temporary-worktree cleanup fail closed before the GitHub create/merge call. After a successful rewrite, the merge path re-reads PR state and fences the request with the refreshed head OID. Manual `fn pr`, dashboard operator actions, and revert PRs remain user-directed and are not changed by this automated lifecycle behavior.
### Per-task direct-merge override ### Per-task direct-merge override
When a project uses `mergeStrategy: "direct"`, an individual task can override the project-level `directMergeCommitStrategy` by adding this line anywhere in `PROMPT.md`: When a project uses `mergeStrategy: "direct"`, an individual task can override the project-level `directMergeCommitStrategy` by adding this line anywhere in `PROMPT.md`:

View File

@@ -0,0 +1,438 @@
import { execFileSync } from "node:child_process";
import { chmodSync, mkdtempSync, rmSync, writeFileSync } from "node:fs";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { afterEach, describe, expect, it, vi } from "vitest";
vi.mock("@fusion/core", async () => {
const actual = await vi.importActual<typeof import("@fusion/core")>("@fusion/core");
return {
...actual,
// The fixture's origin is a local bare repository; the production adapters
// still need a GitHub repository identity to exercise their GitHub boundary.
getCurrentRepo: vi.fn(() => ({ owner: "fixture-owner", repo: "fixture-repo" })),
getPushRepo: vi.fn(() => ({ owner: "fixture-owner", repo: "fixture-repo" })),
};
});
import {
createGroupPrCallback,
createPrNodeGithubOps,
processPullRequestMergeTask,
refreshAutomatedPrHead,
} from "../task-lifecycle.js";
const fixtures: string[] = [];
function git(cwd: string, ...args: string[]): string {
return execFileSync("git", args, { cwd, encoding: "utf8" }).trim();
}
function makeFixture(head = "fusion/fn-refresh-fixture"): { root: string; remote: string; integration: string; head: string } {
const root = mkdtempSync(join(tmpdir(), "fusion-pr-refresh-"));
fixtures.push(root);
const remote = join(root, "remote.git");
const project = join(root, "project");
const integration = join(root, "integration");
git(root, "init", "--bare", remote);
git(root, "clone", remote, project);
git(project, "config", "user.email", "test@example.com");
git(project, "config", "user.name", "Fusion Test");
writeFileSync(join(project, "base.txt"), "base\n");
git(project, "add", "base.txt");
git(project, "commit", "-m", "base");
git(project, "branch", "-M", "main");
git(project, "push", "-u", "origin", "main");
git(project, "checkout", "-b", head);
writeFileSync(join(project, "feature.txt"), "feature\n");
git(project, "add", "feature.txt");
git(project, "commit", "-m", "feature");
git(project, "checkout", "main");
git(root, "clone", remote, integration);
git(integration, "config", "user.email", "test@example.com");
git(integration, "config", "user.name", "Fusion Test");
git(integration, "checkout", "main");
writeFileSync(join(integration, "sentinel.txt"), "late integration security fix\n");
git(integration, "add", "sentinel.txt");
git(integration, "commit", "-m", "security sentinel");
git(integration, "push", "origin", "main");
return { root: project, remote, integration, head };
}
function makeConflictingFixture(head: string) {
const fixture = makeFixture(head);
git(fixture.root, "checkout", head);
writeFileSync(join(fixture.root, "base.txt"), "head conflicts with integration\n");
git(fixture.root, "add", "base.txt");
git(fixture.root, "commit", "-m", "conflicting head change");
git(fixture.root, "checkout", "main");
writeFileSync(join(fixture.integration, "base.txt"), "integration conflicts with head\n");
git(fixture.integration, "add", "base.txt");
git(fixture.integration, "commit", "-m", "conflicting integration change");
git(fixture.integration, "push", "origin", "main");
return fixture;
}
afterEach(() => {
for (const fixture of fixtures.splice(0)) rmSync(fixture, { recursive: true, force: true });
});
function assertPublishedSentinel(remote: string, head: string, sentinel = "sentinel.txt"): void {
expect(git(remote, "show", `refs/heads/${head}:${sentinel}`)).toBe("late integration security fix");
}
function makeLifecycleStore(task: Record<string, unknown>) {
return {
getTask: async () => task,
getSettings: async () => ({ baseBranch: "main" }),
getTaskWorkflowSelection: () => undefined,
getWorkflowDefinition: async () => undefined,
getWorkflowSettingValues: () => ({}),
getWorkflowSettingsProjectId: () => "fixture-project",
getBranchGroup: async () => null,
listTasksByBranchGroup: async () => [],
getActiveMergingTask: async () => null,
updateTask: async () => undefined,
updatePrInfo: async (_id: string, prInfo: unknown) => { Object.assign(task, { prInfo }); },
updateBranchGroup: async () => undefined,
logEntry: async () => undefined,
moveTask: async () => task,
emit: () => undefined,
};
}
describe("refreshAutomatedPrHead local git fixture", () => {
it("publishes a stale automated head only after it contains the late integration sentinel", async () => {
const { root, remote, head } = makeFixture();
/*
FNXC:PullRequestFreshness 2026-08-09-03:20:
An automated PR head created before a later integration security fix must be
rebased and lease-published before any GitHub create boundary can observe it.
*/
const refreshed = await refreshAutomatedPrHead({
projectRoot: root,
headBranch: head,
targetBranch: "main",
});
expect(refreshed.refreshed).toBe(true);
expect(git(remote, "show", `refs/heads/${head}:feature.txt`)).toBe("feature");
expect(git(remote, "show", `refs/heads/${head}:sentinel.txt`)).toBe("late integration security fix");
expect(git(root, "branch", "--show-current")).toBe("main");
expect(git(root, "worktree", "list", "--porcelain")).not.toContain("/.worktrees/pr-refresh-");
});
it("refuses a head checked out at the primary project root", async () => {
const { root, head } = makeFixture("fusion/fn-8838-root-refusal");
git(root, "checkout", head);
await expect(refreshAutomatedPrHead({ projectRoot: root, headBranch: head, targetBranch: "main" }))
.rejects.toThrow(/project root/);
expect(git(root, "branch", "--show-current")).toBe(head);
});
it("uses a verified existing task worktree without creating a temporary checkout", async () => {
const { root, remote, head } = makeFixture("fusion/fn-8838-verified-worktree");
const taskWorktree = join(root, "task-worktree");
git(root, "worktree", "add", taskWorktree, head);
const refreshed = await refreshAutomatedPrHead({
projectRoot: root,
headBranch: head,
targetBranch: "main",
preferredWorktree: taskWorktree,
});
expect(refreshed.refreshed).toBe(true);
assertPublishedSentinel(remote, head);
expect(git(taskWorktree, "branch", "--show-current")).toBe(head);
expect(git(root, "worktree", "list", "--porcelain")).not.toContain("/.worktrees/pr-refresh-");
});
it("materializes a remote-only head and refuses a missing head without GitHub boundaries", async () => {
const remoteOnly = makeFixture("fusion/fn-8838-remote-only");
git(remoteOnly.root, "push", "origin", remoteOnly.head);
git(remoteOnly.root, "branch", "-D", remoteOnly.head);
await expect(refreshAutomatedPrHead({
projectRoot: remoteOnly.root,
headBranch: remoteOnly.head,
targetBranch: "main",
})).resolves.toEqual(expect.objectContaining({ refreshed: true }));
assertPublishedSentinel(remoteOnly.remote, remoteOnly.head);
const missing = makeFixture("fusion/fn-8838-missing");
await expect(refreshAutomatedPrHead({
projectRoot: missing.root,
headBranch: "fusion/fn-8838-does-not-exist",
targetBranch: "main",
})).rejects.toThrow(/missing local and origin head/);
});
it("restores the canonical local head when a force-with-lease publication is rejected", async () => {
const { root, remote, head } = makeFixture("fusion/fn-8838-lease-rejection");
// Make this an existing remote head, so the refresh publication must use a lease.
git(root, "push", "origin", head);
const originalHead = git(root, "rev-parse", `refs/heads/${head}`);
const remoteMain = git(remote, "rev-parse", "refs/heads/main");
const hook = join(root, ".git", "hooks", "pre-push");
writeFileSync(hook, `#!/bin/sh\ngit --git-dir='${remote}' update-ref 'refs/heads/${head}' '${remoteMain}'\n`);
chmodSync(hook, 0o755);
/*
FNXC:PullRequestFreshness 2026-08-09-04:52:
A remote writer can win after refresh observes the old OID but before its
guarded push reaches origin. The rejected lease must restore the shared
local ref, preventing a later unguarded push from publishing the rebase.
*/
await expect(refreshAutomatedPrHead({ projectRoot: root, headBranch: head, targetBranch: "main" }))
.rejects.toThrow(/stale info|lease|failed to push/i);
expect(git(root, "rev-parse", `refs/heads/${head}`)).toBe(originalHead);
expect(git(remote, "rev-parse", `refs/heads/${head}`)).toBe(remoteMain);
expect(git(root, "worktree", "list", "--porcelain")).not.toContain("/.worktrees/pr-refresh-");
});
it("refuses to overwrite a head that advanced before refresh observes its push lease", async () => {
const { root, integration, remote, head } = makeFixture("fusion/fn-8838-pre-observation-race");
git(root, "push", "origin", head);
const originalHead = git(root, "rev-parse", `refs/heads/${head}`);
git(integration, "fetch", "origin", `${head}:${head}`);
git(integration, "checkout", head);
writeFileSync(join(integration, "concurrent.txt"), "concurrent head update\n");
git(integration, "add", "concurrent.txt");
git(integration, "commit", "-m", "concurrent head update");
git(integration, "push", "origin", head);
const concurrentHead = git(remote, "rev-parse", `refs/heads/${head}`);
/*
FNXC:PullRequestFreshness 2026-08-09-05:32:
A head update that reaches origin before refresh reads its lease must survive.
The refresher may not replace that unincorporated commit with its rebased tip.
*/
await expect(refreshAutomatedPrHead({ projectRoot: root, headBranch: head, targetBranch: "main" }))
.rejects.toThrow(/changed before publication/);
expect(git(root, "rev-parse", `refs/heads/${head}`)).toBe(originalHead);
expect(git(remote, "rev-parse", `refs/heads/${head}`)).toBe(concurrentHead);
expect(git(remote, "show", `refs/heads/${head}:concurrent.txt`)).toBe("concurrent head update");
});
it("refreshes each production PR creator before its recording GitHub fake sees the head", async () => {
const groupFixture = makeFixture("fusion/group-refresh-fixture");
const groupGithub = {
findPrForBranch: vi.fn(async () => null),
createPr: vi.fn(async () => {
assertPublishedSentinel(groupFixture.remote, groupFixture.head);
return { number: 1, url: "https://example.test/pr/1", status: "open" as const };
}),
};
await createGroupPrCallback(groupGithub as never)({
cwd: groupFixture.root,
group: { id: "BG-fixture", branchName: groupFixture.head } as never,
members: [{ id: "FN-8838", title: "fixture" }] as never,
headBranch: groupFixture.head,
baseBranch: "main",
});
const workflowFixture = makeFixture("fusion/fn-8838-workflow");
const workflowGithub = {
createPr: vi.fn(async () => {
assertPublishedSentinel(workflowFixture.remote, workflowFixture.head);
return { number: 2, url: "https://example.test/pr/2", status: "open" as const };
}),
getPrStatus: vi.fn(async () => ({ number: 2, url: "https://example.test/pr/2", status: "open" as const })),
mergePr: vi.fn(async () => ({ number: 2, url: "https://example.test/pr/2", status: "merged" as const })),
replyToReviewThread: vi.fn(),
resolveReviewThread: vi.fn(),
getViewerLogin: vi.fn(),
getPrReviewThreadsDetailed: vi.fn(),
};
const workflowOps = createPrNodeGithubOps(workflowGithub as never);
const task = { id: "FN-8838-WORKFLOW", title: "fixture", description: "fixture", worktree: workflowFixture.root };
const entity = { id: "pr-fixture", sourceId: task.id, repo: "fixture-owner/fixture-repo", headBranch: workflowFixture.head, baseBranch: "main", prNumber: 2, headOid: "old" };
await workflowOps.createPr({ task, entity } as never);
writeFileSync(join(workflowFixture.integration, "merge-sentinel.txt"), "late integration security fix\n");
git(workflowFixture.integration, "add", "merge-sentinel.txt");
git(workflowFixture.integration, "commit", "-m", "merge sentinel");
git(workflowFixture.integration, "push", "origin", "main");
workflowGithub.mergePr.mockImplementation(async () => {
assertPublishedSentinel(workflowFixture.remote, workflowFixture.head, "merge-sentinel.txt");
return { number: 2, url: "https://example.test/pr/2", status: "merged" as const };
});
const persisted: string[] = [];
await workflowOps.mergePr({ task, entity, persistRefreshedHead: async (oid) => { persisted.push(oid); } } as never);
expect(workflowGithub.getPrStatus).toHaveBeenCalledTimes(1);
expect(persisted).toHaveLength(1);
const lifecycleFixture = makeFixture("fusion/fn-8838-lifecycle");
const lifecycleTask = {
id: "FN-8838-LIFECYCLE",
title: "fixture",
description: "fixture",
branch: lifecycleFixture.head,
worktree: lifecycleFixture.root,
column: "in-review",
};
let lifecycleMergeReady = false;
const lifecycleGithub = {
findPrForBranch: vi.fn(async () => null),
createPr: vi.fn(async () => {
assertPublishedSentinel(lifecycleFixture.remote, lifecycleFixture.head);
return { number: 3, url: "https://example.test/pr/3", status: "open" as const };
}),
getPrMergeStatus: vi.fn(async () => ({
prInfo: { number: 3, url: "https://example.test/pr/3", status: "open" as const },
reviewDecision: null,
checks: [],
mergeReady: lifecycleMergeReady,
blockingReasons: lifecycleMergeReady ? [] : ["checks pending"],
})),
mergePr: vi.fn(async () => ({ number: 3, url: "https://example.test/pr/3", status: "merged" as const })),
};
const lifecycleResult = await processPullRequestMergeTask(
makeLifecycleStore(lifecycleTask) as never,
lifecycleFixture.root,
lifecycleTask.id,
lifecycleGithub as never,
() => undefined,
);
expect(lifecycleResult).toBe("waiting");
expect(lifecycleGithub.createPr).toHaveBeenCalledTimes(1);
writeFileSync(join(lifecycleFixture.integration, "lifecycle-merge-sentinel.txt"), "late integration security fix\n");
git(lifecycleFixture.integration, "add", "lifecycle-merge-sentinel.txt");
git(lifecycleFixture.integration, "commit", "-m", "lifecycle merge sentinel");
git(lifecycleFixture.integration, "push", "origin", "main");
lifecycleMergeReady = true;
lifecycleGithub.mergePr.mockImplementation(async () => {
assertPublishedSentinel(lifecycleFixture.remote, lifecycleFixture.head, "lifecycle-merge-sentinel.txt");
return { number: 3, url: "https://example.test/pr/3", status: "merged" as const };
});
const lifecycleMergeResult = await processPullRequestMergeTask(
makeLifecycleStore(lifecycleTask) as never,
lifecycleFixture.root,
lifecycleTask.id,
lifecycleGithub as never,
() => undefined,
);
expect(lifecycleMergeResult).toBe("merged");
expect(lifecycleGithub.mergePr).toHaveBeenCalledTimes(1);
});
it("refreshes shared-group creation and merge boundaries before GitHub sees either head", async () => {
const fixture = makeFixture("fusion/groups/fn-8838-refresh");
const task = {
id: "FN-8838-GROUP",
title: "group fixture",
description: "fixture",
branch: fixture.head,
worktree: fixture.root,
column: "in-review",
branchContext: { assignmentMode: "shared", groupId: "BG-8838", source: "planning" },
};
const group = {
id: "BG-8838",
sourceType: "planning",
sourceId: "P-8838",
branchName: fixture.head,
prState: "none",
status: "open",
createdAt: Date.now(),
updatedAt: Date.now(),
};
let mergeReady = false;
const store = {
...makeLifecycleStore(task),
getBranchGroup: async () => group,
listTasksByBranchGroup: async () => [task],
};
const github = {
findPrForBranch: vi.fn(async () => null),
createPr: vi.fn(async () => {
assertPublishedSentinel(fixture.remote, fixture.head);
return { number: 4, url: "https://example.test/pr/4", status: "open" as const };
}),
getPrMergeStatus: vi.fn(async () => ({
prInfo: { number: 4, url: "https://example.test/pr/4", status: "open" as const },
reviewDecision: "APPROVED" as const,
checks: [],
mergeReady,
blockingReasons: mergeReady ? [] : ["checks pending"],
})),
mergePr: vi.fn(async () => ({ number: 4, url: "https://example.test/pr/4", status: "merged" as const })),
};
await expect(processPullRequestMergeTask(store as never, fixture.root, task.id, github as never, () => undefined)).resolves.toBe("waiting");
expect(github.createPr).toHaveBeenCalledTimes(1);
writeFileSync(join(fixture.integration, "group-merge-sentinel.txt"), "late integration security fix\n");
git(fixture.integration, "add", "group-merge-sentinel.txt");
git(fixture.integration, "commit", "-m", "group merge sentinel");
git(fixture.integration, "push", "origin", "main");
mergeReady = true;
github.mergePr.mockImplementation(async () => {
assertPublishedSentinel(fixture.remote, fixture.head, "group-merge-sentinel.txt");
return { number: 4, url: "https://example.test/pr/4", status: "merged" as const };
});
await expect(processPullRequestMergeTask(store as never, fixture.root, task.id, github as never, () => undefined)).resolves.toBe("merged");
expect(github.mergePr).toHaveBeenCalledWith(expect.objectContaining({ expectedHeadOid: expect.any(String) }));
});
it("fails closed at every PR-create adapter when rebase conflicts", async () => {
/*
FNXC:PullRequestFreshness 2026-08-09-04:26:
A rebase conflict is a hard stop at all automated create boundaries. GitHub
must never receive a normal-looking PR whose stale head omits base changes.
*/
const promotionFixture = makeConflictingFixture("fusion/groups/fn-8838-conflict-promotion");
const promotionGithub = { findPrForBranch: vi.fn(async () => null), createPr: vi.fn() };
await expect(createGroupPrCallback(promotionGithub as never)({
cwd: promotionFixture.root,
group: { id: "BG-conflict", sourceType: "planning", sourceId: "P-conflict" } as never,
members: [],
headBranch: promotionFixture.head,
baseBranch: "main",
})).rejects.toThrow(/rebase/);
expect(promotionGithub.createPr).not.toHaveBeenCalled();
const workflowFixture = makeConflictingFixture("fusion/fn-8838-conflict-workflow");
const workflowGithub = {
createPr: vi.fn(), getPrStatus: vi.fn(), mergePr: vi.fn(), replyToReviewThread: vi.fn(),
resolveReviewThread: vi.fn(), getViewerLogin: vi.fn(), getPrReviewThreadsDetailed: vi.fn(),
};
await expect(createPrNodeGithubOps(workflowGithub as never).createPr({
task: { id: "FN-8838-CONFLICT-WORKFLOW", title: "fixture", description: "fixture", worktree: workflowFixture.root },
entity: { id: "pr-conflict", sourceId: "FN-8838-CONFLICT-WORKFLOW", repo: "fixture-owner/fixture-repo", headBranch: workflowFixture.head, baseBranch: "main" },
} as never)).rejects.toThrow(/rebase/);
expect(workflowGithub.createPr).not.toHaveBeenCalled();
const workflowMergeFixture = makeConflictingFixture("fusion/fn-8838-conflict-workflow-merge");
const workflowMergeGithub = {
createPr: vi.fn(), getPrStatus: vi.fn(), mergePr: vi.fn(), replyToReviewThread: vi.fn(),
resolveReviewThread: vi.fn(), getViewerLogin: vi.fn(), getPrReviewThreadsDetailed: vi.fn(),
};
await expect(createPrNodeGithubOps(workflowMergeGithub as never).mergePr({
task: { id: "FN-8838-CONFLICT-WORKFLOW-MERGE", worktree: workflowMergeFixture.root },
entity: { id: "pr-conflict-merge", sourceId: "FN-8838-CONFLICT-WORKFLOW-MERGE", repo: "fixture-owner/fixture-repo", headBranch: workflowMergeFixture.head, baseBranch: "main", prNumber: 5 },
} as never)).rejects.toThrow(/rebase/);
expect(workflowMergeGithub.getPrStatus).not.toHaveBeenCalled();
expect(workflowMergeGithub.mergePr).not.toHaveBeenCalled();
const lifecycleFixture = makeConflictingFixture("fusion/fn-8838-conflict-lifecycle");
const lifecycleGithub = {
findPrForBranch: vi.fn(async () => null), createPr: vi.fn(), getPrMergeStatus: vi.fn(), mergePr: vi.fn(),
};
const lifecycleTask = { id: "FN-8838-CONFLICT-LIFECYCLE", title: "fixture", description: "fixture", branch: lifecycleFixture.head, worktree: lifecycleFixture.root, column: "in-review" };
await expect(processPullRequestMergeTask(
makeLifecycleStore(lifecycleTask) as never, lifecycleFixture.root, lifecycleTask.id, lifecycleGithub as never, () => undefined,
)).rejects.toThrow(/rebase/);
expect(lifecycleGithub.createPr).not.toHaveBeenCalled();
});
});

View File

@@ -10,11 +10,32 @@ const execMock = vi.hoisted(() => vi.fn());
const execFileCalls = vi.hoisted( const execFileCalls = vi.hoisted(
() => [] as Array<{ file: string; args: string[]; cwd: string | undefined }>, () => [] as Array<{ file: string; args: string[]; cwd: string | undefined }>,
); );
const refreshFixture = vi.hoisted(() => ({ next: 0, branch: "fusion/fn-9601" }));
vi.mock("node:fs/promises", async () => {
const actual = await vi.importActual<typeof import("node:fs/promises")>("node:fs/promises");
return {
...actual,
realpath: vi.fn(async (path: string) => path),
mkdir: vi.fn(async () => undefined),
mkdtemp: vi.fn(async (prefix: string) => `${prefix}${++refreshFixture.next}`),
};
});
function defaultGitOutput(command: string, cwd?: string): string {
if (command.includes("worktree list --porcelain")) return "";
if (command.includes(" config --get branch.")) return "origin\n";
if (command.endsWith(" remote")) return "origin\n";
if (command.includes("rev-parse --show-toplevel")) return cwd ?? "";
if (command.includes("symbolic-ref -q HEAD")) return "refs/heads/fusion/fn-9601\n";
if (command.includes("rev-parse HEAD")) return "1111111111111111111111111111111111111111\n";
if (command.includes("rev-parse --verify refs/heads/")) return "0000000000000000000000000000000000000000\n";
if (command.includes("ls-remote origin refs/heads/")) return "1111111111111111111111111111111111111111\trefs/heads/test\n";
return "";
}
vi.mock("node:child_process", () => ({ vi.mock("node:child_process", () => ({
exec: (cmd: string, opts: unknown, cb: (err: Error | null, stdout: string, stderr: string) => void) => { exec: (cmd: string, opts: unknown, cb: (err: Error | null, stdout: string, stderr: string) => void) => {
try { try {
const result = execMock(cmd, opts); const result = execMock(cmd, opts);
cb(null, typeof result === "string" ? result : "", ""); cb(null, typeof result === "string" && result ? result : defaultGitOutput(cmd, (opts as { cwd?: string } | undefined)?.cwd), "");
} catch (err) { } catch (err) {
cb(err as Error, "", (err as Error).message); cb(err as Error, "", (err as Error).message);
} }
@@ -22,8 +43,17 @@ vi.mock("node:child_process", () => ({
execFile: (file: string, args: string[] | undefined, opts: unknown, cb: (err: Error | null, stdout: string, stderr: string) => void) => { execFile: (file: string, args: string[] | undefined, opts: unknown, cb: (err: Error | null, stdout: string, stderr: string) => void) => {
try { try {
execFileCalls.push({ file, args: args ?? [], cwd: (opts as { cwd?: string } | undefined)?.cwd }); execFileCalls.push({ file, args: args ?? [], cwd: (opts as { cwd?: string } | undefined)?.cwd });
const result = execMock(`${file} ${(args ?? []).join(" ")}`.trim(), opts); if (file === "git" && args?.[0] === "worktree" && args[1] === "add") {
cb(null, typeof result === "string" ? result : "", ""); const ref = args.at(-1);
if (ref?.startsWith("refs/heads/")) refreshFixture.branch = ref.slice("refs/heads/".length);
else if (ref?.startsWith("fusion/")) refreshFixture.branch = ref;
}
const command = `${file} ${(args ?? []).join(" ")}`.trim();
const result = execMock(command, opts);
const output = file === "git" && args?.[0] === "symbolic-ref"
? `refs/heads/${refreshFixture.branch}\n`
: typeof result === "string" && result ? result : defaultGitOutput(command, (opts as { cwd?: string } | undefined)?.cwd);
cb(null, output, "");
} catch (err) { } catch (err) {
cb(err as Error, "", (err as Error).message); cb(err as Error, "", (err as Error).message);
} }
@@ -36,16 +66,23 @@ vi.mock("@fusion/core", async () => {
...actual, ...actual,
getCurrentRepo: vi.fn(() => ({ owner: "owner", repo: "repo" })), getCurrentRepo: vi.fn(() => ({ owner: "owner", repo: "repo" })),
getPushRepo: vi.fn(() => ({ owner: "owner", repo: "repo" })), getPushRepo: vi.fn(() => ({ owner: "owner", repo: "repo" })),
acquireWorktreePathReservation: vi.fn(async (input: { canonicalPath: string }) => ({
canonicalPath: input.canonicalPath,
state: "held" as const,
release: vi.fn(async () => undefined),
quarantine: vi.fn(async () => undefined),
})),
}; };
}); });
import { getCurrentRepo, getPushRepo } from "@fusion/core"; import { acquireWorktreePathReservation, getCurrentRepo, getPushRepo } from "@fusion/core";
import { activeSessionRegistry } from "@fusion/engine"; import { activeSessionRegistry } from "@fusion/engine";
import { import {
cleanupMergedTaskArtifacts, cleanupMergedTaskArtifacts,
createGroupPrCallback, createGroupPrCallback,
createPrNodeGithubOps, createPrNodeGithubOps,
processPullRequestMergeTask, processPullRequestMergeTask,
refreshAutomatedPrHead,
getTaskBranchName, getTaskBranchName,
syncGroupPrCallback, syncGroupPrCallback,
} from "../task-lifecycle.js"; } from "../task-lifecycle.js";
@@ -1193,6 +1230,7 @@ describe("processPullRequestMergeTask", () => {
description: "desc", description: "desc",
column: "in-review", column: "in-review",
worktree: "/tmp/worktree-fn-9104", worktree: "/tmp/worktree-fn-9104",
mergeRetries: 2,
prInfo: { prInfo: {
number: 124, number: 124,
url: "https://github.com/x/y/pull/124", url: "https://github.com/x/y/pull/124",
@@ -1248,11 +1286,12 @@ describe("processPullRequestMergeTask", () => {
); );
expect(result).toBe("merged"); expect(result).toBe("merged");
expect(github.mergePr).toHaveBeenCalledWith({ owner: "owner", repo: "repo", number: 124, method: "squash" }); expect(github.mergePr).toHaveBeenCalledWith({ owner: "owner", repo: "repo", number: 124, method: "squash", expectedHeadOid: "1111111111111111111111111111111111111111" });
expect(github.getPrMergeStatus).toHaveBeenCalledTimes(2); expect(github.getPrMergeStatus).toHaveBeenCalledTimes(2);
expect(github.getPrMergeStatus).toHaveBeenNthCalledWith(1, "owner", "repo", 124); expect(github.getPrMergeStatus).toHaveBeenNthCalledWith(1, "owner", "repo", 124);
expect(github.getPrMergeStatus).toHaveBeenNthCalledWith(2, "owner", "repo", 124); expect(github.getPrMergeStatus).toHaveBeenNthCalledWith(2, "owner", "repo", 124);
expect(store.updatePrInfo).toHaveBeenLastCalledWith("FN-9104", expect.objectContaining({ status: "merged" })); expect(store.updatePrInfo).toHaveBeenLastCalledWith("FN-9104", expect.objectContaining({ status: "merged" }));
expect(store.updateTask).toHaveBeenCalledWith("FN-9104", { mergeRetries: 0 });
expect(store.updateTask).toHaveBeenCalledWith("FN-9104", { status: null, mergeRetries: 0 }); expect(store.updateTask).toHaveBeenCalledWith("FN-9104", { status: null, mergeRetries: 0 });
expect(store.moveTask).toHaveBeenCalledWith("FN-9104", "done"); expect(store.moveTask).toHaveBeenCalledWith("FN-9104", "done");
expect(store.logEntry).toHaveBeenCalledWith( expect(store.logEntry).toHaveBeenCalledWith(
@@ -1319,7 +1358,7 @@ describe("processPullRequestMergeTask", () => {
), ),
).rejects.toThrow(mergeError.message); ).rejects.toThrow(mergeError.message);
expect(github.mergePr).toHaveBeenCalledWith({ owner: "owner", repo: "repo", number: 125, method: "squash" }); expect(github.mergePr).toHaveBeenCalledWith({ owner: "owner", repo: "repo", number: 125, method: "squash", expectedHeadOid: "1111111111111111111111111111111111111111" });
expect(github.getPrMergeStatus).toHaveBeenCalledTimes(2); expect(github.getPrMergeStatus).toHaveBeenCalledTimes(2);
expect(store.updatePrInfo).not.toHaveBeenCalledWith("FN-9105", expect.objectContaining({ status: "merged" })); expect(store.updatePrInfo).not.toHaveBeenCalledWith("FN-9105", expect.objectContaining({ status: "merged" }));
expect(store.moveTask).not.toHaveBeenCalled(); expect(store.moveTask).not.toHaveBeenCalled();
@@ -1431,7 +1470,7 @@ describe("processPullRequestMergeTask", () => {
), ),
).rejects.toThrow(mergeError.message); ).rejects.toThrow(mergeError.message);
expect(github.mergePr).toHaveBeenCalledWith({ owner: "owner", repo: "repo", number: 126, method: "squash" }); expect(github.mergePr).toHaveBeenCalledWith({ owner: "owner", repo: "repo", number: 126, method: "squash", expectedHeadOid: "1111111111111111111111111111111111111111" });
expect(github.getPrMergeStatus).toHaveBeenCalledTimes(2); expect(github.getPrMergeStatus).toHaveBeenCalledTimes(2);
expect(store.moveTask).not.toHaveBeenCalled(); expect(store.moveTask).not.toHaveBeenCalled();
}); });
@@ -1598,7 +1637,7 @@ describe("processPullRequestMergeTask", () => {
); );
expect(result).toBe("merged"); expect(result).toBe("merged");
expect(github.mergePr).toHaveBeenCalledWith({ owner: "owner", repo: "repo", number: 100, method: "squash" }); expect(github.mergePr).toHaveBeenCalledWith({ owner: "owner", repo: "repo", number: 100, method: "squash", expectedHeadOid: "1111111111111111111111111111111111111111" });
}); });
it("preserves existing behavior when requirePrApproval is false", async () => { it("preserves existing behavior when requirePrApproval is false", async () => {
@@ -1862,6 +1901,46 @@ describe("createGroupPrCallback", () => {
}); });
}); });
it("passes the configured integration remote to the group-head refresh", async () => {
const github = {
findPrForBranch: vi.fn(async () => null),
createPr: vi.fn(async () => ({ number: 98, url: "https://github.com/owner/repo/pull/98", status: "open" as const })),
};
await createGroupPrCallback(github as never)({
cwd: "/repo",
group: group as never,
members,
headBranch: group.branchName,
baseBranch: "main",
integrationRemote: "upstream",
});
expect(execFileCalls).toContainEqual(expect.objectContaining({
file: "git",
args: ["fetch", "upstream", "main"],
}));
});
it("fails closed when promotion is cancelled before its refresh boundary", async () => {
const controller = new AbortController();
controller.abort(new Error("promotion cancelled"));
const github = {
findPrForBranch: vi.fn(async () => null),
createPr: vi.fn(),
};
await expect(createGroupPrCallback(github as never)({
cwd: "/repo",
group: group as never,
members,
headBranch: group.branchName,
baseBranch: "main",
signal: controller.signal,
})).rejects.toThrow("promotion cancelled");
expect(github.createPr).not.toHaveBeenCalled();
});
it("does not reuse a closed PR from a prior group — creates a fresh one", async () => { it("does not reuse a closed PR from a prior group — creates a fresh one", async () => {
// With state:"open", findPrForBranch returns null for a head whose only PR // With state:"open", findPrForBranch returns null for a head whose only PR
// is closed/merged, so the create path runs instead of resurrecting the // is closed/merged, so the create path runs instead of resurrecting the
@@ -1967,6 +2046,160 @@ describe("createGroupPrCallback", () => {
}); });
}); });
describe("refreshAutomatedPrHead", () => {
beforeEach(() => {
execMock.mockReset();
execFileCalls.length = 0;
vi.mocked(acquireWorktreePathReservation).mockImplementation(async (input: { canonicalPath: string }) => ({
canonicalPath: input.canonicalPath,
state: "held" as const,
release: vi.fn(async () => undefined),
quarantine: vi.fn(async () => undefined),
}) as never);
});
it("refuses a cancelled refresh before any git mutation", async () => {
const controller = new AbortController();
controller.abort(new Error("execution cancelled"));
await expect(refreshAutomatedPrHead({
projectRoot: "/projects/repo-a",
headBranch: "fusion/fn-cancelled",
targetBranch: "main",
signal: controller.signal,
})).rejects.toThrow("execution cancelled");
expect(execFileCalls).toEqual([]);
});
it("cleans a worktree created during cancellation before returning the failure", async () => {
refreshFixture.next = 0;
let addAttempted = false;
execMock.mockImplementation((command: string) => {
if (command.startsWith("git worktree add")) {
addAttempted = true;
throw new Error("The operation was aborted");
}
if (addAttempted && command === "git worktree list --porcelain") {
return "worktree /projects/repo-a/.worktrees/pr-refresh-49e2094e7f117641\nbranch refs/heads/fusion/fn-cancel-cleanup\n";
}
return "";
});
await expect(refreshAutomatedPrHead({
projectRoot: "/projects/repo-a",
headBranch: "fusion/fn-cancel-cleanup",
targetBranch: "main",
})).rejects.toThrow("no verified checkout");
expect(execFileCalls).toContainEqual(expect.objectContaining({
file: "git",
args: ["worktree", "remove", "--force", expect.stringMatching(/\/projects\/repo-a\/\.worktrees\/pr-refresh-/)],
}));
});
it("quarantines a retained temporary checkout when cleanup fails", async () => {
const reservation = {
canonicalPath: "/projects/repo-a/.worktrees/pr-refresh-test",
state: "held" as const,
release: vi.fn(async () => undefined),
quarantine: vi.fn(async () => undefined),
};
vi.mocked(acquireWorktreePathReservation).mockResolvedValue(reservation as never);
let worktreeAdded = false;
execMock.mockImplementation((command: string) => {
if (command.startsWith("git worktree add")) worktreeAdded = true;
if (worktreeAdded && command === "git worktree list --porcelain") {
return "worktree /projects/repo-a/.worktrees/pr-refresh-f7ab243dfb531960\nbranch refs/heads/fusion/fn-cleanup-failure\n";
}
if (command.startsWith("git worktree remove --force")) throw new Error("device busy");
return "";
});
await expect(refreshAutomatedPrHead({
projectRoot: "/projects/repo-a",
headBranch: "fusion/fn-cleanup-failure",
targetBranch: "main",
})).rejects.toThrow("cleanup failed");
expect(reservation.quarantine).toHaveBeenCalledWith("device busy");
expect(reservation.release).not.toHaveBeenCalled();
});
it("does not roll back a completed push when cancellation arrives after publication", async () => {
const controller = new AbortController();
let headReads = 0;
execMock.mockImplementation((command: string) => {
if (command === "git rev-parse HEAD") {
headReads += 1;
return headReads === 1
? "1111111111111111111111111111111111111111\n"
: "2222222222222222222222222222222222222222\n";
}
if (command.startsWith("git push --force-with-lease=")) controller.abort(new Error("cancelled after push"));
return "";
});
await expect(refreshAutomatedPrHead({
projectRoot: "/projects/repo-a",
headBranch: "fusion/fn-cancel-after-push",
targetBranch: "main",
signal: controller.signal,
})).resolves.toEqual(expect.objectContaining({
headOid: "2222222222222222222222222222222222222222",
refreshed: true,
}));
expect(execFileCalls.some((call) => call.args[0] === "update-ref")).toBe(false);
});
it("uses the configured integration remote rather than falling back to origin", async () => {
await refreshAutomatedPrHead({
projectRoot: "/projects/repo-a",
headBranch: "fusion/fn-8838",
targetBranch: "main",
integrationRemote: "upstream",
});
expect(execFileCalls).toContainEqual(expect.objectContaining({
file: "git",
args: ["fetch", "upstream", "main"],
}));
});
it("materializes a remote-only head before attaching an isolated worktree", async () => {
let remoteHeadLookups = 0;
execMock.mockImplementation((command: string) => {
if (command === "git rev-parse HEAD") return "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa\n";
if (command === "git rev-parse --verify refs/heads/fusion/fn-remote") {
if (remoteHeadLookups++ === 0) {
const missing = Object.assign(new Error("missing ref"), { code: 128 });
throw missing;
}
return "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa\n";
}
if (command === "git ls-remote origin refs/heads/fusion/fn-remote") return "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa\trefs/heads/fusion/fn-remote\n";
if (command === "git rev-parse --verify refs/heads/main") return "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb\n";
return "";
});
await refreshAutomatedPrHead({
projectRoot: "/projects/repo-a",
headBranch: "fusion/fn-remote",
targetBranch: "main",
});
expect(execFileCalls).toContainEqual(expect.objectContaining({
file: "git",
args: ["fetch", "origin", "refs/heads/fusion/fn-remote:refs/heads/fusion/fn-remote"],
cwd: "/projects/repo-a",
}));
expect(execFileCalls).toContainEqual(expect.objectContaining({
file: "git",
args: ["worktree", "add", expect.any(String), "fusion/fn-remote"],
}));
});
});
describe("createPrNodeGithubOps repo resolution (gh-4)", () => { describe("createPrNodeGithubOps repo resolution (gh-4)", () => {
beforeEach(() => { beforeEach(() => {
execMock.mockReset(); execMock.mockReset();
@@ -2055,15 +2288,71 @@ describe("createPrNodeGithubOps repo resolution (gh-4)", () => {
); );
}); });
it("createPr uses the engine-provided configured integration remote", async () => {
const github = githubStub();
const ops = createPrNodeGithubOps(github as never);
await ops.createPr({
task: { id: "FN-9601", title: "t", description: "d", worktree: "/projects/repo-a/.worktrees/fn-9601" },
entity: { id: "e1", sourceId: "FN-9601", repo: "central-owner/central-repo", headBranch: "fusion/fn-9601", baseBranch: "main" },
integrationRemote: "upstream",
} as never);
expect(execFileCalls).toContainEqual(expect.objectContaining({
file: "git",
args: ["fetch", "upstream", "main"],
}));
});
it("fails closed when workflow PR creation is cancelled", async () => {
const controller = new AbortController();
controller.abort(new Error("workflow cancelled"));
const github = githubStub();
await expect(createPrNodeGithubOps(github as never).createPr({
task: { id: "FN-9601", title: "t", description: "d", worktree: "/projects/repo-a/.worktrees/fn-9601" },
entity: { id: "e1", sourceId: "FN-9601", repo: "central-owner/central-repo", headBranch: "fusion/fn-9601", baseBranch: "main" },
signal: controller.signal,
} as never)).rejects.toThrow("workflow cancelled");
expect(github.createPr).not.toHaveBeenCalled();
});
it("mergePr passes owner/repo parsed from entity.repo", async () => { it("mergePr passes owner/repo parsed from entity.repo", async () => {
const github = githubStub(); const github = githubStub();
const ops = createPrNodeGithubOps(github as never); const ops = createPrNodeGithubOps(github as never);
const result = await ops.mergePr({ const result = await ops.mergePr({
entity: { id: "e1", sourceId: "FN-9601", repo: "central-owner/central-repo", prNumber: 9, headOid: "abc123" }, entity: { id: "e1", sourceId: "FN-9601", repo: "central-owner/central-repo", prNumber: 9, headOid: "abc123" },
} as never); } as never);
expect(result).toEqual({ status: "merged-requested" }); expect(result).toEqual({ status: "merged-requested", headOid: "1111111111111111111111111111111111111111" });
expect(github.mergePr).toHaveBeenCalledWith( expect(github.mergePr).toHaveBeenCalledWith(
expect.objectContaining({ owner: "central-owner", repo: "central-repo", number: 9, method: "squash", expectedHeadOid: "abc123" }), expect.objectContaining({ owner: "central-owner", repo: "central-repo", number: 9, method: "squash", expectedHeadOid: "1111111111111111111111111111111111111111" }),
); );
}); });
it("fails closed when workflow PR merge is cancelled", async () => {
const controller = new AbortController();
controller.abort(new Error("workflow merge cancelled"));
const github = githubStub();
await expect(createPrNodeGithubOps(github as never).mergePr({
entity: { id: "e1", sourceId: "FN-9601", repo: "central-owner/central-repo", prNumber: 9, headOid: "abc123" },
signal: controller.signal,
} as never)).rejects.toThrow("workflow merge cancelled");
expect(github.mergePr).not.toHaveBeenCalled();
});
it("persists the refreshed workflow head before requesting GitHub merge", async () => {
const calls: string[] = [];
const github = githubStub();
github.mergePr.mockImplementation(async () => {
calls.push("merge");
return { number: 9, url: "https://github.com/central-owner/central-repo/pull/9", status: "merged" };
});
const ops = createPrNodeGithubOps(github as never);
await ops.mergePr({
entity: { id: "e1", sourceId: "FN-9601", repo: "central-owner/central-repo", prNumber: 9, headOid: "old" },
persistRefreshedHead: async () => { calls.push("persist"); },
} as never);
expect(calls).toEqual(["persist", "merge"]);
});
}); });

View File

@@ -365,8 +365,8 @@ export async function runDaemon(opts: DaemonOptions = {}) {
onMigrationProgress: (event) => migrationHoldingServer?.setMigrationProgress(event), onMigrationProgress: (event) => migrationHoldingServer?.setMigrationProgress(event),
cliPackageVersion, cliPackageVersion,
getMergeStrategy, getMergeStrategy,
processPullRequestMerge: (s, wd, taskId, pool) => processPullRequestMerge: (s, wd, taskId, pool, signal) =>
processPullRequestMergeTask(s, wd, taskId, githubClient, getTaskMergeBlocker, pool), processPullRequestMergeTask(s, wd, taskId, githubClient, getTaskMergeBlocker, pool, signal),
createGroupPr: createGroupPrCallback(githubClient), createGroupPr: createGroupPrCallback(githubClient),
syncGroupPr: syncGroupPrCallback(githubClient), syncGroupPr: syncGroupPrCallback(githubClient),
prNodeGithubOps: createPrNodeGithubOps(githubClient), prNodeGithubOps: createPrNodeGithubOps(githubClient),

View File

@@ -2135,8 +2135,8 @@ export async function runDashboard(port: number, opts: { paused?: boolean; dev?:
const engineManager = new ProjectEngineManager(centralCoreForEngine, { const engineManager = new ProjectEngineManager(centralCoreForEngine, {
cliPackageVersion, cliPackageVersion,
getMergeStrategy, getMergeStrategy,
processPullRequestMerge: (s, wd, taskId, pool) => processPullRequestMerge: (s, wd, taskId, pool, signal) =>
processPullRequestMergeTask(s, wd, taskId, githubClient, getTaskMergeBlocker, pool), processPullRequestMergeTask(s, wd, taskId, githubClient, getTaskMergeBlocker, pool, signal),
createGroupPr: createGroupPrCallback(githubClient), createGroupPr: createGroupPrCallback(githubClient),
syncGroupPr: syncGroupPrCallback(githubClient), syncGroupPr: syncGroupPrCallback(githubClient),
prNodeGithubOps: createPrNodeGithubOps(githubClient), prNodeGithubOps: createPrNodeGithubOps(githubClient),

View File

@@ -425,8 +425,8 @@ export async function runServe(
const engineManager = startupEngineManager = new ProjectEngineManager(sharedCentralCore, { const engineManager = startupEngineManager = new ProjectEngineManager(sharedCentralCore, {
cliPackageVersion, cliPackageVersion,
getMergeStrategy, getMergeStrategy,
processPullRequestMerge: (s, wd, taskId, pool) => processPullRequestMerge: (s, wd, taskId, pool, signal) =>
processPullRequestMergeTask(s, wd, taskId, githubClient, getTaskMergeBlocker, pool), processPullRequestMergeTask(s, wd, taskId, githubClient, getTaskMergeBlocker, pool, signal),
createGroupPr: createGroupPrCallback(githubClient), createGroupPr: createGroupPrCallback(githubClient),
syncGroupPr: syncGroupPrCallback(githubClient), syncGroupPr: syncGroupPrCallback(githubClient),
prNodeGithubOps: createPrNodeGithubOps(githubClient), prNodeGithubOps: createPrNodeGithubOps(githubClient),

View File

@@ -13,8 +13,11 @@
* - Full PR lifecycle orchestration (create → status check → merge) * - Full PR lifecycle orchestration (create → status check → merge)
*/ */
import { createHash } from "node:crypto";
import { exec } from "node:child_process"; import { exec } from "node:child_process";
import * as childProcess from "node:child_process"; import * as childProcess from "node:child_process";
import { mkdir, realpath, rm } from "node:fs/promises";
import { join } from "node:path";
import { promisify } from "node:util"; import { promisify } from "node:util";
const execAsync = promisify(exec); const execAsync = promisify(exec);
// `execFile` is resolved lazily through the namespace import so test mocks that // `execFile` is resolved lazily through the namespace import so test mocks that
@@ -34,6 +37,8 @@ import {
assertNotWorkspaceTaskMerge, assertNotWorkspaceTaskMerge,
classifyGhError, classifyGhError,
WorkspaceTaskMergeError, WorkspaceTaskMergeError,
acquireWorktreePathReservation,
type WorktreePathReservation,
} from "@fusion/core"; } from "@fusion/core";
import type { Settings, TaskDetail, PrInfo, MergeResult, BranchGroup, BranchGroupPrState, Task } from "@fusion/core"; import type { Settings, TaskDetail, PrInfo, MergeResult, BranchGroup, BranchGroupPrState, Task } from "@fusion/core";
import { resolveWorkflowIrForTask, resolveCompleteColumn, resolveMergeOrchestrationColumn } from "@fusion/core"; import { resolveWorkflowIrForTask, resolveCompleteColumn, resolveMergeOrchestrationColumn } from "@fusion/core";
@@ -65,7 +70,7 @@ export async function resolveCompleteTargetForTask(store: TaskStore, taskId: str
} catch { /* degraded: legacy id */ } } catch { /* degraded: legacy id */ }
return "done"; return "done";
} }
import { activeSessionRegistry, resolveIntegrationBranch } from "@fusion/engine"; import { activeSessionRegistry, resolveIntegrationBranch, resolveIntegrationRemote } from "@fusion/engine";
import type { import type {
CreateGroupPrFn, CreateGroupPrFn,
SyncGroupPrFn, SyncGroupPrFn,
@@ -187,7 +192,392 @@ async function gitCommandSucceeds(
} }
} }
async function pushTaskBranchToOrigin(cwd: string, branch: string): Promise<void> { /*
FNXC:PullRequestFreshness 2026-08-09-01:17:
Automated PR creation and merge must never rely on the creation-time base. Refresh
only an exact-head checkout, publish any rewritten head with a lease, and complete
temporary-checkout cleanup before a GitHub mutation is allowed.
*/
export interface RefreshAutomatedPrHeadInput {
projectRoot: string;
preferredWorktree?: string;
headBranch: string;
targetBranch: string;
/** Explicit project policy wins over branch/default remote discovery. */
integrationRemote?: string;
/** Cancels git work without allowing a later GitHub mutation. */
signal?: AbortSignal;
}
export interface RefreshAutomatedPrHeadResult {
headOid: string;
refreshed: boolean;
}
function parseWorktreeBranches(output: string): Array<{ path: string; branch?: string }> {
const entries: Array<{ path: string; branch?: string }> = [];
let entry: { path: string; branch?: string } | undefined;
for (const line of output.split("\n")) {
if (line.startsWith("worktree ")) {
entry = { path: line.slice("worktree ".length) };
entries.push(entry);
} else if (line.startsWith("branch ") && entry) {
entry.branch = line.slice("branch ".length).trim();
}
}
return entries;
}
async function gitStdout(cwd: string, args: string[]): Promise<string> {
const result = await execFileAsync("git", args, { cwd, timeout: 60_000, encoding: "utf-8" }) as unknown;
// Node's execFile promisify custom returns `{ stdout, stderr }`; lightweight
// embedders may expose the ordinary promisify string result instead.
const stdout = typeof result === "string"
? result
: (result as { stdout?: string }).stdout ?? "";
return stdout.trim();
}
async function abortRebase(cwd: string): Promise<void> {
try {
await execFileAsync("git", ["rebase", "--abort"], { cwd, timeout: 30_000 });
} catch {
// There may be no active rebase; never obscure the refresh failure with abort cleanup.
}
}
/**
* Restore the canonical head only when it is still the exact ref advanced by
* this refresh. `update-ref <new> <old>` is the CAS fence: a concurrent local
* writer is never overwritten while recovering from a rejected publication.
*/
async function restoreRefreshHead(
root: string,
checkout: string,
ref: string,
before: string | undefined,
): Promise<void> {
if (!before) return;
const current = await gitStdout(root, ["rev-parse", "--verify", ref]).catch(() => "");
if (!current || current === before) return;
await execFileAsync("git", ["update-ref", ref, before, current], { cwd: root, timeout: 30_000 });
// The symbolic checkout now resolves to `before`; reset its index and tree
// without changing the ref again. This is deliberately not cancellable.
const restoredHead = await gitStdout(checkout, ["rev-parse", "HEAD"]);
if (restoredHead !== before) {
throw new Error(`PR head refresh rollback refused: ${ref} changed during recovery`);
}
await execFileAsync("git", ["reset", "--hard", before], { cwd: checkout, timeout: 30_000 });
}
/*
FNXC:PullRequestFreshness 2026-08-09-02:32:
Cancellation is a fail-closed PR boundary. Check it before every mutating git
operation and pass it to child processes, but never pass it to rebase abort or
worktree cleanup: those repairs must finish after the owning run is cancelled.
*/
export function throwIfRefreshAborted(signal: AbortSignal | undefined): void {
if (signal?.aborted) {
throw signal.reason instanceof Error ? signal.reason : new Error("PR head refresh cancelled");
}
}
async function runRefreshGit(
cwd: string,
args: string[],
signal: AbortSignal | undefined,
timeout = 60_000,
checkAfter = true,
): Promise<void> {
throwIfRefreshAborted(signal);
await execFileAsync("git", args, { cwd, timeout, signal });
if (checkAfter) throwIfRefreshAborted(signal);
}
/*
FNXC:PullRequestFreshness 2026-08-09-04:57:
Retry retained cleanup under a fresh reservation claim. The reconciliation is
bounded and path-specific, so it never sweeps an OS temp directory or steals an
active checkout.
*/
async function removeInactiveRetainedRefreshWorktree(root: string, path: string): Promise<void> {
if (activeSessionRegistry.isPathActive(path)) {
throw new Error(`PR head refresh cleanup remains active at ${path}`);
}
const cleanupError = await removeRefreshWorktree(root, path);
if (cleanupError) throw cleanupError;
}
async function reconcileRetainedRefreshWorktree(
root: string,
path: string,
reservationDir: string,
): Promise<void> {
const reservation = await acquireWorktreePathReservation({
canonicalPath: path,
worktreesDir: reservationDir,
rootDir: root,
isLiveWorktree: async (candidate) => activeSessionRegistry.isPathActive(candidate),
reconcileQuarantined: async (candidate) => removeInactiveRetainedRefreshWorktree(root, candidate),
});
await reservation.release();
}
/*
FNXC:PullRequestFreshness 2026-08-09-04:57:
A cleanup quarantine must be retried without waiting for another PR on the same
branch. Keep retries bounded and retain the durable reservation when they fail,
so a stopped daemon still fails closed and a later refresh can reconcile it.
*/
function scheduleRetainedRefreshReconciliation(root: string, path: string, reservationDir: string): void {
const maxAttempts = 3;
let attempt = 0;
const retry = () => {
attempt += 1;
void reconcileRetainedRefreshWorktree(root, path, reservationDir).catch(() => {
if (attempt >= maxAttempts) return;
const timer = setTimeout(retry, attempt * 1_000);
timer.unref?.();
});
};
const timer = setTimeout(retry, 1_000);
timer.unref?.();
}
async function removeRefreshWorktree(root: string, path: string): Promise<Error | undefined> {
try {
const registered = parseWorktreeBranches(await gitStdout(root, ["worktree", "list", "--porcelain"]))
.some((entry) => entry.path === path);
if (registered) {
await execFileAsync("git", ["worktree", "remove", "--force", path], { cwd: root, timeout: 60_000 });
}
await rm(path, { recursive: true, force: true });
return undefined;
} catch (error) {
return error instanceof Error ? error : new Error(String(error));
}
}
/**
* Refresh an automated PR head immediately before a GitHub boundary.
*
* The project root is deliberately never a rebase cwd. A matching task checkout
* is verified against `git worktree list --porcelain`; otherwise a managed,
* detached worktree is created for the local head and removed before returning.
*/
export async function refreshAutomatedPrHead(
input: RefreshAutomatedPrHeadInput,
): Promise<RefreshAutomatedPrHeadResult> {
throwIfRefreshAborted(input.signal);
const root = await realpath(input.projectRoot);
const requestedRef = `refs/heads/${input.headBranch}`;
const listed = await gitStdout(root, ["worktree", "list", "--porcelain"]);
const entries = parseWorktreeBranches(listed);
const primaryWorktree = entries[0] ? await realpath(entries[0].path).catch(() => null) : null;
const preferred = input.preferredWorktree ? await realpath(input.preferredWorktree).catch(() => null) : null;
const candidates = await Promise.all(entries
.filter((entry) => entry.branch === requestedRef)
.map(async (entry) => ({ entry, canonical: await realpath(entry.path).catch(() => null) })));
const rootOwnsHead = candidates.some((candidate) => candidate.canonical === primaryWorktree);
if (rootOwnsHead) {
/*
FNXC:PullRequestFreshness 2026-08-09-02:01:
Automated refresh must never alter an operator's primary checkout. Refuse a
head checked out at the project root rather than using a detached copy that
would publish a rewrite while leaving the canonical local branch stale.
*/
throw new Error(`PR head refresh refused: ${input.headBranch} is checked out in the project root`);
}
const exact = candidates.filter((candidate) => candidate.canonical && candidate.canonical !== primaryWorktree);
let checkout: string | undefined = preferred && exact.find((candidate) => candidate.canonical === preferred)?.canonical || undefined;
if (!checkout && exact.length === 1) checkout = exact[0].canonical ?? undefined;
if (!checkout && exact.length > 1) {
throw new Error(`PR head refresh refused: branch ${input.headBranch} is checked out in multiple worktrees`);
}
let temporary: string | undefined;
let temporaryReservation: WorktreePathReservation | undefined;
if (!checkout) {
/*
FNXC:PullRequestFreshness 2026-08-09-02:01:
A remote-only automated head must be materialized at an explicit local ref
before a temporary worktree attaches it. Compare the fetched OID to the
observed remote OID so a racing remote update cannot rebase an unknown tip.
*/
const localHead = await gitStdout(root, ["rev-parse", "--verify", requestedRef]).catch(() => "");
if (!localHead) {
const remoteHead = await gitStdout(root, ["ls-remote", "origin", requestedRef]);
const observedRemoteHead = remoteHead.split(/\s+/)[0] || "";
if (!observedRemoteHead) {
throw new Error(`PR head refresh refused: missing local and origin head ${input.headBranch}`);
}
await runRefreshGit(root, ["fetch", "origin", `${requestedRef}:${requestedRef}`], input.signal);
const fetchedHead = await gitStdout(root, ["rev-parse", "--verify", requestedRef]);
if (fetchedHead !== observedRemoteHead) {
throw new Error(`PR head refresh refused: origin head ${input.headBranch} changed during fetch`);
}
}
/*
FNXC:PullRequestFreshness 2026-08-09-03:20:
Attach (rather than detach) the exact local branch. Rebase then advances the
canonical ref, so the later guarded push cannot be followed by a stale ordinary
`git push -u origin <branch>`; a full refname would detach this checkout.
*/
// Keep retained failed-cleanup worktrees under the bounded project worktree
// area, where the existing native worktree reconciliation can discover them.
const refreshWorktreesDir = join(root, ".worktrees");
await mkdir(refreshWorktreesDir, { recursive: true });
/*
FNXC:PullRequestFreshness 2026-08-09-02:44:
Refresh worktrees use a deterministic managed path and a durable path reservation.
A failed removal quarantines that path; a later refresh reconciles the exact
inactive checkout before it can reuse the branch or path.
*/
temporary = join(
refreshWorktreesDir,
`pr-refresh-${createHash("sha256").update(input.headBranch).digest("hex").slice(0, 16)}`,
);
const reservationDir = join(root, ".fusion", "pr-refresh-reservations");
await mkdir(reservationDir, { recursive: true });
temporaryReservation = await acquireWorktreePathReservation({
canonicalPath: temporary,
worktreesDir: reservationDir,
rootDir: root,
isLiveWorktree: async (path) => activeSessionRegistry.isPathActive(path),
reconcileQuarantined: async (path) => removeInactiveRetainedRefreshWorktree(root, path),
});
try {
await runRefreshGit(root, ["worktree", "add", temporary, input.headBranch], input.signal);
checkout = temporary;
} catch (error) {
// `execFile` can observe cancellation after git created the worktree. Clean
// both forms without the cancelled signal before surfacing the primary error.
const cleanupError = await removeRefreshWorktree(root, temporary);
if (cleanupError) await temporaryReservation.quarantine(cleanupError.message);
else await temporaryReservation.release();
throw new Error(`PR head refresh refused: no verified checkout for ${input.headBranch}: ${error instanceof Error ? error.message : String(error)}`);
}
}
let primaryError: unknown;
let refreshStartHead: string | undefined;
let successfulResult: RefreshAutomatedPrHeadResult | undefined;
let cleanupFailure: unknown;
try {
const actualRoot = await realpath(await gitStdout(checkout, ["rev-parse", "--show-toplevel"]));
if (actualRoot !== checkout || checkout === primaryWorktree) {
throw new Error(`PR head refresh refused: ${input.headBranch} checkout is not an isolated worktree`);
}
const actualRef = await gitStdout(checkout, ["symbolic-ref", "-q", "HEAD"]);
if (actualRef !== requestedRef) {
throw new Error(`PR head refresh refused: checkout does not own ${requestedRef}`);
}
const remote = input.integrationRemote?.trim() || await resolveIntegrationRemote({
settings: { worktreeRebaseRemote: "" },
rootDir: root,
integrationBranch: input.targetBranch,
});
if (!remote) throw new Error(`PR head refresh refused: no integration remote for ${input.targetBranch}`);
await runRefreshGit(checkout, ["fetch", remote, input.targetBranch], input.signal);
const remoteTarget = `${remote}/${input.targetBranch}`;
const before = await gitStdout(checkout, ["rev-parse", "HEAD"]);
refreshStartHead = before;
const needsRemoteRebase = !(await gitCommandSucceeds(checkout, "git", ["merge-base", "--is-ancestor", remoteTarget, "HEAD"], 1));
if (needsRemoteRebase) await runRefreshGit(checkout, ["rebase", remoteTarget], input.signal);
const localTarget = await gitStdout(root, ["rev-parse", "--verify", `refs/heads/${input.targetBranch}`]).catch(() => "");
if (localTarget) {
const needsLocalRebase = !(await gitCommandSucceeds(checkout, "git", ["merge-base", "--is-ancestor", localTarget, "HEAD"], 1));
if (needsLocalRebase) await runRefreshGit(checkout, ["rebase", localTarget], input.signal);
}
const observedRemote = await gitStdout(checkout, ["ls-remote", "origin", requestedRef]);
const observedOid = observedRemote.split(/\s+/)[0] || "";
/*
FNXC:PullRequestFreshness 2026-08-09-05:32:
A refresh may rewrite only the exact head it began from. A remote head that
advanced before the lease observation contains work this refresh did not
incorporate, so fail closed rather than force-pushing over it.
*/
if (observedOid && observedOid !== before) {
throw new Error(`PR head refresh refused: origin head ${input.headBranch} changed before publication`);
}
const after = await gitStdout(checkout, ["rev-parse", "HEAD"]);
if (after !== before) {
const pushArgs = observedOid
? ["push", `--force-with-lease=${requestedRef}:${observedOid}`, "origin", `HEAD:${requestedRef}`]
: ["push", "origin", `HEAD:${requestedRef}`];
/*
FNXC:PullRequestFreshness 2026-08-09-04:57:
Do not turn a completed push into a rollback merely because cancellation
arrived in the tiny post-publication window. The caller checks cancellation
again before GitHub, while local and remote heads remain identical.
*/
await runRefreshGit(checkout, pushArgs, input.signal, 60_000, false);
}
successfulResult = { headOid: after, refreshed: after !== before };
} catch (error) {
primaryError = error;
await abortRebase(checkout);
try {
const publishedAfterAbort = refreshStartHead
&& await (async () => {
const remoteOid = (await gitStdout(root, ["ls-remote", "origin", requestedRef])).split(/\s+/)[0];
return remoteOid === await gitStdout(checkout, ["rev-parse", "HEAD"]);
})().catch(() => false);
// An aborted push can still have reached origin. In that case preserve the
// canonical rewritten ref; restoring only local would create an unsafe split.
if (!publishedAfterAbort) {
/*
FNXC:PullRequestFreshness 2026-08-09-04:46:
A guarded push can reject after rebase has advanced the shared local head.
Restore that canonical ref with compare-and-swap before surfacing failure,
so a later ordinary push cannot publish the rewrite that lost its lease.
*/
await restoreRefreshHead(root, checkout, requestedRef, refreshStartHead);
}
} catch (rollbackError) {
const original = error instanceof Error ? error.message : String(error);
const rollback = rollbackError instanceof Error ? rollbackError.message : String(rollbackError);
primaryError = new Error(`${original}; PR head refresh rollback failed: ${rollback}`);
}
} finally {
if (temporary && temporaryReservation) {
const removalError = await removeRefreshWorktree(root, temporary);
if (removalError) {
cleanupFailure = removalError;
let quarantined = false;
await temporaryReservation.quarantine(removalError.message).then(() => {
quarantined = true;
}).catch((quarantineError) => {
cleanupFailure = quarantineError;
});
if (quarantined) {
scheduleRetainedRefreshReconciliation(root, temporary, join(root, ".fusion", "pr-refresh-reservations"));
}
} else {
await temporaryReservation.release();
}
}
}
if (primaryError) {
const message = primaryError instanceof Error ? primaryError.message : String(primaryError);
if (cleanupFailure) {
throw new Error(`${message}; retained PR refresh cleanup failed: ${cleanupFailure instanceof Error ? cleanupFailure.message : String(cleanupFailure)}`);
}
throw primaryError;
}
if (cleanupFailure) {
throw new Error(`PR head refresh cleanup failed; GitHub mutation was not attempted: ${cleanupFailure instanceof Error ? cleanupFailure.message : String(cleanupFailure)}`);
}
if (!successfulResult) {
throw new Error(`PR head refresh failed without a result for ${input.headBranch}`);
}
return successfulResult;
}
async function assertTaskBranchAvailable(cwd: string, branch: string): Promise<boolean> {
const localRef = `refs/heads/${branch}`; const localRef = `refs/heads/${branch}`;
const localBranchExists = await gitCommandSucceeds( const localBranchExists = await gitCommandSucceeds(
cwd, cwd,
@@ -205,7 +595,7 @@ async function pushTaskBranchToOrigin(cwd: string, branch: string): Promise<void
); );
if (remoteBranchExists) { if (remoteBranchExists) {
return; return false;
} }
throw new Error( throw new Error(
@@ -213,13 +603,21 @@ async function pushTaskBranchToOrigin(cwd: string, branch: string): Promise<void
); );
} }
return true;
}
async function pushTaskBranchToOrigin(cwd: string, branch: string, signal?: AbortSignal): Promise<void> {
throwIfRefreshAborted(signal);
if (!await assertTaskBranchAvailable(cwd, branch)) return;
try { try {
// No-shell invocation (Fix #11): pass the branch as a discrete argv entry so a // No-shell invocation (Fix #11): pass the branch as a discrete argv entry so a
// crafted branch name (e.g. `$(...)`) cannot be interpreted by a shell. // crafted branch name (e.g. `$(...)`) cannot be interpreted by a shell.
await execFileAsync("git", ["push", "-u", "origin", branch], { await execFileAsync("git", ["push", "-u", "origin", branch], {
cwd, cwd,
timeout: 60_000, timeout: 60_000,
signal,
}); });
throwIfRefreshAborted(signal);
} catch (err: unknown) { } catch (err: unknown) {
const message = err instanceof Error ? err.message : String(err); const message = err instanceof Error ? err.message : String(err);
throw new Error( throw new Error(
@@ -314,7 +712,7 @@ function toBranchGroupPrState(prInfo: PrInfo | null): BranchGroupPrState {
export function createGroupPrCallback( export function createGroupPrCallback(
github: Pick<GitHubOperations, "findPrForBranch" | "createPr">, github: Pick<GitHubOperations, "findPrForBranch" | "createPr">,
): CreateGroupPrFn { ): CreateGroupPrFn {
return async ({ cwd, group, members, headBranch, baseBranch }) => { return async ({ cwd, group, members, headBranch, baseBranch, integrationRemote, signal }) => {
// FNXC:PrMergeAutoMerge 2026-07-17-16:50 (gh-4): // FNXC:PrMergeAutoMerge 2026-07-17-16:50 (gh-4):
// Resolve the repo from the PROJECT cwd, not the process cwd (same T4 // Resolve the repo from the PROJECT cwd, not the process cwd (same T4
// requirement as syncGroupPrCallback below) — in a centrally-installed // requirement as syncGroupPrCallback below) — in a centrally-installed
@@ -329,7 +727,15 @@ export function createGroupPrCallback(
return { prNumber: existing.number, prUrl: existing.url, prState: toBranchGroupPrState(existing) }; return { prNumber: existing.number, prUrl: existing.url, prState: toBranchGroupPrState(existing) };
} }
await pushTaskBranchToOrigin(cwd, headBranch); await refreshAutomatedPrHead({
projectRoot: cwd,
headBranch,
targetBranch: baseBranch,
integrationRemote,
signal,
});
throwIfRefreshAborted(signal);
await pushTaskBranchToOrigin(cwd, headBranch, signal);
const membersWithBranch = members.map((member) => ({ const membersWithBranch = members.map((member) => ({
id: member.id, id: member.id,
title: member.title, title: member.title,
@@ -337,6 +743,7 @@ export function createGroupPrCallback(
})); }));
// FNXC:ForkAwarePrHead 2026-07-26-07:18: group/shared-branch PRs also push via // FNXC:ForkAwarePrHead 2026-07-26-07:18: group/shared-branch PRs also push via
// origin and must qualify head with the fork owner when push ≠ fetch repo. // origin and must qualify head with the fork owner when push ≠ fetch repo.
throwIfRefreshAborted(signal);
const created = await github.createPr({ const created = await github.createPr({
owner: repo.owner, owner: repo.owner,
repo: repo.repo, repo: repo.repo,
@@ -422,19 +829,6 @@ export function syncGroupPrCallback(
}; };
} }
/** Best-effort resolve the head commit OID for a branch (so `pr-merge` can pass
* `expectedHeadOid`). Returns undefined on any failure — the merge then runs
* without the stale-head guard, which the reconcile still corroborates. */
async function resolveBranchHeadOid(cwd: string, branch: string): Promise<string | undefined> {
try {
const { stdout } = await execFileAsync("git", ["rev-parse", branch], { cwd, timeout: 30_000 });
const oid = stdout.trim();
return oid.length > 0 ? oid : undefined;
} catch {
return undefined;
}
}
/** Structural detection of the dashboard `PrStaleHeadError` without importing the /** Structural detection of the dashboard `PrStaleHeadError` without importing the
* class (task-lifecycle.ts deliberately has no @fusion/dashboard dependency). */ * class (task-lifecycle.ts deliberately has no @fusion/dashboard dependency). */
function isStaleHeadError(err: unknown): boolean { function isStaleHeadError(err: unknown): boolean {
@@ -509,13 +903,23 @@ export function createPrNodeGithubOps(
headBranch: getTaskBranchName(task.id), headBranch: getTaskBranchName(task.id),
}; };
}, },
createPr: async ({ task, entity }) => { createPr: async ({ task, entity, integrationRemote, signal }) => {
// FNXC:PrMergeAutoMerge 2026-07-17-19:18 (gh-4): // FNXC:PrMergeAutoMerge 2026-07-17-19:18 (gh-4):
// Git ops run in the task worktree when known; process.cwd() only as the // Git ops run in the task worktree when known; process.cwd() only as the
// single-project fallback. // single-project fallback.
const cwd = options.getTaskWorktree?.(entity.sourceId) ?? task.worktree ?? process.cwd(); const cwd = options.getTaskWorktree?.(entity.sourceId) ?? task.worktree ?? process.cwd();
const headBranch = entity.headBranch || getTaskBranchName(task.id); const headBranch = entity.headBranch || getTaskBranchName(task.id);
await pushTaskBranchToOrigin(cwd, headBranch); const refreshed = await refreshAutomatedPrHead({
projectRoot: cwd,
preferredWorktree: task.worktree,
headBranch,
targetBranch: entity.baseBranch || "main",
integrationRemote: integrationRemote,
signal,
});
throwIfRefreshAborted(signal);
await pushTaskBranchToOrigin(cwd, headBranch, signal);
throwIfRefreshAborted(signal);
const { owner, name } = splitRepoSlug(entity.repo); const { owner, name } = splitRepoSlug(entity.repo);
// FNXC:ForkAwarePrHead 2026-07-26-07:18: qualify head as owner:branch when // FNXC:ForkAwarePrHead 2026-07-26-07:18: qualify head as owner:branch when
// origin pushes to a fork while the PR targets upstream. // origin pushes to a fork while the PR targets upstream.
@@ -527,23 +931,43 @@ export function createPrNodeGithubOps(
head: qualifyForkAwarePrHead(cwd, owner, headBranch), head: qualifyForkAwarePrHead(cwd, owner, headBranch),
base: entity.baseBranch, base: entity.baseBranch,
}); });
const headOid = await resolveBranchHeadOid(cwd, headBranch); return { prNumber: created.number, prUrl: created.url, headOid: refreshed.headOid };
return { prNumber: created.number, prUrl: created.url, headOid };
}, },
mergePr: async ({ entity }) => { mergePr: async ({ task, entity, integrationRemote, persistRefreshedHead, signal }) => {
if (entity.prNumber == null) { if (entity.prNumber == null) {
throw new Error(`pr-merge: entity ${entity.id} has no persisted prNumber`); throw new Error(`pr-merge: entity ${entity.id} has no persisted prNumber`);
} }
const { owner, name } = splitRepoSlug(entity.repo); const { owner, name } = splitRepoSlug(entity.repo);
const cwd = options.getTaskWorktree?.(entity.sourceId) ?? task?.worktree ?? process.cwd();
const refreshed = await refreshAutomatedPrHead({
projectRoot: cwd,
preferredWorktree: task?.worktree,
headBranch: entity.headBranch || getTaskBranchName(task?.id ?? entity.sourceId),
targetBranch: entity.baseBranch || "main",
integrationRemote,
signal,
});
throwIfRefreshAborted(signal);
// Re-read after a rewrite: GitHub's head OID is the merge fence, not a
// locally remembered pre-refresh value.
const status = await github.getPrStatus(owner ?? "", name ?? "", entity.prNumber);
if (status && status.status !== "open" && status.status !== "draft") return { status: "stale-head" };
/*
FNXC:PullRequestFreshness 2026-08-09-02:14:
Persist the locally published OID before invoking GitHub. The workflow
handler owns this write so a crash cannot retain a pre-refresh merge fence.
*/
await persistRefreshedHead?.(refreshed.headOid);
try { try {
throwIfRefreshAborted(signal);
await github.mergePr({ await github.mergePr({
owner, owner,
repo: name, repo: name,
number: entity.prNumber, number: entity.prNumber,
method: "squash", method: "squash",
expectedHeadOid: entity.headOid, expectedHeadOid: refreshed.headOid,
}); });
return { status: "merged-requested" }; return { status: "merged-requested", headOid: refreshed.headOid };
} catch (err) { } catch (err) {
if (isStaleHeadError(err)) return { status: "stale-head" }; if (isStaleHeadError(err)) return { status: "stale-head" };
throw err; throw err;
@@ -811,6 +1235,7 @@ export async function processPullRequestMergeTask(
github: GitHubOperations, github: GitHubOperations,
getTaskMergeBlocker: TaskMergeBlockerFn, getTaskMergeBlocker: TaskMergeBlockerFn,
pool?: WorktreePool, pool?: WorktreePool,
signal?: AbortSignal,
): Promise<ProcessPullRequestResult> { ): Promise<ProcessPullRequestResult> {
const task = await store.getTask(taskId); const task = await store.getTask(taskId);
/* /*
@@ -921,8 +1346,18 @@ export async function processPullRequestMergeTask(
// not-found and fall through to push + createPr for a fresh open PR. // not-found and fall through to push + createPr for a fresh open PR.
groupPrInfo = await github.findPrForBranch({ owner: prRepo.owner, repo: prRepo.repo, head: branchGroup.branchName, state: "open" }); groupPrInfo = await github.findPrForBranch({ owner: prRepo.owner, repo: prRepo.repo, head: branchGroup.branchName, state: "open" });
if (!groupPrInfo) { if (!groupPrInfo) {
await pushTaskBranchToOrigin(cwd, branchGroup.branchName); await refreshAutomatedPrHead({
projectRoot: cwd,
preferredWorktree: task.worktree,
headBranch: branchGroup.branchName,
targetBranch: projectDefaultBranch,
integrationRemote: settings.worktreeRebaseRemote,
signal,
});
throwIfRefreshAborted(signal);
await pushTaskBranchToOrigin(cwd, branchGroup.branchName, signal);
try { try {
throwIfRefreshAborted(signal);
// FNXC:ForkAwarePrHead 2026-07-26-07:18: shared-branch processPullRequest // FNXC:ForkAwarePrHead 2026-07-26-07:18: shared-branch processPullRequest
// path must qualify head for fork push URLs (same as createGroupPrCallback). // path must qualify head for fork push URLs (same as createGroupPrCallback).
groupPrInfo = await github.createPr({ groupPrInfo = await github.createPr({
@@ -995,8 +1430,37 @@ export async function processPullRequestMergeTask(
return "waiting"; return "waiting";
} }
const refreshedHead = await refreshAutomatedPrHead({
projectRoot: cwd,
preferredWorktree: task.worktree,
headBranch: branchGroup.branchName,
targetBranch: projectDefaultBranch,
integrationRemote: settings.worktreeRebaseRemote,
signal,
});
const latestMergeStatus = refreshedHead.refreshed
? await github.getPrMergeStatus(prRepo.owner, prRepo.repo, refreshedPrInfo.number) ?? mergeStatus
: mergeStatus;
if (refreshedHead.refreshed) {
await store.updateBranchGroup(branchGroup.id, {
prNumber: latestMergeStatus.prInfo.number,
prUrl: latestMergeStatus.prInfo.url,
prState: toBranchGroupPrState(latestMergeStatus.prInfo),
});
}
// A rewritten head may invalidate approval or checks; re-admit against the
// authoritative post-publication state rather than the pre-refresh poll.
if (settings.requirePrApproval && latestMergeStatus.reviewDecision !== "APPROVED") {
await store.updateTask(task.id, { status: "awaiting-pr-checks" });
return "waiting";
}
if (!latestMergeStatus.mergeReady) {
await store.updateTask(task.id, { status: "awaiting-pr-checks" });
return "waiting";
}
await store.updateTask(task.id, { status: "merging-pr" }); await store.updateTask(task.id, { status: "merging-pr" });
const mergedPr = await github.mergePr({ owner: prRepo.owner, repo: prRepo.repo, number: refreshedPrInfo.number, method: "squash" }); throwIfRefreshAborted(signal);
const mergedPr = await github.mergePr({ owner: prRepo.owner, repo: prRepo.repo, number: refreshedPrInfo.number, method: "squash", expectedHeadOid: refreshedHead.headOid });
await store.updateBranchGroup(branchGroup.id, { await store.updateBranchGroup(branchGroup.id, {
prNumber: mergedPr.number, prNumber: mergedPr.number,
prUrl: mergedPr.url, prUrl: mergedPr.url,
@@ -1025,12 +1489,29 @@ export async function processPullRequestMergeTask(
const existingPr = await github.findPrForBranch({ owner: prRepo.owner, repo: prRepo.repo, head: branch, state: "all" }); const existingPr = await github.findPrForBranch({ owner: prRepo.owner, repo: prRepo.repo, head: branch, state: "all" });
if (!existingPr) { if (!existingPr) {
// gh pr create / GitHub REST require the head branch to exist on // Refresh before the first external PR mutation, then publish the exact
// origin. Nothing else in the merge path publishes the per-task // rewritten head rather than the creation-time task base.
// branch, so we push it here right before creating the PR. await assertTaskBranchAvailable(cwd, branch);
await pushTaskBranchToOrigin(cwd, branch); await refreshAutomatedPrHead({
projectRoot: cwd,
preferredWorktree: task.worktree,
headBranch: branch,
targetBranch: mergeTarget.branch,
integrationRemote: settings.worktreeRebaseRemote,
signal,
});
throwIfRefreshAborted(signal);
await pushTaskBranchToOrigin(cwd, branch, signal);
/*
FNXC:PullRequestFreshness 2026-08-09-03:32:
A first-time, already-current head is published by pushTaskBranchToOrigin,
not by refreshAutomatedPrHead. Do not erase the retry budget until both
parts of the lifecycle publication boundary have succeeded.
*/
await store.updateTask(task.id, { mergeRetries: 0 });
} }
try { try {
throwIfRefreshAborted(signal);
// FNXC:ForkAwarePrHead 2026-07-26-07:18: per-task processPullRequest path // FNXC:ForkAwarePrHead 2026-07-26-07:18: per-task processPullRequest path
// must qualify head for fork push URLs (same as createPrNodeGithubOps). // must qualify head for fork push URLs (same as createPrNodeGithubOps).
prInfo = existingPr ?? await github.createPr({ prInfo = existingPr ?? await github.createPr({
@@ -1117,10 +1598,46 @@ export async function processPullRequestMergeTask(
await store.updateTask(task.id, { status: "awaiting-pr-checks" }); await store.updateTask(task.id, { status: "awaiting-pr-checks" });
return "waiting"; return "waiting";
} }
const refreshedHead = await refreshAutomatedPrHead({
projectRoot: cwd,
preferredWorktree: task.worktree,
headBranch: branch,
targetBranch: mergeTarget.branch,
integrationRemote: settings.worktreeRebaseRemote,
signal,
});
const latestMergeStatus = refreshedHead.refreshed
? await github.getPrMergeStatus(prRepo.owner, prRepo.repo, prInfo.number) ?? mergeStatus
: mergeStatus;
/*
FNXC:PullRequestFreshness 2026-08-09-03:02:
A completed refresh, guarded publication, and temporary-worktree cleanup are
real lifecycle progress. Reset only after that sequence so stale-base conflict
polls cannot permanently consume the task's retry budget.
*/
await store.updateTask(task.id, { mergeRetries: 0 });
if (refreshedHead.refreshed) {
await store.updatePrInfo(task.id, {
...prInfo,
...latestMergeStatus.prInfo,
lastCheckedAt: new Date().toISOString(),
});
}
// Rebase publication can reset approval/check state. Never merge from the
// pre-refresh admission result.
if (settings.requirePrApproval && latestMergeStatus.reviewDecision !== "APPROVED") {
await store.updateTask(task.id, { status: "awaiting-pr-checks" });
return "waiting";
}
if (!latestMergeStatus.mergeReady) {
await store.updateTask(task.id, { status: "awaiting-pr-checks" });
return "waiting";
}
await store.updateTask(task.id, { status: "merging-pr" }); await store.updateTask(task.id, { status: "merging-pr" });
let mergedPr: PrInfo; let mergedPr: PrInfo;
try { try {
mergedPr = await github.mergePr({ owner: prRepo.owner, repo: prRepo.repo, number: prInfo.number, method: "squash" }); throwIfRefreshAborted(signal);
mergedPr = await github.mergePr({ owner: prRepo.owner, repo: prRepo.repo, number: prInfo.number, method: "squash", expectedHeadOid: refreshedHead.headOid });
} catch (err: unknown) { } catch (err: unknown) {
let refreshedStatus: Awaited<ReturnType<GitHubOperations["getPrMergeStatus"]>>; let refreshedStatus: Awaited<ReturnType<GitHubOperations["getPrMergeStatus"]>>;
try { try {

View File

@@ -566,6 +566,45 @@ describe("promoteBranchGroup PR creation (U5)", () => {
baseBranch: "main", baseBranch: "main",
}; };
it("fails closed before PR-mode promotion mutates the project checkout when cancelled", async () => {
const rootDir = makePrRepo();
const controller = new AbortController();
controller.abort(new Error("promotion cancelled"));
let group = makeGroup();
const createGroupPr = vi.fn();
await expect(promoteBranchGroup({
rootDir,
groupId: group.id,
settings: prSettings,
signal: controller.signal,
store: makeStore(() => group, (g) => { group = g; }, [landedMember("FN-A", group.branchName)]),
createGroupPr,
})).rejects.toThrow("promotion cancelled");
expect(createGroupPr).not.toHaveBeenCalled();
expect(execSync("git branch --show-current", { cwd: rootDir, encoding: "utf8" }).trim()).toBe("main");
expect(execSync("git log --format=%s -1", { cwd: rootDir, encoding: "utf8" }).trim()).not.toBe("Merge branch 'fusion/groups/planning-x'");
});
it("forwards the configured integration remote to automated group PR creation", async () => {
const rootDir = makePrRepo();
let group = makeGroup();
let receivedRemote: string | undefined;
await promoteBranchGroup({
rootDir,
groupId: group.id,
settings: { ...prSettings, worktreeRebaseRemote: "upstream" },
store: makeStore(() => group, (g) => { group = g; }, [landedMember("FN-A", group.branchName)]),
createGroupPr: async ({ integrationRemote }) => {
receivedRemote = integrationRemote;
return { prNumber: 42, prUrl: "https://github.com/x/y/pull/42", prState: "open" };
},
});
expect(receivedRemote).toBe("upstream");
});
it("creates exactly one PR for a complete PR-mode group and persists prNumber/prUrl/prState=open", async () => { it("creates exactly one PR for a complete PR-mode group and persists prNumber/prUrl/prState=open", async () => {
const rootDir = makePrRepo(); const rootDir = makePrRepo();
let group = makeGroup(); let group = makeGroup();
@@ -1286,6 +1325,10 @@ function createInterpreterMergeEngine(repo: string, store: any): any {
const engine = Object.create(ProjectEngine.prototype) as any; const engine = Object.create(ProjectEngine.prototype) as any;
engine.config = { workingDirectory: repo }; engine.config = { workingDirectory: repo };
engine.options = {}; engine.options = {};
// FNXC:PullRequestFreshness 2026-08-09-04:26:
// The production merge admission clears its deduplicated capacity reason on a
// successful release, so this prototype-backed engine must provide that state.
engine.capacityDeferredMergeReasons = new Map();
engine.runtime = { getTaskStore: () => store, getPluginRunner: () => undefined }; engine.runtime = { getTaskStore: () => store, getPluginRunner: () => undefined };
return engine; return engine;
} }

View File

@@ -40,6 +40,10 @@ export type CreateGroupPrFn = (input: {
headBranch: string; headBranch: string;
/** Base branch — the integration/default target. */ /** Base branch — the integration/default target. */
baseBranch: string; baseBranch: string;
/** Explicit project integration remote, when configured. */
integrationRemote?: string;
/** Cancels a queued promotion before it can mutate GitHub. */
signal?: AbortSignal;
}) => Promise<{ prNumber: number; prUrl: string; prState: BranchGroupPrState }>; }) => Promise<{ prNumber: number; prUrl: string; prState: BranchGroupPrState }>;
/** Result shape shared by group-PR sync/close callbacks. */ /** Result shape shared by group-PR sync/close callbacks. */
@@ -204,6 +208,18 @@ async function ensureGroupBranchExists(rootDir: string, branchName: string, star
*/ */
const promotionLocks = new Map<string, Promise<unknown>>(); const promotionLocks = new Map<string, Promise<unknown>>();
/*
FNXC:PullRequestFreshness 2026-08-09-03:20:
PR-mode promotion must not alter the project-root checkout before the verified
head refresh. Cancellation is checked at the coordinator boundary so an aborted
claim cannot reach an injected GitHub creator or the legacy root-checkout merge.
*/
function throwIfPromotionAborted(signal: AbortSignal | undefined): void {
if (signal?.aborted) {
throw signal.reason instanceof Error ? signal.reason : new Error("Branch-group promotion cancelled");
}
}
/** /**
* The only entrypoint allowed to perform shared-branch-group → default-branch promotion. * The only entrypoint allowed to perform shared-branch-group → default-branch promotion.
* Promotion is intentionally idempotent and must never run inline in aiMergeTask. * Promotion is intentionally idempotent and must never run inline in aiMergeTask.
@@ -214,13 +230,15 @@ export interface PromoteBranchGroupInput {
store: Pick<TaskStore, "getBranchGroup" | "getBranchGroupByBranchName" | "listTasksByBranchGroup" | "updateBranchGroup">; store: Pick<TaskStore, "getBranchGroup" | "getBranchGroupByBranchName" | "listTasksByBranchGroup" | "updateBranchGroup">;
rootDir: string; rootDir: string;
groupId: string; groupId: string;
settings: Pick<Settings, "autoMerge" | "globalPause" | "enginePaused"> & Partial<Pick<Settings, "mergeStrategy" | "integrationBranch" | "baseBranch">>; settings: Pick<Settings, "autoMerge" | "globalPause" | "enginePaused"> & Partial<Pick<Settings, "mergeStrategy" | "integrationBranch" | "baseBranch" | "worktreeRebaseRemote">>;
/** /**
* Injected GitHub PR creator (KTD7). When PR mode is active and the group is * Injected GitHub PR creator (KTD7). When PR mode is active and the group is
* complete, the coordinator uses this to create the single managed PR. Omitted * complete, the coordinator uses this to create the single managed PR. Omitted
* for direct-merge mode and in tests that don't exercise PR creation. * for direct-merge mode and in tests that don't exercise PR creation.
*/ */
createGroupPr?: CreateGroupPrFn; createGroupPr?: CreateGroupPrFn;
/** Merge-queue cancellation propagated to the automated PR creation boundary. */
signal?: AbortSignal;
recordAudit?: (event: { recordAudit?: (event: {
domain: string; domain: string;
mutationType: string; mutationType: string;
@@ -253,6 +271,7 @@ export async function promoteBranchGroup(input: PromoteBranchGroupInput): Promis
} }
async function promoteBranchGroupInner(input: PromoteBranchGroupInput): Promise<BranchGroupPromotionResult> { async function promoteBranchGroupInner(input: PromoteBranchGroupInput): Promise<BranchGroupPromotionResult> {
throwIfPromotionAborted(input.signal);
const group = await input.store.getBranchGroup(input.groupId); const group = await input.store.getBranchGroup(input.groupId);
if (!group) { if (!group) {
return { return {
@@ -356,15 +375,19 @@ async function promoteBranchGroupInner(input: PromoteBranchGroupInput): Promise<
} }
const integrationBranch = await resolveIntegrationBranch(input.rootDir, input.settings); const integrationBranch = await resolveIntegrationBranch(input.rootDir, input.settings);
if (!needsPrRepair) { if (!needsPrRepair && !isPrMode) {
throwIfPromotionAborted(input.signal);
await ensureGroupBranchExists(input.rootDir, group.branchName, integrationBranch); await ensureGroupBranchExists(input.rootDir, group.branchName, integrationBranch);
const currentBranch = ( const currentBranch = (
await execFileAsync("git", ["rev-parse", "--abbrev-ref", "HEAD"], { cwd: input.rootDir }) await execFileAsync("git", ["rev-parse", "--abbrev-ref", "HEAD"], { cwd: input.rootDir })
).stdout.trim(); ).stdout.trim();
try { try {
await execFileAsync("git", ["checkout", integrationBranch], { cwd: input.rootDir }); throwIfPromotionAborted(input.signal);
await execFileAsync("git", ["merge", "--no-ff", "--no-edit", group.branchName], { cwd: input.rootDir }); await execFileAsync("git", ["checkout", integrationBranch], { cwd: input.rootDir, signal: input.signal });
throwIfPromotionAborted(input.signal);
await execFileAsync("git", ["merge", "--no-ff", "--no-edit", group.branchName], { cwd: input.rootDir, signal: input.signal });
} finally { } finally {
// Restore the direct-merge caller's checkout even after cancellation.
await execFileAsync("git", ["checkout", currentBranch], { cwd: input.rootDir }); await execFileAsync("git", ["checkout", currentBranch], { cwd: input.rootDir });
} }
} }
@@ -399,12 +422,15 @@ async function promoteBranchGroupInner(input: PromoteBranchGroupInput): Promise<
// GitHub failure must leave the group recoverable: do NOT flip prState to a // GitHub failure must leave the group recoverable: do NOT flip prState to a
// lie. The group is already merged to the integration branch locally; we // lie. The group is already merged to the integration branch locally; we
// surface the error so the caller can retry promotion (which is idempotent). // surface the error so the caller can retry promotion (which is idempotent).
throwIfPromotionAborted(input.signal);
const created = await input.createGroupPr({ const created = await input.createGroupPr({
cwd: input.rootDir, cwd: input.rootDir,
group, group,
members, members,
headBranch: group.branchName, headBranch: group.branchName,
baseBranch: integrationBranch, baseBranch: integrationBranch,
integrationRemote: input.settings.worktreeRebaseRemote,
signal: input.signal,
}); });
prNumber = created.prNumber; prNumber = created.prNumber;
prUrl = created.prUrl; prUrl = created.prUrl;

View File

@@ -42,6 +42,8 @@ import { makePrResponseAgentRunner, makePrResponseGitOps } from "./pr-response-r
* store and the handlers stay trivially fakeable in tests. * store and the handlers stay trivially fakeable in tests.
*/ */
export interface PrNodeStore extends PrResponseRunStore { export interface PrNodeStore extends PrResponseRunStore {
/** The engine supplies project policy without coupling the CLI callback to TaskStore. */
getSettings?: () => Promise<{ worktreeRebaseRemote?: string }>;
/** Create-or-reuse the single non-terminal entity for a source (AE6 idempotency). */ /** Create-or-reuse the single non-terminal entity for a source (AE6 idempotency). */
ensurePrEntityForSource(input: PrEntityCreateInput): Promise<PrEntity>; ensurePrEntityForSource(input: PrEntityCreateInput): Promise<PrEntity>;
getPrEntity(id: string): Promise<PrEntity | null>; getPrEntity(id: string): Promise<PrEntity | null>;
@@ -74,6 +76,10 @@ export interface PrCreateCallInput {
task: TaskDetail; task: TaskDetail;
node: WorkflowIrNode; node: WorkflowIrNode;
entity: PrEntity; entity: PrEntity;
/** Project-configured integration remote resolved by the engine-owned store. */
integrationRemote?: string;
/** Graph cancellation must prevent a PR mutation after a refresh completes. */
signal?: AbortSignal;
} }
/** Result of a successful PR creation — the GitHub-mirror fields the node persists. */ /** Result of a successful PR creation — the GitHub-mirror fields the node persists. */
@@ -91,6 +97,12 @@ export interface PrMergeCallInput {
entity: PrEntity; entity: PrEntity;
/** The head OID the merge is gated on (defeats the push/merge race, U2/U6). */ /** The head OID the merge is gated on (defeats the push/merge race, U2/U6). */
expectedHeadOid?: string; expectedHeadOid?: string;
/** Project-configured integration remote resolved by the engine-owned store. */
integrationRemote?: string;
/** Persist a refreshed head before the injected callback performs GitHub merge. */
persistRefreshedHead?: (headOid: string) => Promise<void>;
/** Graph cancellation must prevent a PR mutation after a refresh completes. */
signal?: AbortSignal;
} }
/** /**
@@ -101,7 +113,7 @@ export interface PrMergeCallInput {
* handler classifies it as a benign retryable outcome. * handler classifies it as a benign retryable outcome.
*/ */
export type PrMergeCallResult = export type PrMergeCallResult =
| { status: "merged-requested" } | { status: "merged-requested"; headOid?: string }
| { status: "stale-head" }; | { status: "stale-head" };
/** Input for the injected `respond` callback (U5 implements the real body). */ /** Input for the injected `respond` callback (U5 implements the real body). */
@@ -386,7 +398,8 @@ export function createPrNodeHandlers(deps: PrNodeDeps): Record<
let created: PrCreateCallResult; let created: PrCreateCallResult;
try { try {
created = await deps.createPr({ task: ctx.task, node, entity: creating }); const integrationRemote = (await store.getSettings?.())?.worktreeRebaseRemote;
created = await deps.createPr({ task: ctx.task, node, entity: creating, integrationRemote, signal: ctx.signal });
} catch (err) { } catch (err) {
const reason = classifyError(err); const reason = classifyError(err);
audit("pr-create-failed", `pr-create node '${node.id}' creation failed: ${reason}`); audit("pr-create-failed", `pr-create node '${node.id}' creation failed: ${reason}`);
@@ -446,11 +459,23 @@ export function createPrNodeHandlers(deps: PrNodeDeps): Record<
let result: PrMergeCallResult; let result: PrMergeCallResult;
try { try {
const integrationRemote = (await store.getSettings?.())?.worktreeRebaseRemote;
result = await deps.mergePr({ result = await deps.mergePr({
task: ctx.task, task: ctx.task,
node, node,
entity, entity,
expectedHeadOid: entity.headOid, expectedHeadOid: entity.headOid,
integrationRemote,
/*
FNXC:PullRequestFreshness 2026-08-09-02:14:
A refreshed head identity is durable state before GitHub receives a merge
request. This keeps workflow fencing and the published branch aligned if
the process dies after refresh but before the provider call returns.
*/
persistRefreshedHead: async (headOid) => {
if (headOid !== entity.headOid) await store.updatePrEntity(entity.id, { headOid });
},
signal: ctx.signal,
}); });
} catch (err) { } catch (err) {
// A non-stale merge error is benign/retryable — never throw out of the // A non-stale merge error is benign/retryable — never throw out of the
@@ -466,6 +491,12 @@ export function createPrNodeHandlers(deps: PrNodeDeps): Record<
return { outcome: "success", value: "stale-head" }; return { outcome: "success", value: "stale-head" };
} }
// FNXC:PullRequestFreshness 2026-08-09-01:17:
// Persist the refreshed GitHub merge fence before reconciliation observes the
// request, so a rewrite cannot leave the entity claiming its stale head.
if (result.headOid && result.headOid !== entity.headOid) {
await store.updatePrEntity(entity.id, { headOid: result.headOid });
}
// Merge requested cleanly. Do NOT write `merged` here — reconcile corroborates. // Merge requested cleanly. Do NOT write `merged` here — reconcile corroborates.
return { outcome: "success", value: "merged-requested" }; return { outcome: "success", value: "merged-requested" };
}; };

View File

@@ -121,6 +121,8 @@ export type ProcessPullRequestMergeFn = (
cwd: string, cwd: string,
taskId: string, taskId: string,
pool?: WorktreePool, pool?: WorktreePool,
/** Propagates merge-queue cancellation into refresh git mutations. */
signal?: AbortSignal,
) => Promise<"merged" | "waiting" | "skipped">; ) => Promise<"merged" | "waiting" | "skipped">;
const execFileAsync = promisify(execFile); const execFileAsync = promisify(execFile);
@@ -2308,6 +2310,7 @@ export class ProjectEngine {
mergeStrategy: settings.mergeStrategy, mergeStrategy: settings.mergeStrategy,
integrationBranch: settings.integrationBranch, integrationBranch: settings.integrationBranch,
baseBranch: settings.baseBranch, baseBranch: settings.baseBranch,
worktreeRebaseRemote: settings.worktreeRebaseRemote,
}; };
return await promoteBranchGroup({ return await promoteBranchGroup({
store, store,
@@ -3858,8 +3861,15 @@ export class ProjectEngine {
mergeStrategy: settings.mergeStrategy, mergeStrategy: settings.mergeStrategy,
integrationBranch: settings.integrationBranch, integrationBranch: settings.integrationBranch,
baseBranch: settings.baseBranch, baseBranch: settings.baseBranch,
worktreeRebaseRemote: settings.worktreeRebaseRemote,
}; };
const attemptBranchGroupPromotion = async (taskForPromotion: Task | null): Promise<void> => { /*
FNXC:PullRequestFreshness 2026-08-09-03:02:
Branch-group promotion is an automated PR producer after a member merge.
Preserve the merge claim's cancellation signal through the coordinator so
a cancelled refresh cannot proceed to GitHub PR creation.
*/
const attemptBranchGroupPromotion = async (taskForPromotion: Task | null, signal?: AbortSignal): Promise<void> => {
// groupId is optional on TaskBranchContext (non-shared members carry none); // groupId is optional on TaskBranchContext (non-shared members carry none);
// isSharedBranchGroupMemberIntegration guarantees it semantically, but capture // isSharedBranchGroupMemberIntegration guarantees it semantically, but capture
// it explicitly so TypeScript narrows. // it explicitly so TypeScript narrows.
@@ -3874,6 +3884,7 @@ export class ProjectEngine {
groupId: promotionGroupId, groupId: promotionGroupId,
settings: promotionSettings, settings: promotionSettings,
createGroupPr: this.options.createGroupPr, createGroupPr: this.options.createGroupPr,
signal,
recordAudit: async (event) => { recordAudit: async (event) => {
await store.recordRunAuditEvent({ await store.recordRunAuditEvent({
domain: event.domain as any, domain: event.domain as any,
@@ -4093,6 +4104,7 @@ export class ProjectEngine {
cwd, cwd,
taskId, taskId,
(this.runtime as any).worktreePool, (this.runtime as any).worktreePool,
abortSignal,
), ),
abortSignal, abortSignal,
taskId, taskId,
@@ -4121,7 +4133,7 @@ export class ProjectEngine {
mergeTargetBranch: mergedTask.mergeDetails?.mergeTargetBranch, mergeTargetBranch: mergedTask.mergeDetails?.mergeTargetBranch,
} as MergeResult); } as MergeResult);
} }
await attemptBranchGroupPromotion(mergedTask); await attemptBranchGroupPromotion(mergedTask, this.mergeAbortController?.signal);
} else if (result === "waiting") { } else if (result === "waiting") {
runtimeLog.log(`${hasManualResolver ? "Manual" : "Auto"}-merge PR waiting: ${taskId}`); runtimeLog.log(`${hasManualResolver ? "Manual" : "Auto"}-merge PR waiting: ${taskId}`);
} }