Addresses the five blocking findings of review 644 on PR #972.
B1 — live ownership and status revalidation. resolve_missing_worktree_binding
now re-reads the authoritative lease, session, checkpoint, and issue-lock rows
immediately before mutating and diffs them against the snapshot the audit
recorded (binding path and identity, lease id/status/session/owner pid,
checkpoint path/status, live-session evidence, trustworthy ownership evidence,
issue-lock state). Any drift fails closed without mutation, and the binding is
reclassified from the live values rather than the audit snapshot. A candidate
carrying no audited snapshot is refused rather than trusted.
B2 — server-enforced cleanup authorization. Apply mode no longer accepts a
client-supplied operator_authorized boolean; it is rejected outright at the MCP
tool and in the module (#709 F1 / review 434). Authorization is now the
project's own reconciliation cleanup gate, required at both the task-capability
boundary (new reconciler-only reconcile_missing_worktree_bindings capability,
gitea.branch.delete, role-exclusive) and the production mutation boundary
(an authorized audit_reconciliation_mode cleanup phase, re-checked at the point
of mutation so a forged authorization mapping cannot stand in for the gate).
Dry-run remains available to any gitea.read profile and stays non-mutating.
Existing role, repository, parity, and provenance gates are unchanged.
B3 — expected-path compare-and-swap. retire_session_checkpoint_worktree_path
now requires expected_path and performs a guarded update keyed on the stored
path, refusing without mutation when the stored path was moved, replaced, or
concurrently changed, when the row is unknown, or when a selector matches more
than one checkpoint. retire_lease_worktree_path gains the same treatment plus
optional status/session/owner-pid compare-and-swap, and its UPDATE is keyed on
the audited path. Both report an idempotent already_retired outcome instead of
falsely reporting a retirement.
B4 — live-session and issue-lock evidence. session_active is now derived from
the control-plane sessions table instead of never being set, along two axes:
genuine liveness (recorded active, PID not dead, heartbeat fresh — the rule
reused from restart_coordinator) and weaker but still trustworthy recorded
ownership. A non-terminal lease now protects its binding regardless of whether
the recorded PID is alive, so dead-PID evidence alone can no longer retire a
lease the control plane still holds. The previously unused issue_lock_store is
now read: a live durable issue lock binding the path or branch blocks cleanup,
and locks whose own paths are missing are reported for release through their
own lifecycle rather than retired here.
B5 — adversarial regression coverage. The suite now drives the registered MCP
tools through mcp_server, the real ControlPlaneDB, and the real cleanup gate,
covering lease status/ownership/session/path drift, expected-path mismatch,
concurrent recreation, an unauthorized caller submitting operator_authorized,
wrong profile and missing capability, live-session and trustworthy-owner
evidence, conflicting issue locks, non-mutating dry-run, exact-binding-only
retirement, preservation of unrelated worktrees and git metadata, idempotent
re-execution, and worktrees-dimension resolution.
All original #970 acceptance criteria are preserved, including the distinctions
between deleted paths, moved paths, unavailable hosts or mounts, transient
filesystem failures, live ownership, and concurrent recreation.
Tests: focused #970 suite 47 passed. Adjacent suites (capability role
invariants, audit reconciliation mode, control plane DB, lease lifecycle,
reconciler cleanup integration, delete-branch capability, restart coordinator,
bootstrap lock contract) 252 passed / 93 subtests. Full suite 28 failed /
6000 passed, an exact match of the pre-change baseline's 28 failing test ids
at 3f584352 (28 failed / 5960 passed).
Closes #970
Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
1332 lines
51 KiB
Python
1332 lines
51 KiB
Python
"""Sanctioned reconciliation for worktree bindings whose paths are missing (#970).
|
|
|
|
When control-plane leases, session checkpoints, or issue locks hold a
|
|
``worktree_path`` pointing to a filesystem path that no longer exists on disk,
|
|
this module provides a fail-closed, auditable classification and resolution
|
|
workflow.
|
|
|
|
Core principles:
|
|
1. **Multi-dimensional Correlation**:
|
|
Correlates the binding with repository, host, branch, issue/PR, session,
|
|
and lease state.
|
|
2. **Safety-first Distinction**:
|
|
- Distinguishes a deleted worktree from a moved path, unavailable host/mount,
|
|
live lease, active session, or temporary filesystem problem.
|
|
- Host / mount failure (e.g. repo root or branches/ directory inaccessible)
|
|
blocks retirement (``unavailable_host_or_mount``).
|
|
- Moved worktree (found under another path in ``git worktree list``)
|
|
blocks retirement and reports the new location (``moved_worktree_path``).
|
|
- Live lease or active session/process blocks retirement (``live_lease_protected``,
|
|
``live_session_protected``).
|
|
3. **Atomic Pre-Mutation Re-Validation**:
|
|
Immediately before committing a mutation, re-evaluates filesystem existence,
|
|
git worktree list, host health, and lease status to prevent recreation races.
|
|
4. **Targeted Retirement**:
|
|
Retires only the exact confirmed stale binding (clearing ``worktree_path`` on
|
|
the lease/checkpoint/lock row and recording durable provenance) while
|
|
preserving unrelated worktrees, leases, branches, and git metadata.
|
|
5. **Dry-Run & Operator Authorization**:
|
|
Supports dry-run inspection and requires explicit cleanup authorization or
|
|
reconciler workflow for mutation.
|
|
6. **Idempotency & Durable Proof**:
|
|
Repeated execution on an already-retired binding is safe, idempotent, and
|
|
produces durable audit evidence.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import json
|
|
import os
|
|
import subprocess
|
|
from datetime import datetime, timezone
|
|
from typing import Any, Mapping, Sequence
|
|
|
|
import audit_reconciliation_mode
|
|
import control_plane_db as cpd
|
|
import issue_lock_store
|
|
import lease_lifecycle
|
|
import restart_coordinator
|
|
from merged_cleanup_reconcile import list_local_worktrees
|
|
|
|
# Single authority for "this session's heartbeat has gone stale". Imported
|
|
# rather than redefined so session liveness cannot drift between the restart
|
|
# coordinator and this workflow.
|
|
SESSION_HEARTBEAT_STALE_SECONDS = restart_coordinator.DEFAULT_SESSION_HEARTBEAT_STALE_SECONDS
|
|
|
|
# Classification vocabulary (#970 AC1-AC4).
|
|
CLASS_CONFIRMED_STALE_DELETED = "confirmed_stale_deleted_worktree"
|
|
CLASS_MOVED_WORKTREE = "moved_worktree_path"
|
|
CLASS_UNAVAILABLE_HOST_MOUNT = "unavailable_host_or_mount"
|
|
CLASS_LIVE_LEASE_PROTECTED = "live_lease_protected"
|
|
CLASS_LIVE_SESSION_PROTECTED = "live_session_protected"
|
|
CLASS_TRUSTED_OWNER_PROTECTED = "trusted_owner_evidence_protected"
|
|
CLASS_ISSUE_LOCK_PROTECTED = "conflicting_issue_lock_protected"
|
|
CLASS_PRESENT_VALID = "present_valid_worktree"
|
|
CLASS_ALREADY_RETIRED = "already_retired_binding"
|
|
|
|
ALL_CLASSIFICATIONS = frozenset({
|
|
CLASS_CONFIRMED_STALE_DELETED,
|
|
CLASS_MOVED_WORKTREE,
|
|
CLASS_UNAVAILABLE_HOST_MOUNT,
|
|
CLASS_LIVE_LEASE_PROTECTED,
|
|
CLASS_LIVE_SESSION_PROTECTED,
|
|
CLASS_TRUSTED_OWNER_PROTECTED,
|
|
CLASS_ISSUE_LOCK_PROTECTED,
|
|
CLASS_PRESENT_VALID,
|
|
CLASS_ALREADY_RETIRED,
|
|
})
|
|
|
|
RETIRE_ELIGIBLE_CLASSES = frozenset({
|
|
CLASS_CONFIRMED_STALE_DELETED,
|
|
})
|
|
|
|
# A lease stops protecting its binding only once it reaches a terminal status.
|
|
# Anything else — including an unrecognised status — keeps the binding
|
|
# protected, so a dead recorded PID can never on its own retire a lease the
|
|
# control plane still considers live (#970 review 644 B4).
|
|
TERMINAL_LEASE_STATUSES = frozenset({
|
|
lease_lifecycle.LEASE_STATUS_RELEASED,
|
|
lease_lifecycle.LEASE_STATUS_EXPIRED,
|
|
lease_lifecycle.LEASE_STATUS_ABANDONED,
|
|
})
|
|
|
|
TERMINAL_SESSION_STATUSES = frozenset({
|
|
"exited", "dead", "released", "terminated", "closed", "completed",
|
|
})
|
|
|
|
# Server-enforced cleanup authorization (#970 review 644 B2).
|
|
#
|
|
# Apply mode is authorized by the project's existing reconciliation cleanup
|
|
# gate — ``audit_reconciliation_mode.authorize_cleanup_phase``, minted through
|
|
# ``gitea_authorize_reconciliation_cleanup_phase``, which itself requires the
|
|
# reconciler-only ``gitea.branch.delete`` capability, safety proof, and a
|
|
# before/after snapshot. A caller-supplied ``operator_authorized`` boolean is
|
|
# self-assertable and is therefore never authorization evidence
|
|
# (#709 F1 / review 434).
|
|
CLEANUP_AUTHORIZATION_SOURCE = "audit_reconciliation_mode.authorize_cleanup_phase"
|
|
CLEANUP_REQUIRED_ROLE = "reconciler"
|
|
CLEANUP_REQUIRED_PERMISSION = "gitea.branch.delete"
|
|
CLEANUP_TASK = "reconcile_missing_worktree_bindings"
|
|
|
|
OPERATOR_AUTHORIZED_REJECTION = (
|
|
"operator_authorized is not accepted as authorization evidence "
|
|
"(#709 F1 / review 434); authorize cleanup server-side via "
|
|
"gitea_authorize_reconciliation_cleanup_phase"
|
|
)
|
|
|
|
|
|
def _norm_path(p: str | None) -> str:
|
|
if not p or not str(p).strip():
|
|
return ""
|
|
try:
|
|
return os.path.realpath(os.path.abspath(str(p).strip()))
|
|
except Exception:
|
|
return str(p).strip()
|
|
|
|
|
|
def check_host_mount_health(project_root: str) -> dict[str, Any]:
|
|
"""Check if the repository root and branches/ mount are healthy and accessible."""
|
|
root = _norm_path(project_root)
|
|
if not root or not os.path.isdir(root):
|
|
return {
|
|
"healthy": False,
|
|
"reason": f"project_root '{project_root}' does not exist or is not a directory",
|
|
}
|
|
|
|
# Verify project_root is inside a valid git repository
|
|
try:
|
|
res = subprocess.run(
|
|
["git", "-C", root, "rev-parse", "--git-dir"],
|
|
capture_output=True,
|
|
text=True,
|
|
check=False,
|
|
)
|
|
if res.returncode != 0:
|
|
return {
|
|
"healthy": False,
|
|
"reason": f"project_root '{root}' is not a valid git repository: {res.stderr.strip()}",
|
|
}
|
|
except Exception as exc:
|
|
return {
|
|
"healthy": False,
|
|
"reason": f"git execution failed at project_root '{root}': {exc}",
|
|
}
|
|
|
|
branches_dir = os.path.join(root, "branches")
|
|
if os.path.exists(branches_dir) and not os.path.isdir(branches_dir):
|
|
return {
|
|
"healthy": False,
|
|
"reason": f"branches path '{branches_dir}' exists but is not a directory",
|
|
}
|
|
|
|
return {"healthy": True, "project_root": root, "branches_dir": branches_dir}
|
|
|
|
|
|
def find_moved_worktree_path(
|
|
project_root: str,
|
|
recorded_path: str,
|
|
branch: str | None = None,
|
|
) -> str | None:
|
|
"""Check if a missing recorded worktree path is actually present at another location.
|
|
|
|
Returns the new path if found, or None if truly absent.
|
|
"""
|
|
norm_rec = _norm_path(recorded_path)
|
|
if not norm_rec:
|
|
return None
|
|
|
|
wt_list = []
|
|
try:
|
|
wt_list = list_local_worktrees(project_root)
|
|
except Exception:
|
|
wt_list = []
|
|
|
|
for entry in wt_list:
|
|
p = _norm_path(entry.get("path"))
|
|
b = entry.get("branch")
|
|
if not p:
|
|
continue
|
|
# If git worktree list shows a worktree matching the branch at a different path
|
|
if branch and b and b.strip() == branch.strip() and p != norm_rec:
|
|
if os.path.exists(p):
|
|
return p
|
|
# If the basename matches and path exists
|
|
if os.path.basename(p) == os.path.basename(norm_rec) and p != norm_rec:
|
|
if os.path.exists(p):
|
|
return p
|
|
|
|
# Check potential filesystem candidates under branches/
|
|
branches_dir = os.path.join(_norm_path(project_root), "branches")
|
|
if os.path.isdir(branches_dir):
|
|
candidates_to_check: list[str] = []
|
|
if branch:
|
|
candidates_to_check.append(os.path.join(branches_dir, branch.replace("/", "-")))
|
|
candidates_to_check.append(os.path.join(branches_dir, branch))
|
|
if recorded_path:
|
|
candidates_to_check.append(os.path.join(branches_dir, os.path.basename(recorded_path)))
|
|
|
|
for cand in candidates_to_check:
|
|
norm_cand = _norm_path(cand)
|
|
if norm_cand and norm_cand != norm_rec and os.path.isdir(norm_cand):
|
|
return norm_cand
|
|
|
|
return None
|
|
|
|
|
|
def _heartbeat_age_seconds(value: str | None, *, now: datetime | None = None) -> float | None:
|
|
"""Seconds since *value*, or None when it cannot be parsed."""
|
|
text = (value or "").strip()
|
|
if not text:
|
|
return None
|
|
if text.endswith("Z"):
|
|
text = text[:-1] + "+00:00"
|
|
try:
|
|
parsed = datetime.fromisoformat(text)
|
|
except ValueError:
|
|
return None
|
|
if parsed.tzinfo is None:
|
|
parsed = parsed.replace(tzinfo=timezone.utc)
|
|
moment = now or datetime.now(timezone.utc)
|
|
return (moment - parsed).total_seconds()
|
|
|
|
|
|
def collect_session_evidence(db: cpd.ControlPlaneDB) -> dict[str, dict[str, Any]]:
|
|
"""Return live-session evidence per session id (#970 review 644 B4).
|
|
|
|
``session_active`` was previously never derived, so every binding was
|
|
classified as if no session owned it. The control-plane ``sessions`` table
|
|
is the authoritative record of which process claimed a worktree, so it is
|
|
read here and reported along two distinct axes:
|
|
|
|
* ``session_live`` — the project's existing liveness rule, reused verbatim
|
|
from ``restart_coordinator._classify_session``: recorded ``active``, the
|
|
PID not known-dead, and the heartbeat not stale. A live session is real
|
|
ownership and blocks cleanup outright.
|
|
* ``session_recorded_active`` — the row still says ``active`` but its
|
|
heartbeat has gone stale. That is weaker, still-trustworthy ownership
|
|
evidence: it blocks cleanup unless a terminal lease contradicts it.
|
|
|
|
A merely-alive PID never establishes liveness on its own (#790 AC-N2): the
|
|
recorded PID is the long-lived MCP daemon, so it proves the daemon is up and
|
|
nothing about the task. It is carried as corroboration only.
|
|
"""
|
|
evidence: dict[str, dict[str, Any]] = {}
|
|
try:
|
|
rows = db.list_sessions(limit=500)
|
|
except Exception: # noqa: BLE001
|
|
return evidence
|
|
|
|
now = datetime.now(timezone.utc)
|
|
for row in rows:
|
|
sid = str(row.get("session_id") or "").strip()
|
|
if not sid:
|
|
continue
|
|
status = str(row.get("status") or "").strip().lower()
|
|
pid = row.get("pid")
|
|
pid_alive = lease_lifecycle.is_process_alive(pid) if pid else None
|
|
recorded_active = bool(status) and status not in TERMINAL_SESSION_STATUSES
|
|
age = _heartbeat_age_seconds(row.get("last_heartbeat_at"), now=now)
|
|
heartbeat_stale = bool(
|
|
age is not None and age > SESSION_HEARTBEAT_STALE_SECONDS
|
|
)
|
|
session_live = bool(
|
|
recorded_active and pid_alive is not False and not heartbeat_stale
|
|
)
|
|
evidence[sid] = {
|
|
"session_id": sid,
|
|
"session_status": status,
|
|
"session_pid": pid,
|
|
"session_pid_alive": bool(pid_alive),
|
|
"session_live": session_live,
|
|
"session_recorded_active": recorded_active,
|
|
"session_active": session_live or recorded_active,
|
|
"heartbeat_stale": heartbeat_stale,
|
|
"heartbeat_age_seconds": age,
|
|
"last_heartbeat_at": row.get("last_heartbeat_at"),
|
|
"role": row.get("role"),
|
|
"profile": row.get("profile"),
|
|
}
|
|
return evidence
|
|
|
|
|
|
def collect_issue_lock_bindings(lock_dir: str | None = None) -> list[dict[str, Any]]:
|
|
"""Return durable author issue-lock bindings (#970 review 644 B4).
|
|
|
|
``issue_lock_store`` was imported but never used, so a worktree still owned
|
|
by a live author issue lock could be classified stale and retired. Every
|
|
durable lock file is read here and reported with its recorded path, branch,
|
|
owner pid, and lease liveness, so a conflicting lock can block cleanup.
|
|
"""
|
|
bindings: list[dict[str, Any]] = []
|
|
try:
|
|
paths = issue_lock_store.iter_lock_files(lock_dir)
|
|
except Exception: # noqa: BLE001
|
|
return bindings
|
|
|
|
for lock_path in paths:
|
|
try:
|
|
lock = issue_lock_store.read_lock_file(lock_path)
|
|
except Exception: # noqa: BLE001
|
|
continue
|
|
if not isinstance(lock, dict):
|
|
continue
|
|
wt = str(lock.get("worktree_path") or "").strip()
|
|
if not wt:
|
|
continue
|
|
pid = lock.get("session_pid") or lock.get("pid")
|
|
try:
|
|
pid_alive = issue_lock_store.is_process_alive(pid)
|
|
except Exception: # noqa: BLE001
|
|
pid_alive = False
|
|
try:
|
|
lease_live = issue_lock_store.is_lease_live(lock)
|
|
except Exception: # noqa: BLE001
|
|
lease_live = False
|
|
bindings.append({
|
|
"source": "issue_lock",
|
|
"lock_file_path": lock_path,
|
|
"issue_number": lock.get("issue_number"),
|
|
"branch": lock.get("branch_name") or lock.get("branch"),
|
|
"worktree_path": wt,
|
|
"normalized_path": _norm_path(wt),
|
|
"owner_pid": pid,
|
|
"owner_pid_alive": pid_alive,
|
|
"lease_live": lease_live,
|
|
"claimant": (lock.get("claimant") or {}).get("username")
|
|
if isinstance(lock.get("claimant"), Mapping)
|
|
else lock.get("claimant"),
|
|
# A lock is trustworthy ownership evidence while its lease is live
|
|
# or its recorded process is still running.
|
|
"conflicting": bool(lease_live or pid_alive),
|
|
})
|
|
return bindings
|
|
|
|
|
|
def find_conflicting_issue_lock(
|
|
recorded_path: str,
|
|
issue_locks: Sequence[Mapping[str, Any]] | None,
|
|
*,
|
|
branch: str | None = None,
|
|
) -> dict[str, Any] | None:
|
|
"""Return a live issue lock that still binds *recorded_path* (or *branch*)."""
|
|
if not issue_locks:
|
|
return None
|
|
norm_p = _norm_path(recorded_path)
|
|
for lock in issue_locks:
|
|
if not lock.get("conflicting"):
|
|
continue
|
|
if norm_p and lock.get("normalized_path") == norm_p:
|
|
return dict(lock)
|
|
lock_branch = str(lock.get("branch") or "").strip()
|
|
if branch and lock_branch and lock_branch == str(branch).strip():
|
|
return dict(lock)
|
|
return None
|
|
|
|
|
|
def classify_worktree_binding(
|
|
*,
|
|
recorded_path: str,
|
|
lease_status: str | None = None,
|
|
lease_phase: str | None = None,
|
|
owner_pid: int | None = None,
|
|
session_id: str | None = None,
|
|
session_active: bool = False,
|
|
session_live: bool = False,
|
|
session_pid_alive: bool = False,
|
|
session_status: str | None = None,
|
|
issue_lock: Mapping[str, Any] | None = None,
|
|
branch: str | None = None,
|
|
project_root: str,
|
|
host_health: dict[str, Any] | None = None,
|
|
path_exists_override: bool | None = None,
|
|
) -> dict[str, Any]:
|
|
"""Classify a single worktree binding record for missing-path safety.
|
|
|
|
Fail-closed: returns detailed reasons and classification.
|
|
"""
|
|
norm_p = _norm_path(recorded_path)
|
|
reasons: list[str] = []
|
|
|
|
if not norm_p:
|
|
return {
|
|
"classification": CLASS_ALREADY_RETIRED,
|
|
"retire_eligible": False,
|
|
"recorded_path": "",
|
|
"reasons": ["worktree_path is empty (already retired/unbound)"],
|
|
}
|
|
|
|
# Check 1: Host / mount health
|
|
hh = host_health or check_host_mount_health(project_root)
|
|
if not hh.get("healthy"):
|
|
reasons.append(
|
|
f"host or filesystem mount unavailable: {hh.get('reason')}; "
|
|
"refusing to classify missing path as deleted (fail closed)"
|
|
)
|
|
return {
|
|
"classification": CLASS_UNAVAILABLE_HOST_MOUNT,
|
|
"retire_eligible": False,
|
|
"recorded_path": norm_p,
|
|
"reasons": reasons,
|
|
"host_health": hh,
|
|
}
|
|
|
|
# Check 2: Path existence
|
|
exists = (
|
|
path_exists_override
|
|
if path_exists_override is not None
|
|
else os.path.exists(norm_p)
|
|
)
|
|
|
|
if exists:
|
|
return {
|
|
"classification": CLASS_PRESENT_VALID,
|
|
"retire_eligible": False,
|
|
"recorded_path": norm_p,
|
|
"reasons": [f"worktree path '{norm_p}' exists on disk"],
|
|
}
|
|
|
|
# Path is missing on disk. Run safety checks before declaring confirmed stale.
|
|
|
|
# Check 3: Moved path
|
|
moved_path = find_moved_worktree_path(project_root, norm_p, branch=branch)
|
|
if moved_path:
|
|
reasons.append(
|
|
f"recorded path '{norm_p}' is missing, but worktree for branch '{branch}' "
|
|
f"was found moved to '{moved_path}'"
|
|
)
|
|
return {
|
|
"classification": CLASS_MOVED_WORKTREE,
|
|
"retire_eligible": False,
|
|
"recorded_path": norm_p,
|
|
"moved_to_path": moved_path,
|
|
"reasons": reasons,
|
|
}
|
|
|
|
# Check 4: Conflicting durable issue lock.
|
|
#
|
|
# Checked before the lease, because an author issue lock can still own a
|
|
# worktree after its control-plane lease row has gone terminal.
|
|
if issue_lock:
|
|
reasons.append(
|
|
f"durable issue lock for issue #{issue_lock.get('issue_number')} "
|
|
f"(branch '{issue_lock.get('branch')}', owner pid "
|
|
f"{issue_lock.get('owner_pid')}, lease_live="
|
|
f"{issue_lock.get('lease_live')}, pid_alive="
|
|
f"{issue_lock.get('owner_pid_alive')}) still binds this worktree; "
|
|
"a conflicting issue lock blocks cleanup"
|
|
)
|
|
return {
|
|
"classification": CLASS_ISSUE_LOCK_PROTECTED,
|
|
"retire_eligible": False,
|
|
"recorded_path": norm_p,
|
|
"issue_lock": dict(issue_lock),
|
|
"reasons": reasons,
|
|
}
|
|
|
|
# Check 5: Lease protection.
|
|
#
|
|
# A non-terminal lease protects its binding regardless of whether the
|
|
# recorded PID is still alive: the recorded PID is the long-lived MCP
|
|
# daemon in normal operation, so its absence is weak evidence, while the
|
|
# lease status is the control plane's own record of ownership. Dead-PID
|
|
# evidence alone must never retire a lease the control plane still holds
|
|
# (#970 review 644 B4).
|
|
st = (lease_status or "").strip().lower()
|
|
pid_alive = lease_lifecycle.is_process_alive(owner_pid) if owner_pid else False
|
|
if st and st not in TERMINAL_LEASE_STATUSES:
|
|
reasons.append(
|
|
f"lease status '{st}' is not terminal (owner_pid={owner_pid}, "
|
|
f"alive={pid_alive}); a non-terminal lease protects its worktree "
|
|
"binding"
|
|
)
|
|
return {
|
|
"classification": CLASS_LIVE_LEASE_PROTECTED,
|
|
"retire_eligible": False,
|
|
"recorded_path": norm_p,
|
|
"lease_status": st,
|
|
"owner_pid_alive": pid_alive,
|
|
"reasons": reasons,
|
|
}
|
|
|
|
# Check 6: Live session protection.
|
|
#
|
|
# "Live" is the project's existing rule (recorded active, PID not dead,
|
|
# heartbeat fresh) — not bare PID liveness, which is only the daemon.
|
|
if session_live:
|
|
reasons.append(
|
|
f"session '{session_id}' is live (session_status='{session_status}', "
|
|
f"heartbeat fresh, session_pid_alive={session_pid_alive}); a live "
|
|
"session protects its worktree binding"
|
|
)
|
|
return {
|
|
"classification": CLASS_LIVE_SESSION_PROTECTED,
|
|
"retire_eligible": False,
|
|
"recorded_path": norm_p,
|
|
"session_status": session_status,
|
|
"reasons": reasons,
|
|
}
|
|
|
|
# Check 7: Trustworthy ownership evidence without a live process.
|
|
#
|
|
# The owning process may be gone while the control plane still records the
|
|
# session as active. That is a valid ownership signal and must block
|
|
# cleanup, even though nothing is running.
|
|
#
|
|
# One signal outranks it: an explicitly terminal lease. Session rows are
|
|
# opened when work is assigned and are not closed on release, so a
|
|
# recorded-active session routinely outlives the work it owned. A lease
|
|
# that was explicitly released, expired, or abandoned is the later and more
|
|
# specific statement about *this* binding, so it wins. Without such a lease
|
|
# — a checkpoint binding, or a lease row that is simply gone — the session
|
|
# record is the best ownership evidence available and blocks cleanup.
|
|
lease_is_terminal = bool(st) and st in TERMINAL_LEASE_STATUSES
|
|
if session_active and not lease_is_terminal:
|
|
reasons.append(
|
|
f"session '{session_id}' is still recorded active "
|
|
f"(session_status='{session_status}') even though its process is "
|
|
"not running; trustworthy ownership evidence blocks cleanup"
|
|
)
|
|
return {
|
|
"classification": CLASS_TRUSTED_OWNER_PROTECTED,
|
|
"retire_eligible": False,
|
|
"recorded_path": norm_p,
|
|
"session_status": session_status,
|
|
"reasons": reasons,
|
|
}
|
|
|
|
# Confirmed stale deleted
|
|
reasons.append(
|
|
f"worktree path '{norm_p}' does not exist on disk, host is healthy, "
|
|
"no moved path was found, the lease status is terminal "
|
|
f"('{st or 'none recorded'}'), no live or recorded-active session owns "
|
|
"it, and no durable issue lock binds it"
|
|
)
|
|
return {
|
|
"classification": CLASS_CONFIRMED_STALE_DELETED,
|
|
"retire_eligible": True,
|
|
"recorded_path": norm_p,
|
|
"lease_status": st,
|
|
"owner_pid_alive": pid_alive,
|
|
"reasons": reasons,
|
|
}
|
|
|
|
|
|
def audit_missing_worktree_bindings(
|
|
db: cpd.ControlPlaneDB | None = None,
|
|
*,
|
|
project_root: str | None = None,
|
|
remote: str = "prgs",
|
|
org: str | None = None,
|
|
repo: str | None = None,
|
|
host: str | None = None,
|
|
) -> dict[str, Any]:
|
|
"""Audit all recorded worktree bindings and classify missing-path entries.
|
|
|
|
Read-only / side-effect free.
|
|
"""
|
|
root = _norm_path(project_root or os.getcwd())
|
|
hh = check_host_mount_health(root)
|
|
cp_db = db or cpd.ControlPlaneDB()
|
|
|
|
# Ownership evidence the classifier needs. Populated once per audit and
|
|
# re-read live before any mutation (#970 review 644 B1/B4).
|
|
session_evidence = collect_session_evidence(cp_db)
|
|
issue_locks = collect_issue_lock_bindings()
|
|
|
|
bindings: list[dict[str, Any]] = []
|
|
|
|
# 1. Audit control-plane leases
|
|
try:
|
|
leases_res = lease_lifecycle.list_active_leases(
|
|
cp_db,
|
|
remote=remote,
|
|
org=org,
|
|
repo=repo,
|
|
include_non_active=True,
|
|
limit=500,
|
|
)
|
|
for L in leases_res.get("leases") or []:
|
|
wt = (L.get("worktree_path") or "").strip()
|
|
if not wt:
|
|
continue
|
|
owner_pid = L.get("owner_pid") or L.get("session_pid")
|
|
sess_id = L.get("session_id")
|
|
br = L.get("branch")
|
|
if not br and isinstance(L.get("provenance"), dict):
|
|
br = L["provenance"].get("branch")
|
|
|
|
sess = session_evidence.get(str(sess_id or "").strip(), {})
|
|
conflicting_lock = find_conflicting_issue_lock(wt, issue_locks, branch=br)
|
|
|
|
cls_info = classify_worktree_binding(
|
|
recorded_path=wt,
|
|
lease_status=L.get("status"),
|
|
lease_phase=L.get("phase"),
|
|
owner_pid=owner_pid,
|
|
session_id=sess_id,
|
|
session_active=bool(sess.get("session_active")),
|
|
session_live=bool(sess.get("session_live")),
|
|
session_pid_alive=bool(sess.get("session_pid_alive")),
|
|
session_status=sess.get("session_status"),
|
|
issue_lock=conflicting_lock,
|
|
branch=br,
|
|
project_root=root,
|
|
host_health=hh,
|
|
)
|
|
|
|
bindings.append({
|
|
"source": "lease",
|
|
"lease_id": L.get("lease_id"),
|
|
"session_id": sess_id,
|
|
"work_kind": L.get("work_kind"),
|
|
"work_number": L.get("work_number"),
|
|
"branch": br,
|
|
"worktree_path": wt,
|
|
"lease_status": L.get("status"),
|
|
"lease_phase": L.get("phase"),
|
|
"owner_pid": owner_pid,
|
|
"session_status": sess.get("session_status"),
|
|
"session_active": bool(sess.get("session_active")),
|
|
"session_pid_alive": bool(sess.get("session_pid_alive")),
|
|
"issue_lock_conflict": conflicting_lock,
|
|
"remote": L.get("remote") or remote,
|
|
"org": L.get("org") or org,
|
|
"repo": L.get("repo") or repo,
|
|
"host": L.get("host") or host,
|
|
**cls_info,
|
|
"audited_state": _safety_snapshot(
|
|
source="lease",
|
|
worktree_path=wt,
|
|
lease_id=L.get("lease_id"),
|
|
lease_status=L.get("status"),
|
|
lease_session_id=sess_id,
|
|
lease_owner_pid=owner_pid,
|
|
session_evidence=sess,
|
|
issue_lock=conflicting_lock,
|
|
),
|
|
})
|
|
except Exception as exc: # noqa: BLE001
|
|
bindings.append({
|
|
"source": "lease_query_error",
|
|
"error": str(exc),
|
|
})
|
|
|
|
# 2. Audit session checkpoints
|
|
try:
|
|
checkpoints = cp_db.list_session_checkpoints(limit=500)
|
|
for cp in checkpoints:
|
|
wt = (cp.get("worktree_path") or "").strip()
|
|
if not wt:
|
|
continue
|
|
sess_id = cp.get("session_id")
|
|
sess = session_evidence.get(str(sess_id or "").strip(), {})
|
|
conflicting_lock = find_conflicting_issue_lock(
|
|
wt, issue_locks, branch=cp.get("branch")
|
|
)
|
|
cls_info = classify_worktree_binding(
|
|
recorded_path=wt,
|
|
session_id=sess_id,
|
|
session_active=bool(sess.get("session_active")),
|
|
session_live=bool(sess.get("session_live")),
|
|
session_pid_alive=bool(sess.get("session_pid_alive")),
|
|
session_status=sess.get("session_status"),
|
|
issue_lock=conflicting_lock,
|
|
branch=cp.get("branch"),
|
|
project_root=root,
|
|
host_health=hh,
|
|
)
|
|
bindings.append({
|
|
"source": "session_checkpoint",
|
|
"checkpoint_id": cp.get("checkpoint_id"),
|
|
"session_id": sess_id,
|
|
"work_kind": cp.get("work_kind"),
|
|
"work_number": cp.get("work_number"),
|
|
"branch": cp.get("branch"),
|
|
"worktree_path": wt,
|
|
"status": cp.get("status"),
|
|
"checkpoint_status": cp.get("status"),
|
|
"session_status": sess.get("session_status"),
|
|
"session_active": bool(sess.get("session_active")),
|
|
"session_pid_alive": bool(sess.get("session_pid_alive")),
|
|
"issue_lock_conflict": conflicting_lock,
|
|
"remote": cp.get("remote") or remote,
|
|
"org": cp.get("org") or org,
|
|
"repo": cp.get("repo") or repo,
|
|
**cls_info,
|
|
"audited_state": _safety_snapshot(
|
|
source="session_checkpoint",
|
|
worktree_path=wt,
|
|
checkpoint_id=cp.get("checkpoint_id"),
|
|
checkpoint_status=cp.get("status"),
|
|
lease_session_id=sess_id,
|
|
session_evidence=sess,
|
|
issue_lock=conflicting_lock,
|
|
),
|
|
})
|
|
except Exception: # noqa: BLE001
|
|
pass
|
|
|
|
# 3. Report durable issue-lock bindings whose paths are missing.
|
|
#
|
|
# These are reported, never retired: a durable author lock is released
|
|
# through its own sanctioned lifecycle. They are kept out of
|
|
# ``missing_bindings`` for exactly that reason — counting a binding this
|
|
# workflow must not touch would leave the worktrees dimension permanently
|
|
# unresolvable — but they are still surfaced so a missing path a lock
|
|
# claims is visible rather than silently ignored (#970 AC1).
|
|
issue_lock_missing: list[dict[str, Any]] = []
|
|
for lock in issue_locks:
|
|
wt = lock.get("worktree_path") or ""
|
|
if not wt or os.path.exists(_norm_path(wt)):
|
|
continue
|
|
issue_lock_missing.append({
|
|
"source": "issue_lock",
|
|
"lock_file_path": lock.get("lock_file_path"),
|
|
"issue_number": lock.get("issue_number"),
|
|
"branch": lock.get("branch"),
|
|
"worktree_path": wt,
|
|
"recorded_path": _norm_path(wt),
|
|
"owner_pid": lock.get("owner_pid"),
|
|
"owner_pid_alive": lock.get("owner_pid_alive"),
|
|
"lease_live": lock.get("lease_live"),
|
|
"claimant": lock.get("claimant"),
|
|
"conflicting": lock.get("conflicting"),
|
|
"classification": CLASS_ISSUE_LOCK_PROTECTED
|
|
if lock.get("conflicting")
|
|
else CLASS_ALREADY_RETIRED,
|
|
"retire_eligible": False,
|
|
"reasons": [
|
|
"durable issue lock records a worktree path that is missing on "
|
|
"disk; release it through the sanctioned issue-lock lifecycle, "
|
|
"never by retiring the binding here"
|
|
],
|
|
})
|
|
|
|
missing_bindings = [
|
|
b for b in bindings if b.get("classification") != CLASS_PRESENT_VALID
|
|
]
|
|
confirmed_stale = [
|
|
b for b in missing_bindings if b.get("classification") == CLASS_CONFIRMED_STALE_DELETED
|
|
]
|
|
moved = [
|
|
b for b in missing_bindings if b.get("classification") == CLASS_MOVED_WORKTREE
|
|
]
|
|
live_protected = [
|
|
b for b in missing_bindings
|
|
if b.get("classification") in (
|
|
CLASS_LIVE_LEASE_PROTECTED,
|
|
CLASS_LIVE_SESSION_PROTECTED,
|
|
CLASS_TRUSTED_OWNER_PROTECTED,
|
|
)
|
|
]
|
|
issue_lock_protected = [
|
|
b for b in missing_bindings if b.get("classification") == CLASS_ISSUE_LOCK_PROTECTED
|
|
]
|
|
unavailable_host = [
|
|
b for b in missing_bindings if b.get("classification") == CLASS_UNAVAILABLE_HOST_MOUNT
|
|
]
|
|
|
|
return {
|
|
"project_root": root,
|
|
"host_mount_healthy": hh.get("healthy", False),
|
|
"host_health_reason": hh.get("reason"),
|
|
"total_bindings_audited": len(bindings),
|
|
"present_count": len(bindings) - len(missing_bindings),
|
|
"missing_count": len(missing_bindings),
|
|
"confirmed_stale_count": len(confirmed_stale),
|
|
"moved_count": len(moved),
|
|
"live_protected_count": len(live_protected),
|
|
"unavailable_host_count": len(unavailable_host),
|
|
"issue_lock_protected_count": len(issue_lock_protected),
|
|
"missing_bindings": missing_bindings,
|
|
"confirmed_stale_candidates": confirmed_stale,
|
|
"issue_lock_missing_bindings": issue_lock_missing,
|
|
"issue_lock_missing_count": len(issue_lock_missing),
|
|
"live_session_evidence_count": sum(
|
|
1 for e in session_evidence.values() if e.get("session_active")
|
|
),
|
|
"worktrees_dimension_resolved": len(missing_bindings) == 0,
|
|
}
|
|
|
|
|
|
def _safety_snapshot(
|
|
*,
|
|
source: str,
|
|
worktree_path: str | None,
|
|
lease_id: str | None = None,
|
|
lease_status: str | None = None,
|
|
lease_session_id: str | None = None,
|
|
lease_owner_pid: Any = None,
|
|
checkpoint_id: str | None = None,
|
|
checkpoint_status: str | None = None,
|
|
session_evidence: Mapping[str, Any] | None = None,
|
|
issue_lock: Mapping[str, Any] | None = None,
|
|
) -> dict[str, Any]:
|
|
"""Normalize the safety-relevant fields cleanup must re-verify (#970 B1).
|
|
|
|
Every value cleanup relies on is captured here in one comparable shape, so
|
|
the apply path can diff a fresh authoritative read against what the audit
|
|
saw instead of trusting the audit's own snapshot.
|
|
"""
|
|
sess = dict(session_evidence or {})
|
|
lock = dict(issue_lock or {})
|
|
return {
|
|
"source": source,
|
|
"binding_path": _norm_path(worktree_path),
|
|
"binding_identity": f"{source}:{lease_id or checkpoint_id or lease_session_id or ''}",
|
|
"lease_id": lease_id,
|
|
"lease_status": (str(lease_status).strip().lower() if lease_status else None),
|
|
"lease_session_id": lease_session_id,
|
|
"lease_owner_pid": (int(lease_owner_pid) if isinstance(lease_owner_pid, int) else lease_owner_pid),
|
|
"checkpoint_id": checkpoint_id,
|
|
"checkpoint_status": (str(checkpoint_status).strip().lower() if checkpoint_status else None),
|
|
"checkpoint_path": _norm_path(worktree_path) if source == "session_checkpoint" else None,
|
|
"session_status": sess.get("session_status"),
|
|
"session_active": bool(sess.get("session_active")),
|
|
"session_live": bool(sess.get("session_live")),
|
|
"session_pid_alive": bool(sess.get("session_pid_alive")),
|
|
"issue_lock_conflicting": bool(lock.get("conflicting")),
|
|
"issue_lock_file": lock.get("lock_file_path"),
|
|
"issue_lock_issue_number": lock.get("issue_number"),
|
|
}
|
|
|
|
|
|
# Fields whose drift between audit and apply must fail closed without mutation.
|
|
SAFETY_SNAPSHOT_FIELDS = (
|
|
"binding_path",
|
|
"binding_identity",
|
|
"lease_id",
|
|
"lease_status",
|
|
"lease_session_id",
|
|
"lease_owner_pid",
|
|
"checkpoint_id",
|
|
"checkpoint_status",
|
|
"checkpoint_path",
|
|
"session_status",
|
|
"session_active",
|
|
"session_live",
|
|
"session_pid_alive",
|
|
"issue_lock_conflicting",
|
|
"issue_lock_file",
|
|
"issue_lock_issue_number",
|
|
)
|
|
|
|
|
|
def read_live_safety_state(
|
|
db: cpd.ControlPlaneDB,
|
|
binding: Mapping[str, Any],
|
|
) -> dict[str, Any]:
|
|
"""Re-read the authoritative safety state for one binding (#970 B1).
|
|
|
|
Reads the lease row, session row, checkpoint row, and durable issue locks
|
|
from their authoritative stores — never from the audit snapshot — so
|
|
cleanup decides against current state.
|
|
"""
|
|
source = binding.get("source") or "lease"
|
|
lease_id = binding.get("lease_id")
|
|
checkpoint_id = binding.get("checkpoint_id")
|
|
session_id = binding.get("session_id")
|
|
|
|
lease_row: dict[str, Any] = {}
|
|
if lease_id:
|
|
try:
|
|
state = db.get_lease_workflow_state(str(lease_id))
|
|
except Exception: # noqa: BLE001
|
|
state = None
|
|
if isinstance(state, Mapping):
|
|
lease_row = dict(state.get("lease") or {})
|
|
|
|
checkpoint_row: dict[str, Any] = {}
|
|
if source == "session_checkpoint":
|
|
try:
|
|
rows = db.list_session_checkpoints(
|
|
session_id=str(session_id) if session_id else None,
|
|
limit=500,
|
|
)
|
|
except Exception: # noqa: BLE001
|
|
rows = []
|
|
for row in rows or []:
|
|
if checkpoint_id and row.get("checkpoint_id") != checkpoint_id:
|
|
continue
|
|
checkpoint_row = dict(row)
|
|
break
|
|
|
|
session_evidence = collect_session_evidence(db)
|
|
live_session = session_evidence.get(str(session_id or "").strip(), {})
|
|
|
|
issue_locks = collect_issue_lock_bindings()
|
|
live_path = (
|
|
lease_row.get("worktree_path")
|
|
if source == "lease"
|
|
else checkpoint_row.get("worktree_path")
|
|
)
|
|
if live_path is None:
|
|
live_path = binding.get("worktree_path")
|
|
conflicting_lock = find_conflicting_issue_lock(
|
|
str(live_path or ""), issue_locks, branch=binding.get("branch")
|
|
)
|
|
|
|
snapshot = _safety_snapshot(
|
|
source=source,
|
|
worktree_path=str(live_path or ""),
|
|
lease_id=lease_id,
|
|
lease_status=lease_row.get("status") if lease_row else binding.get("lease_status"),
|
|
lease_session_id=lease_row.get("session_id") if lease_row else session_id,
|
|
lease_owner_pid=lease_row.get("owner_pid") if lease_row else binding.get("owner_pid"),
|
|
checkpoint_id=checkpoint_row.get("checkpoint_id") if checkpoint_row else checkpoint_id,
|
|
checkpoint_status=checkpoint_row.get("status") if checkpoint_row else binding.get("checkpoint_status"),
|
|
session_evidence=live_session,
|
|
issue_lock=conflicting_lock,
|
|
)
|
|
return {
|
|
"snapshot": snapshot,
|
|
"lease_row": lease_row,
|
|
"checkpoint_row": checkpoint_row,
|
|
"session_evidence": live_session,
|
|
"issue_lock": conflicting_lock,
|
|
"live_worktree_path": str(live_path or ""),
|
|
}
|
|
|
|
|
|
def compare_safety_state(
|
|
audited: Mapping[str, Any] | None,
|
|
live: Mapping[str, Any],
|
|
) -> dict[str, Any]:
|
|
"""Diff an audited safety snapshot against a fresh one (#970 B1)."""
|
|
if not audited:
|
|
return {
|
|
"comparable": False,
|
|
"changed": True,
|
|
"changed_fields": ["audited_state"],
|
|
"reasons": [
|
|
"binding carries no audited safety snapshot; cleanup cannot "
|
|
"prove nothing changed since the audit (fail closed)"
|
|
],
|
|
"differences": {},
|
|
}
|
|
|
|
differences: dict[str, Any] = {}
|
|
for field in SAFETY_SNAPSHOT_FIELDS:
|
|
before = audited.get(field)
|
|
after = live.get(field)
|
|
if before != after:
|
|
differences[field] = {"audited": before, "live": after}
|
|
|
|
reasons = [
|
|
f"safety-relevant value '{field}' changed between audit and apply "
|
|
f"(audited={diff['audited']!r}, live={diff['live']!r})"
|
|
for field, diff in differences.items()
|
|
]
|
|
return {
|
|
"comparable": True,
|
|
"changed": bool(differences),
|
|
"changed_fields": sorted(differences),
|
|
"differences": differences,
|
|
"reasons": reasons,
|
|
}
|
|
|
|
|
|
def assess_cleanup_authorization(
|
|
*,
|
|
role: str | None = None,
|
|
capability_blockers: Sequence[str] | None = None,
|
|
operator_authorized: bool = False,
|
|
task_capability_resolved: bool = False,
|
|
) -> dict[str, Any]:
|
|
"""Server-enforced authorization for apply-mode cleanup (#970 B2).
|
|
|
|
Authorization is never derived from a caller-supplied boolean. It requires,
|
|
together:
|
|
|
|
* the reconciler role — the only role the project grants cleanup authority;
|
|
* proven ``gitea.branch.delete`` capability, evaluated server-side by the
|
|
profile gate and passed in as its blocker list;
|
|
* a resolved cleanup task capability (task-capability resolution boundary);
|
|
* an active, authorized cleanup phase minted through
|
|
``gitea_authorize_reconciliation_cleanup_phase`` (mutation boundary).
|
|
|
|
An unauthorized or wrongly-profiled caller is rejected even when it submits
|
|
``operator_authorized=True``.
|
|
"""
|
|
reasons: list[str] = []
|
|
if operator_authorized:
|
|
reasons.append(OPERATOR_AUTHORIZED_REJECTION)
|
|
|
|
normalized_role = (role or "").strip().lower()
|
|
if normalized_role != CLEANUP_REQUIRED_ROLE:
|
|
reasons.append(
|
|
f"cleanup requires role '{CLEANUP_REQUIRED_ROLE}'; active role is "
|
|
f"'{normalized_role or 'unknown'}' (fail closed)"
|
|
)
|
|
|
|
blockers = list(capability_blockers or [])
|
|
if blockers:
|
|
reasons.append(
|
|
f"active profile lacks {CLEANUP_REQUIRED_PERMISSION}: "
|
|
+ "; ".join(str(b) for b in blockers)
|
|
)
|
|
|
|
if not task_capability_resolved:
|
|
reasons.append(
|
|
f"task capability '{CLEANUP_TASK}' is not resolved for this session; "
|
|
"resolve it before apply-mode cleanup (fail closed)"
|
|
)
|
|
|
|
phase_allowed, phase_reasons = audit_reconciliation_mode.check_cleanup_execution_allowed()
|
|
if not phase_allowed:
|
|
reasons.extend(phase_reasons)
|
|
|
|
return {
|
|
"authorized": not reasons,
|
|
"source": CLEANUP_AUTHORIZATION_SOURCE,
|
|
"required_role": CLEANUP_REQUIRED_ROLE,
|
|
"required_permission": CLEANUP_REQUIRED_PERMISSION,
|
|
"required_task": CLEANUP_TASK,
|
|
"role": normalized_role or None,
|
|
"phase": audit_reconciliation_mode.current_phase(),
|
|
"operator_authorized_rejected": bool(operator_authorized),
|
|
"reasons": reasons,
|
|
}
|
|
|
|
|
|
def _authorization_valid_for_apply(
|
|
cleanup_authorization: Mapping[str, Any] | None,
|
|
) -> tuple[bool, list[str]]:
|
|
"""Validate a minted authorization at the mutation boundary (#970 B2).
|
|
|
|
The authorization artifact alone is not trusted: the process-local cleanup
|
|
phase is re-checked here, so a forged mapping cannot authorize a mutation
|
|
that the server never approved.
|
|
"""
|
|
reasons: list[str] = []
|
|
if not isinstance(cleanup_authorization, Mapping):
|
|
reasons.append(
|
|
"apply mode requires a server-minted cleanup authorization from "
|
|
f"{CLEANUP_AUTHORIZATION_SOURCE} (fail closed)"
|
|
)
|
|
return False, reasons
|
|
if not cleanup_authorization.get("authorized"):
|
|
reasons.extend(
|
|
list(cleanup_authorization.get("reasons") or [])
|
|
or ["cleanup authorization is not authorized (fail closed)"]
|
|
)
|
|
if cleanup_authorization.get("source") != CLEANUP_AUTHORIZATION_SOURCE:
|
|
reasons.append(
|
|
"cleanup authorization did not come from "
|
|
f"{CLEANUP_AUTHORIZATION_SOURCE} (fail closed)"
|
|
)
|
|
if (cleanup_authorization.get("role") or "").strip().lower() != CLEANUP_REQUIRED_ROLE:
|
|
reasons.append(
|
|
f"cleanup authorization role must be '{CLEANUP_REQUIRED_ROLE}' "
|
|
"(fail closed)"
|
|
)
|
|
# Re-check the server-side phase at the mutation boundary itself.
|
|
phase_allowed, phase_reasons = audit_reconciliation_mode.check_cleanup_execution_allowed()
|
|
if not phase_allowed:
|
|
reasons.extend(phase_reasons)
|
|
return (not reasons), reasons
|
|
|
|
|
|
def resolve_missing_worktree_binding(
|
|
db: cpd.ControlPlaneDB | None = None,
|
|
*,
|
|
binding: dict[str, Any],
|
|
dry_run: bool = True,
|
|
operator_authorized: bool = False,
|
|
cleanup_authorization: Mapping[str, Any] | None = None,
|
|
project_root: str | None = None,
|
|
) -> dict[str, Any]:
|
|
"""Safely resolve (retire) one confirmed stale missing worktree binding (#970).
|
|
|
|
Apply mode requires a server-minted cleanup authorization; a client-supplied
|
|
``operator_authorized`` is rejected outright (B2). Immediately before
|
|
mutating, the authoritative lease, session, checkpoint, and issue-lock state
|
|
is re-read and diffed against what the audit saw, and the binding is
|
|
reclassified from those live values — so a lease that was re-activated, an
|
|
owner that changed, or a path that moved after the audit fails closed
|
|
without mutation (B1/B4).
|
|
"""
|
|
root = _norm_path(project_root or os.getcwd())
|
|
cp_db = db or cpd.ControlPlaneDB()
|
|
|
|
rec_path = binding.get("worktree_path") or binding.get("recorded_path") or ""
|
|
if not rec_path:
|
|
return {
|
|
"success": True,
|
|
"performed": False,
|
|
"outcome": "already_retired",
|
|
"message": "worktree_path is already empty (no action needed)",
|
|
}
|
|
|
|
# 1. Authorization (apply mode only; dry-run is always non-mutating).
|
|
if operator_authorized:
|
|
return {
|
|
"success": False,
|
|
"performed": False,
|
|
"reason": "operator_authorized_rejected",
|
|
"reasons": [OPERATOR_AUTHORIZED_REJECTION],
|
|
"message": OPERATOR_AUTHORIZED_REJECTION,
|
|
}
|
|
|
|
if not dry_run:
|
|
auth_ok, auth_reasons = _authorization_valid_for_apply(cleanup_authorization)
|
|
if not auth_ok:
|
|
return {
|
|
"success": False,
|
|
"performed": False,
|
|
"reason": "authorization_required",
|
|
"reasons": auth_reasons,
|
|
"cleanup_authorization": dict(cleanup_authorization or {}),
|
|
"message": (
|
|
"Mutation refused: server-enforced cleanup authorization is "
|
|
"required to retire a stale binding"
|
|
),
|
|
}
|
|
|
|
# 2. Fresh authoritative read of every safety-relevant record (B1).
|
|
live_state = read_live_safety_state(cp_db, binding)
|
|
live_snapshot = live_state["snapshot"]
|
|
drift = compare_safety_state(binding.get("audited_state"), live_snapshot)
|
|
if drift["changed"]:
|
|
return {
|
|
"success": False,
|
|
"performed": False,
|
|
"reason": "safety_state_changed",
|
|
"reasons": drift["reasons"],
|
|
"changed_fields": drift["changed_fields"],
|
|
"audited_state": dict(binding.get("audited_state") or {}),
|
|
"live_state": live_snapshot,
|
|
"differences": drift["differences"],
|
|
"message": (
|
|
"Pre-mutation revalidation refused: safety-relevant state "
|
|
f"changed since the audit ({', '.join(drift['changed_fields'])})"
|
|
),
|
|
}
|
|
|
|
# 3. Reclassify from the live values, not the audit snapshot (B1/B4).
|
|
reval = classify_worktree_binding(
|
|
recorded_path=live_state["live_worktree_path"] or rec_path,
|
|
lease_status=live_snapshot.get("lease_status"),
|
|
lease_phase=binding.get("lease_phase"),
|
|
owner_pid=live_snapshot.get("lease_owner_pid"),
|
|
session_id=live_snapshot.get("lease_session_id"),
|
|
session_active=bool(live_snapshot.get("session_active")),
|
|
session_live=bool(live_snapshot.get("session_live")),
|
|
session_pid_alive=bool(live_snapshot.get("session_pid_alive")),
|
|
session_status=live_snapshot.get("session_status"),
|
|
issue_lock=live_state.get("issue_lock"),
|
|
branch=binding.get("branch"),
|
|
project_root=root,
|
|
)
|
|
|
|
if reval.get("classification") != CLASS_CONFIRMED_STALE_DELETED:
|
|
return {
|
|
"success": False,
|
|
"performed": False,
|
|
"reason": "revalidation_failed",
|
|
"revalidation": reval,
|
|
"live_state": live_snapshot,
|
|
"message": f"Pre-mutation re-validation failed: binding classified as '{reval.get('classification')}' instead of '{CLASS_CONFIRMED_STALE_DELETED}'",
|
|
}
|
|
|
|
source = binding.get("source", "lease")
|
|
lease_id = binding.get("lease_id")
|
|
session_id = binding.get("session_id")
|
|
checkpoint_id = live_snapshot.get("checkpoint_id") or binding.get("checkpoint_id")
|
|
|
|
planned_action = {
|
|
"source": source,
|
|
"lease_id": lease_id,
|
|
"session_id": session_id,
|
|
"checkpoint_id": checkpoint_id,
|
|
"retired_worktree_path": rec_path,
|
|
"classification": CLASS_CONFIRMED_STALE_DELETED,
|
|
}
|
|
|
|
if dry_run:
|
|
return {
|
|
"success": True,
|
|
"performed": False,
|
|
"dry_run": True,
|
|
"action": "retire_stale_binding",
|
|
"planned_action": planned_action,
|
|
"live_state": live_snapshot,
|
|
"message": f"Dry-run: would retire stale missing worktree_path '{rec_path}' for {source} (lease={lease_id})",
|
|
}
|
|
|
|
# 4. Apply mode: targeted, compare-and-swapped retirement.
|
|
audit_proof: dict[str, Any] = {
|
|
"retired_worktree_path": rec_path,
|
|
"source": source,
|
|
"lease_id": lease_id,
|
|
"session_id": session_id,
|
|
"checkpoint_id": checkpoint_id,
|
|
"classification": CLASS_CONFIRMED_STALE_DELETED,
|
|
"authorization_source": CLEANUP_AUTHORIZATION_SOURCE,
|
|
"revalidated_against": live_snapshot,
|
|
}
|
|
|
|
try:
|
|
if source == "lease" and lease_id:
|
|
audit_proof["db_lease_update"] = cp_db.retire_lease_worktree_path(
|
|
lease_id,
|
|
expected_path=rec_path,
|
|
expected_status=live_snapshot.get("lease_status"),
|
|
expected_session_id=live_snapshot.get("lease_session_id"),
|
|
expected_owner_pid=live_snapshot.get("lease_owner_pid"),
|
|
reason="missing_worktree_path_retired_by_reconciler",
|
|
)
|
|
elif source == "session_checkpoint" and (checkpoint_id or session_id):
|
|
audit_proof["db_checkpoint_update"] = (
|
|
cp_db.retire_session_checkpoint_worktree_path(
|
|
session_id=session_id or "",
|
|
checkpoint_id=checkpoint_id,
|
|
expected_path=rec_path,
|
|
expected_status=live_snapshot.get("checkpoint_status"),
|
|
reason="missing_worktree_path_retired_by_reconciler",
|
|
)
|
|
)
|
|
else:
|
|
return {
|
|
"success": False,
|
|
"performed": False,
|
|
"reason": "unsupported_binding_source",
|
|
"message": (
|
|
f"binding source '{source}' has no sanctioned retirement "
|
|
"path; refusing to mutate"
|
|
),
|
|
}
|
|
except cpd.ControlPlaneError as exc:
|
|
# The compare-and-swap lost a race. Nothing was mutated.
|
|
return {
|
|
"success": False,
|
|
"performed": False,
|
|
"reason": "compare_and_swap_failed",
|
|
"reasons": [str(exc)],
|
|
"live_state": live_snapshot,
|
|
"message": (
|
|
"Retirement refused by compare-and-swap: the stored binding "
|
|
"changed concurrently"
|
|
),
|
|
}
|
|
|
|
return {
|
|
"success": True,
|
|
"performed": True,
|
|
"dry_run": False,
|
|
"action": "retired_stale_binding",
|
|
"audit_proof": audit_proof,
|
|
"message": f"Successfully retired stale missing worktree_path '{rec_path}' for lease {lease_id}",
|
|
}
|
|
|
|
|
|
def reconcile_missing_worktree_bindings(
|
|
db: cpd.ControlPlaneDB | None = None,
|
|
*,
|
|
project_root: str | None = None,
|
|
remote: str = "prgs",
|
|
org: str | None = None,
|
|
repo: str | None = None,
|
|
host: str | None = None,
|
|
dry_run: bool = True,
|
|
operator_authorized: bool = False,
|
|
cleanup_authorization: Mapping[str, Any] | None = None,
|
|
) -> dict[str, Any]:
|
|
"""Audit and safely reconcile missing worktree path bindings (#970).
|
|
|
|
Audits all recorded bindings, identifies confirmed stale missing paths,
|
|
retires eligible bindings (when ``dry_run=False`` and a server-minted
|
|
``cleanup_authorization`` proves authority), and returns durable
|
|
before/after audit evidence. A caller-supplied ``operator_authorized`` is
|
|
rejected, never honoured (#970 review 644 B2).
|
|
"""
|
|
cp_db = db or cpd.ControlPlaneDB()
|
|
root = _norm_path(project_root or os.getcwd())
|
|
|
|
initial_audit = audit_missing_worktree_bindings(
|
|
cp_db,
|
|
project_root=root,
|
|
remote=remote,
|
|
org=org,
|
|
repo=repo,
|
|
host=host,
|
|
)
|
|
|
|
candidates = initial_audit.get("confirmed_stale_candidates") or []
|
|
resolutions: list[dict[str, Any]] = []
|
|
|
|
for candidate in candidates:
|
|
res = resolve_missing_worktree_binding(
|
|
cp_db,
|
|
binding=candidate,
|
|
dry_run=dry_run,
|
|
operator_authorized=operator_authorized,
|
|
cleanup_authorization=cleanup_authorization,
|
|
project_root=root,
|
|
)
|
|
resolutions.append(res)
|
|
|
|
post_audit = audit_missing_worktree_bindings(
|
|
cp_db,
|
|
project_root=root,
|
|
remote=remote,
|
|
org=org,
|
|
repo=repo,
|
|
host=host,
|
|
)
|
|
|
|
refusals = [r for r in resolutions if not r.get("success")]
|
|
return {
|
|
"success": True,
|
|
"dry_run": dry_run,
|
|
"operator_authorized_rejected": bool(operator_authorized),
|
|
"cleanup_authorization": dict(cleanup_authorization or {}),
|
|
"before_audit": initial_audit,
|
|
"resolutions": resolutions,
|
|
"refused_count": len(refusals),
|
|
"after_audit": post_audit,
|
|
"worktrees_dimension_resolved": post_audit.get("worktrees_dimension_resolved", False),
|
|
"stale_candidates_count": len(candidates),
|
|
"resolved_count": sum(1 for r in resolutions if r.get("performed") or (dry_run and r.get("success"))),
|
|
}
|