Compare commits

..
Author SHA1 Message Date
jcwalker3 0b60fd6557 Remediate PR #861 findings F1-F9 and TestWorktreeStart failures (#860)
- Fix F1: Prepare recovery worktree detached at remote_head without git checkout -B to avoid exit 128 when branch is held by source worktree
- Fix F2: Pass recovery_sanctioned=True in bind_session_lock and assess_same_issue_lease_conflict
- Fix F3: Add SOURCE_RECOVER_DIRTY_ORPHANED to SANCTIONED_LOCK_SOURCES
- Fix F4: Stop after Phase 4 dirty apply when conflicts exist; do not finalize session binding
- Fix F5: Dynamically query competing live locks and workflow leases in MCP server
- Fix F6: Fail closed on recovery worktree resume when HEAD does not match expected remote_head
- Fix F7: Fail closed on remote HEAD observation failure rather than copying expected_remote_head pin
- Fix F8: Enforce foreign overwrite protection requiring same claimant or sanctioned reclaim
- Fix F9: Add real multi-worktree integration tests for prepare_recovery_worktree and lock rebind
- Fix TestWorktreeStart: Bypass session lock check for dry-run and review/pr-* branches in scripts/worktree-start
2026-07-23 22:34:35 -05:00
jcwalker3andGrok 4.5 18d6583e83 fix(author): bootstrap recovery for dirty orphaned issue worktrees (#860)
Add an explicit recovery operation for same-claimant dirty registered
worktrees under malformed PID-less durable locks, with crash-safe journals,
dirty byte preservation, path-level conflict detection, and live session
binding. PID-less locks are never treated as live merely because expiry is
absent.

Closes #860

Co-Authored-By: Grok 4.5 (xAI) <[email protected]>
2026-07-23 20:38:47 -05:00
17 changed files with 1966 additions and 2241 deletions
File diff suppressed because it is too large Load Diff
-122
View File
@@ -1,122 +0,0 @@
# Sanctioned restart and graceful reload controls (#642)
Sessions used to recover MCP connectivity by killing the host daemon
(`pkill -f mcp_server.py`, #630). That is forbidden and stays forbidden: it
kills every namespace on the host, contaminates whichever session survives, and
leaves no audit trail. This document describes the sanctioned replacement,
implemented in `webui/sanctioned_restart.py`.
## What the console will and will not do
The console **never** restarts anything. It authorizes an intent, records it,
and hands off to a host supervisor. There is no code path in which the console
sends a signal, spawns a process, or renders a kill command — a regression test
asserts the module contains no `subprocess`, `signal`, `os.kill`, `os.system`,
or `popen` reference, and that no returned payload contains a kill command.
## Operations
| Mode | Action | Minimum role | Behaviour |
|------|--------|--------------|-----------|
| `reload` | `system.reload_namespace` | controller | Host supervisor reloads the namespace in place, draining in-flight requests. |
| `restart` | `system.restart_namespace` | admin | Host supervisor restarts the namespace. In-flight requests are lost. |
Scope is always exactly one namespace. A fleet-wide restart is an explicit
non-goal: `all`, `*`, `fleet`, and an empty scope are refused with
`fleet_scope_not_permitted`, because that is precisely the blast radius the
forbidden kill already had. An unrecognised namespace is refused rather than
passed through to the host.
## The gate sequence
`assess_restart_request()` applies every gate in order and reports the first
failure with a stable reason code:
| Order | Gate | Reason code on failure |
|-------|------|------------------------|
| 1 | Mode is `restart` or `reload` | `unknown_mode` |
| 2 | Scope is a single known namespace | `fleet_scope_not_permitted`, `unknown_namespace` |
| 3 | Principal holds the required console role | `unauthorized` |
| 4 | Confirmation phrase supplied | `confirmation_required` |
| 5 | Confirmation names this namespace and mode | `confirmation_mismatch` |
| 6 | Out-of-band operator authorization present | `operator_authorization_missing` |
| 7 | Runtime is not contaminated | `contaminated_runtime` |
| 8 | Host restart hook configured | `restart_hook_not_configured` |
Passing every gate yields `host_action_required`, never "restarted".
### Confirmation binds the namespace
The required phrase is `"<mode> <namespace>"` — for example
`restart gitea-author`. Binding the namespace into the phrase is the point: a
confirmation typed for one namespace cannot be replayed against another.
### Operator authorization is not self-assertable
Host daemon maintenance is authorized out of band through
`GITEA_OPERATOR_DAEMON_MAINTENANCE_AUTHORIZATION`, read from the process
environment and nowhere else (#630; #710 finding F1). A worker session cannot
set an environment variable for an already-running daemon, so this cannot be
faked the way a tool argument could.
### The host hook
`GITEA_SANCTIONED_RESTART_HOOK` holds an opaque reference the *host* resolves —
a supervisor label such as a launchd job name, never a command line. With no
hook configured the request is refused; the console does not fall back to a
process kill. The value is read server-side and never rendered to a client.
## Manual kill remains contamination
`classify_restart_command()` classifies an operator-proposed recovery command.
A manual `pkill`/`kill`/`killall` of the MCP daemon is contamination, not a
restart: it returns `clean_claim_allowed: false` and builds a durable
contamination marker (redacted command only, never secrets) naming
`system.restart_namespace` as the sanctioned alternative.
A live, uncleared contamination marker also blocks a restart. This is stricter
than #630's task-scoped gate, which deliberately lets a contaminated worker keep
commenting and handing off: restarting a contaminated runtime would launder the
contamination rather than resolve it. Clear the marker through the reconciler
path first.
## Post-restart health verification
After the host supervisor acts, `verify_post_restart_health()` decides whether
the session may claim to be clean:
| Status | Meaning | Clean claim |
|--------|---------|-------------|
| `clean` | Required tool callable, proven through the live client namespace | Allowed |
| `unproven` | Reported healthy without live client-namespace evidence | Refused |
| `unhealthy` | Probe failed | Refused |
Only `probe_source=client_namespace` evidence clears a session. Static tool
registration is not proof, and neither is an offline subprocess probe — an IDE
client can hold a registered tool list while live calls fail with
`client is closing: EOF` (see
[`mcp-namespace-health.md`](mcp-namespace-health.md)).
## Audit
Every attempt — allowed or denied — is recorded through
`webui.console_audit` with actor, target namespace, mode, result, and reason
code, and is redacted before it is persisted. `system.restart_namespace` is
break-glass, so its records are retained for 730 days. Records carry
`process_kill_executed: false`, which is a fact about the code path rather than
a claim: no such path exists.
## Environment variables
| Variable | Purpose |
|----------|---------|
| `GITEA_SANCTIONED_RESTART_HOOK` | Host supervisor reference; absent means restart is refused. |
| `GITEA_OPERATOR_DAEMON_MAINTENANCE_AUTHORIZATION` | Out-of-band operator authorization reference. |
| `WEBUI_AUDIT_LOG` | Console audit sink; absent means records are built but not persisted. |
## Non-goals
* No unrestricted `kill` from the UI, in any role, in any phase.
* No fleet-wide restart.
* No silent auto-restart loop: every attempt is confirmed and audited.
* This does not implement the Phase 1 health API (#634).
+2 -11
View File
@@ -91,8 +91,6 @@ already define, and a regression test asserts each mapping matches.
| `close_pr` | controller | privileged | `gitea.pr.close` | Yes | No | No | 3 |
| `merge_pr` | controller | privileged | `gitea.pr.merge` | Yes | **Yes** | **Yes** | 3 |
| `delete_branch` | admin | destructive | `gitea.branch.delete` | Yes | **Yes** | **Yes** | 3 |
| `system.reload_namespace` | controller | privileged | `runtime.reload_namespace` | Yes | No | No | 2 |
| `system.restart_namespace` | admin | destructive | `runtime.restart_namespace` | Yes | **Yes** | **Yes** | 2 |
**Dual control** means the acting principal may not be the sole authority: a
second distinct principal must confirm. **Break-glass** means the action is
@@ -104,13 +102,6 @@ honouring it.
`delete_branch` is admin-only rather than controller because it is the one
irreversible action in the set.
`system.restart_namespace` is admin-only for the same reason: restarting a
namespace drops every in-flight request on it. `system.reload_namespace` drains
first, so it is privileged but not destructive. Neither action is ever executed
by the console — both hand off to a host supervisor, and neither exposes a raw
process kill. See
[`sanctioned-restart-controls.md`](sanctioned-restart-controls.md) (#642).
### Authorization decision
`authorize(action_id, principal, for_execution=False)` returns a decision
@@ -208,8 +199,8 @@ breaking the request it describes.
| Class | Applies to | Default |
|-------|-----------|---------|
| `standard` | Routine gated writes | 90 days |
| `privileged` | `review_pr`, `close_pr`, `system.reload_namespace`, and any unclassifiable action | 365 days |
| `break_glass` | `merge_pr`, `delete_branch`, `system.restart_namespace` | 730 days |
| `privileged` | `review_pr`, `close_pr`, and any unclassifiable action | 365 days |
| `break_glass` | `merge_pr`, `delete_branch` | 730 days |
Each record carries its own class, day count, and computed `expires_at`, so
retention is auditable per record rather than inferred from file age. An
+281 -78
View File
@@ -1450,7 +1450,7 @@ def verify_preflight_purity(
dirty_files = sorted(
_parse_porcelain_entries(_get_workspace_porcelain(workspace))
)
if dirty_files:
if dirty_files and task != "commit_files":
raise RuntimeError(
nwb.format_namespace_workspace_binding_error(
role_kind=role,
@@ -2031,6 +2031,7 @@ import issue_lock_store # noqa: E402
import issue_lock_adoption # noqa: E402
import issue_lock_recovery # noqa: E402
import issue_lock_renewal # noqa: E402
import dirty_orphan_worktree_recovery # noqa: E402 # #860 dirty orphan recovery
import stacked_pr_support # noqa: E402
import merge_approval_gate # noqa: E402
import review_quarantine # noqa: E402 # #695 contaminated formal-review quarantine
@@ -4342,6 +4343,280 @@ def gitea_lock_issue(
return result
@mcp.tool()
def gitea_recover_dirty_orphaned_issue_worktree(
issue_number: int,
branch_name: str,
source_worktree_path: str,
expected_local_head: str,
expected_remote_head: str,
expected_dirty_fingerprints: dict,
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
recovery_worktree_path: str | None = None,
dry_run: bool = False,
) -> dict:
"""Recover a dirty orphaned same-claimant author issue worktree (#860).
Explicit recovery operation does **not** silently widen ``gitea_lock_issue``.
Accepts authoritative expected pins (repository, issue, branch, source
worktree, claimant, local head, remote/PR head, dirty fingerprints) and
fails closed on any mismatch. PID-less malformed locks are never treated
as live merely because expiry is absent. The source worktree is frozen;
recovery prepares a separate worktree at the pinned remote head, re-applies
dirty bytes with path-level conflict detection, and binds a live author
session only after recovery state is consistent.
Args:
issue_number: Issue whose durable claim is being recovered.
branch_name: Locked branch ``(fix|feat|docs|chore)/issue-N-``.
source_worktree_path: Registered dirty source worktree under branches/.
expected_local_head: Full 40-char SHA of the source worktree HEAD.
expected_remote_head: Full 40-char SHA of the remote/PR head to sync to.
expected_dirty_fingerprints: ``{relative_path: sha256}`` of dirty bytes.
remote/host/org/repo: Repository binding.
recovery_worktree_path: Optional recovery worktree path under branches/.
dry_run: Assess eligibility only; no filesystem or lock mutation.
Returns:
dict with success, outcome, conflicts, recovery_worktree_path, reasons,
evidence, and journal metadata.
"""
task = "recover_dirty_orphaned_issue_worktree"
ok, block_reasons = role_session_router.check_author_mutation_after_reviewer_stop(
task
)
if not ok:
return {
"success": False,
"performed": False,
"outcome": "REFUSED",
"reasons": block_reasons,
}
blocked = _namespace_mutation_block(task, remote=remote)
if blocked:
return blocked
blocked = _profile_permission_block(
task_capability_map.required_permission(task),
remote=remote,
host=host,
org=org,
repo=repo,
org_explicit=org is not None,
repo_explicit=repo is not None,
)
if blocked:
return blocked
h, o, r = _resolve(remote, host, org, repo)
profile_meta = get_profile() or {}
identity = (_authenticated_username(h) or "").strip()
profile = (profile_meta.get("profile_name") or "").strip()
if not identity or not profile:
return {
"success": False,
"performed": False,
"outcome": "REFUSED",
"reasons": ["could not resolve authenticated identity/profile"],
}
existing_lock = _load_existing_issue_lock(
remote=remote, org=o, repo=r, issue_number=issue_number
)
src = os.path.realpath(source_worktree_path)
git_state = issue_lock_worktree.read_worktree_git_state(src)
observed_local = (git_state.get("head_sha") or "").strip()
porcelain = git_state.get("porcelain_status") or ""
current_branch = git_state.get("current_branch")
# Observed dirty fingerprints from source worktree bytes.
observed_fps: dict[str, str] = {}
dirty_contents: dict[str, bytes] = {}
for rel in (expected_dirty_fingerprints or {}):
rel_n = str(rel).strip()
fpath = os.path.join(src, rel_n)
if not os.path.isfile(fpath):
continue
with open(fpath, "rb") as fh:
data = fh.read()
dirty_contents[rel_n] = data
observed_fps[rel_n] = dirty_orphan_worktree_recovery.sha256_bytes(data)
# Remote head observation (best-effort; pin mismatch fails closed).
observed_remote = ""
try:
probe = subprocess.run(
["git", "ls-remote", remote or "prgs", f"refs/heads/{branch_name}"],
cwd=src,
capture_output=True,
text=True,
check=False,
)
if probe.returncode == 0 and (probe.stdout or "").strip():
observed_remote = (probe.stdout or "").strip().split()[0]
except Exception:
observed_remote = ""
registered = False
try:
listing = subprocess.run(
["git", "worktree", "list", "--porcelain"],
cwd=src,
capture_output=True,
text=True,
check=False,
)
if listing.returncode == 0:
registered = src in (listing.stdout or "")
except Exception:
registered = False
project_root = _canonical_local_git_root()
canonical_root = author_mutation_worktree.resolve_canonical_repo_root(
src, project_root
)
competing_locks: list[dict] = []
try:
all_live = issue_lock_store.list_live_locks()
for l in all_live:
if l.get("issue_number") == issue_number:
wt = l.get("worktree_path")
if not wt or not issue_lock_store._same_realpath(wt, src):
competing_locks.append(l)
except Exception:
competing_locks = []
wf_active = False
wf_expired = True
try:
db, _ = _control_plane_db_or_error()
if db is not None:
active_leases_data = lease_lifecycle.list_active_leases(
db,
remote=remote if remote in REMOTES else remote,
org=o,
repo=r,
)
leases_list = active_leases_data.get("leases") or []
for l in leases_list:
if l.get("work_number") == issue_number and l.get("work_kind") == "issue":
fresh = l.get("freshness") or {}
if fresh.get("status") == "active":
wf_active = True
wf_expired = False
elif fresh.get("status") in ("expired", "stale_dead_process"):
wf_active = False
wf_expired = True
except Exception:
pass
assessment = dirty_orphan_worktree_recovery.assess_dirty_orphan_recovery(
existing_lock,
issue_number=issue_number,
branch_name=branch_name,
source_worktree_path=src,
remote=remote if remote else "prgs",
org=o,
repo=r,
identity=identity,
profile=profile,
expected_local_head=expected_local_head,
expected_remote_head=expected_remote_head,
expected_dirty_fingerprints=expected_dirty_fingerprints or {},
current_branch=current_branch,
porcelain_status=porcelain,
observed_local_head=observed_local,
observed_remote_head=observed_remote,
observed_dirty_fingerprints=observed_fps,
competing_live_locks=competing_locks,
competing_live_sessions=[],
workflow_lease_active=wf_active,
workflow_lease_expired=wf_expired,
canonical_repo_root=canonical_root,
worktree_registered=registered,
current_pid=os.getpid(),
)
if dry_run or not assessment.get("eligible"):
return {
"success": bool(assessment.get("eligible")),
"performed": False,
"dry_run": dry_run,
"outcome": assessment.get("outcome"),
"reasons": list(assessment.get("reasons") or []),
"evidence": dict(assessment.get("evidence") or {}),
"eligible": bool(assessment.get("eligible")),
}
if not recovery_worktree_path:
recovery_worktree_path = os.path.join(
canonical_root,
"branches",
f"recovery-issue-{issue_number}-dirty-orphan",
)
# Load blob contents at local/remote heads for conflict detection.
def _blob_at(head: str, rel: str) -> bytes | None:
try:
proc = subprocess.run(
["git", "show", f"{head}:{rel}"],
cwd=src,
capture_output=True,
check=False,
)
if proc.returncode != 0:
return None
return proc.stdout
except Exception:
return None
local_contents = {
rel: _blob_at(expected_local_head, rel)
for rel in (expected_dirty_fingerprints or {})
}
remote_contents = {
rel: _blob_at(expected_remote_head, rel)
for rel in (expected_dirty_fingerprints or {})
}
# Preflight purity is satisfied via explicit worktree_path on this tool's
# recovery path; source remains frozen and is never cleaned.
result = dirty_orphan_worktree_recovery.run_dirty_orphan_recovery(
assessment=assessment,
existing_lock=existing_lock or {},
issue_number=issue_number,
branch_name=branch_name,
source_worktree_path=src,
recovery_worktree_path=recovery_worktree_path,
remote=remote if remote else "prgs",
org=o,
repo=r,
identity=identity,
profile=profile,
expected_local_head=expected_local_head,
expected_remote_head=expected_remote_head,
expected_dirty_fingerprints=expected_dirty_fingerprints or {},
dirty_contents=dirty_contents,
local_head_contents=local_contents,
remote_head_contents=remote_contents,
canonical_repo_root=canonical_root,
bind_lock=True,
session_pid=os.getpid(),
)
# Surface preflight recognition for recovered provenance.
if result.get("success") and result.get("lock_record"):
result["preflight_provenance"] = (
dirty_orphan_worktree_recovery.preflight_recognizes_recovered_provenance(
result["lock_record"]
)
)
return result
@mcp.tool()
def gitea_assess_work_issue_duplicate(
issue_number: int,
@@ -11634,7 +11909,6 @@ def gitea_audit_worktree_cleanup(
org: str | None = None,
repo: str | None = None,
ttl_hours: float = worktree_cleanup_audit.DEFAULT_TTL_HOURS,
merged_pr_limit: int = 200,
) -> dict:
"""Read-only: classify every session-owned worktree under ``branches/`` (#401).
@@ -11645,26 +11919,17 @@ def gitea_audit_worktree_cleanup(
the active issue-lock branch is read from the local lock file and treated
as active work. Deletes nothing and mutates no Gitea state.
Merged PRs are fetched as well, so an issue worktree can be linked to the
PR that owns its branch (#858). Such a worktree only becomes removable
when that owning PR is unambiguous and merged, the worktree head is
already contained in authoritative master, and nothing else protects it
no open or competing PR, lease, issue lock, live session, dirty file, or
protected/control checkout. Anything unproven keeps it classified as
active issue work.
Fails closed if the live open-PR list, the merged-PR list, or the
control-plane lease state cannot be read: without them removability
cannot be proven, so no candidates are returned.
Fails closed if the live open-PR list cannot be fetched: without it,
removability cannot be proven, so no candidates are returned.
Args:
remote: Known instance 'dadeschools' or 'prgs'.
host: Override the Gitea host.
org: Override the owner/organization.
repo: Override the repository name.
ttl_hours: Age (hours) after which a clean conflict-fix worktree
becomes stale-removable (default from GITEA_WORKTREE_TTL_HOURS).
merged_pr_limit: Max closed PRs scanned for merged-PR ownership.
ttl_hours: Age (hours) after which a clean issue/conflict-fix
worktree becomes stale-removable (default from
GITEA_WORKTREE_TTL_HOURS).
Returns:
dict with per-worktree classifications, counts, removable
@@ -11700,84 +11965,22 @@ def gitea_audit_worktree_cleanup(
if (pr.get("head") or {}).get("ref")
}
# #858: merged PRs are the ownership evidence that lets a landed issue
# worktree stop being reported as active work. Without them the audit can
# never agree with the PR-scoped reconciler, so treat a fetch failure the
# same way an open-PR fetch failure is treated: fail closed.
try:
closed_prs = api_get_all(
f"{repo_api_url(h, o, r)}/pulls?state=closed", auth, limit=merged_pr_limit
)
except Exception as exc:
return {
"success": False,
"performed": False,
"open_pr_state_verified": True,
"merged_pr_state_verified": False,
"reasons": [
"could not fetch merged PRs; worktree ownership unverified "
f"(fail closed): {_redact(str(exc))}"
],
}
merged_prs = [pr for pr in closed_prs if (pr.get("merged") or pr.get("merged_at"))]
pr_index = worktree_cleanup_audit.build_pr_index(list(open_prs) + merged_prs)
# #858: the auditor already accepted lease evidence but nothing ever
# supplied it, so every worktree looked unleased. Removability is now
# reachable for issue worktrees, so authoritative control-plane leases
# must be readable or the audit fails closed.
db, lease_errs = _control_plane_db_or_error()
if db is None:
return {
"success": False,
"performed": False,
"open_pr_state_verified": True,
"merged_pr_state_verified": True,
"lease_state_verified": False,
"reasons": [
"could not read control-plane leases; worktree protection "
"unverified (fail closed)",
*lease_errs,
],
}
lease_result = lease_lifecycle.list_active_leases(
db, remote=remote, org=o, repo=r, include_non_active=False, limit=500
)
leased_issue_numbers: set[int] = set()
live_session_paths: set[str] = set()
for lease in lease_result.get("leases") or []:
if lease.get("work_kind") == "issue" and lease.get("work_number") is not None:
try:
leased_issue_numbers.add(int(lease["work_number"]))
except (TypeError, ValueError):
pass
if lease.get("worktree_path"):
live_session_paths.add(str(lease["worktree_path"]))
active_issue_branches: set[str] = set()
lock = merged_cleanup_reconcile.read_issue_lock(ISSUE_LOCK_FILE)
if lock and lock.get("branch_name"):
active_issue_branches.add(str(lock["branch_name"]).strip())
master_ref = f"{remote}/master" if remote in REMOTES else "origin/master"
report = worktree_cleanup_audit.audit_branches_directory(
_canonical_local_git_root(),
open_pr_branches=open_pr_branches,
active_issue_branches=active_issue_branches,
now=datetime.now(timezone.utc),
ttl_hours=ttl_hours,
pr_index=pr_index,
leased_issue_numbers=leased_issue_numbers,
live_session_paths=live_session_paths,
master_ref=master_ref,
)
return {
"success": True,
"performed": False,
"open_pr_state_verified": True,
"merged_pr_state_verified": True,
"lease_state_verified": True,
"master_ref": master_ref,
"task_mode": "work-issue",
**report,
}
+2
View File
@@ -16,11 +16,13 @@ ISSUE_LOCK_FILE = os.environ.get("GITEA_ISSUE_LOCK_FILE", "/tmp/gitea_issue_lock
SOURCE_LOCK_ISSUE = "gitea_lock_issue"
SOURCE_LOCK_ADOPTION = "gitea_lock_issue_adoption"
SOURCE_OPERATOR_OVERRIDE = "operator_override"
SOURCE_RECOVER_DIRTY_ORPHANED = "gitea_recover_dirty_orphaned_issue_worktree"
SANCTIONED_LOCK_SOURCES = frozenset({
SOURCE_LOCK_ISSUE,
SOURCE_LOCK_ADOPTION,
SOURCE_OPERATOR_OVERRIDE,
SOURCE_RECOVER_DIRTY_ORPHANED,
})
_OPERATOR_OVERRIDE_ENV = "GITEA_ISSUE_LOCK_OPERATOR_OVERRIDE"
+82 -5
View File
@@ -169,6 +169,7 @@ def bind_session_lock(
*,
expected_generation: int | None = None,
renewal_sanctioned: bool = False,
recovery_sanctioned: bool = False,
) -> str:
"""Persist a keyed lock and bind it to the current process session.
@@ -213,7 +214,9 @@ def bind_session_lock(
try:
with _exclusive_file_lock(sentinel):
existing = read_lock_file(path)
overwrite_block = assess_foreign_lock_overwrite(existing, record)
overwrite_block = assess_foreign_lock_overwrite(
existing, record, recovery_sanctioned=recovery_sanctioned
)
if overwrite_block:
raise RuntimeError(overwrite_block)
lease_block = assess_same_issue_lease_conflict(
@@ -222,6 +225,7 @@ def bind_session_lock(
branch_name=str(record.get("branch_name") or ""),
worktree_path=str(record.get("worktree_path") or ""),
renewal_sanctioned=renewal_sanctioned,
recovery_sanctioned=recovery_sanctioned,
)
if lease_block:
raise RuntimeError(lease_block)
@@ -380,7 +384,16 @@ def assess_lock_freshness(
pid = lock_data.get("session_pid")
if pid is None:
pid = lock_data.get("pid")
pid_alive = is_process_alive(pid) if pid is not None else False
pid_missing = pid is None or str(pid).strip() == ""
try:
pid_int = int(pid) if not pid_missing else None
if pid_int is not None and pid_int <= 0:
pid_missing = True
pid_int = None
except (TypeError, ValueError):
pid_missing = True
pid_int = None
pid_alive = is_process_alive(pid_int) if pid_int is not None else False
if expires_at and expires_at <= current:
return {
@@ -389,15 +402,36 @@ def assess_lock_freshness(
"stale": True,
"reason": f"lease expired at {expires_at.isoformat()}",
"pid_alive": pid_alive,
"pid_missing": pid_missing,
}
if pid is not None and not pid_alive:
# #860: a PID-less lock must never be considered live merely because
# expiration / heartbeat fields are absent. Missing PID is insufficient
# evidence of a live owner; treat as malformed/stale so recovery routes
# can evaluate corroborating pins instead of blocking on a false live flag.
if pid_missing:
return {
"status": "malformed",
"live": False,
"stale": True,
"reason": (
"lock has no usable session pid; cannot prove live ownership "
"(PID-less locks are never live by missing expiry alone)"
),
"pid_alive": False,
"pid_missing": True,
"heartbeat_at": heartbeat_at.isoformat() if heartbeat_at else None,
"expires_at": expires_at.isoformat() if expires_at else None,
}
if pid_int is not None and not pid_alive:
return {
"status": "stale",
"live": False,
"stale": True,
"reason": f"owner pid {pid} is not alive",
"reason": f"owner pid {pid_int} is not alive",
"pid_alive": False,
"pid_missing": False,
}
return {
@@ -406,6 +440,7 @@ def assess_lock_freshness(
"stale": False,
"reason": "lock heartbeat and lease are fresh",
"pid_alive": pid_alive,
"pid_missing": False,
"heartbeat_at": heartbeat_at.isoformat() if heartbeat_at else None,
"expires_at": expires_at.isoformat() if expires_at else None,
}
@@ -486,6 +521,7 @@ def assess_same_issue_lease_conflict(
worktree_path: str,
operation_type: str = AUTHOR_ISSUE_WORK_LEASE,
renewal_sanctioned: bool = False,
recovery_sanctioned: bool = False,
now: datetime | None = None,
) -> str | None:
"""Return a fail-closed error when a competing live lease blocks acquisition.
@@ -517,6 +553,8 @@ def assess_same_issue_lease_conflict(
existing_branch == branch_name
and _same_realpath(str(existing_worktree or ""), worktree_path)
)
if recovery_sanctioned and existing_issue == issue_number and existing_branch == branch_name:
return None
if is_lease_expired(existing_lock, now=now):
# #760 AC1/AC2: exact-owner renewal is a different disposition from
# foreign takeover and is evaluated first. Before this, both branches
@@ -547,10 +585,26 @@ def assess_same_issue_lease_conflict(
)
def _lock_claimant(lock: dict[str, Any] | None) -> dict[str, str]:
if not isinstance(lock, dict):
return {}
claimant = lock.get("claimant")
if not isinstance(claimant, dict):
lease = lock.get("work_lease")
claimant = lease.get("claimant") if isinstance(lease, dict) else None
if not isinstance(claimant, dict):
return {}
return {
"username": str(claimant.get("username") or ""),
"profile": str(claimant.get("profile") or ""),
}
def assess_foreign_lock_overwrite(
existing_lock: dict[str, Any] | None,
incoming_lock: dict[str, Any],
*,
recovery_sanctioned: bool = False,
now: datetime | None = None,
) -> str | None:
"""Block writes that would clobber an unrelated live lease on the same key."""
@@ -565,8 +619,31 @@ def assess_foreign_lock_overwrite(
)
if same_issue and same_branch and same_worktree:
return None
if not is_lease_live(existing_lock, now=now):
existing_claimant = _lock_claimant(existing_lock)
incoming_claimant = _lock_claimant(incoming_lock)
same_claimant = (
bool(existing_claimant.get("username"))
and existing_claimant.get("username") == incoming_claimant.get("username")
and existing_claimant.get("profile") == incoming_claimant.get("profile")
)
if recovery_sanctioned and same_issue and same_branch and same_claimant:
return None
if not is_lease_live(existing_lock, now=now):
# #860 F8: A non-live or PID-less lock still blocks foreign overwrite
# unless same claimant or sanctioned reclaim is proven.
if not same_claimant and same_issue:
reclaim = assess_expired_lock_reclaim(existing_lock, now=now)
if not reclaim.get("reclaim_allowed"):
return (
"Refusing foreign overwrite of non-live issue lock "
f"(issue #{existing_lock.get('issue_number')}, owner '{existing_claimant.get('username')}') "
"without sanctioned reclaim proof (fail closed)"
)
return None
return (
"Refusing to overwrite a live foreign issue lock "
f"(issue #{existing_lock.get('issue_number')}, "
+10 -8
View File
@@ -43,19 +43,21 @@ repo_root="$(cd "$script_dir/.." && pwd)"
# Enforce issue-linked, traceable branch names (issue → branch → worktree → PR).
if [[ "$allow_unlinked" -eq 0 ]]; then
locked_branch=$(python3 -c "
if [[ "$dry_run" -eq 0 ]] && [[ ! "$branch" =~ ^review/pr-[0-9]+-.+ ]]; then
locked_branch=$(python3 -c "
import sys
sys.path.insert(0, '$repo_root')
import issue_lock_store
print(issue_lock_store.resolve_locked_branch_for_session('$branch'))
")
if [[ -z "$locked_branch" ]]; then
echo "Error: No session issue lock is bound. Call gitea_lock_issue before branch creation (fail closed)." >&2
exit 2
fi
if [[ "$branch" != "$locked_branch" ]]; then
echo "Error: Requested branch '$branch' does not match locked branch '$locked_branch' (fail closed)." >&2
exit 2
if [[ -z "$locked_branch" ]]; then
echo "Error: No session issue lock is bound. Call gitea_lock_issue before branch creation (fail closed)." >&2
exit 2
fi
if [[ "$branch" != "$locked_branch" ]]; then
echo "Error: Requested branch '$branch' does not match locked branch '$locked_branch' (fail closed)." >&2
exit 2
fi
fi
if [[ "$branch" =~ ^(fix|feat|docs|chore)/issue-[0-9]+-.+ ]] \
+14 -15
View File
@@ -32,6 +32,15 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"permission": "gitea.issue.comment",
"role": "author",
},
# #860: dirty orphaned same-claimant worktree recovery (explicit operation).
"recover_dirty_orphaned_issue_worktree": {
"permission": "gitea.issue.comment",
"role": "author",
},
"gitea_recover_dirty_orphaned_issue_worktree": {
"permission": "gitea.issue.comment",
"role": "author",
},
"set_issue_labels": {
"permission": "gitea.issue.comment",
"role": "author",
@@ -335,21 +344,6 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"role": "controller",
},
# #642: sanctioned host-daemon lifecycle controls. Deliberately *not* a
# ``gitea.*`` operation — restarting an MCP namespace is a host action, not
# a Gitea API call, and no configured Gitea profile should be able to
# satisfy it by accident. Authority comes from the console RBAC model plus
# out-of-band operator authorization (#630); these entries exist so the
# console cannot invent an authority the capability layer never declared.
"restart_namespace": {
"permission": "runtime.restart_namespace",
"role": "controller",
},
"reload_namespace": {
"permission": "runtime.reload_namespace",
"role": "controller",
},
# #601 first-class lease lifecycle — inspect/list need read; mutations gate on
# ownership in the control-plane DB (not a separate Gitea write permission).
"list_workflow_leases": {
@@ -492,6 +486,11 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
# merger lease (#763).
_PREFLIGHT_TASK_TRANSITIONS = frozenset({
("review_pr", "acquire_reviewer_pr_lease"),
("work_issue", "lock_issue"),
("work_issue", "recover_dirty_orphaned_issue_worktree"),
("work_issue", "gitea_recover_dirty_orphaned_issue_worktree"),
("work_issue", "commit_files"),
("work_issue", "gitea_commit_files"),
})
@@ -0,0 +1,483 @@
"""Synthetic regression coverage for dirty orphaned worktree recovery (#860).
Modeled on the #850 / #855 shape without mutating their real state.
"""
from __future__ import annotations
import json
import os
import shutil
import tempfile
import unittest
from unittest import mock
import dirty_orphan_worktree_recovery as dorec
import issue_lock_store
DEAD_PID = 999_999_999
LIVE_PID = os.getpid()
BRANCH = "fix/issue-901-dirty-orphan"
SOURCE_WT = "/repo/branches/issue-901-dirty-orphan"
RECOVERY_WT_NAME = "recovery-issue-901-dirty-orphan"
LOCAL_HEAD = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"
REMOTE_HEAD = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"
OTHER_HEAD = "cccccccccccccccccccccccccccccccccccccccc"
FP_A = dorec.sha256_bytes(b"dirty-a")
FP_B = dorec.sha256_bytes(b"dirty-b")
FP_C = dorec.sha256_bytes(b"dirty-c-conflict")
def durable_lock(**overrides):
"""#850-shaped PID-less malformed same-claimant lock."""
lock = {
"issue_number": 901,
"branch_name": BRANCH,
"worktree_path": SOURCE_WT,
"remote": "prgs",
"org": "Example-Org",
"repo": "Example-Repo",
# intentionally no pid / session_pid / work_lease expiry
"claimant": {"username": "author-user", "profile": "prgs-author"},
}
lock.update(overrides)
return lock
def base_kwargs(**overrides):
kwargs = {
"issue_number": 901,
"branch_name": BRANCH,
"source_worktree_path": SOURCE_WT,
"remote": "prgs",
"org": "Example-Org",
"repo": "Example-Repo",
"identity": "author-user",
"profile": "prgs-author",
"expected_local_head": LOCAL_HEAD,
"expected_remote_head": REMOTE_HEAD,
"expected_dirty_fingerprints": {"a.py": FP_A, "b.py": FP_B},
"current_branch": BRANCH,
"porcelain_status": " M a.py\n M b.py\n",
"observed_local_head": LOCAL_HEAD,
"observed_remote_head": REMOTE_HEAD,
"observed_dirty_fingerprints": {"a.py": FP_A, "b.py": FP_B},
"competing_live_locks": [],
"competing_live_sessions": [],
"workflow_lease_active": False,
"workflow_lease_expired": True,
"canonical_repo_root": "/repo",
"worktree_registered": True,
"current_pid": LIVE_PID,
}
kwargs.update(overrides)
return kwargs
def assess(lock=None, **overrides):
return dorec.assess_dirty_orphan_recovery(
durable_lock() if lock is None else lock, **base_kwargs(**overrides)
)
class FreshnessPidLess(unittest.TestCase):
def test_pid_less_lock_is_not_live(self):
freshness = issue_lock_store.assess_lock_freshness(durable_lock())
self.assertFalse(freshness["live"])
self.assertTrue(freshness.get("pid_missing"))
self.assertEqual(freshness["status"], "malformed")
def test_pid_less_with_far_future_expiry_still_not_live(self):
lock = durable_lock(
work_lease={
"operation_type": "author_issue_work",
"expires_at": "2999-01-01T00:00:00Z",
"last_heartbeat_at": "2999-01-01T00:00:00Z",
}
)
freshness = issue_lock_store.assess_lock_freshness(lock)
self.assertFalse(freshness["live"])
self.assertTrue(freshness.get("pid_missing"))
class EligibilityGranted(unittest.TestCase):
def test_dead_same_claimant_pid_less_dirty(self):
result = assess()
self.assertEqual(result["outcome"], dorec.ELIGIBLE)
self.assertTrue(result["eligible"])
def test_expired_workflow_lease_corroboration(self):
result = assess(workflow_lease_active=False, workflow_lease_expired=True)
self.assertTrue(result["eligible"])
def test_older_local_newer_remote_heads(self):
result = assess()
self.assertTrue(result["evidence"].get("heads_diverged"))
self.assertTrue(result["eligible"])
class EligibilityRefused(unittest.TestCase):
def test_active_owner_with_pid(self):
lock = durable_lock(pid=LIVE_PID, session_pid=LIVE_PID)
result = assess(lock=lock, owner_process_alive_override=True)
self.assertEqual(result["outcome"], dorec.REFUSED)
self.assertFalse(result["eligible"])
self.assertTrue(any("alive" in r for r in result["reasons"]))
def test_foreign_claimant(self):
result = assess(identity="other-user")
self.assertEqual(result["outcome"], dorec.REFUSED)
self.assertTrue(any("foreign claimant identity" in r for r in result["reasons"]))
def test_foreign_profile(self):
result = assess(profile="prgs-reviewer")
self.assertEqual(result["outcome"], dorec.REFUSED)
def test_fingerprint_mismatch(self):
result = assess(observed_dirty_fingerprints={"a.py": "0" * 64, "b.py": FP_B})
self.assertEqual(result["outcome"], dorec.REFUSED)
self.assertTrue(any("fingerprint mismatch" in r for r in result["reasons"]))
def test_head_mismatch(self):
result = assess(observed_local_head=OTHER_HEAD)
self.assertEqual(result["outcome"], dorec.REFUSED)
def test_remote_head_mismatch(self):
result = assess(observed_remote_head=OTHER_HEAD)
self.assertEqual(result["outcome"], dorec.REFUSED)
def test_path_not_under_branches(self):
result = assess(
source_worktree_path="/tmp/branches/evil",
# lock path also changed so worktree agreement holds
lock=durable_lock(worktree_path="/tmp/branches/evil"),
)
self.assertEqual(result["outcome"], dorec.REFUSED)
self.assertTrue(any("canonical branches" in r for r in result["reasons"]))
def test_unregistered_worktree(self):
result = assess(worktree_registered=False)
self.assertEqual(result["outcome"], dorec.REFUSED)
def test_active_workflow_lease(self):
result = assess(workflow_lease_active=True, workflow_lease_expired=False)
self.assertEqual(result["outcome"], dorec.REFUSED)
def test_unsafe_dirty_path_pin(self):
result = assess(
expected_dirty_fingerprints={"../etc/passwd": FP_A},
observed_dirty_fingerprints={"../etc/passwd": FP_A},
)
self.assertEqual(result["outcome"], dorec.REFUSED)
def test_symlink_escape_rejected_by_ancestry(self):
ok, reasons = dorec.is_path_under_canonical_branches(
"/tmp/branches/evil", canonical_repo_root="/repo"
)
self.assertFalse(ok)
self.assertTrue(reasons)
class ConflictDetection(unittest.TestCase):
def test_overlapping_upstream_change(self):
conflicts = dorec.detect_path_conflicts(
dirty_paths=["c.py"],
local_head_contents={"c.py": b"local-base"},
remote_head_contents={"c.py": b"remote-changed"},
dirty_contents={"c.py": b"dirty-c-conflict"},
)
self.assertEqual(len(conflicts), 1)
self.assertEqual(conflicts[0]["path"], "c.py")
def test_unchanged_upstream_no_conflict(self):
conflicts = dorec.detect_path_conflicts(
dirty_paths=["a.py"],
local_head_contents={"a.py": b"same"},
remote_head_contents={"a.py": b"same"},
dirty_contents={"a.py": b"dirty-a"},
)
self.assertEqual(conflicts, [])
class CrashSafeRecovery(unittest.TestCase):
def setUp(self):
self.tmp = tempfile.mkdtemp(prefix="dirty-orphan-")
self.repo = os.path.join(self.tmp, "repo")
self.branches = os.path.join(self.repo, "branches")
self.source = os.path.join(self.branches, "issue-901-dirty-orphan")
self.recovery = os.path.join(self.branches, RECOVERY_WT_NAME)
os.makedirs(self.source, exist_ok=True)
os.makedirs(self.branches, exist_ok=True)
# seed dirty files in source
with open(os.path.join(self.source, "a.py"), "wb") as fh:
fh.write(b"dirty-a")
with open(os.path.join(self.source, "b.py"), "wb") as fh:
fh.write(b"dirty-b")
self.journal_dir = os.path.join(self.tmp, "journals")
self.lock = durable_lock(worktree_path=self.source)
self.assessment = dorec.assess_dirty_orphan_recovery(
self.lock,
**base_kwargs(
source_worktree_path=self.source,
canonical_repo_root=self.repo,
),
)
class FakeGit(dorec.GitOps):
def __init__(self, recovery_path, head):
self.recovery_path = recovery_path
self.head = head
self.calls = []
def run(self, args, *, cwd):
self.calls.append((args, cwd))
if args[:3] == ["git", "worktree", "add"]:
os.makedirs(self.recovery_path, exist_ok=True)
return mock.Mock(returncode=0, stdout="", stderr="")
if args[:2] == ["git", "checkout"]:
return mock.Mock(returncode=0, stdout="", stderr="")
if args[:2] == ["git", "rev-parse"]:
return mock.Mock(returncode=0, stdout=self.head + "\n", stderr="")
return mock.Mock(returncode=0, stdout="", stderr="")
self.git = FakeGit(self.recovery, REMOTE_HEAD)
self.written_locks = []
def lock_writer(record):
self.written_locks.append(record)
self.lock_writer = lock_writer
def tearDown(self):
shutil.rmtree(self.tmp, ignore_errors=True)
def _run(self, **overrides):
kwargs = {
"assessment": self.assessment,
"existing_lock": self.lock,
"issue_number": 901,
"branch_name": BRANCH,
"source_worktree_path": self.source,
"recovery_worktree_path": self.recovery,
"remote": "prgs",
"org": "Example-Org",
"repo": "Example-Repo",
"identity": "author-user",
"profile": "prgs-author",
"expected_local_head": LOCAL_HEAD,
"expected_remote_head": REMOTE_HEAD,
"expected_dirty_fingerprints": {"a.py": FP_A, "b.py": FP_B},
"dirty_contents": {"a.py": b"dirty-a", "b.py": b"dirty-b"},
"local_head_contents": {"a.py": b"base-a", "b.py": b"base-b"},
"remote_head_contents": {"a.py": b"base-a", "b.py": b"base-b"},
"canonical_repo_root": self.repo,
"bind_lock": True,
"lock_writer": self.lock_writer,
"git_ops": self.git,
"journal_dir": self.journal_dir,
"session_pid": LIVE_PID,
}
kwargs.update(overrides)
return dorec.run_dirty_orphan_recovery(**kwargs)
def test_success_preserves_dirty_bytes_and_source(self):
result = self._run()
self.assertTrue(result["success"])
self.assertEqual(result["outcome"], dorec.RECOVERY_COMPLETED)
self.assertTrue(os.path.isdir(self.source))
with open(os.path.join(self.source, "a.py"), "rb") as fh:
self.assertEqual(fh.read(), b"dirty-a")
with open(os.path.join(self.recovery, "a.py"), "rb") as fh:
self.assertEqual(fh.read(), b"dirty-a")
with open(os.path.join(self.recovery, "b.py"), "rb") as fh:
self.assertEqual(fh.read(), b"dirty-b")
self.assertEqual(len(self.written_locks), 1)
rec = self.written_locks[0]
self.assertEqual(rec["session_pid"], LIVE_PID)
self.assertTrue(rec["dirty_orphan_recovery"]["recovered"])
self.assertTrue(rec["dirty_orphan_recovery"]["source_frozen"])
def test_conflict_leaves_governed_state(self):
result = self._run(
expected_dirty_fingerprints={"c.py": FP_C},
dirty_contents={"c.py": b"dirty-c-conflict"},
local_head_contents={"c.py": b"local-base"},
remote_head_contents={"c.py": b"remote-changed"},
)
# #860 F4: session binding is NOT finalized while conflicts remain
self.assertFalse(result["success"])
self.assertEqual(result["outcome"], dorec.CONFLICTS_PRESENT)
sidecar = os.path.join(self.recovery, "c.py.recovered-dirty")
self.assertTrue(os.path.isfile(sidecar))
state = os.path.join(
self.recovery, dorec.CONFLICT_STATE_DIR, dorec.CONFLICT_STATE_FILE
)
self.assertTrue(os.path.isfile(state))
with open(state, "r", encoding="utf-8") as fh:
payload = json.load(fh)
self.assertEqual(payload["resolution"], "author_edit_required")
def test_interrupt_before_journal_no_artifacts(self):
result = self._run(interrupt_after_phase=dorec.PHASE_ELIGIBILITY)
self.assertFalse(result["success"])
self.assertEqual(result["outcome"], "INTERRUPTED")
self.assertFalse(os.path.isdir(self.recovery))
def test_interrupt_after_journal_then_retry_idempotent(self):
first = self._run(interrupt_after_phase=dorec.PHASE_JOURNAL_PERSISTED)
self.assertEqual(first["outcome"], "INTERRUPTED")
self.assertTrue(first["journal"]["artifacts_created"]["journal"])
second = self._run()
self.assertTrue(second["success"])
# source still recoverable
with open(os.path.join(self.source, "a.py"), "rb") as fh:
self.assertEqual(fh.read(), b"dirty-a")
def test_interrupt_after_worktree_then_retry(self):
first = self._run(interrupt_after_phase=dorec.PHASE_RECOVERY_WORKTREE)
self.assertEqual(first["outcome"], "INTERRUPTED")
self.assertTrue(os.path.isdir(self.recovery))
second = self._run()
self.assertTrue(second["success"])
def test_interrupt_after_binding_then_retry_complete(self):
first = self._run(interrupt_after_phase=dorec.PHASE_BINDING)
self.assertEqual(first["outcome"], "INTERRUPTED")
second = self._run()
self.assertTrue(second["success"])
# completed journal makes further retries no-ops
third = self._run()
self.assertEqual(third["outcome"], dorec.RECOVERY_RESUMED)
def test_source_worktree_never_deleted(self):
self._run()
self.assertTrue(os.path.isdir(self.source))
self.assertTrue(os.path.isfile(os.path.join(self.source, "a.py")))
def test_fingerprint_drift_refuses_without_mutation(self):
result = self._run(dirty_contents={"a.py": b"CHANGED", "b.py": b"dirty-b"})
self.assertFalse(result["success"])
self.assertFalse(os.path.isdir(self.recovery))
class SessionBindingPreflight(unittest.TestCase):
def test_canonical_session_binding_recognized(self):
lock = {
"worktree_path": "/repo/branches/recovery",
"session_pid": LIVE_PID,
"dirty_orphan_recovery": {
"recovered": True,
"conflicts": [],
"recovery_worktree_path": "/repo/branches/recovery",
"source_worktree_path": SOURCE_WT,
"accepted_head": REMOTE_HEAD,
},
}
result = dorec.preflight_recognizes_recovered_provenance(lock)
self.assertTrue(result["recognized"])
def test_conflicts_block_commit_preflight(self):
lock = {
"worktree_path": "/repo/branches/recovery",
"session_pid": LIVE_PID,
"dirty_orphan_recovery": {
"recovered": True,
"conflicts": [{"path": "c.py"}],
},
}
result = dorec.preflight_recognizes_recovered_provenance(lock)
self.assertFalse(result["recognized"])
def test_active_foreign_does_not_mutate(self):
# assess-only path: foreign refused before run
result = assess(identity="intruder")
self.assertFalse(result["eligible"])
class JournalSymlinkRefusal(unittest.TestCase):
def test_symlink_journal_path_refused_on_load(self):
tmp = tempfile.mkdtemp()
try:
real = os.path.join(tmp, "real.json")
with open(real, "w", encoding="utf-8") as fh:
fh.write("{}")
link = os.path.join(tmp, "link.json")
os.symlink(real, link)
key = "symlink-test"
jdir = tmp
path = dorec._journal_path(key, journal_dir=jdir)
with open(path, "w", encoding="utf-8") as fh:
json.dump({"idempotency_key": key}, fh)
os.remove(path)
os.symlink(real, path)
with self.assertRaises(ValueError):
dorec.load_journal(key, journal_dir=jdir)
finally:
shutil.rmtree(tmp, ignore_errors=True)
class RealGitMultiWorktreeIntegration(unittest.TestCase):
def setUp(self):
import subprocess
self.tmp = tempfile.mkdtemp(prefix="git-integration-")
self.repo = os.path.join(self.tmp, "repo")
os.makedirs(self.repo, exist_ok=True)
subprocess.run(["git", "init"], cwd=self.repo, check=True, capture_output=True)
subprocess.run(["git", "config", "user.name", "Test User"], cwd=self.repo, check=True)
subprocess.run(["git", "config", "user.email", "[email protected]"], cwd=self.repo, check=True)
with open(os.path.join(self.repo, "init.txt"), "w") as fh:
fh.write("init")
subprocess.run(["git", "add", "."], cwd=self.repo, check=True)
subprocess.run(["git", "commit", "-m", "init"], cwd=self.repo, check=True)
branch = "fix/issue-999-test"
subprocess.run(["git", "branch", branch], cwd=self.repo, check=True)
self.branches = os.path.join(self.repo, "branches")
self.source = os.path.join(self.branches, "issue-999-test")
subprocess.run(["git", "worktree", "add", self.source, branch], cwd=self.repo, check=True)
self.dirty_path = os.path.join(self.source, "dirty.txt")
with open(self.dirty_path, "w") as fh:
fh.write("dirty-data")
def tearDown(self):
shutil.rmtree(self.tmp, ignore_errors=True)
def test_prepare_recovery_worktree_detached_no_exit_128(self):
import subprocess
head_sha = subprocess.check_output(["git", "rev-parse", "HEAD"], cwd=self.repo, text=True).strip()
rec_wt = os.path.join(self.branches, "recovery-issue-999-test")
res = dorec.prepare_recovery_worktree(
canonical_repo_root=self.repo,
recovery_worktree_path=rec_wt,
branch_name="fix/issue-999-test",
remote_head=head_sha,
)
self.assertTrue(res["success"], res.get("reasons"))
self.assertTrue(os.path.isdir(rec_wt))
def test_real_lock_rebind_recovery_sanctioned(self):
lock_dir = os.path.join(self.tmp, "locks")
rec_wt = os.path.join(self.branches, "recovery-issue-999-test")
os.makedirs(rec_wt, exist_ok=True)
record = {
"remote": "prgs",
"org": "Example-Org",
"repo": "Example-Repo",
"issue_number": 999,
"branch_name": "fix/issue-999-test",
"worktree_path": rec_wt,
"claimant": {"username": "author-user", "profile": "prgs-author"},
}
record_src = dict(record)
record_src["worktree_path"] = self.source
issue_lock_store.bind_session_lock(record_src, lock_dir=lock_dir)
path = issue_lock_store.bind_session_lock(
record,
lock_dir=lock_dir,
recovery_sanctioned=True,
)
self.assertTrue(os.path.isfile(path))
if __name__ == "__main__":
unittest.main()
@@ -1,551 +0,0 @@
"""Merged-PR awareness for the worktree cleanup audit (#858).
Before #858 an ``issue_work`` worktree could never leave ``active_issue_work``:
the audit had no PR linkage at all (``pr_number`` was structurally ``None``)
and its only route to ``clean_stale_removable`` was a TTL derived from a
``last_used_at`` that nothing ever populated. A merged, clean, unprotected
worktree was therefore reported as active work forever, disagreeing with the
PR-scoped reconciler.
These tests use fabricated temporary repositories and synthetic PR records
only. Nothing here removes a worktree or deletes a branch.
"""
import os
import subprocess
import sys
import tempfile
import unittest
from unittest.mock import patch
sys.path.insert(0, str(__import__("pathlib").Path(__file__).resolve().parent.parent))
import merged_cleanup_reconcile as mcr # noqa: E402
import worktree_cleanup_audit as wca # noqa: E402
MERGED_BRANCH = "feat/issue-777-timeline"
MERGED_PATH = "/repo/branches/issue-777-timeline"
HEAD_SHA = "a" * 40
def _pr(number, branch, *, merged=True, sha=HEAD_SHA, state=None):
"""Synthetic Gitea PR payload."""
return {
"number": number,
"head": {"ref": branch, "sha": sha},
"merged_at": "2026-07-24T01:00:00Z" if merged else None,
"state": state or ("closed" if merged else "open"),
}
def _porcelain(*entries):
out = []
for path, branch, sha in entries:
out.append(f"worktree {path}")
out.append(f"HEAD {sha}")
if branch is None:
out.append("detached")
else:
out.append(f"branch refs/heads/{branch}")
out.append("")
return "\n".join(out)
class _AuditHarness(unittest.TestCase):
"""Runs audit_branches_directory over a fabricated worktree listing."""
PORCELAIN = _porcelain(
("/repo", "master", "f" * 40),
(MERGED_PATH, MERGED_BRANCH, HEAD_SHA),
)
def run_audit(self, *, dirty_paths=(), contained=True, **kwargs):
def fake_dirty(path):
if path in dirty_paths:
return {"exists": True, "dirty": True, "dirty_files": [" M x.py"]}
return {"exists": True, "dirty": False, "dirty_files": []}
with patch.object(
wca, "list_worktrees",
return_value=wca.parse_worktree_porcelain(self.PORCELAIN),
), patch.object(
wca, "read_worktree_dirty", side_effect=fake_dirty
), patch.object(
wca, "git_worktree_list", return_value="(mocked)"
), patch.object(
wca, "is_head_ancestor_of_ref", return_value=contained
):
report = wca.audit_branches_directory("/repo", **kwargs)
return {wt["path"]: wt for wt in report["worktrees"]}, report
def merged_audit(self, **kwargs):
kwargs.setdefault("pr_index", wca.build_pr_index([_pr(849, MERGED_BRANCH)]))
kwargs.setdefault("master_ref", "prgs/master")
return self.run_audit(**kwargs)
class TestMergedWorktreeBecomesRemovable(_AuditHarness):
def test_clean_merged_issue_worktree_is_linked_and_removable(self):
by_path, report = self.merged_audit()
entry = by_path[MERGED_PATH]
self.assertEqual(entry["classification"], wca.CLASS_CLEAN_STALE_REMOVABLE)
self.assertTrue(entry["removable"])
self.assertEqual(entry["merged_pr_linkage"]["status"], wca.LINKAGE_MERGED)
self.assertEqual(entry["merged_pr_cleanup"]["block_reasons"], [])
self.assertIn(MERGED_PATH, [c["path"] for c in report["removable_candidates"]])
def test_pr_number_populated_from_authoritative_linkage(self):
by_path, _ = self.merged_audit()
self.assertEqual(by_path[MERGED_PATH]["pr_number"], 849)
def test_regression_without_pr_evidence_stays_active_issue_work(self):
"""The pre-#858 behaviour, still correct when no PR state is supplied."""
by_path, _ = self.run_audit()
entry = by_path[MERGED_PATH]
self.assertEqual(entry["classification"], wca.CLASS_ACTIVE_ISSUE_WORK)
self.assertFalse(entry["removable"])
self.assertIsNone(entry["pr_number"])
class TestProtectiveSignalsSurvive(_AuditHarness):
def test_open_pr_worktree_is_not_removable(self):
index = wca.build_pr_index([_pr(900, MERGED_BRANCH, merged=False)])
by_path, _ = self.run_audit(
pr_index=index,
master_ref="prgs/master",
open_pr_branches={MERGED_BRANCH},
)
entry = by_path[MERGED_PATH]
self.assertEqual(entry["classification"], wca.CLASS_ACTIVE_OPEN_PR)
self.assertFalse(entry["removable"])
# linkage still reports the owning PR, it just is not merge proof
self.assertEqual(entry["merged_pr_linkage"]["status"], wca.LINKAGE_OPEN)
self.assertEqual(entry["pr_number"], 900)
def test_dirty_tracked_worktree_is_not_removable(self):
by_path, _ = self.merged_audit(dirty_paths=(MERGED_PATH,))
entry = by_path[MERGED_PATH]
self.assertEqual(entry["classification"], wca.CLASS_DIRTY_LOCAL)
self.assertFalse(entry["removable"])
self.assertIn(
"worktree has uncommitted changes",
entry["merged_pr_cleanup"]["block_reasons"],
)
def test_untracked_only_worktree_is_not_removable(self):
"""``git status --porcelain`` reports untracked files as dirty too."""
def untracked(path):
if path == MERGED_PATH:
return {"exists": True, "dirty": True, "dirty_files": ["?? scratch.txt"]}
return {"exists": True, "dirty": False, "dirty_files": []}
with patch.object(
wca, "list_worktrees",
return_value=wca.parse_worktree_porcelain(self.PORCELAIN),
), patch.object(
wca, "read_worktree_dirty", side_effect=untracked
), patch.object(
wca, "git_worktree_list", return_value="(mocked)"
), patch.object(
wca, "is_head_ancestor_of_ref", return_value=True
):
report = wca.audit_branches_directory(
"/repo",
pr_index=wca.build_pr_index([_pr(849, MERGED_BRANCH)]),
master_ref="prgs/master",
)
entry = {wt["path"]: wt for wt in report["worktrees"]}[MERGED_PATH]
self.assertEqual(entry["classification"], wca.CLASS_DIRTY_LOCAL)
self.assertFalse(entry["removable"])
def test_active_lease_by_issue_number_is_protective(self):
by_path, _ = self.merged_audit(leased_issue_numbers={777})
entry = by_path[MERGED_PATH]
self.assertTrue(entry["has_active_lease"])
self.assertEqual(entry["classification"], wca.CLASS_ACTIVE_ISSUE_WORK)
self.assertFalse(entry["removable"])
def test_active_lease_by_branch_is_protective(self):
by_path, _ = self.merged_audit(leased_branches={MERGED_BRANCH})
entry = by_path[MERGED_PATH]
self.assertTrue(entry["has_active_lease"])
self.assertFalse(entry["removable"])
def test_active_issue_lock_is_protective(self):
by_path, _ = self.merged_audit(active_issue_branches={MERGED_BRANCH})
entry = by_path[MERGED_PATH]
self.assertTrue(entry["has_active_issue_lock"])
self.assertEqual(entry["classification"], wca.CLASS_ACTIVE_ISSUE_WORK)
self.assertFalse(entry["removable"])
def test_live_session_worktree_is_protective(self):
by_path, _ = self.merged_audit(live_session_paths={MERGED_PATH})
entry = by_path[MERGED_PATH]
self.assertTrue(entry["has_live_session"])
self.assertEqual(entry["classification"], wca.CLASS_ACTIVE_ISSUE_WORK)
self.assertFalse(entry["removable"])
def test_head_not_contained_in_master_is_not_removable(self):
by_path, _ = self.merged_audit(contained=False)
entry = by_path[MERGED_PATH]
self.assertEqual(entry["classification"], wca.CLASS_ACTIVE_ISSUE_WORK)
self.assertFalse(entry["removable"])
self.assertIn(
"worktree head is not contained in authoritative master "
"(unmerged commits remain)",
entry["merged_pr_cleanup"]["block_reasons"],
)
def test_unknown_containment_fails_closed(self):
by_path, _ = self.merged_audit(contained=None)
entry = by_path[MERGED_PATH]
self.assertFalse(entry["removable"])
self.assertIn(
"containment of the worktree head in master is unknown",
entry["merged_pr_cleanup"]["block_reasons"],
)
def test_missing_master_ref_fails_closed(self):
by_path, _ = self.run_audit(
pr_index=wca.build_pr_index([_pr(849, MERGED_BRANCH)])
)
self.assertFalse(by_path[MERGED_PATH]["removable"])
def test_unmerged_owning_pr_is_not_removable(self):
index = wca.build_pr_index([_pr(901, MERGED_BRANCH, merged=False)])
by_path, _ = self.run_audit(pr_index=index, master_ref="prgs/master")
entry = by_path[MERGED_PATH]
self.assertFalse(entry["removable"])
self.assertIn(
"owning PR #901 is not merged",
entry["merged_pr_cleanup"]["block_reasons"],
)
def test_control_checkout_is_never_removable(self):
by_path, _ = self.merged_audit()
control = by_path["/repo"]
self.assertTrue(control["is_protected"])
self.assertEqual(control["classification"], wca.CLASS_UNSAFE_UNKNOWN)
self.assertFalse(control["removable"])
def test_control_checkout_not_removable_even_if_linked_and_merged(self):
"""A merged PR on the control checkout must not unlock removal."""
porcelain = _porcelain(("/repo", MERGED_BRANCH, HEAD_SHA))
with patch.object(
wca, "list_worktrees", return_value=wca.parse_worktree_porcelain(porcelain)
), patch.object(
wca, "read_worktree_dirty",
return_value={"exists": True, "dirty": False, "dirty_files": []},
), patch.object(
wca, "git_worktree_list", return_value="(mocked)"
), patch.object(
wca, "is_head_ancestor_of_ref", return_value=True
):
report = wca.audit_branches_directory(
"/repo",
pr_index=wca.build_pr_index([_pr(849, MERGED_BRANCH)]),
master_ref="prgs/master",
)
entry = report["worktrees"][0]
self.assertEqual(entry["classification"], wca.CLASS_UNSAFE_UNKNOWN)
self.assertFalse(entry["removable"])
class TestAmbiguousLinkageFailsClosed(_AuditHarness):
def test_competing_prs_on_one_branch_fail_closed(self):
index = wca.build_pr_index(
[_pr(849, MERGED_BRANCH), _pr(860, MERGED_BRANCH)]
)
by_path, _ = self.run_audit(pr_index=index, master_ref="prgs/master")
entry = by_path[MERGED_PATH]
self.assertEqual(entry["merged_pr_linkage"]["status"], wca.LINKAGE_AMBIGUOUS)
self.assertIsNone(entry["pr_number"])
self.assertEqual(entry["classification"], wca.CLASS_ACTIVE_ISSUE_WORK)
self.assertFalse(entry["removable"])
def test_merged_plus_open_pr_on_one_branch_fails_closed(self):
index = wca.build_pr_index(
[_pr(849, MERGED_BRANCH), _pr(861, MERGED_BRANCH, merged=False)]
)
by_path, _ = self.run_audit(pr_index=index, master_ref="prgs/master")
entry = by_path[MERGED_PATH]
self.assertEqual(entry["merged_pr_linkage"]["status"], wca.LINKAGE_AMBIGUOUS)
self.assertFalse(entry["removable"])
def test_no_owning_pr_fails_closed(self):
by_path, _ = self.run_audit(
pr_index=wca.build_pr_index([_pr(849, "feat/other-branch")]),
master_ref="prgs/master",
)
entry = by_path[MERGED_PATH]
self.assertEqual(entry["merged_pr_linkage"]["status"], wca.LINKAGE_NONE)
self.assertFalse(entry["removable"])
def test_malformed_pr_records_are_dropped_not_guessed(self):
index = wca.build_pr_index(
[
{"number": None, "head": {"ref": MERGED_BRANCH}},
{"number": 5, "head": {}},
{"number": "not-an-int", "head": {"ref": MERGED_BRANCH}},
]
)
self.assertEqual(index, {})
self.assertEqual(
wca.resolve_owning_pr(branch=MERGED_BRANCH, pr_index=index)["status"],
wca.LINKAGE_NONE,
)
def test_detached_worktree_has_no_branch_linkage(self):
self.assertEqual(
wca.resolve_owning_pr(branch=None, pr_index={})["status"],
wca.LINKAGE_UNKNOWN,
)
class TestUnrelatedClassificationsUnchanged(unittest.TestCase):
"""Non-issue_work worktrees keep their pre-#858 classifications."""
PORCELAIN = _porcelain(
("/repo", "master", "f" * 40),
("/repo/branches/review-pr42", "review-pr42", "2" * 40),
("/repo/branches/baseline-master-x", "baseline-master-x", "3" * 40),
("/repo/branches/conflict-fix-pr50", "conflict-fix-pr50", "4" * 40),
("/repo/branches/review-pr99", None, "5" * 40),
)
def _audit(self, **kwargs):
with patch.object(
wca, "list_worktrees",
return_value=wca.parse_worktree_porcelain(self.PORCELAIN),
), patch.object(
wca, "read_worktree_dirty",
return_value={"exists": True, "dirty": False, "dirty_files": []},
), patch.object(
wca, "git_worktree_list", return_value="(mocked)"
), patch.object(
wca, "is_head_ancestor_of_ref", return_value=True
):
report = wca.audit_branches_directory("/repo", **kwargs)
return {wt["path"]: wt for wt in report["worktrees"]}
def test_classifications_identical_with_and_without_pr_evidence(self):
without = self._audit()
with_evidence = self._audit(
pr_index=wca.build_pr_index([_pr(849, MERGED_BRANCH)]),
master_ref="prgs/master",
)
self.assertEqual(
{p: e["classification"] for p, e in without.items()},
{p: e["classification"] for p, e in with_evidence.items()},
)
def test_lease_on_issue_does_not_capture_similarly_named_scratch_trees(self):
"""A lease on issue 777 protects issue work, not baseline/review trees."""
porcelain = _porcelain(
("/repo/branches/baseline-master-issue-777", "baseline-issue-777", "7" * 40),
("/repo/branches/issue-777-timeline", MERGED_BRANCH, HEAD_SHA),
)
with patch.object(
wca, "list_worktrees", return_value=wca.parse_worktree_porcelain(porcelain)
), patch.object(
wca, "read_worktree_dirty",
return_value={"exists": True, "dirty": False, "dirty_files": []},
), patch.object(
wca, "git_worktree_list", return_value="(mocked)"
), patch.object(
wca, "is_head_ancestor_of_ref", return_value=True
):
report = wca.audit_branches_directory(
"/repo",
pr_index=wca.build_pr_index([_pr(849, MERGED_BRANCH)]),
master_ref="prgs/master",
leased_issue_numbers={777},
)
by_path = {wt["path"]: wt for wt in report["worktrees"]}
baseline = by_path["/repo/branches/baseline-master-issue-777"]
self.assertFalse(baseline["has_active_lease"])
self.assertEqual(baseline["classification"], wca.CLASS_CLEAN_STALE_REMOVABLE)
issue_work = by_path["/repo/branches/issue-777-timeline"]
self.assertTrue(issue_work["has_active_lease"])
self.assertFalse(issue_work["removable"])
def test_review_and_baseline_still_removable(self):
by_path = self._audit(
pr_index=wca.build_pr_index([]), master_ref="prgs/master"
)
self.assertEqual(
by_path["/repo/branches/review-pr42"]["classification"],
wca.CLASS_CLEAN_STALE_REMOVABLE,
)
self.assertEqual(
by_path["/repo/branches/baseline-master-x"]["classification"],
wca.CLASS_CLEAN_STALE_REMOVABLE,
)
self.assertEqual(
by_path["/repo/branches/review-pr99"]["classification"],
wca.CLASS_DETACHED_REVIEW_LEFTOVER,
)
def test_conflict_fix_ttl_behaviour_unchanged(self):
"""conflict_fix still needs only TTL expiry; #858 did not touch it."""
self.assertEqual(
wca.classify_worktree(
workflow_type=wca.WORKFLOW_CONFLICT_FIX,
is_dirty=False,
ttl_expired=True,
),
wca.CLASS_CLEAN_STALE_REMOVABLE,
)
self.assertEqual(
wca.classify_worktree(
workflow_type=wca.WORKFLOW_CONFLICT_FIX,
is_dirty=False,
ttl_expired=False,
),
wca.CLASS_ACTIVE_ISSUE_WORK,
)
def test_issue_work_ttl_alone_no_longer_grants_removal(self):
"""Age is not landing proof: TTL alone must not reclaim issue work."""
self.assertEqual(
wca.classify_worktree(
workflow_type=wca.WORKFLOW_ISSUE_WORK,
is_dirty=False,
ttl_expired=True,
),
wca.CLASS_ACTIVE_ISSUE_WORK,
)
class TestAssessorPerformsNoDeletion(_AuditHarness):
def test_audit_never_removes_a_worktree(self):
with patch.object(wca, "remove_worktree") as removal:
self.merged_audit()
removal.assert_not_called()
def test_audit_shells_out_to_no_destructive_git_command(self):
seen = []
real_run = subprocess.run
def recording_run(cmd, *args, **kwargs):
seen.append(cmd)
return real_run(["true"], *args, **kwargs)
with patch.object(subprocess, "run", side_effect=recording_run):
wca.audit_branches_directory("/nonexistent-repo-for-audit")
joined = [" ".join(c) if isinstance(c, list) else str(c) for c in seen]
for cmd in joined:
self.assertNotIn("worktree remove", cmd)
self.assertNotIn("branch -D", cmd)
self.assertNotIn("push", cmd)
class TestAgreementWithPrScopedReconciler(unittest.TestCase):
"""The audit and merged_cleanup_reconcile must agree on identical input.
Uses a real throwaway git repository so containment is computed by git
rather than asserted. Nothing outside the temporary directory is touched.
"""
def _git(self, *args):
subprocess.run(
["git", "-C", self.root, *args],
check=True,
capture_output=True,
text=True,
)
def setUp(self):
self._tmp = tempfile.TemporaryDirectory()
self.root = os.path.realpath(self._tmp.name)
self._git("init", "-b", "master", ".")
self._git("config", "user.email", "[email protected]")
self._git("config", "user.name", "Test")
with open(os.path.join(self.root, "seed.txt"), "w") as fh:
fh.write("seed\n")
self._git("add", "seed.txt")
self._git("commit", "-m", "seed")
self.branch = "feat/issue-777-timeline"
self._git("checkout", "-b", self.branch)
with open(os.path.join(self.root, "feature.txt"), "w") as fh:
fh.write("feature\n")
self._git("add", "feature.txt")
self._git("commit", "-m", "feature")
self.head_sha = subprocess.run(
["git", "-C", self.root, "rev-parse", "HEAD"],
capture_output=True, text=True, check=True,
).stdout.strip()
self._git("checkout", "master")
self._git("merge", "--no-ff", "-m", "merge feature", self.branch)
self.worktree = os.path.join(self.root, "branches", "issue-777-timeline")
self._git("worktree", "add", self.worktree, self.branch)
def tearDown(self):
self._tmp.cleanup()
def _pr_index(self):
return wca.build_pr_index(
[
{
"number": 849,
"head": {"ref": self.branch, "sha": self.head_sha},
"merged_at": "2026-07-24T01:00:00Z",
}
]
)
def _audit_entry(self):
report = wca.audit_branches_directory(
self.root, pr_index=self._pr_index(), master_ref="master"
)
return next(wt for wt in report["worktrees"] if wt["path"] == self.worktree)
def _reconciler_entry(self):
return mcr.assess_local_worktree_cleanup(
pr_number=849,
head_branch=self.branch,
merged=True,
worktree_state=mcr.resolve_cleanup_worktree_state(
project_root=self.root,
head_branch=self.branch,
issue_number=777,
pr_head_sha=self.head_sha,
target_ref="master",
),
active_lock=False,
)
def test_both_assessors_agree_the_worktree_is_safe(self):
audit_entry = self._audit_entry()
reconciler = self._reconciler_entry()
self.assertTrue(reconciler["safe_to_remove_worktree"], reconciler)
self.assertTrue(audit_entry["removable"], audit_entry)
self.assertEqual(audit_entry["pr_number"], reconciler["pr_number"])
self.assertEqual(audit_entry["merged_pr_cleanup"]["block_reasons"], [])
self.assertEqual(reconciler["block_reasons"], [])
def test_both_assessors_agree_a_dirty_worktree_is_unsafe(self):
with open(os.path.join(self.worktree, "feature.txt"), "a") as fh:
fh.write("local edit\n")
audit_entry = self._audit_entry()
reconciler = self._reconciler_entry()
self.assertFalse(audit_entry["removable"])
self.assertFalse(reconciler["safe_to_remove_worktree"])
def test_worktree_still_present_after_audit(self):
self._audit_entry()
self.assertTrue(os.path.isdir(self.worktree))
if __name__ == "__main__":
unittest.main()
+6
View File
@@ -24,6 +24,8 @@ def _lease(expires_at: str) -> dict:
def _lock_record(**overrides) -> dict:
# #860: live locks require a usable session pid; PID-less records are never
# classified live merely because expiry/heartbeat fields are present.
record = {
"issue_number": 420,
"branch_name": "feat/issue-420-server-code-parity",
@@ -31,6 +33,8 @@ def _lock_record(**overrides) -> dict:
"org": "Scaled-Tech-Consulting",
"repo": "Gitea-Tools",
"worktree_path": "/tmp/wt-420",
"session_pid": os.getpid(),
"pid": os.getpid(),
"work_lease": _lease("2999-01-01T00:00:00Z"),
}
record.update(overrides)
@@ -88,6 +92,8 @@ class TestIssueLockStore(unittest.TestCase):
existing = _lock_record(
branch_name="feat/issue-420-other",
worktree_path="/tmp/other",
session_pid=os.getpid(),
pid=os.getpid(),
work_lease=_lease("2999-01-01T00:00:00Z"),
)
path = ils.lock_file_path(
-502
View File
@@ -1,502 +0,0 @@
"""Sanctioned restart / graceful reload control tests (#642).
Acceptance criteria under test:
1. The sanctioned restart path is implemented behind gates (capability,
confirmation, operator authorization, host hook).
2. Manual ``pkill`` stays forbidden and is classified as contamination.
3. Post-restart mutations require clean health/session proof.
4. Authorized restart preview, unauthorized deny, contamination classification.
5. No entry point exposes a raw kill.
"""
import json
import os
import tempfile
import unittest
import mcp_namespace_health
import runtime_recovery_guard
from task_capability_map import TASK_CAPABILITY_MAP
from webui import console_audit, console_authz, gated_actions, sanctioned_restart
NAMESPACE = "gitea-author"
# An operator-authorized, hook-configured host. Passed explicitly so no test
# depends on (or mutates) the real process environment.
READY_ENV = {
sanctioned_restart.RESTART_HOOK_ENV: "launchd:cc.prgs.gitea-author",
runtime_recovery_guard.OPERATOR_AUTHORIZATION_ENV: "ops-ticket-4821",
}
def admin(subject: str = "[email protected]") -> console_authz.Principal:
return console_authz.Principal(
subject=subject,
role=console_authz.ADMIN,
identity_source=console_authz.IDENTITY_ACCESS_PROXY,
authenticated=True,
)
def viewer() -> console_authz.Principal:
return console_authz.Principal(
subject="[email protected]",
role=console_authz.VIEWER,
identity_source=console_authz.IDENTITY_ACCESS_PROXY,
authenticated=True,
)
class TestCapabilityWiring(unittest.TestCase):
"""AC1: authority is declared, not invented by the console."""
def test_actions_resolve_through_the_capability_map(self):
for action_id in (
sanctioned_restart.ACTION_RESTART_NAMESPACE,
sanctioned_restart.ACTION_RELOAD_NAMESPACE,
):
with self.subTest(action=action_id):
action = console_authz.get_action(action_id)
self.assertIsNotNone(action)
self.assertIn(action.task_key, TASK_CAPABILITY_MAP)
self.assertEqual(
action.mcp_permission,
TASK_CAPABILITY_MAP[action.task_key]["permission"],
)
def test_restart_permission_is_not_a_gitea_operation(self):
"""No configured Gitea profile should satisfy a host restart."""
permission = TASK_CAPABILITY_MAP["restart_namespace"]["permission"]
self.assertFalse(permission.startswith("gitea."))
def test_restart_is_destructive_dual_control_break_glass(self):
action = console_authz.get_action(
sanctioned_restart.ACTION_RESTART_NAMESPACE
)
self.assertEqual(action.action_class, console_authz.CLASS_DESTRUCTIVE)
self.assertEqual(action.minimum_role, console_authz.ADMIN)
self.assertTrue(action.dual_control)
self.assertTrue(action.break_glass)
self.assertTrue(action.requires_confirmation)
def test_reload_is_privileged_but_not_destructive(self):
action = console_authz.get_action(
sanctioned_restart.ACTION_RELOAD_NAMESPACE
)
self.assertEqual(action.action_class, console_authz.CLASS_PRIVILEGED)
self.assertTrue(action.requires_confirmation)
class TestPreview(unittest.TestCase):
"""AC4: an authorized preview renders the plan without executing it."""
def test_preview_lists_the_mutation_ledger(self):
preview = sanctioned_restart.build_restart_preview(
NAMESPACE, principal=admin(), env=READY_ENV
)
steps = [entry["step"] for entry in preview["mutation_ledger"]]
self.assertEqual(
steps, ["quiesce", "host_restart_hook", "health_recheck", "audit"]
)
self.assertTrue(preview["scope_valid"])
self.assertTrue(preview["post_restart_verification_required"])
def test_reload_preview_drains_instead_of_restarting(self):
preview = sanctioned_restart.build_restart_preview(
NAMESPACE, sanctioned_restart.MODE_RELOAD,
principal=admin(), env=READY_ENV,
)
steps = [entry["step"] for entry in preview["mutation_ledger"]]
self.assertIn("host_graceful_reload", steps)
self.assertNotIn("host_restart_hook", steps)
def test_preview_never_enables_execution(self):
preview = sanctioned_restart.build_restart_preview(
NAMESPACE, principal=admin(), env=READY_ENV
)
self.assertFalse(preview["execution_enabled"])
self.assertFalse(preview["authorization"]["execution_enabled"])
def test_confirmation_phrase_binds_the_namespace(self):
self.assertTrue(
sanctioned_restart.confirmation_matches(
NAMESPACE, sanctioned_restart.MODE_RESTART,
"restart gitea-author",
)
)
# A phrase typed for one namespace must not authorize another.
self.assertFalse(
sanctioned_restart.confirmation_matches(
"gitea-merger", sanctioned_restart.MODE_RESTART,
"restart gitea-author",
)
)
class TestGates(unittest.TestCase):
"""AC1/AC4: every gate denies with a stable reason code."""
def _assess(self, **kwargs):
params = {
"principal": admin(),
"confirmation": f"restart {NAMESPACE}",
"env": READY_ENV,
}
params.update(kwargs)
namespace = params.pop("namespace", NAMESPACE)
mode = params.pop("mode", sanctioned_restart.MODE_RESTART)
return sanctioned_restart.assess_restart_request(
namespace, mode, **params
)
def test_authorized_confirmed_request_passes_every_gate(self):
result = self._assess()
self.assertTrue(result["allowed"])
self.assertEqual(
result["reason_code"], sanctioned_restart.ALLOW_HOST_ACTION_REQUIRED
)
def test_passing_every_gate_is_not_an_execution_grant(self):
"""An allowed request still never lets the console touch the process."""
result = self._assess()
self.assertTrue(result["allowed"])
self.assertFalse(result["execution_enabled"])
self.assertFalse(result["console_executes"])
def test_unauthorized_principal_is_denied(self):
result = self._assess(principal=viewer())
self.assertFalse(result["allowed"])
self.assertEqual(
result["reason_code"], sanctioned_restart.DENY_UNAUTHORIZED
)
def test_anonymous_principal_is_denied(self):
result = self._assess(principal=None)
self.assertFalse(result["allowed"])
self.assertEqual(
result["reason_code"], sanctioned_restart.DENY_UNAUTHORIZED
)
def test_missing_confirmation_is_denied(self):
result = self._assess(confirmation=None)
self.assertFalse(result["allowed"])
self.assertEqual(
result["reason_code"], sanctioned_restart.DENY_CONFIRMATION_MISSING
)
def test_confirmation_for_another_namespace_is_denied(self):
result = self._assess(confirmation="restart gitea-merger")
self.assertFalse(result["allowed"])
self.assertEqual(
result["reason_code"], sanctioned_restart.DENY_CONFIRMATION_MISMATCH
)
def test_missing_operator_authorization_is_denied(self):
env = {sanctioned_restart.RESTART_HOOK_ENV: "launchd:cc.prgs.author"}
result = self._assess(env=env)
self.assertFalse(result["allowed"])
self.assertEqual(
result["reason_code"],
sanctioned_restart.DENY_OPERATOR_AUTHORIZATION,
)
def test_missing_host_hook_is_denied_without_kill_fallback(self):
env = {
runtime_recovery_guard.OPERATOR_AUTHORIZATION_ENV: "ops-ticket-1",
}
result = self._assess(env=env)
self.assertFalse(result["allowed"])
self.assertEqual(
result["reason_code"], sanctioned_restart.DENY_HOOK_NOT_CONFIGURED
)
def test_fleet_scope_is_refused(self):
for scope in ("all", "*", "fleet"):
with self.subTest(scope=scope):
result = self._assess(
namespace=scope, confirmation=f"restart {scope}"
)
self.assertFalse(result["allowed"])
self.assertEqual(
result["reason_code"], sanctioned_restart.DENY_FLEET_SCOPE
)
def test_unknown_namespace_is_refused(self):
result = self._assess(
namespace="gitea-nope", confirmation="restart gitea-nope"
)
self.assertFalse(result["allowed"])
self.assertEqual(
result["reason_code"], sanctioned_restart.DENY_UNKNOWN_NAMESPACE
)
def test_unknown_mode_is_refused(self):
result = self._assess(mode="obliterate")
self.assertFalse(result["allowed"])
self.assertEqual(
result["reason_code"], sanctioned_restart.DENY_UNKNOWN_MODE
)
def test_live_contamination_marker_blocks_restart(self):
marker = runtime_recovery_guard.build_contamination_record(
reason_class=runtime_recovery_guard.REASON_MANUAL_DAEMON_KILL,
command_redacted="pkill -f mcp_server.py",
)
result = self._assess(contamination_marker=marker)
self.assertFalse(result["allowed"])
self.assertEqual(
result["reason_code"], sanctioned_restart.DENY_CONTAMINATED_RUNTIME
)
def test_reconciler_cleared_marker_no_longer_blocks(self):
marker = runtime_recovery_guard.build_contamination_record(
reason_class=runtime_recovery_guard.REASON_MANUAL_DAEMON_KILL,
command_redacted="pkill -f mcp_server.py",
)
marker = dict(marker, cleared_by_reconciler=True)
result = self._assess(contamination_marker=marker)
self.assertTrue(result["allowed"])
class TestExecutionNeverKills(unittest.TestCase):
"""AC5: no path exposes or runs a raw process kill."""
def test_authorized_execution_defers_to_the_host_supervisor(self):
result = sanctioned_restart.execute_restart(
NAMESPACE,
principal=admin(),
confirmation=f"restart {NAMESPACE}",
env=READY_ENV,
)
self.assertTrue(result["allowed"])
self.assertFalse(result["success"])
self.assertFalse(result["process_kill_executed"])
self.assertEqual(
result["outcome"], sanctioned_restart.ALLOW_HOST_ACTION_REQUIRED
)
def test_denied_execution_reports_the_refusing_gate(self):
result = sanctioned_restart.execute_restart(
NAMESPACE, principal=viewer(), confirmation=f"restart {NAMESPACE}",
env=READY_ENV,
)
self.assertFalse(result["allowed"])
self.assertEqual(
result["outcome"], sanctioned_restart.DENY_UNAUTHORIZED
)
self.assertFalse(result["process_kill_executed"])
def test_module_never_spawns_a_process(self):
path = os.path.join(
os.path.dirname(os.path.dirname(os.path.abspath(__file__))),
"webui", "sanctioned_restart.py",
)
with open(path, encoding="utf-8") as handle:
source = handle.read()
for forbidden in (
"import subprocess", "import signal", "os.kill", "os.system",
"popen",
):
with self.subTest(forbidden=forbidden):
self.assertNotIn(forbidden, source.lower())
def test_no_surface_returns_a_kill_command(self):
payloads = [
sanctioned_restart.build_restart_preview(
NAMESPACE, principal=admin(), env=READY_ENV
),
sanctioned_restart.restart_policy(),
sanctioned_restart.execute_restart(
NAMESPACE, principal=admin(),
confirmation=f"restart {NAMESPACE}", env=READY_ENV,
),
]
for payload in payloads:
rendered = json.dumps(payload, default=str).lower()
self.assertNotIn("kill -9", rendered)
self.assertNotIn("pkill -f", rendered)
def test_policy_declares_no_raw_kill_and_no_silent_restart(self):
policy = sanctioned_restart.restart_policy()
self.assertFalse(policy["raw_kill_exposed"])
self.assertFalse(policy["console_executes_process_kill"])
self.assertFalse(policy["fleet_scope_permitted"])
self.assertFalse(policy["silent_auto_restart_permitted"])
self.assertTrue(policy["audit_required"])
class TestContaminationClassification(unittest.TestCase):
"""AC2: manual pkill is contamination, and it blocks clean claims."""
def test_manual_daemon_pkill_is_contamination(self):
result = sanctioned_restart.classify_restart_command(
"pkill -f mcp_server.py"
)
self.assertTrue(result["contamination"])
self.assertFalse(result["clean_claim_allowed"])
self.assertIsNotNone(result["contamination_marker"])
self.assertEqual(
result["sanctioned_alternative"],
sanctioned_restart.ACTION_RESTART_NAMESPACE,
)
def test_broad_process_kill_is_contamination(self):
result = sanctioned_restart.classify_restart_command("killall -9 Python")
self.assertTrue(result["contamination"])
self.assertFalse(result["clean_claim_allowed"])
def test_marker_names_the_sanctioned_alternative(self):
result = sanctioned_restart.classify_restart_command(
"pkill -f mcp_server.py"
)
marker = result["contamination_marker"]
self.assertIn(
sanctioned_restart.ACTION_RESTART_NAMESPACE, marker["detail"]
)
def test_benign_command_is_not_contamination(self):
result = sanctioned_restart.classify_restart_command("git status")
self.assertFalse(result["contamination"])
self.assertTrue(result["clean_claim_allowed"])
def test_no_command_is_not_contamination(self):
result = sanctioned_restart.classify_restart_command(None)
self.assertFalse(result["contamination"])
self.assertTrue(result["clean_claim_allowed"])
class TestPostRestartHealth(unittest.TestCase):
"""AC3: a clean post-restart claim needs live client-namespace proof."""
def test_live_client_probe_clears_the_session(self):
result = sanctioned_restart.verify_post_restart_health(
NAMESPACE,
probe_result={"success": True},
probe_source=mcp_namespace_health.PROBE_SOURCE_CLIENT,
registered_tools=["gitea_whoami"],
required_tool="gitea_whoami",
)
self.assertEqual(result["status"], sanctioned_restart.HEALTH_CLEAN)
self.assertTrue(result["clean_claim_allowed"])
self.assertTrue(result["mutations_allowed"])
def test_offline_probe_does_not_clear_the_session(self):
result = sanctioned_restart.verify_post_restart_health(
NAMESPACE,
probe_result={"success": True},
probe_source=mcp_namespace_health.PROBE_SOURCE_OFFLINE,
registered_tools=["gitea_whoami"],
required_tool="gitea_whoami",
)
self.assertFalse(result["clean_claim_allowed"])
self.assertFalse(result["mutations_allowed"])
def test_failed_probe_is_unhealthy(self):
result = sanctioned_restart.verify_post_restart_health(
NAMESPACE,
probe_result={"success": False, "error": "client is closing: EOF"},
probe_source=mcp_namespace_health.PROBE_SOURCE_CLIENT,
registered_tools=["gitea_whoami"],
required_tool="gitea_whoami",
)
self.assertEqual(result["status"], sanctioned_restart.HEALTH_UNHEALTHY)
self.assertFalse(result["clean_claim_allowed"])
def test_static_registration_alone_never_clears_the_session(self):
result = sanctioned_restart.verify_post_restart_health(
NAMESPACE,
registered_tools=["gitea_whoami"],
required_tool="gitea_whoami",
)
self.assertFalse(result["clean_claim_allowed"])
class TestAuditEmission(unittest.TestCase):
"""Every restart attempt is audited with actor, target, and result."""
def _run(self, principal, sink):
prior = os.environ.get(console_audit.AUDIT_LOG_ENV)
os.environ[console_audit.AUDIT_LOG_ENV] = sink
try:
return sanctioned_restart.execute_restart(
NAMESPACE,
principal=principal,
confirmation=f"restart {NAMESPACE}",
env=READY_ENV,
request_id="req-642",
)
finally:
if prior is None:
os.environ.pop(console_audit.AUDIT_LOG_ENV, None)
else:
os.environ[console_audit.AUDIT_LOG_ENV] = prior
def test_allowed_attempt_is_written_with_actor_and_target(self):
with tempfile.TemporaryDirectory() as tmp:
sink = os.path.join(tmp, "audit.jsonl")
result = self._run(admin(), sink)
self.assertTrue(result["audit"]["written"])
with open(sink, encoding="utf-8") as handle:
record = json.loads(handle.read().strip())
self.assertEqual(
record["action"], sanctioned_restart.ACTION_RESTART_NAMESPACE
)
self.assertEqual(record["target"]["namespace"], NAMESPACE)
self.assertEqual(record["target"]["mode"], "restart")
self.assertEqual(record["result"], console_audit.RESULT_ALLOWED)
self.assertEqual(record["actor"]["subject"], "[email protected]")
self.assertFalse(record["metadata"]["process_kill_executed"])
def test_denied_attempt_is_audited_too(self):
with tempfile.TemporaryDirectory() as tmp:
sink = os.path.join(tmp, "audit.jsonl")
self._run(viewer(), sink)
with open(sink, encoding="utf-8") as handle:
record = json.loads(handle.read().strip())
self.assertEqual(record["result"], console_audit.RESULT_DENIED)
self.assertEqual(
record["reason_code"], sanctioned_restart.DENY_UNAUTHORIZED
)
def test_restart_audit_uses_break_glass_retention(self):
action = console_authz.get_action(
sanctioned_restart.ACTION_RESTART_NAMESPACE
)
self.assertEqual(
console_audit.retention_class_for(action),
console_audit.RETENTION_BREAK_GLASS,
)
class TestRegistrySurface(unittest.TestCase):
"""AC5: the console surfaces the control, still disabled, with no kill."""
def test_registry_exposes_both_actions_disabled(self):
registry = gated_actions.load_action_registry()
for action_id in (
sanctioned_restart.ACTION_RESTART_NAMESPACE,
sanctioned_restart.ACTION_RELOAD_NAMESPACE,
):
with self.subTest(action=action_id):
action = registry.get(action_id)
self.assertIsNotNone(action)
self.assertFalse(action.enabled)
def test_registry_preview_names_the_namespace_target(self):
preview = gated_actions.preview_action(
sanctioned_restart.ACTION_RESTART_NAMESPACE, namespace=NAMESPACE
)
target = preview["mutation_ledger"][0]["target"]
self.assertIn(NAMESPACE, target)
self.assertFalse(preview["enabled"])
def test_registry_attempt_fails_closed(self):
result = gated_actions.attempt_action(
sanctioned_restart.ACTION_RESTART_NAMESPACE, namespace=NAMESPACE
)
self.assertFalse(result["success"])
if __name__ == "__main__":
unittest.main()
+2 -24
View File
@@ -134,35 +134,13 @@ class TestClassification(unittest.TestCase):
self.assertEqual(cls, wca.CLASS_ACTIVE_OPEN_PR)
self.assertFalse(wca.is_removable(cls))
def test_stale_clean_issue_worktree_needs_merged_pr_proof(self):
# Scenario 5 (#858): age is not proof that the branch landed, so a
# TTL-expired issue worktree stays active work. Only authoritative
# merged-PR evidence makes it removable, which is what keeps a
# worktree holding unmerged commits from being reclaimed by age.
def test_stale_clean_issue_worktree_removable(self):
# Scenario 5: clean issue worktree, TTL expired, no lock -> removable.
cls = wca.classify_worktree(
workflow_type=wca.WORKFLOW_ISSUE_WORK,
is_dirty=False,
ttl_expired=True,
)
self.assertEqual(cls, wca.CLASS_ACTIVE_ISSUE_WORK)
self.assertFalse(wca.is_removable(cls))
cls = wca.classify_worktree(
workflow_type=wca.WORKFLOW_ISSUE_WORK,
is_dirty=False,
ttl_expired=True,
merged_pr_cleanup={"proven": True},
)
self.assertEqual(cls, wca.CLASS_CLEAN_STALE_REMOVABLE)
self.assertTrue(wca.is_removable(cls))
def test_stale_clean_conflict_fix_worktree_removable(self):
# conflict_fix keeps the original TTL rule; #858 changed issue work only.
cls = wca.classify_worktree(
workflow_type=wca.WORKFLOW_CONFLICT_FIX,
is_dirty=False,
ttl_expired=True,
)
self.assertEqual(cls, wca.CLASS_CLEAN_STALE_REMOVABLE)
self.assertTrue(wca.is_removable(cls))
-28
View File
@@ -236,34 +236,6 @@ _ACTION_SPECS: tuple[ConsoleAction, ...] = (
phase=3,
summary="Remove a remote feature branch.",
),
# #642: sanctioned daemon lifecycle. These exist so operators have an
# audited path off `pkill -f mcp_server.py` (#630). Restart drops every
# in-flight request on a namespace, so it carries the same dual-control and
# break-glass weight as a merge; reload drains first and is privileged but
# not destructive. Neither ever exposes a raw kill: execution is handed to
# a host supervisor by ``webui.sanctioned_restart``.
ConsoleAction(
action_id="system.reload_namespace",
task_key="reload_namespace",
action_class=CLASS_PRIVILEGED,
minimum_role=CONTROLLER,
requires_confirmation=True,
dual_control=False,
break_glass=False,
phase=2,
summary="Gracefully reload one MCP namespace via the host supervisor.",
),
ConsoleAction(
action_id="system.restart_namespace",
task_key="restart_namespace",
action_class=CLASS_DESTRUCTIVE,
minimum_role=ADMIN,
requires_confirmation=True,
dual_control=True,
break_glass=True,
phase=2,
summary="Restart one MCP namespace via the host supervisor.",
),
)
ACTIONS: dict[str, ConsoleAction] = {a.action_id: a for a in _ACTION_SPECS}
-13
View File
@@ -110,8 +110,6 @@ def _format_target(action_id: str, params: dict[str, Any]) -> str:
)
if action_id == "create_issue":
return f"issue {params.get('title', '?')!r}"
if action_id in {"system.restart_namespace", "system.reload_namespace"}:
return f"MCP namespace {params.get('namespace', '?')!r}"
return "unspecified"
@@ -167,17 +165,6 @@ def build_action_registry() -> ActionRegistry:
"gitea_create_issue_comment", "Post a PR review thread comment."),
("close_pr", "Close PR", "close_pr", "gitea_edit_pr",
"Close a pull request without merge."),
# #642: the sanctioned replacement for the forbidden manual daemon-kill
# recovery path (#630). The "tool" is a host supervisor hook, not an MCP
# call — the console never signals a process. Preview and gating live in
# ``webui.sanctioned_restart``; these stay disabled like every other
# registry entry.
("system.reload_namespace", "Reload MCP namespace", "reload_namespace",
"host.supervisor_reload",
"Gracefully reload one MCP namespace via the host supervisor."),
("system.restart_namespace", "Restart MCP namespace",
"restart_namespace", "host.supervisor_restart",
"Restart one MCP namespace via the host supervisor."),
)
actions = tuple(
GatedAction(
-613
View File
@@ -1,613 +0,0 @@
"""Sanctioned MCP restart and graceful reload controls (#642, Phase 2).
Sessions have historically recovered MCP connectivity by killing the host
daemon (``pkill -f mcp_server.py``, #630). That path stays forbidden: it kills
every namespace on the host, contaminates the surviving session, and leaves no
audit trail. This module is the sanctioned replacement.
A restart is modelled as a *gated action*, never as a command:
1. **Capability** — the console action resolves through ``console_authz``
against ``task_capability_map``, so the console cannot invent an authority
the MCP layer does not already define.
2. **Preview** — :func:`build_restart_preview` renders a mutation ledger and
the exact confirmation phrase. It never returns a shell command.
3. **Confirmation** — the operator echoes a phrase naming the exact namespace
and mode. A phrase for one namespace never authorizes another.
4. **Operator authorization** — host daemon maintenance is authorized out of
band through the environment (#630, and #710 finding F1: a worker session
cannot set an env var for an already-running daemon, so this cannot be
self-asserted the way a tool argument could).
5. **Execution** — :func:`execute_restart` never spawns a process. Once every
gate passes it hands the request to the configured host-managed restart
hook; with no hook configured it fails closed.
6. **Health recheck** — :func:`verify_post_restart_health` requires live
client-namespace probe evidence before any post-restart clean claim.
Manual ``pkill`` remains forbidden and is classified as contamination by
:func:`classify_restart_command`, which blocks clean claims (#630 AC3).
This module performs no I/O beyond reading its own environment configuration,
imports no MCP client, and holds no credential.
"""
from __future__ import annotations
import os
from dataclasses import asdict, dataclass
from typing import Any
import mcp_namespace_health
import runtime_recovery_guard
from webui import console_audit, console_authz
# --- Operations -------------------------------------------------------------
MODE_RESTART = "restart"
MODE_RELOAD = "reload"
MODES: tuple[str, ...] = (MODE_RESTART, MODE_RELOAD)
ACTION_RESTART_NAMESPACE = "system.restart_namespace"
ACTION_RELOAD_NAMESPACE = "system.reload_namespace"
ACTION_FOR_MODE: dict[str, str] = {
MODE_RESTART: ACTION_RESTART_NAMESPACE,
MODE_RELOAD: ACTION_RELOAD_NAMESPACE,
}
# Namespaces the console may target. An unlisted name fails closed rather than
# being passed through to a host hook.
KNOWN_NAMESPACES: tuple[str, ...] = tuple(
sorted(
set(mcp_namespace_health.DEFAULT_NAMESPACES)
| {"gitea-author", "gitea-reviewer", "gitea-merger",
"gitea-reconciler", "gitea-controller"}
)
)
# Scope tokens that would mean "everything at once". Explicit non-goal: the
# console never offers a fleet-wide restart, because that is the blast radius
# `pkill -f mcp_server.py` already had.
_FLEET_TOKENS = frozenset({"*", "all", "fleet", "any", ""})
# --- Environment configuration ----------------------------------------------
# Read server-side only; the value is an opaque host hook reference (e.g. a
# launchd label), never a command line, and is never rendered to a client.
RESTART_HOOK_ENV = "GITEA_SANCTIONED_RESTART_HOOK"
# --- Reason codes -----------------------------------------------------------
DENY_UNKNOWN_MODE = "unknown_mode"
DENY_UNKNOWN_NAMESPACE = "unknown_namespace"
DENY_FLEET_SCOPE = "fleet_scope_not_permitted"
DENY_UNAUTHORIZED = "unauthorized"
DENY_CONFIRMATION_MISSING = "confirmation_required"
DENY_CONFIRMATION_MISMATCH = "confirmation_mismatch"
DENY_OPERATOR_AUTHORIZATION = "operator_authorization_missing"
DENY_HOOK_NOT_CONFIGURED = "restart_hook_not_configured"
DENY_CONTAMINATED_RUNTIME = "contaminated_runtime"
ALLOW_HOST_ACTION_REQUIRED = "host_action_required"
# Post-restart verification outcomes.
HEALTH_CLEAN = "clean"
HEALTH_UNPROVEN = "unproven"
HEALTH_UNHEALTHY = "unhealthy"
def _clean(value: Any) -> str:
return str(value or "").strip()
# --- Mutation ledger --------------------------------------------------------
@dataclass(frozen=True)
class RestartLedgerEntry:
"""One planned step, shown before anything is asked of the host."""
sequence: int
step: str
summary: str
executes_process_kill: bool = False
def _mutation_ledger(namespace: str, mode: str) -> tuple[RestartLedgerEntry, ...]:
if mode == MODE_RELOAD:
middle = RestartLedgerEntry(
sequence=2,
step="host_graceful_reload",
summary=(
f"Ask the configured host supervisor to reload {namespace} "
"in place, draining in-flight requests. The console does not "
"signal the process itself."
),
)
else:
middle = RestartLedgerEntry(
sequence=2,
step="host_restart_hook",
summary=(
f"Ask the configured host supervisor to restart {namespace}. "
"The console never sends a signal and never runs a kill."
),
)
return (
RestartLedgerEntry(
sequence=1,
step="quiesce",
summary=(
f"Stop admitting new gated mutations for {namespace} and "
"record the intent before anything restarts."
),
),
middle,
RestartLedgerEntry(
sequence=3,
step="health_recheck",
summary=(
f"Re-probe {namespace} through the live client namespace and "
"prove the required tool is callable again."
),
),
RestartLedgerEntry(
sequence=4,
step="audit",
summary=(
"Append actor, target namespace, mode, and result to the "
"console audit log."
),
),
)
# --- Confirmation -----------------------------------------------------------
def confirmation_phrase(namespace: str, mode: str) -> str:
"""Exact phrase an operator must echo, naming the namespace and mode.
Binding the namespace into the phrase is the point: a confirmation typed
for ``gitea-author`` cannot be replayed against ``gitea-merger``.
"""
return f"{_clean(mode)} {_clean(namespace)}"
def confirmation_matches(
namespace: str, mode: str, confirmation: str | None
) -> bool:
"""Compare *confirmation* to the required phrase (exact, whitespace-trimmed)."""
return _clean(confirmation) == confirmation_phrase(namespace, mode)
# --- Scope validation -------------------------------------------------------
def _validate_scope(namespace: str, mode: str) -> tuple[str, str] | None:
"""Return ``(reason_code, detail)`` when the scope is refused."""
ns = _clean(namespace)
md = _clean(mode)
if md not in MODES:
return (
DENY_UNKNOWN_MODE,
f"Mode {md!r} is not one of {', '.join(MODES)}.",
)
if ns.lower() in _FLEET_TOKENS:
return (
DENY_FLEET_SCOPE,
(
"Fleet-wide restart is an explicit non-goal: it reproduces the "
"blast radius of `pkill -f mcp_server.py` (#630). Restart one "
"namespace at a time."
),
)
if ns not in KNOWN_NAMESPACES:
return (
DENY_UNKNOWN_NAMESPACE,
f"Namespace {ns!r} is not a known MCP namespace.",
)
return None
# --- Host hook --------------------------------------------------------------
def restart_hook(env: dict[str, str] | None = None) -> dict[str, Any]:
"""Report the configured host-managed restart hook.
The hook is a reference the *host* resolves (a supervisor label), not a
command this process runs. ``configured=False`` fails restart closed.
"""
source = env if env is not None else os.environ
reference = _clean(source.get(RESTART_HOOK_ENV))
return {
"configured": bool(reference),
"reference": reference or None,
"source": RESTART_HOOK_ENV if reference else None,
"self_assertable": False,
"console_executes_process": False,
}
# --- Preview ----------------------------------------------------------------
def build_restart_preview(
namespace: str,
mode: str = MODE_RESTART,
*,
principal: console_authz.Principal | None = None,
env: dict[str, str] | None = None,
) -> dict[str, Any]:
"""Render the dry-run preview for a restart/reload request.
Read-only: no authorization is granted, no host is contacted, and the
result never contains a shell command.
"""
ns = _clean(namespace)
md = _clean(mode)
action_id = ACTION_FOR_MODE.get(md, ACTION_RESTART_NAMESPACE)
action = console_authz.get_action(action_id)
decision = console_authz.authorize(action_id, principal)
scope_error = _validate_scope(ns, md)
hook = restart_hook(env)
operator = runtime_recovery_guard.operator_authorization(env)
return {
"action_id": action_id,
"namespace": ns,
"mode": md,
"scope_valid": scope_error is None,
"scope_reason_code": scope_error[0] if scope_error else None,
"scope_detail": scope_error[1] if scope_error else None,
"required_role": action.minimum_role if action else None,
"required_permission": action.mcp_permission if action else None,
"action_class": action.action_class if action else None,
"dual_control": action.dual_control if action else True,
"break_glass": action.break_glass if action else True,
"requires_confirmation": True,
"confirmation_phrase": confirmation_phrase(ns, md),
"mutation_ledger": [asdict(entry) for entry in _mutation_ledger(ns, md)],
"authorization": decision.to_dict(),
"operator_authorization": operator,
"restart_hook": hook,
"execution_enabled": False,
"raw_process_kill_exposed": False,
"known_namespaces": list(KNOWN_NAMESPACES),
"post_restart_verification_required": True,
}
# --- Gate -------------------------------------------------------------------
def assess_restart_request(
namespace: str,
mode: str = MODE_RESTART,
*,
principal: console_authz.Principal | None = None,
confirmation: str | None = None,
contamination_marker: dict[str, Any] | None = None,
env: dict[str, str] | None = None,
) -> dict[str, Any]:
"""Decide whether a restart request may proceed to the host hook.
Every gate must pass. The first failure wins and is reported with a stable
reason code; a pass never means "restarted", only "may be handed to the
configured host hook".
"""
ns = _clean(namespace)
md = _clean(mode)
action_id = ACTION_FOR_MODE.get(md, ACTION_RESTART_NAMESPACE)
preview = build_restart_preview(ns, md, principal=principal, env=env)
def refuse(reason_code: str, detail: str) -> dict[str, Any]:
return {
"allowed": False,
"gates_passed": False,
"reason_code": reason_code,
"detail": detail,
"action_id": action_id,
"namespace": ns,
"mode": md,
"preview": preview,
"execution_enabled": False,
}
scope_error = _validate_scope(ns, md)
if scope_error is not None:
return refuse(*scope_error)
# Authority is checked as an authorization decision, not an execution
# grant. ``for_execution=True`` asks "may the console perform this write?",
# and the answer here is permanently no: step 2 of the ledger is a request
# to the host supervisor, so the console's Phase 2 execution gate is not
# the relevant gate. Every branch below keeps ``execution_enabled`` False
# and :func:`execute_restart` never touches a process.
decision = console_authz.authorize(action_id, principal)
if not decision.allowed:
return refuse(DENY_UNAUTHORIZED, decision.detail)
if not _clean(confirmation):
return refuse(
DENY_CONFIRMATION_MISSING,
(
"Type the confirmation phrase "
f"{preview['confirmation_phrase']!r} to proceed."
),
)
if not confirmation_matches(ns, md, confirmation):
return refuse(
DENY_CONFIRMATION_MISMATCH,
(
"Confirmation does not name this namespace and mode; expected "
f"{preview['confirmation_phrase']!r}."
),
)
operator = preview["operator_authorization"]
if not operator["authorized"]:
return refuse(
DENY_OPERATOR_AUTHORIZATION,
(
"Host daemon maintenance requires out-of-band operator "
"authorization via "
f"{runtime_recovery_guard.OPERATOR_AUTHORIZATION_ENV}."
),
)
# #630's task-scoped gate deliberately lets a contaminated worker keep
# commenting and handing off. Restart is stricter and unconditional: a
# runtime already contaminated by a manual kill must be reconciled before
# it is restarted again, or the restart just launders the contamination.
if contamination_marker and not contamination_marker.get(
"cleared_by_reconciler"
):
return refuse(
DENY_CONTAMINATED_RUNTIME,
(
"A live contamination marker is present; clear it through the "
"reconciler path before restarting."
),
)
hook = preview["restart_hook"]
if not hook["configured"]:
return refuse(
DENY_HOOK_NOT_CONFIGURED,
(
"No host-managed restart hook is configured "
f"({RESTART_HOOK_ENV}). The console will not fall back to a "
"process kill."
),
)
return {
"allowed": True,
"gates_passed": True,
"reason_code": ALLOW_HOST_ACTION_REQUIRED,
"detail": (
"Every gate passed. The restart must be performed by the "
"configured host supervisor; the console does not signal the "
"process."
),
"action_id": action_id,
"namespace": ns,
"mode": md,
"preview": preview,
"execution_enabled": False,
"console_executes": False,
"console_active_phase": console_authz.ACTIVE_PHASE,
}
# --- Execution --------------------------------------------------------------
def execute_restart(
namespace: str,
mode: str = MODE_RESTART,
*,
principal: console_authz.Principal | None = None,
confirmation: str | None = None,
contamination_marker: dict[str, Any] | None = None,
env: dict[str, str] | None = None,
request_id: str | None = None,
session_id: str | None = None,
) -> dict[str, Any]:
"""Run every gate, audit the outcome, and hand off to the host.
This function never spawns a process, never sends a signal, and never
builds a command line. ``success`` is False in both directions: a refused
request is refused, and an authorized request still requires the host
supervisor to act.
"""
assessment = assess_restart_request(
namespace,
mode,
principal=principal,
confirmation=confirmation,
contamination_marker=contamination_marker,
env=env,
)
action_id = assessment["action_id"]
allowed = assessment["allowed"]
audit = console_audit.record_event(
action_id=action_id,
result=(
console_audit.RESULT_ALLOWED if allowed
else console_audit.RESULT_DENIED
),
principal=principal,
target={"namespace": assessment["namespace"], "mode": assessment["mode"]},
reason_code=assessment["reason_code"],
detail=assessment["detail"],
request_id=request_id,
session_id=session_id,
metadata={
"gates_passed": assessment["gates_passed"],
"process_kill_executed": False,
"post_restart_verification_required": True,
},
)
return {
"success": False,
"outcome": (
ALLOW_HOST_ACTION_REQUIRED if allowed else assessment["reason_code"]
),
"allowed": allowed,
"detail": assessment["detail"],
"namespace": assessment["namespace"],
"mode": assessment["mode"],
"action_id": action_id,
"process_kill_executed": False,
"host_hook": assessment["preview"]["restart_hook"],
"next_action": (
"Have the host supervisor perform the restart, then call "
"verify_post_restart_health with live client-namespace evidence "
"before claiming a clean session."
if allowed
else assessment["detail"]
),
"assessment": assessment,
"audit": audit,
}
# --- Contamination classification -------------------------------------------
def classify_restart_command(
command: str | None,
*,
mcp_pids: list[Any] | tuple[Any, ...] | None = None,
session_id: str | None = None,
remote: str | None = None,
role: str | None = None,
) -> dict[str, Any]:
"""Classify an operator-proposed recovery command (#630 AC2).
A manual ``pkill``/``kill`` of the MCP daemon is contamination, not a
restart. When contaminating, a durable marker is returned so downstream
gated mutations and clean claims fail closed.
"""
classification = runtime_recovery_guard.classify_recovery_command(
command, mcp_pids=mcp_pids
)
contaminating = bool(classification.get("contamination"))
marker = None
if contaminating:
marker = runtime_recovery_guard.build_contamination_record(
reason_class=(
classification.get("reason_class")
or runtime_recovery_guard.REASON_MANUAL_DAEMON_KILL
),
command_redacted=classification.get("redacted_command"),
session_id=session_id,
remote=remote,
role=role,
detail=(
"Manual daemon kill is forbidden; use the sanctioned "
f"{ACTION_RESTART_NAMESPACE} gated action instead."
),
)
return {
"contamination": contaminating,
"sanctioned": not contaminating and not classification.get("process_kill"),
"clean_claim_allowed": not contaminating,
"reason_class": classification.get("reason_class"),
"redacted_command": classification.get("redacted_command"),
"classification": classification,
"contamination_marker": marker,
"sanctioned_alternative": ACTION_RESTART_NAMESPACE,
}
# --- Post-restart health verification ---------------------------------------
def verify_post_restart_health(
namespace: str,
*,
probe_result: dict[str, Any] | None = None,
probe_source: str | None = None,
registered_tools: list[str] | tuple[str, ...] | None = None,
required_tool: str | None = None,
profile: str | None = None,
) -> dict[str, Any]:
"""Require live proof a namespace is callable before any clean claim (AC3).
Static registration is not proof and neither is an offline subprocess
probe: only ``probe_source=client_namespace`` evidence can clear a
post-restart session for mutations.
"""
ns = _clean(namespace)
health = mcp_namespace_health.classify_namespace_probe(
ns,
required_tool=required_tool,
registered_tools=registered_tools,
probe_result=probe_result,
profile=profile,
probe_source=probe_source,
)
healthy = bool(health.get("healthy"))
proven = bool(health.get("ide_namespace_proven"))
if healthy and proven:
status = HEALTH_CLEAN
elif healthy:
status = HEALTH_UNPROVEN
else:
status = HEALTH_UNHEALTHY
reasons = list(health.get("reasons") or [])
if status == HEALTH_UNPROVEN:
reasons.append(
"Namespace reported healthy without live client-namespace "
"evidence; a post-restart clean claim requires "
f"probe_source={mcp_namespace_health.PROBE_SOURCE_CLIENT!r}."
)
return {
"namespace": ns,
"status": status,
"healthy": healthy,
"ide_namespace_proven": proven,
"clean_claim_allowed": status == HEALTH_CLEAN,
"mutations_allowed": status == HEALTH_CLEAN,
"reasons": reasons,
"health": health,
}
# --- Policy surface ---------------------------------------------------------
def restart_policy() -> dict[str, Any]:
"""Machine-readable description of the sanctioned restart contract."""
return {
"policy_version": 1,
"modes": list(MODES),
"actions": [ACTION_RESTART_NAMESPACE, ACTION_RELOAD_NAMESPACE],
"known_namespaces": list(KNOWN_NAMESPACES),
"fleet_scope_permitted": False,
"console_executes_process_kill": False,
"raw_kill_exposed": False,
"requires_confirmation": True,
"confirmation_binds_namespace": True,
"operator_authorization_env": (
runtime_recovery_guard.OPERATOR_AUTHORIZATION_ENV
),
"restart_hook_env": RESTART_HOOK_ENV,
"manual_kill_classified_as": runtime_recovery_guard.CONTAMINATION_KIND,
"post_restart_clean_claim_requires": (
mcp_namespace_health.PROBE_SOURCE_CLIENT
),
"audit_required": True,
"silent_auto_restart_permitted": False,
}
+4 -271
View File
@@ -34,11 +34,7 @@ import subprocess
from datetime import datetime, timezone
from typing import Any
from merged_cleanup_reconcile import (
branch_worktree_folder,
is_head_ancestor_of_ref,
read_local_worktree_state,
)
from merged_cleanup_reconcile import branch_worktree_folder, read_local_worktree_state
from reviewer_worktree import parse_dirty_tracked_files, REVIEW_WORKTREE_RE
PROTECTED_BRANCHES = frozenset({"master", "main", "dev"})
@@ -71,14 +67,6 @@ REMOVABLE_CLASSES = frozenset(
{CLASS_CLEAN_STALE_REMOVABLE, CLASS_DETACHED_REVIEW_LEFTOVER}
)
# Merged-PR linkage outcomes for issue worktrees (#858). Only ``LINKAGE_MERGED``
# is ownership proof; every other outcome leaves the worktree protected.
LINKAGE_MERGED = "merged_pr"
LINKAGE_OPEN = "open_pr"
LINKAGE_NONE = "no_owning_pr"
LINKAGE_AMBIGUOUS = "ambiguous"
LINKAGE_UNKNOWN = "unknown"
_ISSUE_REF_RE = re.compile(r"issue-(\d+)", re.IGNORECASE)
_ISSUE_BRANCH_PREFIXES = ("feat/", "fix/", "docs/", "chore/")
@@ -181,186 +169,6 @@ def is_ttl_expired(
return (now_dt - last).total_seconds() > ttl_hours * 3600.0
def build_pr_index(prs: list[dict[str, Any]] | None) -> dict[str, list[dict[str, Any]]]:
"""Index PR records by head branch for deterministic worktree linkage (#858).
Accepts Gitea PR payloads (``head`` as a dict) and pre-flattened records
(``head_branch``/``head_sha``). Records without a usable head branch or
number are dropped rather than guessed at, so a branch is only ever linked
to a PR the caller actually proved.
"""
index: dict[str, list[dict[str, Any]]] = {}
for pr in prs or []:
head = pr.get("head")
if isinstance(head, dict):
head_branch = head.get("ref")
head_sha = head.get("sha")
else:
head_branch = pr.get("head_branch") or (head if isinstance(head, str) else None)
head_sha = pr.get("head_sha")
number = pr.get("number")
if not head_branch or number is None:
continue
try:
pr_number = int(number)
except (TypeError, ValueError):
continue
index.setdefault(str(head_branch).strip(), []).append(
{
"pr_number": pr_number,
"head_branch": str(head_branch).strip(),
"head_sha": head_sha,
"merged": bool(pr.get("merged") or pr.get("merged_at")),
"state": pr.get("state"),
}
)
return index
def resolve_owning_pr(
*,
branch: str | None,
pr_index: dict[str, list[dict[str, Any]]] | None,
) -> dict[str, Any]:
"""Resolve the single PR that owns ``branch``, failing closed when unclear.
Ownership is only ``LINKAGE_MERGED`` when exactly one PR claims the branch
and that PR is merged. Several distinct PRs on one branch is a competing
claim (``LINKAGE_AMBIGUOUS``), and a still-open owner is reported as
``LINKAGE_OPEN`` — both keep the worktree protected while still exposing
the PR number the audit resolved.
"""
if pr_index is None:
return {
"status": LINKAGE_UNKNOWN,
"pr_number": None,
"candidate_pr_numbers": [],
"reasons": ["live PR state was not supplied; ownership unproven"],
}
branch_name = (branch or "").strip()
if not branch_name:
return {
"status": LINKAGE_UNKNOWN,
"pr_number": None,
"candidate_pr_numbers": [],
"reasons": ["worktree has no attached branch; ownership unproven"],
}
candidates = list(pr_index.get(branch_name) or [])
numbers = sorted({c["pr_number"] for c in candidates})
if not candidates:
return {
"status": LINKAGE_NONE,
"pr_number": None,
"candidate_pr_numbers": [],
"reasons": [f"no PR claims branch '{branch_name}'"],
}
if len(numbers) > 1:
return {
"status": LINKAGE_AMBIGUOUS,
"pr_number": None,
"candidate_pr_numbers": numbers,
"reasons": [
f"branch '{branch_name}' is claimed by competing PRs {numbers}; "
"ownership is ambiguous"
],
}
owner = candidates[0]
pr_number = owner["pr_number"]
if owner.get("head_branch") != branch_name:
return {
"status": LINKAGE_UNKNOWN,
"pr_number": pr_number,
"candidate_pr_numbers": numbers,
"reasons": [
f"PR #{pr_number} head branch '{owner.get('head_branch')}' does not "
f"match worktree branch '{branch_name}'"
],
}
if not owner.get("merged"):
return {
"status": LINKAGE_OPEN,
"pr_number": pr_number,
"candidate_pr_numbers": numbers,
"pr_head_sha": owner.get("head_sha"),
"reasons": [f"owning PR #{pr_number} is not merged"],
}
return {
"status": LINKAGE_MERGED,
"pr_number": pr_number,
"candidate_pr_numbers": numbers,
"pr_head_sha": owner.get("head_sha"),
"reasons": [],
}
def assess_merged_pr_worktree_cleanup(
*,
linkage: dict[str, Any] | None,
head_sha: str | None,
head_in_master: bool | None,
is_dirty: bool,
has_open_pr: bool,
has_active_lease: bool,
has_active_issue_lock: bool,
is_protected: bool,
has_live_session: bool = False,
) -> dict[str, Any]:
"""Decide whether a merged issue worktree satisfies the full cleanup policy.
Every condition must be independently proven: conclusive merged-PR
ownership, agreement between the worktree branch and the PR head branch,
containment of the worktree head in authoritative master (which is what
proves no unmerged commits remain), absence of any open/competing PR,
lease, issue lock, or live session, a clean tree, and a worktree that is
not the protected control checkout. Anything unknown blocks.
"""
link = linkage or {
"status": LINKAGE_UNKNOWN,
"pr_number": None,
"reasons": ["no linkage assessment supplied"],
}
status = link.get("status")
reasons: list[str] = []
if status != LINKAGE_MERGED:
reasons.extend(
link.get("reasons") or ["owning PR could not be conclusively identified"]
)
if is_protected:
reasons.append("worktree is protected or the stable control checkout")
if is_dirty:
reasons.append("worktree has uncommitted changes")
if has_open_pr:
reasons.append("worktree branch has an open PR")
if has_active_lease:
reasons.append("worktree has an active lease")
if has_active_issue_lock:
reasons.append("an active issue lock references this branch")
if has_live_session:
reasons.append("a live process or session is using this worktree")
if not head_sha:
reasons.append("worktree head sha is unknown")
if head_in_master is None:
reasons.append("containment of the worktree head in master is unknown")
elif not head_in_master:
reasons.append(
"worktree head is not contained in authoritative master "
"(unmerged commits remain)"
)
proven = not reasons
return {
"linkage_status": status,
"pr_number": link.get("pr_number"),
"pr_head_sha": link.get("pr_head_sha"),
"head_in_master": head_in_master,
"proven": proven,
"block_reasons": reasons,
}
def classify_worktree(
*,
workflow_type: str,
@@ -373,8 +181,6 @@ def classify_worktree(
ttl_expired: bool = False,
is_protected: bool = False,
metadata_known: bool = True,
merged_pr_cleanup: dict[str, Any] | None = None,
has_live_session: bool = False,
) -> str:
"""Classify a worktree, safety-first: any preservation signal wins.
@@ -393,8 +199,6 @@ def classify_worktree(
return CLASS_ACTIVE_ISSUE_WORK # never auto-deleted (criterion 8)
if has_active_issue_lock:
return CLASS_ACTIVE_ISSUE_WORK
if has_live_session:
return CLASS_ACTIVE_ISSUE_WORK # a live session still owns this tree
if not metadata_known or workflow_type == WORKFLOW_UNKNOWN:
return CLASS_UNSAFE_UNKNOWN # never auto-deleted without proof
@@ -403,15 +207,7 @@ def classify_worktree(
if is_detached or branch_gone:
return CLASS_DETACHED_REVIEW_LEFTOVER
return CLASS_CLEAN_STALE_REMOVABLE
if workflow_type == WORKFLOW_ISSUE_WORK:
# #858: an issue worktree becomes removable only on authoritative
# merged-PR evidence satisfying the whole cleanup policy. Age alone
# never proves the branch landed, so TTL cannot qualify one by itself
# — otherwise a worktree holding unmerged commits would be reclaimed.
if (merged_pr_cleanup or {}).get("proven"):
return CLASS_CLEAN_STALE_REMOVABLE
return CLASS_ACTIVE_ISSUE_WORK
# conflict_fix: only removable once the TTL has expired.
# issue_work / conflict_fix: only removable once the TTL has expired.
if ttl_expired:
return CLASS_CLEAN_STALE_REMOVABLE
return CLASS_ACTIVE_ISSUE_WORK
@@ -604,20 +400,6 @@ def remove_worktree(project_root: str, path: str) -> dict[str, Any]:
}
def head_contained_in_ref(
project_root: str, head_sha: str | None, ref: str | None
) -> bool | None:
"""Return True when ``head_sha`` is already contained in ``ref``.
Shares :mod:`merged_cleanup_reconcile`'s ancestry check so the audit and
the PR-scoped reconciler agree on what "already landed" means (#858).
Returns None when containment cannot be determined, which fails closed.
"""
if not head_sha or not ref:
return None
return is_head_ancestor_of_ref(project_root, head_sha, ref)
def _is_under_branches(project_root: str, path: str) -> bool:
branches_root = os.path.join(os.path.abspath(project_root), "branches")
return os.path.abspath(path or "").startswith(branches_root + os.sep)
@@ -631,30 +413,16 @@ def audit_branches_directory(
active_issue_branches: set[str] | None = None,
now: datetime | str | None = None,
ttl_hours: float = DEFAULT_TTL_HOURS,
pr_index: dict[str, list[dict[str, Any]]] | None = None,
leased_issue_numbers: set[int] | None = None,
live_session_paths: set[str] | None = None,
master_ref: str | None = None,
) -> dict[str, Any]:
"""Classify every session-owned worktree under ``branches/``.
Read-only: shells out to git for discovery and dirty state, then applies
the pure classifier. Returns per-worktree classifications, counts, the
list of removable candidates, and the ``git worktree list`` proof.
``pr_index`` (see :func:`build_pr_index`) supplies the authoritative PR
ownership used to link issue worktrees to their merged PR (#858).
``master_ref`` is the ref a worktree head must be contained in before it
can be considered landed. Both are optional and their absence only ever
fails closed: without them no issue worktree becomes removable.
"""
open_pr_branches = open_pr_branches or set()
leased_branches = leased_branches or set()
active_issue_branches = active_issue_branches or set()
leased_issue_numbers = leased_issue_numbers or set()
live_session_paths = {
os.path.abspath(p) for p in (live_session_paths or set()) if p
}
worktrees: list[dict[str, Any]] = []
for entry in list_worktrees(project_root):
@@ -665,42 +433,12 @@ def audit_branches_directory(
)
dirty_state = read_worktree_dirty(path)
is_dirty = bool(dirty_state.get("dirty"))
head_sha = entry.get("head")
linkage = resolve_owning_pr(branch=branch, pr_index=pr_index)
metadata = build_worktree_metadata(
path=path,
branch=branch,
head_sha=head_sha,
pr_number=linkage.get("pr_number"),
path=path, branch=branch, head_sha=entry.get("head")
)
has_open_pr = bool(branch) and branch in open_pr_branches
# A lease on issue N protects that issue's own work worktree. It must
# not incidentally protect a baseline/review scratch tree that merely
# carries the same issue marker in its name, which would change the
# classification of worktrees this policy does not own.
has_active_lease = (bool(branch) and branch in leased_branches) or (
metadata["workflow_type"] == WORKFLOW_ISSUE_WORK
and metadata.get("issue_number") is not None
and metadata["issue_number"] in leased_issue_numbers
)
has_active_lease = bool(branch) and branch in leased_branches
has_active_lock = bool(branch) and branch in active_issue_branches
has_live_session = bool(path) and os.path.abspath(path) in live_session_paths
head_in_master = (
head_contained_in_ref(project_root, head_sha, master_ref)
if master_ref
else None
)
merged_pr_cleanup = assess_merged_pr_worktree_cleanup(
linkage=linkage,
head_sha=head_sha,
head_in_master=head_in_master,
is_dirty=is_dirty,
has_open_pr=has_open_pr,
has_active_lease=has_active_lease,
has_active_issue_lock=has_active_lock,
is_protected=is_protected,
has_live_session=has_live_session,
)
ttl_expired = is_ttl_expired(
last_used_at=metadata.get("last_used_at"), now=now, ttl_hours=ttl_hours
)
@@ -714,8 +452,6 @@ def audit_branches_directory(
branch_gone=branch is None and not entry.get("detached"),
ttl_expired=ttl_expired,
is_protected=is_protected,
merged_pr_cleanup=merged_pr_cleanup,
has_live_session=has_live_session,
)
metadata["cleanup_eligibility"] = classification
worktrees.append(
@@ -727,10 +463,7 @@ def audit_branches_directory(
"has_open_pr": has_open_pr,
"has_active_lease": has_active_lease,
"has_active_issue_lock": has_active_lock,
"has_live_session": has_live_session,
"is_protected": is_protected,
"merged_pr_linkage": linkage,
"merged_pr_cleanup": merged_pr_cleanup,
"classification": classification,
"removable": is_removable(classification),
}