diff --git a/.changeset/review-lease-node-attribution.md b/.changeset/review-lease-node-attribution.md new file mode 100644 index 0000000000..98d25f2abe --- /dev/null +++ b/.changeset/review-lease-node-attribution.md @@ -0,0 +1,7 @@ +--- +"@runfusion/fusion": patch +--- + +summary: Review-gate leases now record which node holds them, so a restarted engine can tell its own dead leases from a peer's. +category: internal +dev: Adds `WorkflowStepResult.leaseNodeId` and an optional `LocalNodeLeaseIdentity` argument to `classifyReviewLease`. A pending lease stamped with the caller's own node id whose `startedAt` predates the current process boot now classifies as `reclaim` immediately instead of waiting out `PLAN_REVIEW_LEASE_STALENESS_MS`; peer-owned and legacy unattributed leases are unchanged. `InProcessRuntime.start()` resolves the local node id from CentralCore and passes it to SelfHealingManager. The graph executor stamps the field when `deps.localNodeId` is set — that dep is not yet threaded from the runners, so the field is not written in production yet and behavior is unchanged end to end. diff --git a/packages/core/src/types/workflow-steps.ts b/packages/core/src/types/workflow-steps.ts index 72930d10cc..7c1bf40569 100644 --- a/packages/core/src/types/workflow-steps.ts +++ b/packages/core/src/types/workflow-steps.ts @@ -265,6 +265,23 @@ export interface WorkflowStepResult { * on every terminal (passed/failed/…) record — a lease only exists while pending. */ leaseOwner?: string; + /* + * FNXC:PlanReviewLease 2026-07-26-20:05: + * Node that owns this lease. `leaseOwner` alone identifies a run but not WHERE it runs, and in a + * multi-node deployment (several engines on one central database) every node's self-healing sweep + * sees every other node's leases. Without attribution the only safe liveness test is the + * 15-minute staleness floor, because a lease that is fresh-but-unknown might be running on a peer. + * + * With it, a node can prove one specific case: a lease stamped with ITS OWN id whose `startedAt` + * predates its current process boot cannot have a live owner — the process that took it is gone. + * That is the restart-orphan case (FN-8603: an engine restart killed a Code Review session 34s in + * and the card then waited out the full floor before anything re-ran it). Peer-owned and + * unattributed (legacy) leases keep the floor. + * + * Absent on legacy rows written before this field existed; treat absence as "unknown node", never + * as "this node". + */ + leaseNodeId?: string; /* * FNXC:ReviewLaneBypass 2026-07-09-00:00: * A privileged operator can bypass a `status:"failed"` pre-merge review step diff --git a/packages/core/src/workflow-step-results.ts b/packages/core/src/workflow-step-results.ts index 64111a19b8..ac867f47d9 100644 --- a/packages/core/src/workflow-step-results.ts +++ b/packages/core/src/workflow-step-results.ts @@ -117,6 +117,17 @@ graph executor and unit tests share one lease implementation. * FN-6736 staleness-floor standard for durable single-owner leases. */ export const PLAN_REVIEW_LEASE_STALENESS_MS = 15 * 60 * 1000; +/** + * Identity a caller supplies so {@link classifyReviewLease} can recognize leases left behind by a + * PREVIOUS process on the SAME node. `nodeId` must be the cluster node id stamped into + * `WorkflowStepResult.leaseNodeId`; `processBootAt` is this process's start time (epoch ms). + * Omit it entirely to keep pure staleness-floor semantics. + */ +export interface LocalNodeLeaseIdentity { + nodeId: string; + processBootAt: number; +} + /** Classification of a review-gate's current lease state for a re-entering run. */ export type ReviewLeaseDisposition = /** No prior result — this run should claim the lease and dispatch the reviewer. */ @@ -158,6 +169,7 @@ export function classifyReviewLease( stepId: string, now: number, stalenessMs: number = PLAN_REVIEW_LEASE_STALENESS_MS, + localNode?: LocalNodeLeaseIdentity, ): ReviewLeaseDisposition { const existing = results?.find((r) => r.workflowStepId === stepId); if (!existing) return { kind: "claim" }; @@ -165,6 +177,29 @@ export function classifyReviewLease( // existing.status === "pending": it is a lease. const startedMs = existing.startedAt ? Date.parse(existing.startedAt) : Number.NaN; const ageMs = Number.isFinite(startedMs) ? now - startedMs : Number.POSITIVE_INFINITY; + /* + FNXC:PlanReviewLease 2026-07-26-20:12: + Pre-boot reclaim. A lease stamped with THIS node's id whose `startedAt` predates this process's + boot is provably dead: the process that could have owned it no longer exists. Reclaim it + immediately instead of waiting out the staleness floor — the floor exists to protect leases we + cannot attribute, and this one we can. + + Deliberately narrow, because every widening is a double-dispatch risk: + - `leaseNodeId` must be PRESENT and EQUAL to ours. Absent (legacy rows) or a peer's id both keep + the floor — under multi-node, a fresh peer lease is very likely genuinely running. + - `startedAt` must parse and be STRICTLY before boot. A lease taken by this process after boot is + a live in-process claim and must still be adopted. + Motivating incident FN-8603: an engine restart killed a Code Review session 34s in; the lease then + read "fresh" for the remaining ~14 minutes of the floor, so nothing re-ran the gate until it aged + out and was marked failed. + */ + const ownedByDeadLocalProcess = + localNode !== undefined && + existing.leaseNodeId !== undefined && + existing.leaseNodeId === localNode.nodeId && + Number.isFinite(startedMs) && + startedMs < localNode.processBootAt; + if (ownedByDeadLocalProcess) return { kind: "reclaim", priorOwner: existing.leaseOwner }; const stale = !existing.leaseOwner || !Number.isFinite(startedMs) || ageMs >= stalenessMs; if (stale) return { kind: "reclaim", priorOwner: existing.leaseOwner }; // Not stale ⇒ `leaseOwner` is guaranteed set (the stale check requires it). diff --git a/packages/engine/src/runtimes/in-process-runtime.ts b/packages/engine/src/runtimes/in-process-runtime.ts index d5d225695e..9e00c3fd85 100644 --- a/packages/engine/src/runtimes/in-process-runtime.ts +++ b/packages/engine/src/runtimes/in-process-runtime.ts @@ -325,6 +325,8 @@ export class InProcessRuntime */ private cliAgentRuntime?: BootstrappedCliAgentRuntime; private usageLimitPauser?: UsageLimitPauser; + /** FNXC:PlanReviewLease 2026-07-26-20:42: cluster node id stamped onto review-gate leases; undefined until start() resolves it, or if resolution fails. */ + private localNodeId?: string; private selfHealingManager?: SelfHealingManager; private leaseManager?: MeshLeaseManager; private leaseCentralClaimStore?: AsyncCentralClaimStore; @@ -1201,8 +1203,26 @@ export class InProcessRuntime if (!chatLayer2) throw new Error("Self-healing ChatStore requires the project PostgreSQL AsyncDataLayer"); this.chatStore ??= new ChatStore(chatLayer2); } + /* + FNXC:PlanReviewLease 2026-07-26-20:40: + Resolve this engine's cluster node id once at start so review-gate leases can be attributed. + Attribution is what lets self-healing tell "a lease my own dead process left behind" from "a + peer node's lease that is genuinely running" — the former is reclaimed immediately, the latter + keeps the 15-minute staleness floor. Fail-soft: on any error the id stays undefined, leases are + written unattributed, and floor-only semantics (the pre-existing behavior) apply. + */ + let localNodeId: string | undefined; + try { + const registeredNodes = await this.centralCore.listNodes(); + localNodeId = registeredNodes.find((node) => node.type === "local")?.id; + } catch (error) { + runtimeLog.warn(`Could not resolve local node id for review-gate lease attribution: ${error instanceof Error ? error.message : String(error)}`); + } + this.localNodeId = localNodeId; + this.selfHealingManager = new SelfHealingManager(this.taskStore, { rootDir: this.config.workingDirectory, + localNodeId, agentStore: this.agentStore, isWorktreeResumeReserved: this.cliAgentRuntime?.isWorktreeResumeReserved, recoverCompletedTask: (task) => this.executor.recoverCompletedTask(task), diff --git a/packages/engine/src/self-healing.ts b/packages/engine/src/self-healing.ts index c839149db2..6531af8076 100644 --- a/packages/engine/src/self-healing.ts +++ b/packages/engine/src/self-healing.ts @@ -272,6 +272,14 @@ const PRE_EXECUTION_WORKTREE_MAX_IDLE_MS = 30 * 24 * 60 * 60 * 1000; export interface SelfHealingOptions { /** Project root directory (parent of .worktrees/) */ rootDir: string; + /* + * FNXC:PlanReviewLease 2026-07-26-20:30: + * This engine's cluster node id, matching what the graph stamps into + * `WorkflowStepResult.leaseNodeId`. Only used to recognize review-gate leases left by a PREVIOUS + * process on this same node so they can be reclaimed immediately rather than waiting out the + * 15-minute staleness floor. Unset (single-node, tests) keeps floor-only semantics. + */ + localNodeId?: string; /** Optional callback to release TaskExecutor in-memory worktree ownership for a task. */ releaseExecutorWorktreeOwnership?: (taskId: string) => void; /** @@ -7161,9 +7169,28 @@ export class SelfHealingManager { re-attaches an in-review graph run after a restart, so those leases simply age out and are then marked failed as before. */ + /* + FNXC:PlanReviewLease 2026-07-26-20:26: + FN-8603 follow-up. The paragraph above accepts that "restart-orphaned gates simply age + out" — that acceptance is what cost FN-8603 ~14 minutes of dead wait after an engine + restart killed its Code Review session 34s in. Passing this node's identity lets + classifyReviewLease reclaim a lease THIS node's previous process took (proven dead: it + predates our boot) without touching the floor that protects peer-owned and legacy + unattributed leases. When localNodeId is unset the argument is undefined and behavior is + exactly as before. + */ + const localNodeLeaseIdentity = this.options.localNodeId + ? { nodeId: this.options.localNodeId, processBootAt: this.processBootStartedAt } + : undefined; const hasLiveReviewLease = (result: WorkflowStepResult): boolean => { if (!result.leaseOwner || !result.startedAt) return false; - return classifyReviewLease([result], result.workflowStepId, Date.now(), PLAN_REVIEW_LEASE_STALENESS_MS).kind === "adopt"; + return classifyReviewLease( + [result], + result.workflowStepId, + Date.now(), + PLAN_REVIEW_LEASE_STALENESS_MS, + localNodeLeaseIdentity, + ).kind === "adopt"; }; const { results, orphanedCount } = resolveOrphanedPendingStepResults( fresh.workflowStepResults, diff --git a/packages/engine/src/workflow-graph-executor.ts b/packages/engine/src/workflow-graph-executor.ts index 47ebc94872..d8c5dd6c81 100644 --- a/packages/engine/src/workflow-graph-executor.ts +++ b/packages/engine/src/workflow-graph-executor.ts @@ -127,6 +127,15 @@ export interface WorkflowNodePreparationRequirement { } export interface WorkflowGraphExecutorDeps { + /* + * FNXC:PlanReviewLease 2026-07-26-20:18: + * Cluster node id stamped onto review-gate leases (`WorkflowStepResult.leaseNodeId`). It lets a + * node later recognize a lease its OWN previous process left behind and reclaim it without + * waiting out the staleness floor; peer-owned leases stay protected by the floor. Optional — + * when unset the lease is written unattributed and keeps pure floor semantics, which is the + * pre-existing behavior. + */ + localNodeId?: string; handlers?: Partial>; /* * FNXC:WorkflowNodeRunners 2026-07-01-00:00: @@ -818,6 +827,14 @@ export class WorkflowGraphExecutor { // U3/KTD-4: stamp the lease owner so a concurrent/crashed re-entry // adopts this pending gate instead of dispatching a second reviewer. leaseOwner: runId, + /* + FNXC:PlanReviewLease 2026-07-26-20:20: + Stamp WHERE the lease runs, not just which run holds it. `runId` cannot distinguish + "this node's dead previous process" from "a peer node running right now", so without + this every restart-orphaned gate had to wait out the full staleness floor. Omitted when + the node id is unknown, which preserves the previous floor-only behavior. + */ + ...(this.deps.localNodeId ? { leaseNodeId: this.deps.localNodeId } : {}), }); this.deps.logTaskEntry?.(`${logPrefix} Starting workflow step: ${groupName}`);