fix(FN-4949): defer stale active branch reclaim during live executor work
- Add stale-branch reclaim guards for active executor sessions, recent execution starts, and worktrees with uncommitted changes - Emit branch:stale-active-reclaim-deferred run-audit telemetry with deferral metadata for self-healing decisions - Cover reclaim deferral and legitimate reclaim behavior with reliability interaction tests - Document the new stale active branch deferral contract in AGENTS.md - Add a patch changeset for the published CLI package Fusion-Task-Id: FN-4949
This commit is contained in:
committed by
gsxdsm
parent
22729135e8
commit
287673c261
5
.changeset/fn-4949-reclaim-guards.md
Normal file
5
.changeset/fn-4949-reclaim-guards.md
Normal file
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@runfusion/fusion": patch
|
||||
---
|
||||
|
||||
Engine: `reclaim-stale-active-branches` now defers reclaim when the task has a registered active session, a recent `executionStartedAt`, or a worktree with uncommitted changes. Emits a new `branch:stale-active-reclaim-deferred` run-audit event per deferral. Fixes FN-4924-class loops where executor work was wiped because per-step commits were absent.
|
||||
@@ -169,7 +169,7 @@ Port 4040 is the production dashboard port. A user's live session is typically r
|
||||
Detailed mechanism logs live in `docs/architecture.md` and `docs/design/`. The contracts agents must respect:
|
||||
|
||||
- **Orphan `fusion/*` branches**: prune-or-rescue, never force-delete. Subsumed branches pruned; unique-commit branches rescued into triage tasks.
|
||||
- **Stale active branches**: self-healing's `reclaim-stale-active-branches` stage prunes a `fusion/<task-id>` branch with zero unique commits when no usable worktree mapping exists, then clears `task.branch`/`task.worktree`/`task.baseCommitSha`.
|
||||
- **Stale active branches**: self-healing's `reclaim-stale-active-branches` stage prunes a `fusion/<task-id>` branch with zero unique commits when no usable worktree mapping exists, then clears `task.branch`/`task.worktree`/`task.baseCommitSha`. It must defer reclaim (emit `branch:stale-active-reclaim-deferred`) when the task worktree is in `activeSessionRegistry`, when `executionStartedAt` is within `STALE_ACTIVE_BRANCH_EXECUTION_GRACE_MS` (10 minutes), or when the mapped worktree has uncommitted changes.
|
||||
- **Worktree metadata reconcile ordering (FN-4962)**: `reconcile-task-worktree-metadata` must run before `reclaim-stale-active-branches`; stale `task.worktree` metadata is rebound to live `fusion/<task-id>` worktrees when present (`task:auto-recover-worktree-metadata-rebound`) or cleared (`task:auto-recover-worktree-metadata-cleared`) when absent.
|
||||
- **Completion fan-out is synchronous**: `SelfHealingManager.reconcileCompletedTask()` runs on `in-review → done`. Downstream stale `blockedBy` links and residual `fusion/<task-id>` branch/worktree artifacts are reconciled immediately, not on a periodic sweep.
|
||||
- **In-review stall deadlock**: identical stalls (same code + reason) repeated past `inReviewStallDeadlockThreshold` (default 3) auto-pause with `pausedReason: "in-review-stall-deadlock"` and `status: "failed"`.
|
||||
|
||||
@@ -0,0 +1,207 @@
|
||||
import { describe, it, expect, vi, beforeEach, afterEach } from "vitest";
|
||||
import { EventEmitter } from "node:events";
|
||||
import { mkdtempSync, mkdirSync, rmSync, writeFileSync } from "node:fs";
|
||||
import { tmpdir } from "node:os";
|
||||
import { join } from "node:path";
|
||||
import { execSync } from "node:child_process";
|
||||
import type { Settings, Task, TaskStore } from "@fusion/core";
|
||||
import { SelfHealingManager, STALE_ACTIVE_BRANCH_EXECUTION_GRACE_MS } from "../../self-healing.js";
|
||||
import { activeSessionRegistry } from "../../active-session-registry.js";
|
||||
|
||||
function sh(command: string, cwd: string): string {
|
||||
return String(execSync(command, { cwd, encoding: "utf-8", stdio: ["pipe", "pipe", "pipe"] }) ?? "");
|
||||
}
|
||||
|
||||
function makeRepo(): string {
|
||||
const root = mkdtempSync(join(tmpdir(), "fn-4949-"));
|
||||
sh("git init", root);
|
||||
sh("git config user.email 'test@example.com'", root);
|
||||
sh("git config user.name 'Test User'", root);
|
||||
writeFileSync(join(root, "README.md"), "base\n", "utf-8");
|
||||
sh("git add README.md", root);
|
||||
sh("git commit -m 'init'", root);
|
||||
sh("git branch -M main", root);
|
||||
return root;
|
||||
}
|
||||
|
||||
function createFusionBranch(repo: string, taskId: string): string {
|
||||
const branch = `fusion/${taskId.toLowerCase()}`;
|
||||
sh(`git checkout -b ${branch}`, repo);
|
||||
sh("git checkout main", repo);
|
||||
return branch;
|
||||
}
|
||||
|
||||
function makeTask(taskId: string, branch: string, worktree: string | null, executionStartedAt: string): Task {
|
||||
return {
|
||||
id: taskId,
|
||||
title: "test",
|
||||
description: "test",
|
||||
column: "in-progress",
|
||||
branch,
|
||||
worktree: worktree ?? undefined,
|
||||
executionStartedAt,
|
||||
paused: false,
|
||||
userPaused: false,
|
||||
checkedOutBy: undefined,
|
||||
pausedReason: undefined,
|
||||
dependencies: [],
|
||||
steps: Array.from({ length: 7 }, (_, index) => ({ name: `step-${index + 1}`, status: "done" })),
|
||||
currentStep: 7,
|
||||
log: [],
|
||||
createdAt: new Date().toISOString(),
|
||||
updatedAt: new Date().toISOString(),
|
||||
} as unknown as Task;
|
||||
}
|
||||
|
||||
function makeStore(task: Task): TaskStore & EventEmitter & { auditEvents: any[] } {
|
||||
const emitter = new EventEmitter();
|
||||
const settings = {
|
||||
autoMerge: true,
|
||||
globalPause: false,
|
||||
enginePaused: false,
|
||||
baseBranch: "main",
|
||||
mergeStrategy: "direct",
|
||||
autoRecovery: { mode: "deterministic-only", maxRetries: 3 },
|
||||
} as unknown as Settings;
|
||||
const auditEvents: any[] = [];
|
||||
return Object.assign(emitter, {
|
||||
auditEvents,
|
||||
getSettings: vi.fn(async () => settings),
|
||||
getTask: vi.fn(async () => task),
|
||||
listTasks: vi.fn(async () => [task]),
|
||||
updateTask: vi.fn(async (_id: string, updates: Partial<Task>) => Object.assign(task, updates)),
|
||||
moveTask: vi.fn(async (_id: string, column: Task["column"]) => {
|
||||
task.column = column;
|
||||
return task;
|
||||
}),
|
||||
logEntry: vi.fn(async () => undefined),
|
||||
appendAgentLog: vi.fn(async () => undefined),
|
||||
updateSettings: vi.fn(async () => settings),
|
||||
clearStaleExecutionStartBranchReferences: vi.fn(() => []),
|
||||
recordRunAuditEvent: vi.fn(async (event: any) => {
|
||||
auditEvents.push(event);
|
||||
}),
|
||||
walCheckpoint: vi.fn(() => ({ busy: 0, log: 0, checkpointed: 0 })),
|
||||
archiveTaskAndCleanup: vi.fn(async () => ({})),
|
||||
mergeTask: vi.fn(async () => undefined),
|
||||
getRootDir: vi.fn(() => "/tmp/test"),
|
||||
}) as unknown as TaskStore & EventEmitter & { auditEvents: any[] };
|
||||
}
|
||||
|
||||
describe("FN-4924 / FN-4949: reclaim-stale-active-branches defers in-flight executor work", () => {
|
||||
const tempRoots: string[] = [];
|
||||
|
||||
beforeEach(() => {
|
||||
activeSessionRegistry.clear();
|
||||
vi.restoreAllMocks();
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
activeSessionRegistry.clear();
|
||||
for (const root of tempRoots.splice(0)) {
|
||||
rmSync(root, { recursive: true, force: true });
|
||||
}
|
||||
});
|
||||
|
||||
it("defers when execution started recently (before reclaim)", async () => {
|
||||
const repo = makeRepo();
|
||||
tempRoots.push(repo);
|
||||
const branch = createFusionBranch(repo, "FN-4924");
|
||||
const worktree = join(repo, ".worktrees", "fn-4924");
|
||||
mkdirSync(worktree, { recursive: true });
|
||||
writeFileSync(join(worktree, "dirty.txt"), "dirty\n", "utf-8");
|
||||
|
||||
const task = makeTask("FN-4924", branch, worktree, new Date(Date.now() - 60_000).toISOString());
|
||||
const store = makeStore(task);
|
||||
const manager = new SelfHealingManager(store as any, { rootDir: repo } as any);
|
||||
vi.spyOn(manager as any, "inspectOrphanedBranch").mockResolvedValue({
|
||||
branch,
|
||||
tipSha: "abc123def456",
|
||||
uniqueCommitCount: 0,
|
||||
subjects: [],
|
||||
});
|
||||
|
||||
const recovered = await manager.reclaimStaleActiveBranches();
|
||||
|
||||
expect(recovered).toBe(0);
|
||||
expect((store.updateTask as any).mock.calls.some((call: any[]) => call[1]?.worktree === null && call[1]?.branch === null)).toBe(false);
|
||||
expect(sh(`git rev-parse --verify ${branch}`, repo).trim().length).toBeGreaterThan(0);
|
||||
expect(store.auditEvents.some((event) => event.mutationType === "branch:stale-active-reclaim-deferred" && event.metadata?.reason === "recent-execution-started")).toBe(true);
|
||||
|
||||
manager.stop();
|
||||
});
|
||||
|
||||
it("defers when worktree is in activeSessionRegistry", async () => {
|
||||
const repo = makeRepo();
|
||||
tempRoots.push(repo);
|
||||
const branch = createFusionBranch(repo, "FN-4925");
|
||||
const worktree = join(repo, ".worktrees", "fn-4925");
|
||||
mkdirSync(worktree, { recursive: true });
|
||||
|
||||
const oldStart = new Date(Date.now() - STALE_ACTIVE_BRANCH_EXECUTION_GRACE_MS - 60_000).toISOString();
|
||||
const task = makeTask("FN-4925", branch, worktree, oldStart);
|
||||
const store = makeStore(task);
|
||||
activeSessionRegistry.registerPath(worktree, { taskId: task.id, kind: "executor", ownerKey: task.id });
|
||||
|
||||
const manager = new SelfHealingManager(store as any, { rootDir: repo } as any);
|
||||
const inspectSpy = vi.spyOn(manager as any, "inspectOrphanedBranch");
|
||||
|
||||
const recovered = await manager.reclaimStaleActiveBranches();
|
||||
|
||||
expect(recovered).toBe(0);
|
||||
expect(inspectSpy).not.toHaveBeenCalled();
|
||||
expect(store.auditEvents.some((event) => event.mutationType === "branch:stale-active-reclaim-deferred" && event.metadata?.reason === "active-session")).toBe(true);
|
||||
manager.stop();
|
||||
});
|
||||
|
||||
it("still reclaims legitimate stale active branch without in-flight signals", async () => {
|
||||
const repo = makeRepo();
|
||||
tempRoots.push(repo);
|
||||
const branch = createFusionBranch(repo, "FN-4926");
|
||||
|
||||
const oldStart = new Date(Date.now() - STALE_ACTIVE_BRANCH_EXECUTION_GRACE_MS - 60_000).toISOString();
|
||||
const task = makeTask("FN-4926", branch, null, oldStart);
|
||||
const store = makeStore(task);
|
||||
const manager = new SelfHealingManager(store as any, { rootDir: repo } as any);
|
||||
vi.spyOn(manager as any, "inspectOrphanedBranch").mockResolvedValue({
|
||||
branch,
|
||||
tipSha: "abc123def456",
|
||||
uniqueCommitCount: 0,
|
||||
subjects: [],
|
||||
});
|
||||
|
||||
const recovered = await manager.reclaimStaleActiveBranches();
|
||||
|
||||
expect(recovered).toBe(1);
|
||||
expect((store.updateTask as any).mock.calls.some((call: any[]) => call[1]?.worktree === null && call[1]?.branch === null && call[1]?.baseCommitSha === null)).toBe(true);
|
||||
expect(() => sh(`git rev-parse --verify ${branch}`, repo)).toThrow();
|
||||
expect(store.auditEvents.some((event) => event.mutationType === "branch:stale-active-reclaim" && event.target === branch)).toBe(true);
|
||||
manager.stop();
|
||||
});
|
||||
|
||||
it("reproduces FN-4924 loop conditions and defers reclaim", async () => {
|
||||
const repo = makeRepo();
|
||||
tempRoots.push(repo);
|
||||
const branch = createFusionBranch(repo, "FN-4924");
|
||||
const worktree = join(repo, ".worktrees", "fn-4924-loop");
|
||||
mkdirSync(worktree, { recursive: true });
|
||||
writeFileSync(join(worktree, "session-progress.txt"), "pending commit\n", "utf-8");
|
||||
|
||||
const task = makeTask("FN-4924", branch, worktree, new Date(Date.now() - 2 * 60_000).toISOString());
|
||||
const store = makeStore(task);
|
||||
const manager = new SelfHealingManager(store as any, { rootDir: repo } as any);
|
||||
vi.spyOn(manager as any, "inspectOrphanedBranch").mockResolvedValue({
|
||||
branch,
|
||||
tipSha: "abc123def456",
|
||||
uniqueCommitCount: 0,
|
||||
subjects: [],
|
||||
});
|
||||
|
||||
const recovered = await manager.reclaimStaleActiveBranches();
|
||||
|
||||
expect(recovered).toBe(0);
|
||||
expect((store.updateTask as any).mock.calls.some((call: any[]) => call[1]?.worktree === null && call[1]?.branch === null)).toBe(false);
|
||||
expect(store.auditEvents.some((event) => event.mutationType === "branch:stale-active-reclaim-deferred")).toBe(true);
|
||||
manager.stop();
|
||||
});
|
||||
});
|
||||
@@ -138,6 +138,7 @@ export type GitMutationType =
|
||||
| "merge:audit-failure"
|
||||
| "branch:auto-reclaim"
|
||||
| "branch:stale-active-reclaim"
|
||||
| "branch:stale-active-reclaim-deferred"
|
||||
| "branch:orphan-prune"
|
||||
| "branch:orphan-rescued"
|
||||
| "branch:reanchor"
|
||||
|
||||
@@ -50,6 +50,7 @@ const log = createLogger("self-healing");
|
||||
const worktreeMetadataReconcileLog = createLogger("worktree-metadata-reconcile");
|
||||
const execAsync = promisify(exec);
|
||||
const DONE_TASK_INTEGRITY_SWEEP_LIMIT = 50;
|
||||
export const STALE_ACTIVE_BRANCH_EXECUTION_GRACE_MS = 10 * 60_000;
|
||||
|
||||
async function classifyOwnedLandedEvidenceForSelfHealing(rootDir: string, task: Task, mergeTargetBranch: string): Promise<OwnedLandedClassification> {
|
||||
const { classifyOwnedLandedEvidence } = await import("./merger.js");
|
||||
@@ -2008,6 +2009,64 @@ export class SelfHealingManager {
|
||||
}
|
||||
if (activeTaskIds.has(task.id.toUpperCase())) continue;
|
||||
|
||||
const emitDeferredReclaimAudit = async (reason: "active-session" | "recent-execution-started" | "worktree-has-uncommitted-changes", hasActiveSession: boolean, hasUncommittedChanges: boolean): Promise<void> => {
|
||||
log.log(`[self-healing] deferring stale-active-branch reclaim for ${task.id}: reason=${reason}`);
|
||||
try {
|
||||
const auditor = createRunAuditor(this.store, {
|
||||
runId: generateSyntheticRunId("self-heal", task.id),
|
||||
agentId: "self-healing",
|
||||
taskId: task.id,
|
||||
taskLineageId: task.lineageId,
|
||||
phase: "reclaim-stale-active-branches",
|
||||
});
|
||||
await auditor.git({
|
||||
type: "branch:stale-active-reclaim-deferred",
|
||||
target: branch,
|
||||
metadata: {
|
||||
taskId: task.id,
|
||||
branch,
|
||||
reason,
|
||||
executionStartedAt: task.executionStartedAt ?? null,
|
||||
hasActiveSession,
|
||||
hasUncommittedChanges,
|
||||
},
|
||||
});
|
||||
} catch (auditErr: unknown) {
|
||||
log.warn(`Failed to write branch:stale-active-reclaim-deferred run-audit event for ${task.id}: ${auditErr instanceof Error ? auditErr.message : String(auditErr)}`);
|
||||
}
|
||||
};
|
||||
|
||||
const hasActiveSession = Boolean(task.worktree && activeSessionRegistry.isPathActive(task.worktree));
|
||||
if (hasActiveSession) {
|
||||
await emitDeferredReclaimAudit("active-session", true, false);
|
||||
continue;
|
||||
}
|
||||
|
||||
const executionStartedAtMs = task.executionStartedAt ? Date.parse(task.executionStartedAt) : Number.NaN;
|
||||
const isRecentlyStarted = Number.isFinite(executionStartedAtMs) && Date.now() - executionStartedAtMs <= STALE_ACTIVE_BRANCH_EXECUTION_GRACE_MS;
|
||||
if (isRecentlyStarted) {
|
||||
await emitDeferredReclaimAudit("recent-execution-started", false, false);
|
||||
continue;
|
||||
}
|
||||
|
||||
if (task.worktree && existsSync(task.worktree)) {
|
||||
try {
|
||||
if (statSync(task.worktree).isDirectory()) {
|
||||
const { stdout } = await execAsync(`git -C ${JSON.stringify(task.worktree)} status --porcelain`, {
|
||||
cwd: this.options.rootDir,
|
||||
timeout: 30_000,
|
||||
maxBuffer: 10 * 1024 * 1024,
|
||||
});
|
||||
if ((stdout ?? "").trim().length > 0) {
|
||||
await emitDeferredReclaimAudit("worktree-has-uncommitted-changes", false, true);
|
||||
continue;
|
||||
}
|
||||
}
|
||||
} catch (statusErr: unknown) {
|
||||
log.warn(`[self-healing] stale-active-branch reclaim could not determine worktree status for ${task.id}: ${statusErr instanceof Error ? statusErr.message : String(statusErr)}`);
|
||||
}
|
||||
}
|
||||
|
||||
if (task.worktree && await isUsableTaskWorktree(this.options.rootDir, task.worktree)) continue;
|
||||
|
||||
const inspection = await this.inspectOrphanedBranch(branch);
|
||||
|
||||
Reference in New Issue
Block a user