From b3504f01a4894f0707e08630472dcf75c2a7796b Mon Sep 17 00:00:00 2001 From: gsxdsm Date: Sat, 8 Aug 2026 22:35:59 -0700 Subject: [PATCH] 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) --- .changeset/fn-8838-refresh-pr-heads.md | 7 + docs/settings-reference.md | 6 + ...task-lifecycle-refresh.integration.test.ts | 438 +++++++++++++ .../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 +++++++++++++++++- .../__tests__/group-merge-coordinator.test.ts | 43 ++ .../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(-) create mode 100644 .changeset/fn-8838-refresh-pr-heads.md create mode 100644 packages/cli/src/commands/__tests__/task-lifecycle-refresh.integration.test.ts diff --git a/.changeset/fn-8838-refresh-pr-heads.md b/.changeset/fn-8838-refresh-pr-heads.md new file mode 100644 index 0000000000..fae5273cc4 --- /dev/null +++ b/.changeset/fn-8838-refresh-pr-heads.md @@ -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. diff --git a/docs/settings-reference.md b/docs/settings-reference.md index 28332b4482..15fc9b31de 100644 --- a/docs/settings-reference.md +++ b/docs/settings-reference.md @@ -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`. | | `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/`; 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 When a project uses `mergeStrategy: "direct"`, an individual task can override the project-level `directMergeCommitStrategy` by adding this line anywhere in `PROMPT.md`: diff --git a/packages/cli/src/commands/__tests__/task-lifecycle-refresh.integration.test.ts b/packages/cli/src/commands/__tests__/task-lifecycle-refresh.integration.test.ts new file mode 100644 index 0000000000..a49636710c --- /dev/null +++ b/packages/cli/src/commands/__tests__/task-lifecycle-refresh.integration.test.ts @@ -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("@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) { + 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(); + }); +}); diff --git a/packages/cli/src/commands/__tests__/task-lifecycle.test.ts b/packages/cli/src/commands/__tests__/task-lifecycle.test.ts index 4f042373e9..fc2bb43ec1 100644 --- a/packages/cli/src/commands/__tests__/task-lifecycle.test.ts +++ b/packages/cli/src/commands/__tests__/task-lifecycle.test.ts @@ -10,11 +10,32 @@ const execMock = vi.hoisted(() => vi.fn()); const execFileCalls = vi.hoisted( () => [] 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("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", () => ({ exec: (cmd: string, opts: unknown, cb: (err: Error | null, stdout: string, stderr: string) => void) => { try { 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) { 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) => { try { execFileCalls.push({ file, args: args ?? [], cwd: (opts as { cwd?: string } | undefined)?.cwd }); - const result = execMock(`${file} ${(args ?? []).join(" ")}`.trim(), opts); - cb(null, typeof result === "string" ? result : "", ""); + if (file === "git" && args?.[0] === "worktree" && args[1] === "add") { + 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) { cb(err as Error, "", (err as Error).message); } @@ -36,16 +66,23 @@ vi.mock("@fusion/core", async () => { ...actual, getCurrentRepo: 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 { cleanupMergedTaskArtifacts, createGroupPrCallback, createPrNodeGithubOps, processPullRequestMergeTask, + refreshAutomatedPrHead, getTaskBranchName, syncGroupPrCallback, } from "../task-lifecycle.js"; @@ -1193,6 +1230,7 @@ describe("processPullRequestMergeTask", () => { description: "desc", column: "in-review", worktree: "/tmp/worktree-fn-9104", + mergeRetries: 2, prInfo: { number: 124, url: "https://github.com/x/y/pull/124", @@ -1248,11 +1286,12 @@ describe("processPullRequestMergeTask", () => { ); 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).toHaveBeenNthCalledWith(1, "owner", "repo", 124); expect(github.getPrMergeStatus).toHaveBeenNthCalledWith(2, "owner", "repo", 124); 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.moveTask).toHaveBeenCalledWith("FN-9104", "done"); expect(store.logEntry).toHaveBeenCalledWith( @@ -1319,7 +1358,7 @@ describe("processPullRequestMergeTask", () => { ), ).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(store.updatePrInfo).not.toHaveBeenCalledWith("FN-9105", expect.objectContaining({ status: "merged" })); expect(store.moveTask).not.toHaveBeenCalled(); @@ -1431,7 +1470,7 @@ describe("processPullRequestMergeTask", () => { ), ).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(store.moveTask).not.toHaveBeenCalled(); }); @@ -1598,7 +1637,7 @@ describe("processPullRequestMergeTask", () => { ); 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 () => { @@ -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 () => { // With state:"open", findPrForBranch returns null for a head whose only PR // 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)", () => { beforeEach(() => { 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 () => { const github = githubStub(); const ops = createPrNodeGithubOps(github as never); const result = await ops.mergePr({ entity: { id: "e1", sourceId: "FN-9601", repo: "central-owner/central-repo", prNumber: 9, headOid: "abc123" }, } as never); - expect(result).toEqual({ status: "merged-requested" }); + expect(result).toEqual({ status: "merged-requested", headOid: "1111111111111111111111111111111111111111" }); 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"]); + }); }); diff --git a/packages/cli/src/commands/daemon.ts b/packages/cli/src/commands/daemon.ts index 3db97a2e48..ff0732a077 100644 --- a/packages/cli/src/commands/daemon.ts +++ b/packages/cli/src/commands/daemon.ts @@ -365,8 +365,8 @@ export async function runDaemon(opts: DaemonOptions = {}) { onMigrationProgress: (event) => migrationHoldingServer?.setMigrationProgress(event), cliPackageVersion, getMergeStrategy, - processPullRequestMerge: (s, wd, taskId, pool) => - processPullRequestMergeTask(s, wd, taskId, githubClient, getTaskMergeBlocker, pool), + processPullRequestMerge: (s, wd, taskId, pool, signal) => + processPullRequestMergeTask(s, wd, taskId, githubClient, getTaskMergeBlocker, pool, signal), createGroupPr: createGroupPrCallback(githubClient), syncGroupPr: syncGroupPrCallback(githubClient), prNodeGithubOps: createPrNodeGithubOps(githubClient), diff --git a/packages/cli/src/commands/dashboard.ts b/packages/cli/src/commands/dashboard.ts index 556f05940b..1010c274d6 100644 --- a/packages/cli/src/commands/dashboard.ts +++ b/packages/cli/src/commands/dashboard.ts @@ -2135,8 +2135,8 @@ export async function runDashboard(port: number, opts: { paused?: boolean; dev?: const engineManager = new ProjectEngineManager(centralCoreForEngine, { cliPackageVersion, getMergeStrategy, - processPullRequestMerge: (s, wd, taskId, pool) => - processPullRequestMergeTask(s, wd, taskId, githubClient, getTaskMergeBlocker, pool), + processPullRequestMerge: (s, wd, taskId, pool, signal) => + processPullRequestMergeTask(s, wd, taskId, githubClient, getTaskMergeBlocker, pool, signal), createGroupPr: createGroupPrCallback(githubClient), syncGroupPr: syncGroupPrCallback(githubClient), prNodeGithubOps: createPrNodeGithubOps(githubClient), diff --git a/packages/cli/src/commands/serve.ts b/packages/cli/src/commands/serve.ts index e0c3552e92..d0243500bf 100644 --- a/packages/cli/src/commands/serve.ts +++ b/packages/cli/src/commands/serve.ts @@ -425,8 +425,8 @@ export async function runServe( const engineManager = startupEngineManager = new ProjectEngineManager(sharedCentralCore, { cliPackageVersion, getMergeStrategy, - processPullRequestMerge: (s, wd, taskId, pool) => - processPullRequestMergeTask(s, wd, taskId, githubClient, getTaskMergeBlocker, pool), + processPullRequestMerge: (s, wd, taskId, pool, signal) => + processPullRequestMergeTask(s, wd, taskId, githubClient, getTaskMergeBlocker, pool, signal), createGroupPr: createGroupPrCallback(githubClient), syncGroupPr: syncGroupPrCallback(githubClient), prNodeGithubOps: createPrNodeGithubOps(githubClient), diff --git a/packages/cli/src/commands/task-lifecycle.ts b/packages/cli/src/commands/task-lifecycle.ts index 5b343cbee1..db6d877964 100644 --- a/packages/cli/src/commands/task-lifecycle.ts +++ b/packages/cli/src/commands/task-lifecycle.ts @@ -13,8 +13,11 @@ * - Full PR lifecycle orchestration (create → status check → merge) */ +import { createHash } from "node:crypto"; import { exec } 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"; const execAsync = promisify(exec); // `execFile` is resolved lazily through the namespace import so test mocks that @@ -34,6 +37,8 @@ import { assertNotWorkspaceTaskMerge, classifyGhError, WorkspaceTaskMergeError, + acquireWorktreePathReservation, + type WorktreePathReservation, } from "@fusion/core"; import type { Settings, TaskDetail, PrInfo, MergeResult, BranchGroup, BranchGroupPrState, Task } 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 */ } return "done"; } -import { activeSessionRegistry, resolveIntegrationBranch } from "@fusion/engine"; +import { activeSessionRegistry, resolveIntegrationBranch, resolveIntegrationRemote } from "@fusion/engine"; import type { CreateGroupPrFn, SyncGroupPrFn, @@ -187,7 +192,392 @@ async function gitCommandSucceeds( } } -async function pushTaskBranchToOrigin(cwd: string, branch: string): Promise { +/* +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 { + 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 { + 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 ` 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 { + 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 { + 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 { + 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 { + 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 { + 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 { + 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 `; 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 { const localRef = `refs/heads/${branch}`; const localBranchExists = await gitCommandSucceeds( cwd, @@ -205,7 +595,7 @@ async function pushTaskBranchToOrigin(cwd: string, branch: string): Promise { + throwIfRefreshAborted(signal); + if (!await assertTaskBranchAvailable(cwd, branch)) return; try { // 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. await execFileAsync("git", ["push", "-u", "origin", branch], { cwd, timeout: 60_000, + signal, }); + throwIfRefreshAborted(signal); } catch (err: unknown) { const message = err instanceof Error ? err.message : String(err); throw new Error( @@ -314,7 +712,7 @@ function toBranchGroupPrState(prInfo: PrInfo | null): BranchGroupPrState { export function createGroupPrCallback( github: Pick, ): 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): // Resolve the repo from the PROJECT cwd, not the process cwd (same T4 // 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) }; } - 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) => ({ id: member.id, title: member.title, @@ -337,6 +743,7 @@ export function createGroupPrCallback( })); // 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. + throwIfRefreshAborted(signal); const created = await github.createPr({ owner: repo.owner, 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 { - 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 * class (task-lifecycle.ts deliberately has no @fusion/dashboard dependency). */ function isStaleHeadError(err: unknown): boolean { @@ -509,13 +903,23 @@ export function createPrNodeGithubOps( headBranch: getTaskBranchName(task.id), }; }, - createPr: async ({ task, entity }) => { + createPr: async ({ task, entity, integrationRemote, signal }) => { // FNXC:PrMergeAutoMerge 2026-07-17-19:18 (gh-4): // Git ops run in the task worktree when known; process.cwd() only as the // single-project fallback. const cwd = options.getTaskWorktree?.(entity.sourceId) ?? task.worktree ?? process.cwd(); 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); // FNXC:ForkAwarePrHead 2026-07-26-07:18: qualify head as owner:branch when // origin pushes to a fork while the PR targets upstream. @@ -527,23 +931,43 @@ export function createPrNodeGithubOps( head: qualifyForkAwarePrHead(cwd, owner, headBranch), base: entity.baseBranch, }); - const headOid = await resolveBranchHeadOid(cwd, headBranch); - return { prNumber: created.number, prUrl: created.url, headOid }; + return { prNumber: created.number, prUrl: created.url, headOid: refreshed.headOid }; }, - mergePr: async ({ entity }) => { + mergePr: async ({ task, entity, integrationRemote, persistRefreshedHead, signal }) => { if (entity.prNumber == null) { throw new Error(`pr-merge: entity ${entity.id} has no persisted prNumber`); } 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 { + throwIfRefreshAborted(signal); await github.mergePr({ owner, repo: name, number: entity.prNumber, method: "squash", - expectedHeadOid: entity.headOid, + expectedHeadOid: refreshed.headOid, }); - return { status: "merged-requested" }; + return { status: "merged-requested", headOid: refreshed.headOid }; } catch (err) { if (isStaleHeadError(err)) return { status: "stale-head" }; throw err; @@ -811,6 +1235,7 @@ export async function processPullRequestMergeTask( github: GitHubOperations, getTaskMergeBlocker: TaskMergeBlockerFn, pool?: WorktreePool, + signal?: AbortSignal, ): Promise { 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. groupPrInfo = await github.findPrForBranch({ owner: prRepo.owner, repo: prRepo.repo, head: branchGroup.branchName, state: "open" }); 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 { + throwIfRefreshAborted(signal); // FNXC:ForkAwarePrHead 2026-07-26-07:18: shared-branch processPullRequest // path must qualify head for fork push URLs (same as createGroupPrCallback). groupPrInfo = await github.createPr({ @@ -995,8 +1430,37 @@ export async function processPullRequestMergeTask( 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" }); - 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, { prNumber: mergedPr.number, 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" }); if (!existingPr) { - // gh pr create / GitHub REST require the head branch to exist on - // origin. Nothing else in the merge path publishes the per-task - // branch, so we push it here right before creating the PR. - await pushTaskBranchToOrigin(cwd, branch); + // Refresh before the first external PR mutation, then publish the exact + // rewritten head rather than the creation-time task base. + await assertTaskBranchAvailable(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 { + throwIfRefreshAborted(signal); // FNXC:ForkAwarePrHead 2026-07-26-07:18: per-task processPullRequest path // must qualify head for fork push URLs (same as createPrNodeGithubOps). prInfo = existingPr ?? await github.createPr({ @@ -1117,10 +1598,46 @@ export async function processPullRequestMergeTask( await store.updateTask(task.id, { status: "awaiting-pr-checks" }); 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" }); let mergedPr: PrInfo; 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) { let refreshedStatus: Awaited>; try { diff --git a/packages/engine/src/__tests__/group-merge-coordinator.test.ts b/packages/engine/src/__tests__/group-merge-coordinator.test.ts index 0ef667452d..e4828e8c75 100644 --- a/packages/engine/src/__tests__/group-merge-coordinator.test.ts +++ b/packages/engine/src/__tests__/group-merge-coordinator.test.ts @@ -566,6 +566,45 @@ describe("promoteBranchGroup PR creation (U5)", () => { 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 () => { const rootDir = makePrRepo(); let group = makeGroup(); @@ -1286,6 +1325,10 @@ function createInterpreterMergeEngine(repo: string, store: any): any { const engine = Object.create(ProjectEngine.prototype) as any; engine.config = { workingDirectory: repo }; 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 }; return engine; } diff --git a/packages/engine/src/merge/group-merge-coordinator.ts b/packages/engine/src/merge/group-merge-coordinator.ts index ce0eeb1ab0..649b5a5c0d 100644 --- a/packages/engine/src/merge/group-merge-coordinator.ts +++ b/packages/engine/src/merge/group-merge-coordinator.ts @@ -40,6 +40,10 @@ export type CreateGroupPrFn = (input: { headBranch: string; /** Base branch — the integration/default target. */ 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 }>; /** 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>(); +/* +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. * Promotion is intentionally idempotent and must never run inline in aiMergeTask. @@ -214,13 +230,15 @@ export interface PromoteBranchGroupInput { store: Pick; rootDir: string; groupId: string; - settings: Pick & Partial>; + settings: Pick & Partial>; /** * 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 * for direct-merge mode and in tests that don't exercise PR creation. */ createGroupPr?: CreateGroupPrFn; + /** Merge-queue cancellation propagated to the automated PR creation boundary. */ + signal?: AbortSignal; recordAudit?: (event: { domain: string; mutationType: string; @@ -253,6 +271,7 @@ export async function promoteBranchGroup(input: PromoteBranchGroupInput): Promis } async function promoteBranchGroupInner(input: PromoteBranchGroupInput): Promise { + throwIfPromotionAborted(input.signal); const group = await input.store.getBranchGroup(input.groupId); if (!group) { return { @@ -356,15 +375,19 @@ async function promoteBranchGroupInner(input: PromoteBranchGroupInput): Promise< } const integrationBranch = await resolveIntegrationBranch(input.rootDir, input.settings); - if (!needsPrRepair) { + if (!needsPrRepair && !isPrMode) { + throwIfPromotionAborted(input.signal); await ensureGroupBranchExists(input.rootDir, group.branchName, integrationBranch); const currentBranch = ( await execFileAsync("git", ["rev-parse", "--abbrev-ref", "HEAD"], { cwd: input.rootDir }) ).stdout.trim(); try { - await execFileAsync("git", ["checkout", integrationBranch], { cwd: input.rootDir }); - await execFileAsync("git", ["merge", "--no-ff", "--no-edit", group.branchName], { cwd: input.rootDir }); + throwIfPromotionAborted(input.signal); + 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 { + // Restore the direct-merge caller's checkout even after cancellation. 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 // lie. The group is already merged to the integration branch locally; we // surface the error so the caller can retry promotion (which is idempotent). + throwIfPromotionAborted(input.signal); const created = await input.createGroupPr({ cwd: input.rootDir, group, members, headBranch: group.branchName, baseBranch: integrationBranch, + integrationRemote: input.settings.worktreeRebaseRemote, + signal: input.signal, }); prNumber = created.prNumber; prUrl = created.prUrl; diff --git a/packages/engine/src/merge/pr-nodes.ts b/packages/engine/src/merge/pr-nodes.ts index c11c6ba78f..3661fb4860 100644 --- a/packages/engine/src/merge/pr-nodes.ts +++ b/packages/engine/src/merge/pr-nodes.ts @@ -42,6 +42,8 @@ import { makePrResponseAgentRunner, makePrResponseGitOps } from "./pr-response-r * store and the handlers stay trivially fakeable in tests. */ 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). */ ensurePrEntityForSource(input: PrEntityCreateInput): Promise; getPrEntity(id: string): Promise; @@ -74,6 +76,10 @@ export interface PrCreateCallInput { task: TaskDetail; node: WorkflowIrNode; 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. */ @@ -91,6 +97,12 @@ export interface PrMergeCallInput { entity: PrEntity; /** The head OID the merge is gated on (defeats the push/merge race, U2/U6). */ 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; + /** 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. */ export type PrMergeCallResult = - | { status: "merged-requested" } + | { status: "merged-requested"; headOid?: string } | { status: "stale-head" }; /** Input for the injected `respond` callback (U5 implements the real body). */ @@ -386,7 +398,8 @@ export function createPrNodeHandlers(deps: PrNodeDeps): Record< let created: PrCreateCallResult; 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) { const reason = classifyError(err); 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; try { + const integrationRemote = (await store.getSettings?.())?.worktreeRebaseRemote; result = await deps.mergePr({ task: ctx.task, node, entity, 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) { // 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" }; } + // 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. return { outcome: "success", value: "merged-requested" }; }; diff --git a/packages/engine/src/project-engine.ts b/packages/engine/src/project-engine.ts index 3bc748cf48..728ba5dac6 100644 --- a/packages/engine/src/project-engine.ts +++ b/packages/engine/src/project-engine.ts @@ -121,6 +121,8 @@ export type ProcessPullRequestMergeFn = ( cwd: string, taskId: string, pool?: WorktreePool, + /** Propagates merge-queue cancellation into refresh git mutations. */ + signal?: AbortSignal, ) => Promise<"merged" | "waiting" | "skipped">; const execFileAsync = promisify(execFile); @@ -2308,6 +2310,7 @@ export class ProjectEngine { mergeStrategy: settings.mergeStrategy, integrationBranch: settings.integrationBranch, baseBranch: settings.baseBranch, + worktreeRebaseRemote: settings.worktreeRebaseRemote, }; return await promoteBranchGroup({ store, @@ -3858,8 +3861,15 @@ export class ProjectEngine { mergeStrategy: settings.mergeStrategy, integrationBranch: settings.integrationBranch, baseBranch: settings.baseBranch, + worktreeRebaseRemote: settings.worktreeRebaseRemote, }; - const attemptBranchGroupPromotion = async (taskForPromotion: Task | null): Promise => { + /* + 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 => { // groupId is optional on TaskBranchContext (non-shared members carry none); // isSharedBranchGroupMemberIntegration guarantees it semantically, but capture // it explicitly so TypeScript narrows. @@ -3874,6 +3884,7 @@ export class ProjectEngine { groupId: promotionGroupId, settings: promotionSettings, createGroupPr: this.options.createGroupPr, + signal, recordAudit: async (event) => { await store.recordRunAuditEvent({ domain: event.domain as any, @@ -4093,6 +4104,7 @@ export class ProjectEngine { cwd, taskId, (this.runtime as any).worktreePool, + abortSignal, ), abortSignal, taskId, @@ -4121,7 +4133,7 @@ export class ProjectEngine { mergeTargetBranch: mergedTask.mergeDetails?.mergeTargetBranch, } as MergeResult); } - await attemptBranchGroupPromotion(mergedTask); + await attemptBranchGroupPromotion(mergedTask, this.mergeAbortController?.signal); } else if (result === "waiting") { runtimeLog.log(`${hasManualResolver ? "Manual" : "Auto"}-merge PR waiting: ${taskId}`); }