Files
fusion/packages/engine/src/merger-integration-worktree.ts
gsxdsm f9f06a4fb7 engine tail: nine files to zero — incl. a census-INVISIBLE requeue into a lane that does not exist (−13) (#2797)
Follow-up to #2785. Nine engine files to **zero**, all one question —
*"is this task finished?"* — asked in nine places, wrong in every one on
a renamed board.

## Census, per file (measured, `--strict` verified)

| file | main | here |
| --- | ---: | ---: |
| `agent-reflection.ts` | 2 | **0** |
| `merger-scope-auto-widen.ts` | 2 | **0** |
| `worktree-pool.ts` | 2 | **0** |
| `cli-agent/state-machine.ts` | 2 | **0** (reclassified — see below) |
| `auto-recovery-handlers/branch-worktree.ts` | 1 | **0** |
| `cli-agent/task-session.ts` | 1 | **0** (reclassified) |
| `merger-integration-worktree.ts` | 1 | **0** |
| `merger-orphan-rehome.ts` | 1 | **0** |
| `plugin-runner.ts` | 1 | **0** |
| **net** | | **−13** |

## The one worth reading: `branch-worktree` had TWO defects, and the
census could only see one

```ts
if (task.column === "in-progress") { …clear branch… }        // counted
await this.deps.taskStore.moveTask(task.id, "todo", { … });  // INVISIBLE
```

The census scores **comparisons**. The requeue *destination* is a call
argument, so nothing in the backlog ever pointed at it — and it is the
worse of the two: a board with no `todo` column was requeued into a lane
**that does not exist**. The counted literal is the smaller half (a
renamed wip lane meant the stale branch was never cleared, so the card
carried a dead branch back into execution).

Converting the comparison alone would have dropped a census count and
left the board requeuing into nowhere. Destination now resolves through
`resolveReboundTarget` (KTD-10 ordering: hold → intake → first column).

Reverted **independently**: destination restored → 2 fail (`moveTask`
called with `"todo"`, not `"backlog"`); wip test restored → 1 fail
(`updateTask` never called).

## The rest

- **`plugin-runner`** — `onTaskCompleted` never fired on a renamed
board. Every plugin that closes an issue, posts a notification, or
records a metric on completion **silently stopped**, with nothing
logged. Resolved *asynchronously* inside the existing fire-and-forget
seam, not via `resolveTaskWorkflowIrSync` — per
`sync-workflow-ir-callsite-allowlist` that reader returns the DEFAULT
workflow for every task in production, so a sync guard here would read
as converted and still be wrong. The listener is already
`void`-dispatched, so awaiting inside it changes no observable ordering
(the shape `NotificationService` already uses).
- **`merger-orphan-rehome`** — a renamed complete lane made every source
task read as unfinished, so orphaned commits were never rehomed and
stayed stranded off the integration branch. Resolves by the **trailer
id**, not `sourceTask.id`, which the fake store does not populate.
- **`agent-reflection`** — `classifyOutcome` returned `null` for every
finished task, so both callers treated completed work as nothing to
reflect on: one recorded `reflection:skipped` with reason
`"not-completed"`, the other silently `continue`d. Reflection captured
**nothing at all** on a custom board.
- **`worktree-pool`** — shipped tasks' worktrees stayed in the ACTIVE
set, so the reclaim pass never returned them and the board walks into
worktree exhaustion — a stall whose cause is invisible from the symptom.
- **`merger-scope-auto-widen`** — finished cards counted as active
claimants, so a merge was blocked by a task that no longer exists in any
meaningful sense.
- **`merger-integration-worktree`** — a shipped task still counted as a
live worktree user, so the integration worktree could never be reused
and the merge path took the slower rebuild every time.

## The census OVERSTATED the engine backlog by 3

`cli-agent/state-machine.ts` and `cli-agent/task-session.ts` compare
against `done` — but that is a **`CliMachineState`**
(`ready`/`busy`/`waitingOnInput`/`done`/`resuming`/`idle`) tracking one
CLI agent process. It never reads a board column. The census matches the
bare string.

Marked `DELIBERATE-LITERAL` rather than left for a later sweep to
"convert" a process state into a workflow role. Worth flagging
fleet-wide: the backlog total includes at least these three non-columns.

## Revert results (measured, each run)

| conversion | reverted → |
| --- | --- |
| `plugin-runner` complete gate | RENAMED case fails — `onTaskCompleted`
never invoked |
| `merger-orphan-rehome` source gate | RENAMED case fails —
`orphan:false, reason:"source-task-not-done"` |
| `branch-worktree` destination | 2 fail — `moveTask` called with
`"todo"` |
| `branch-worktree` wip test | 1 fail — `updateTask` never called |

Each has a **non-vacuous companion** (renamed board, non-complete lane /
mid-flight source / non-wip column) so a guard that fired
unconditionally would not pass.

**Four are NOT revert-proven, and I am not claiming otherwise:**
`agent-reflection`, `merger-scope-auto-widen`,
`merger-integration-worktree`, `worktree-pool`. Their suites omit a
workflow and therefore assert the legacy fallback — they pass before and
after. `merger-scope-auto-widen` has no test file at all;
`scanIdleWorktrees` is mocked in every suite that touches it and driving
it for real needs git worktrees on disk. All four strictly **widen** the
finished set (resolved roles ∪ the legacy ids), so default boards are
byte-identical. That is the argument for shipping them, not a substitute
for coverage.

## Examined and deliberately NOT converted

- **`backlog-pressure-reporter:173`** — fed by `listTasks({ column:
"todo" })`, a hardcoded **query** filter. On a renamed board `todoFull`
is empty and the predicate never runs. Converting it drops a census
count and changes nothing observable; the fix belongs at the query
layer.
- **`auto-merge-finalization:28`** — the catch-arm legacy fallback,
which must stay for the same reason `columnRoles.ts` keeps its id
fallback.
- **`auto-merge-finalization:84`** — only selects between two diagnostic
reason strings that are **both** `ok: false`, on a pure validator with
no store in scope. Converting it would thread a store through a pure
function to change a label.

## Merge resolution note

Merging main brought conflicts in `agent-assignment.ts` and
`ephemeral-worker-manager.ts`. **Main's versions won both** and mine are
dropped: main threads an optional `activeColumns` from
`scheduler.ts:2340` (a cleaner seam than widening the store type to
resolve internally), and its `isAgentIdle` carries a greptile P1 fix
mine lacked — `columnsWithFlag` membership rather than first-per-role,
so a workflow declaring two wip lanes has both recognised. That is the
fourth time in this program main's version of a contested file was the
better one.

## Verification

- `pnpm test:gate` — 161 + 487 + 13 + 71, green
- `npx tsc -p packages/engine/tsconfig.json --noEmit` — clean
- `pnpm lint` — clean
- `--strict` exits 0


<!-- This is an auto-generated comment: release notes by coderabbit.ai
-->

## Summary by CodeRabbit

* **Bug Fixes**
* Workflow-dependent task completion now recognizes custom lifecycle
columns, including renamed boards.
* Recovery requeues tasks to the configured destination and clears
branch details only from the appropriate work-in-progress column.
* Improved handling of completed tasks, orphaned work, shared worktrees,
and scope evaluation across custom workflows.
* Plugin completion hooks now trigger for any column configured as
complete.
* **Tests**
* Added coverage for renamed workflow columns and custom completion,
recovery, and rehoming behavior.
* **Documentation**
  * Clarified CLI state terminology in internal developer comments.

<!-- end of auto-generated comment: release notes by coderabbit.ai -->

---------

Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-07-30 11:28:54 -07:00

765 lines
26 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
import { createHash } from "node:crypto";
import { exec, execFile } from "node:child_process";
import { promisify } from "node:util";
import {
normalizeMergeIntegrationWorktreeMode,
resolveWorkflowIrForTask,
columnsWithFlag,
} from "@fusion/core";
import type {
MergeIntegrationWorktreeMode,
MergeQueueReleaseOutcome,
ProjectSettings,
Task,
TaskStore,
WorkflowIr,
} from "@fusion/core";
import {
activeSessionRegistry,
executingTaskLock,
reconcileSelfOwnedActiveSessionForRemoval,
} from "./active-session-registry.js";
import { attemptBranchAutocorrect } from "./branch-autocorrect.js";
import { isBranchAuthoritativeForTask } from "./branch-conflicts.js";
import { MeshLeaseManager } from "./mesh-lease-manager.js";
import {
canonicalizePath,
classifyTaskWorktree,
getRegisteredWorktreeBranchMap,
PoolDoubleLeaseError,
} from "./worktree-pool.js";
import { canonicalFusionBranchName } from "./worktree-names.js";
const execAsync = promisify(exec);
const execFileAsync = promisify(execFile);
const MERGE_HANDOFF_WORKER_ID = "merger-reuse-handoff";
/** Shell-quote a value for safe interpolation into `git` command strings.
* Mirrors the `quoteArg` helper in merger.ts; kept local to avoid an
* import cycle. */
function quoteAutostashArg(value: string): string {
return `"${value.replace(/(["\\$`])/g, "\\$1")}"`;
}
export interface MergeIntegrationRootResolution {
mode: MergeIntegrationWorktreeMode;
// Sentinel: empty string means reuse mode is requested but no reusable
// task.worktree is currently recorded; caller must reacquire before use.
rootDir: string;
branchName: string;
}
export interface ResolveMergeIntegrationRootInput {
task: Pick<Task, "id" | "branch" | "worktree">;
settings: Pick<ProjectSettings, "mergeIntegrationWorktree" | "worktrunk">;
projectRoot: string;
}
export function resolveMergeIntegrationRoot(
input: ResolveMergeIntegrationRootInput,
): MergeIntegrationRootResolution {
const branchName = canonicalFusionBranchName(input.task.id);
const mode = normalizeMergeIntegrationWorktreeMode(
input.settings.mergeIntegrationWorktree,
);
const reusablePath = input.task.worktree?.trim() || "";
return {
mode,
rootDir: mode === "reuse-task-worktree"
? reusablePath
: input.projectRoot,
branchName,
};
}
export type MergeIntegrationRootPreflightResult =
| { ok: true; resolution: MergeIntegrationRootResolution; checked: false | "reuse-task-worktree" }
| {
ok: false;
resolution: MergeIntegrationRootResolution;
checked: "reuse-task-worktree";
reason: "missing-task-worktree" | "unusable-task-worktree";
classification?: Awaited<ReturnType<typeof classifyTaskWorktree>>;
};
export interface EnsureUsableMergeIntegrationRootInput {
resolution: MergeIntegrationRootResolution;
projectRoot: string;
}
/**
* Preflight the merge integration cwd before any merge-runner git spawn.
*
* In cwd/project-root mode this is intentionally a no-op: the project root is
* the stable repository checkout and the common path should not pay an extra
* stat or git-worktree-list probe. In reuse-task-worktree mode the resolved
* `task.worktree` is user/session mutable, so classify it before it can be
* used as `cwd`; callers can then reacquire/recreate the worktree instead of
* letting Node surface a misleading `spawn git ENOENT` for a vanished cwd.
*/
export async function ensureUsableMergeIntegrationRoot(
input: EnsureUsableMergeIntegrationRootInput,
): Promise<MergeIntegrationRootPreflightResult> {
if (input.resolution.mode !== "reuse-task-worktree") {
return { ok: true, resolution: input.resolution, checked: false };
}
const reusableRoot = input.resolution.rootDir.trim();
if (!reusableRoot) {
return {
ok: false,
resolution: input.resolution,
checked: "reuse-task-worktree",
reason: "missing-task-worktree",
};
}
const classification = await classifyTaskWorktree(input.projectRoot, reusableRoot);
if (!classification.ok) {
return {
ok: false,
resolution: input.resolution,
checked: "reuse-task-worktree",
reason: "unusable-task-worktree",
classification,
};
}
return { ok: true, resolution: input.resolution, checked: "reuse-task-worktree" };
}
export interface ResolveIntegrationRemoteInput {
settings: Pick<ProjectSettings, "worktreeRebaseRemote">;
rootDir: string;
integrationBranch: string;
}
export async function resolveIntegrationRemote(
input: ResolveIntegrationRemoteInput,
): Promise<string | undefined> {
const configured = input.settings.worktreeRebaseRemote?.trim();
if (configured) {
return configured;
}
try {
const { stdout } = await execAsync(
`git config --get branch.${input.integrationBranch}.remote`,
{ cwd: input.rootDir, encoding: "utf-8" },
);
const branchRemote = stdout.trim();
if (branchRemote) {
return branchRemote;
}
} catch {
// Fall through to repo remote discovery.
}
try {
const { stdout } = await execAsync("git remote", {
cwd: input.rootDir,
encoding: "utf-8",
});
const remotes = stdout.trim().split(/\s+/).filter(Boolean);
if (remotes.length === 1) {
return remotes[0];
}
if (remotes.includes("origin")) {
return "origin";
}
} catch {
// No remote resolvable.
}
return "origin";
}
export class MergeHandoffRefusedError extends Error {
readonly reason: string;
readonly gate: string;
readonly payload: Record<string, unknown>;
constructor(gate: string, reason: string, payload: Record<string, unknown> = {}) {
super(`Merge handoff refused (${gate}): ${reason}`);
this.name = "MergeHandoffRefusedError";
this.gate = gate;
this.reason = reason;
this.payload = payload;
}
}
export interface ReuseHandoffSuccess {
ok: true;
taskId: string;
worktreePath: string;
branch: string;
workerId: string;
releaseLease: (outcome: MergeQueueReleaseOutcome) => void;
}
export type HandoffResult = ReuseHandoffSuccess;
export interface ReuseHandoffInput {
task: Pick<
Task,
| "id"
| "branch"
| "worktree"
| "checkedOutBy"
| "checkedOutAt"
| "checkoutLeaseRenewedAt"
| "checkoutNodeId"
| "checkoutRunId"
| "checkoutLeaseEpoch"
>;
store: TaskStore;
projectRoot: string;
settings: ProjectSettings;
worktreePath: string;
auditEmit?: (event: { type: string; target?: string; metadata?: Record<string, unknown> }) => Promise<void> | void;
}
export async function snapshotDirtyFilesLocal(rootDir: string): Promise<Set<string>> {
const paths = new Set<string>();
try {
const [unstagedOut, stagedOut, porcelainOut] = await Promise.all([
execFileAsync("git", ["diff", "-z", "--name-only"], { cwd: rootDir, encoding: "utf-8" }).then(
(r) => r.stdout,
() => "",
),
execFileAsync("git", ["diff", "-z", "--cached", "--name-only"], { cwd: rootDir, encoding: "utf-8" }).then(
(r) => r.stdout,
() => "",
),
execFileAsync("git", ["status", "-z", "--porcelain"], { cwd: rootDir, encoding: "utf-8" }).then(
(r) => r.stdout,
() => "",
),
]);
for (const entry of unstagedOut.split("\0")) {
const path = entry.trim();
if (path) paths.add(path);
}
for (const entry of stagedOut.split("\0")) {
const path = entry.trim();
if (path) paths.add(path);
}
for (const entry of porcelainOut.split("\0")) {
if (!entry.startsWith("?? ")) continue;
const path = entry.slice(3);
if (path) paths.add(path);
}
} catch {
// Best-effort gate input.
}
return paths;
}
export async function gitDirtyFingerprintLocal(rootDir: string): Promise<string> {
try {
const [diffOut, statusOut] = await Promise.all([
execFileAsync("git", ["diff", "HEAD"], {
cwd: rootDir,
encoding: "utf-8",
maxBuffer: 64 * 1024 * 1024,
}).then((r) => r.stdout, () => ""),
execFileAsync("git", ["status", "-z", "--porcelain"], { cwd: rootDir, encoding: "utf-8" }).then(
(r) => r.stdout,
() => "",
),
]);
if (!diffOut && !statusOut) return "";
return createHash("sha256").update(diffOut).update("\0").update(statusOut).digest("hex");
} catch {
return "";
}
}
export interface IntegrationWorktreeProbeResult {
userCheckout: {
worktreePath: string;
dirty: boolean;
untrackedCount: number;
dirtyPathSample: string[];
} | null;
dirtyFingerprint: string | null;
}
export interface ProbeIntegrationWorktreeStateInput {
rootDir: string;
integrationBranch: string;
projectRoot: string;
}
export async function probeIntegrationWorktreeState(
input: ProbeIntegrationWorktreeStateInput,
): Promise<IntegrationWorktreeProbeResult> {
try {
const branchMap = await getRegisteredWorktreeBranchMap(input.projectRoot);
const caseInsensitiveMatches = Array.from(branchMap.entries())
.filter(([branch]) => branch.toLowerCase() === input.integrationBranch.toLowerCase())
.map(([, worktreePath]) => worktreePath);
const registeredPath = branchMap.get(input.integrationBranch)
?? caseInsensitiveMatches.find((worktreePath) => canonicalizePath(worktreePath) === canonicalizePath(input.rootDir))
?? caseInsensitiveMatches[0]
?? null;
if (!registeredPath) {
return { userCheckout: null, dirtyFingerprint: null };
}
const dirtyPaths = Array.from(await snapshotDirtyFilesLocal(registeredPath)).sort();
const dirtyFingerprint = await gitDirtyFingerprintLocal(registeredPath);
let untrackedCount = 0;
try {
const { stdout } = await execFileAsync("git", ["status", "-z", "--porcelain"], {
cwd: registeredPath,
encoding: "utf-8",
});
untrackedCount = stdout.split("\0").filter((entry) => entry.startsWith("?? ")).length;
} catch {
// best-effort
}
return {
userCheckout: {
worktreePath: registeredPath,
dirty: dirtyPaths.length > 0 || Boolean(dirtyFingerprint),
untrackedCount,
dirtyPathSample: dirtyPaths.slice(0, 20),
},
dirtyFingerprint: dirtyFingerprint || null,
};
} catch {
return { userCheckout: null, dirtyFingerprint: null };
}
}
/*
FNXC:WorkflowResolvedColumns 2026-07-30-13:10 (batch-engine tail):
"Someone else is still using this worktree" excludes tasks that have FINISHED. Keyed on the literal, a
renamed complete lane meant a shipped task still counted as a live user, so the integration worktree
could never be reused and the merge path took the slower rebuild every time.
Resolved per CANDIDATE task (each may run its own workflow), one IR cache for the scan, and only for
the tasks that actually share the worktree path — the common case resolves nothing at all.
*/
async function findOtherWorktreeUser(store: TaskStore, worktreePath: string, excludeTaskId: string): Promise<string | null> {
const tasks = await store.listTasks({ slim: true, includeArchived: false } as never);
const irCache = new Map<string, WorkflowIr>();
for (const task of tasks) {
if (task.id === excludeTaskId) continue;
if (task.worktree !== worktreePath) continue;
const completeColumns = new Set<string>(["done"]);
try {
const ir = await resolveWorkflowIrForTask(store, task.id, irCache);
if (ir) for (const id of columnsWithFlag(ir, "complete")) completeColumns.add(id);
} catch { /* degraded: legacy id only */ }
if (!completeColumns.has(task.column)) {
return task.id;
}
}
return null;
}
function asCentralClaimAccessor(store: TaskStore): {
projectId?: string;
getTaskClaim?: (projectId: string, taskId: string) => { ownerAgentId?: string | null } | null;
} {
const candidate = store as TaskStore & {
projectId?: string;
getTaskClaim?: (projectId: string, taskId: string) => { ownerAgentId?: string | null } | null;
};
return {
projectId: typeof candidate.projectId === "string" ? candidate.projectId : undefined,
getTaskClaim: typeof candidate.getTaskClaim === "function" ? candidate.getTaskClaim.bind(candidate) : undefined,
};
}
export async function acquireReuseHandoff(input: ReuseHandoffInput): Promise<HandoffResult> {
const expectedBranch = canonicalFusionBranchName(input.task.id);
const worktreePath = input.worktreePath;
const preflight = await ensureUsableMergeIntegrationRoot({
resolution: {
mode: "reuse-task-worktree",
rootDir: worktreePath,
branchName: expectedBranch,
},
projectRoot: input.projectRoot,
});
if (!preflight.ok) {
throw new MergeHandoffRefusedError("integration-root-preflight", preflight.reason, {
taskId: input.task.id,
worktreePath,
classification: preflight.classification ?? null,
});
}
if (canonicalizePath(worktreePath) === canonicalizePath(input.projectRoot)) {
throw new MergeHandoffRefusedError("reuse-misconfigured", "worktree-equals-project-root", {
taskId: input.task.id,
projectRoot: input.projectRoot,
worktreePath,
});
}
const dirtyPaths = Array.from(await snapshotDirtyFilesLocal(worktreePath)).sort();
const dirtyFingerprint = await gitDirtyFingerprintLocal(worktreePath);
if (dirtyPaths.length > 0 || dirtyFingerprint) {
// Previously this refused the handoff and parked the task as
// in-review:failed. Instead, autostash the dirty state so the merge can
// proceed; the stash survives in the repo's stash list even after the
// worktree is later torn down, so the developer can always recover.
const stashLabel = `fusion-reuse-handoff-autostash:${input.task.id}:${Date.now()}`;
let stashSha: string | null = null;
let stashError: string | null = null;
try {
// Stage everything (including untracked) so `git stash create`
// captures the full dirty tree.
await execAsync("git add -A", { cwd: worktreePath });
const { stdout: createOut } = await execAsync("git stash create", {
cwd: worktreePath,
encoding: "utf-8",
});
stashSha = String(createOut).trim() || null;
if (stashSha) {
await execAsync(
`git stash store -m ${quoteAutostashArg(stashLabel)} ${stashSha}`,
{ cwd: worktreePath },
);
// Only reset/clean once the dirty content is safely captured by the
// stash. If the stash failed (no SHA), leaving the working tree as-is
// preserves the user's edits for manual recovery.
await execAsync("git reset --hard HEAD", { cwd: worktreePath });
await execAsync("git clean -fd", { cwd: worktreePath });
}
} catch (err: unknown) {
stashError = err instanceof Error ? err.message : String(err);
}
if (!stashSha || stashError) {
// Stash creation failed: do NOT proceed (the merge's destructive ops
// would wipe the user's edits). Best-effort unstage so the worktree
// isn't left with a half-staged index from the `git add -A` above.
try {
await execAsync("git reset", { cwd: worktreePath });
} catch {
// Nothing more we can do.
}
throw new MergeHandoffRefusedError("working-tree-dirty", "dirty-worktree-autostash-failed", {
taskId: input.task.id,
worktreePath,
dirtyPaths,
dirtyFingerprint,
stashError,
});
}
await input.auditEmit?.({
type: "merge:reuse-handoff-autostash",
target: worktreePath,
metadata: {
taskId: input.task.id,
worktreePath,
stashSha,
stashLabel,
dirtyPathCount: dirtyPaths.length,
dirtyPathSample: dirtyPaths.slice(0, 20),
recoverCommand: `cd ${worktreePath} && git stash apply ${stashSha}`,
},
});
}
const { stdout: headStdout } = await execAsync("git rev-parse --abbrev-ref HEAD", {
cwd: worktreePath,
encoding: "utf-8",
timeout: 10_000,
maxBuffer: 1024 * 1024,
});
let observedBranch = headStdout.trim();
if (observedBranch && observedBranch !== expectedBranch && observedBranch.toLowerCase() === expectedBranch.toLowerCase()) {
const autocorrectResult = await attemptBranchAutocorrect({
worktreePath,
observedBranch,
expectedBranch,
rootDir: input.projectRoot,
});
if (autocorrectResult.status !== "failed") {
await input.auditEmit?.({
type: "branch:auto-canonicalize-case",
target: worktreePath,
metadata: {
taskId: input.task.id,
observed: observedBranch,
expected: expectedBranch,
worktreePath,
mode: autocorrectResult.status,
},
});
const { stdout: correctedHead } = await execAsync("git rev-parse --abbrev-ref HEAD", {
cwd: worktreePath,
encoding: "utf-8",
timeout: 10_000,
maxBuffer: 1024 * 1024,
});
observedBranch = correctedHead.trim();
}
}
if (observedBranch !== expectedBranch) {
// The worktree's HEAD points elsewhere (detached or different branch) but
// the expected branch ref may still hold this task's authoritative work.
// If the branch tip carries the task's Fusion-Task-Id trailer and the
// range against base is contamination-free, re-attach via plain checkout
// (worktree was already asserted clean above, so this is safe and
// non-destructive — unlike `checkout -B` which would clobber the ref).
const authority = await isBranchAuthoritativeForTask(
input.projectRoot,
expectedBranch,
input.task.id,
);
if (authority.ok) {
const reattach = await execAsync(`git checkout ${expectedBranch}`, {
cwd: worktreePath,
encoding: "utf-8",
timeout: 30_000,
maxBuffer: 1024 * 1024,
}).then(
() => ({ ok: true as const }),
(err: unknown) => ({ ok: false as const, reason: err instanceof Error ? err.message : String(err) }),
);
if (reattach.ok) {
const { stdout: reattachedHead } = await execAsync("git rev-parse --abbrev-ref HEAD", {
cwd: worktreePath,
encoding: "utf-8",
timeout: 10_000,
maxBuffer: 1024 * 1024,
});
observedBranch = reattachedHead.trim();
await input.auditEmit?.({
type: "branch:auto-reattach-authoritative",
target: worktreePath,
metadata: {
taskId: input.task.id,
previousHead: observedBranch === expectedBranch ? undefined : observedBranch,
expectedBranch,
worktreePath,
},
});
}
}
if (observedBranch !== expectedBranch) {
throw new MergeHandoffRefusedError("head-branch-mismatch", "unexpected-branch", {
taskId: input.task.id,
worktreePath,
observedBranch,
expectedBranch,
authorityProbe: authority.ok ? "ok" : authority.reason,
});
}
}
const activeRecord = activeSessionRegistry.lookupByPath(worktreePath);
if (activeRecord) {
if (activeRecord.taskId === input.task.id) {
// FN-5256: route through the hardened helper so the minIdleMs window and
// processActiveProbe (executingTaskLock) gates apply — bare
// `reconcileStaleSelfOwned` would race a warming-down session.
const outcome = reconcileSelfOwnedActiveSessionForRemoval(
activeSessionRegistry,
worktreePath,
input.task.id,
() => false,
{ processActiveProbe: (probeTaskId) => executingTaskLock.has(probeTaskId) },
);
if (outcome.action !== "reconciled") {
throw new MergeHandoffRefusedError("active-session-binding", "active-session-present", {
taskId: input.task.id,
worktreePath,
activeRecord,
executingTaskLockHeld: executingTaskLock.has(input.task.id),
reconcileOutcome: outcome.action,
});
}
} else {
throw new MergeHandoffRefusedError("active-session-binding", "active-session-present", {
taskId: input.task.id,
worktreePath,
activeRecord,
executingTaskLockHeld: executingTaskLock.has(input.task.id),
});
}
}
const classification = await classifyTaskWorktree(input.projectRoot, worktreePath);
if (!classification.ok) {
throw new MergeHandoffRefusedError("branch-worktree-mapping", classification.classification, {
taskId: input.task.id,
worktreePath,
classification,
});
}
const otherTaskId = await findOtherWorktreeUser(input.store, worktreePath, input.task.id);
if (otherTaskId) {
throw new MergeHandoffRefusedError("branch-worktree-mapping", "foreign-task-worktree-owner", {
taskId: input.task.id,
worktreePath,
otherTaskId,
});
}
if (input.task.branch?.trim() && input.task.branch.trim().toLowerCase() !== expectedBranch.toLowerCase()) {
throw new MergeHandoffRefusedError("branch-worktree-mapping", "task-branch-metadata-mismatch", {
taskId: input.task.id,
taskBranch: input.task.branch,
expectedBranch,
});
}
const branchMap = await getRegisteredWorktreeBranchMap(input.projectRoot);
const registeredBranchPath = branchMap.get(expectedBranch);
const canonicalWorktreePath = canonicalizePath(worktreePath);
if (!registeredBranchPath || canonicalizePath(registeredBranchPath) !== canonicalWorktreePath) {
throw new MergeHandoffRefusedError("branch-worktree-mapping", "registered-branch-mismatch", {
taskId: input.task.id,
worktreePath,
expectedBranch,
registeredBranchPath: registeredBranchPath ?? null,
});
}
const staleCheck = new MeshLeaseManager({
taskStore: input.store,
getExecutingTaskIds: () => {
const active = new Set<string>();
if (executingTaskLock.has(input.task.id)) {
active.add(input.task.id);
}
return active;
},
});
if (input.task.checkedOutBy) {
const recoverable = await staleCheck.isLeaseRecoverable(input.task as Task);
if (!recoverable.recoverable) {
throw new MergeHandoffRefusedError("lease-handoff-failed", "executor-lease-active", {
taskId: input.task.id,
checkedOutBy: input.task.checkedOutBy,
checkoutNodeId: input.task.checkoutNodeId ?? null,
checkoutRunId: input.task.checkoutRunId ?? null,
reason: recoverable.reason ?? null,
});
}
}
if (executingTaskLock.has(input.task.id)) {
throw new MergeHandoffRefusedError("lease-handoff-failed", "executor-lease-active", {
taskId: input.task.id,
worktreePath,
reason: "active_local_execution",
});
}
const centralAccessor = asCentralClaimAccessor(input.store);
if (centralAccessor.projectId && centralAccessor.getTaskClaim) {
const claim = centralAccessor.getTaskClaim(centralAccessor.projectId, input.task.id);
if (claim?.ownerAgentId) {
throw new MergeHandoffRefusedError("lease-handoff-failed", "central-conflict", {
taskId: input.task.id,
projectId: centralAccessor.projectId,
ownerAgentId: claim.ownerAgentId,
});
}
}
let lease;
try {
// Non-atomic fallback: executor lease checks above race with mergeQueue lease acquisition.
// TaskStore does not yet expose an atomic executor-lease absence check inside mergeQueue leasing.
lease = await (input.store as TaskStore & {
acquireMergeQueueLease(workerId: string, opts: { leaseDurationMs: number; now?: string; targetTaskId?: string }): Promise<unknown>;
}).acquireMergeQueueLease(MERGE_HANDOFF_WORKER_ID, {
leaseDurationMs: 15 * 60 * 1000,
targetTaskId: input.task.id,
});
} catch (error) {
if (error instanceof PoolDoubleLeaseError) {
throw new MergeHandoffRefusedError("lease-handoff-failed", "pool-double-lease", {
taskId: input.task.id,
worktreePath,
path: error.path,
existingHolder: error.existingHolder,
requestingTaskId: error.requestingTaskId,
phase: error.phase,
});
}
throw error;
}
if (!lease) {
throw new MergeHandoffRefusedError("lease-handoff-failed", "target-not-queued", {
taskId: input.task.id,
worktreePath,
});
}
if (!("taskId" in lease) || lease.taskId !== input.task.id) {
const queueHead = await (input.store as TaskStore & {
peekMergeQueueHead?: () => Promise<{ taskId: string; leasedBy: string | null; column: string | null } | null>;
}).peekMergeQueueHead?.();
throw new MergeHandoffRefusedError("lease-handoff-failed", "no-lease", {
taskId: input.task.id,
worktreePath,
acquiredTaskId: "taskId" in lease ? lease.taskId : null,
queueHeadTaskId: queueHead?.taskId ?? null,
queueHeadLeasedBy: queueHead?.leasedBy ?? null,
});
}
// Re-check executor lease after acquiring the merge-queue lease: the
// checks above (lines ~362–391) are non-atomic with acquisition, so a
// local executor could grab the task between them. Releasing here gives
// a precise diagnostic instead of letting the merge proceed with a
// conflicting executor lease and surfacing as a generic failure later.
if (executingTaskLock.has(input.task.id)) {
(input.store as TaskStore & {
releaseMergeQueueLease(taskId: string, workerId: string, outcome: MergeQueueReleaseOutcome): Promise<void>;
}).releaseMergeQueueLease(input.task.id, MERGE_HANDOFF_WORKER_ID, {
kind: "failure",
error: "executor-lease-acquired-after-queue-lease",
});
throw new MergeHandoffRefusedError("lease-handoff-failed", "executor-lease-race-detected", {
taskId: input.task.id,
worktreePath,
reason: "executor_lease_acquired_after_queue_lease",
});
}
return {
ok: true,
taskId: input.task.id,
worktreePath,
branch: expectedBranch,
workerId: MERGE_HANDOFF_WORKER_ID,
releaseLease: (outcome) => {
void (input.store as TaskStore & {
releaseMergeQueueLease(taskId: string, workerId: string, outcome: MergeQueueReleaseOutcome): Promise<void>;
}).releaseMergeQueueLease(input.task.id, MERGE_HANDOFF_WORKER_ID, outcome);
},
};
}
export async function releaseReuseHandoff(input: {
handoff: ReuseHandoffSuccess;
outcome: string;
auditEmit?: (event: { type: string; target?: string; metadata?: Record<string, unknown> }) => Promise<void> | void;
}): Promise<void> {
input.handoff.releaseLease(
input.outcome === "success"
? { kind: "success" }
: { kind: "failure", error: input.outcome },
);
await input.auditEmit?.({
type: "merge:reuse-handoff-released",
target: input.handoff.worktreePath,
metadata: {
taskId: input.handoff.taskId,
outcome: input.outcome,
branch: input.handoff.branch,
worktreePath: input.handoff.worktreePath,
},
});
}