Merge branch 'master' into feat/issue-646-policy-guardrail-visibility

This commit is contained in:
2026-07-24 07:58:16 -04:00
47 changed files with 14712 additions and 190 deletions
+724
View File
@@ -1266,6 +1266,730 @@ class TestSecondRemediationIntegration(unittest.TestCase):
self.assertIn("delete_acknowledged", delete_actions[0])
self.assertTrue(delete_actions[0].get("verified_absent"))
def test_issue_851_worktree_removed_when_remote_blocked_only_by_worktree_binding(self):
"""#851: remote blocked by worktree_binding must not skip safe worktree removal.
Lifecycle: remove clean owned worktree → reassess ownership → delete
remote only if independently safe. Unrelated entries stay untouched.
"""
from mcp_server import gitea_reconcile_merged_cleanups
target_branch = "fix/issue-844-exclude-epic-containers"
foreign_branch = "fix/issue-999-unrelated-active"
worktree_path = "/tmp/branches/fix-issue-844-exclude-epic-containers"
ownership_calls = []
remove_calls = []
delete_api_calls = []
def fake_collect(**kwargs):
ownership_calls.append(dict(kwargs))
# Ownership is reassessed *after* independent worktree removal (#851).
# Target worktree is already gone → no worktree_binding remains.
# Foreign branch keeps an active author lease → remote delete blocked.
if kwargs.get("branch") == foreign_branch:
# Match session-bound org/repo + host used by the tool resolve path.
return {
"records": [
{
"category": guard.OWNERSHIP_CATEGORY_AUTHOR_LEASE,
"status": "active",
"remote": kwargs.get("remote") or "prgs",
"host": kwargs.get("host") or "gitea.example.com",
"org": kwargs.get("org") or "Scaled-Tech-Consulting",
"repo": kwargs.get("repo") or "Gitea-Tools",
"branch": foreign_branch,
"reclaim_allowed": False,
}
],
"inventory_error": False,
}
return {"records": [], "inventory_error": False}
def fake_remove(project_root, branch, worktree_path=None):
remove_calls.append(
{"branch": branch, "worktree_path": worktree_path}
)
return {
"success": True,
"performed": True,
"message": f"removed worktree {worktree_path}",
"worktree_path": worktree_path,
}
def fake_probe(h, o, r, auth, br):
return guard.classify_branch_readback_http_status(
404, not_found_scope=guard.NOT_FOUND_SCOPE_BRANCH
)
def fake_api(method, url, auth, **kwargs):
if method == "DELETE":
delete_api_calls.append(url)
return {}
report = {
"entries": [
{
"pr_number": 848,
"head_branch": target_branch,
"remote_branch": {"safe_to_delete_remote": True},
"local_worktree": {
"safe_to_remove_worktree": True,
"worktree_path": worktree_path,
},
},
{
"pr_number": 999,
"head_branch": foreign_branch,
"remote_branch": {"safe_to_delete_remote": True},
"local_worktree": {
"safe_to_remove_worktree": False,
"worktree_path": None,
},
},
],
"reviewer_scratch_entries": [],
}
patch(
"mcp_server.get_profile",
return_value={
"profile_name": "prgs-reconciler",
"role": "reconciler",
"allowed_operations": [
"gitea.read",
"gitea.branch.delete",
"gitea.pr.close",
],
"forbidden_operations": [],
},
).start()
patch("mcp_server.api_get_all", return_value=[]).start()
patch(
"mcp_server.merged_cleanup_reconcile.build_reconciliation_report",
return_value=report,
).start()
patch(
"mcp_server.merged_cleanup_reconcile.discover_reviewer_scratch_worktrees",
return_value=[],
).start()
patch(
"mcp_server.audit_reconciliation_mode.check_cleanup_execution_allowed",
return_value=(True, []),
).start()
patch("mcp_server.verify_preflight_purity", return_value=None).start()
patch(
"mcp_server._collect_branch_ownership_records",
side_effect=fake_collect,
).start()
patch("mcp_server._probe_remote_branch", side_effect=fake_probe).start()
patch(
"mcp_server.merged_cleanup_reconcile.remove_local_worktree",
side_effect=fake_remove,
).start()
self.mock_api.side_effect = fake_api
res = gitea_reconcile_merged_cleanups(
dry_run=False,
execute_confirmed=True,
remote="prgs",
)
self.assertTrue(res.get("performed") or res.get("executed"))
actions = res.get("actions") or []
remove_actions = [
a for a in actions if a.get("action") == "remove_local_worktree"
]
self.assertEqual(len(remove_actions), 1, actions)
self.assertTrue(remove_actions[0].get("success"))
self.assertEqual(remove_calls[0]["branch"], target_branch)
self.assertEqual(remove_calls[0]["worktree_path"], worktree_path)
# Target remote delete succeeds after worktree removal + reassessment.
target_deletes = [
a
for a in actions
if a.get("action") == "delete_remote_branch"
and a.get("branch") == target_branch
]
self.assertEqual(len(target_deletes), 1, actions)
self.assertTrue(target_deletes[0].get("success"))
self.assertTrue(target_deletes[0].get("after_worktree_removal"))
self.assertTrue(target_deletes[0].get("ownership_reassessed"))
self.assertTrue(target_deletes[0].get("verified_absent"))
# Foreign branch remains protected (author lease) and is not deleted.
foreign_deletes = [
a
for a in actions
if a.get("action") == "delete_remote_branch"
and a.get("branch") == foreign_branch
]
self.assertEqual(len(foreign_deletes), 1, actions)
self.assertFalse(foreign_deletes[0].get("success"))
self.assertEqual(
foreign_deletes[0].get("blocker_kind"), "active_branch_ownership"
)
self.assertIn(
guard.OWNERSHIP_CATEGORY_AUTHOR_LEASE,
foreign_deletes[0].get("blocking_categories") or [],
)
# Only the target branch should hit the DELETE API.
self.assertEqual(len(delete_api_calls), 1)
# Ownership collected for target (post-removal) and foreign; worktree
# removal happened before target remote delete in the action log.
target_idx = next(
i
for i, a in enumerate(actions)
if a.get("action") == "remove_local_worktree"
)
delete_idx = next(
i
for i, a in enumerate(actions)
if a.get("action") == "delete_remote_branch"
and a.get("branch") == target_branch
and a.get("success")
)
self.assertLess(target_idx, delete_idx)
def test_issue_851_dirty_worktree_not_removed_and_remote_stays_protected(self):
"""#851: dirty/foreign worktrees remain protected; no unsafe cleanup."""
from mcp_server import gitea_reconcile_merged_cleanups
branch = "fix/issue-851-dirty"
remove_calls = []
def fake_collect(**kwargs):
return {
"records": [
{
"category": guard.OWNERSHIP_CATEGORY_WORKTREE_BINDING,
"status": "active",
"remote": kwargs.get("remote") or "prgs",
"host": kwargs.get("host") or "gitea.example.com",
"org": kwargs.get("org") or "Scaled-Tech-Consulting",
"repo": kwargs.get("repo") or "Gitea-Tools",
"branch": branch,
"reclaim_allowed": False,
}
],
"inventory_error": False,
}
report = {
"entries": [
{
"pr_number": 851,
"head_branch": branch,
"remote_branch": {"safe_to_delete_remote": True},
"local_worktree": {
"safe_to_remove_worktree": False,
"worktree_path": "/tmp/dirty-wt",
},
}
],
"reviewer_scratch_entries": [],
}
patch(
"mcp_server.get_profile",
return_value={
"profile_name": "prgs-reconciler",
"role": "reconciler",
"allowed_operations": [
"gitea.read",
"gitea.branch.delete",
],
"forbidden_operations": [],
},
).start()
patch("mcp_server.api_get_all", return_value=[]).start()
patch(
"mcp_server.merged_cleanup_reconcile.build_reconciliation_report",
return_value=report,
).start()
patch(
"mcp_server.merged_cleanup_reconcile.discover_reviewer_scratch_worktrees",
return_value=[],
).start()
patch(
"mcp_server.audit_reconciliation_mode.check_cleanup_execution_allowed",
return_value=(True, []),
).start()
patch("mcp_server.verify_preflight_purity", return_value=None).start()
patch(
"mcp_server._collect_branch_ownership_records",
side_effect=fake_collect,
).start()
patch(
"mcp_server.merged_cleanup_reconcile.remove_local_worktree",
side_effect=lambda *a, **k: remove_calls.append(k) or {
"success": True,
"performed": True,
},
).start()
self.mock_api.side_effect = lambda *a, **k: {}
res = gitea_reconcile_merged_cleanups(
dry_run=False,
execute_confirmed=True,
remote="prgs",
)
actions = res.get("actions") or []
self.assertEqual(remove_calls, [])
self.assertFalse(
any(a.get("action") == "remove_local_worktree" for a in actions)
)
deletes = [
a for a in actions if a.get("action") == "delete_remote_branch"
]
self.assertEqual(len(deletes), 1)
self.assertFalse(deletes[0].get("success"))
self.assertEqual(deletes[0].get("blocker_kind"), "active_branch_ownership")
self.assertIn(
guard.OWNERSHIP_CATEGORY_WORKTREE_BINDING,
deletes[0].get("blocking_categories") or [],
)
def test_issue_851_idempotent_resume_when_worktree_already_absent(self):
"""#851: partial failures remain resumable and idempotent."""
from mcp_server import gitea_reconcile_merged_cleanups
branch = "fix/issue-851-resume"
ownership_calls = []
def fake_collect(**kwargs):
ownership_calls.append(kwargs)
return {"records": [], "inventory_error": False}
def fake_remove(project_root, branch, worktree_path=None):
return {
"success": False,
"performed": False,
"message": f"worktree not found: {worktree_path}",
}
def fake_probe(h, o, r, auth, br):
return guard.classify_branch_readback_http_status(
404, not_found_scope=guard.NOT_FOUND_SCOPE_BRANCH
)
report = {
"entries": [
{
"pr_number": 851,
"head_branch": branch,
"remote_branch": {"safe_to_delete_remote": True},
"local_worktree": {
"safe_to_remove_worktree": True,
"worktree_path": "/tmp/already-gone",
},
}
],
"reviewer_scratch_entries": [],
}
patch(
"mcp_server.get_profile",
return_value={
"profile_name": "prgs-reconciler",
"role": "reconciler",
"allowed_operations": [
"gitea.read",
"gitea.branch.delete",
],
"forbidden_operations": [],
},
).start()
patch("mcp_server.api_get_all", return_value=[]).start()
patch(
"mcp_server.merged_cleanup_reconcile.build_reconciliation_report",
return_value=report,
).start()
patch(
"mcp_server.merged_cleanup_reconcile.discover_reviewer_scratch_worktrees",
return_value=[],
).start()
patch(
"mcp_server.audit_reconciliation_mode.check_cleanup_execution_allowed",
return_value=(True, []),
).start()
patch("mcp_server.verify_preflight_purity", return_value=None).start()
patch(
"mcp_server._collect_branch_ownership_records",
side_effect=fake_collect,
).start()
patch("mcp_server._probe_remote_branch", side_effect=fake_probe).start()
patch(
"mcp_server.merged_cleanup_reconcile.remove_local_worktree",
side_effect=fake_remove,
).start()
self.mock_api.side_effect = lambda *a, **k: {}
res = gitea_reconcile_merged_cleanups(
dry_run=False,
execute_confirmed=True,
remote="prgs",
)
actions = res.get("actions") or []
removes = [a for a in actions if a.get("action") == "remove_local_worktree"]
deletes = [a for a in actions if a.get("action") == "delete_remote_branch"]
self.assertEqual(len(removes), 1)
self.assertFalse(removes[0].get("success"))
self.assertEqual(len(deletes), 1)
self.assertTrue(deletes[0].get("success"))
self.assertTrue(deletes[0].get("after_worktree_removal"))
self.assertTrue(ownership_calls)
class TestIssue855ExactPrSelector(unittest.TestCase):
"""#855: exact pr_number pin for reconcile_merged_cleanups (#851 lifecycle)."""
def setUp(self):
self._remotes = patch.dict(
mcp_server.REMOTES,
{
"prgs": {
"host": "gitea.example.com",
"org": "Scaled-Tech-Consulting",
"repo": "Gitea-Tools",
}
},
)
self._remotes.start()
patch("gitea_audit.audit_enabled", return_value=False).start()
self.mock_api = patch("mcp_server.api_request").start()
self.mock_all = patch("mcp_server.api_get_all", return_value=[]).start()
patch("mcp_server.get_auth_header", return_value=FAKE_AUTH).start()
patch(
"mcp_server.merged_cleanup_reconcile.is_head_ancestor_of_ref",
return_value=True,
).start()
patch(
"mcp_server.get_profile",
return_value=dict(RECONCILER_WITH_DELETE),
).start()
patch(
"mcp_server._profile_operation_gate",
return_value=[],
).start()
patch(
"mcp_server._collect_branch_ownership_records",
return_value={"records": [], "inventory_error": False},
).start()
patch(
"mcp_server.merged_cleanup_reconcile.discover_reviewer_scratch_worktrees",
return_value=[],
).start()
patch("mcp_server.verify_preflight_purity", return_value=None).start()
patch(
"mcp_server.audit_reconciliation_mode.check_cleanup_execution_allowed",
return_value=(True, []),
).start()
def tearDown(self):
patch.stopall()
def _merged_pr(self, number, branch, sha="c" * 40):
return {
"number": number,
"title": f"PR {number}",
"body": f"Closes #{number - 4}",
"merged": True,
"merged_at": "2026-07-23T12:00:00Z",
"merge_commit_sha": "f" * 40,
"state": "closed",
"head": {"ref": branch, "sha": sha},
"base": {"ref": "master"},
}
def test_exact_pr_848_ignores_newer_852_in_batch_queue(self):
"""pr_number=848 selects only #848 even when #852 is newer/first."""
from mcp_server import gitea_reconcile_merged_cleanups
pr_848 = self._merged_pr(
848, "fix/issue-844-exclude-epic-containers", sha="c3f282ba" + "0" * 32
)
# Closed list would rank #852 first in batch mode; exact pin must ignore it.
closed_batch = [
self._merged_pr(852, "fix/issue-851-cleanup-worktree-before-remote-delete"),
pr_848,
self._merged_pr(849, "fix/issue-849-other"),
self._merged_pr(846, "fix/issue-846-other"),
self._merged_pr(845, "fix/issue-845-other"),
]
batch_fetch_calls = []
def fake_api(method, url, *args, **kwargs):
if method == "GET" and url.rstrip("/").endswith("/pulls/848"):
return dict(pr_848)
if method == "GET" and "/pulls/" in url:
raise AssertionError(f"unexpected PR fetch: {url}")
if method == "GET" and "/branches/" in url:
return {"name": "present"}
return {}
def fake_all(url, auth, limit=None):
batch_fetch_calls.append((url, limit))
if "state=open" in url:
return []
if "state=closed" in url:
# Exact mode must not use the closed batch list.
raise AssertionError(
"exact pr_number mode must not page closed PRs: " + url
)
return []
self.mock_api.side_effect = fake_api
self.mock_all.side_effect = fake_all
patch(
"mcp_server._remote_branch_exists",
return_value=True,
).start()
patch(
"mcp_server.merged_cleanup_reconcile.build_reconciliation_report",
side_effect=lambda **kwargs: {
"entries": [
{
"pr_number": int(pr["number"]),
"head_branch": (pr.get("head") or {}).get("ref"),
"issue_number": 844,
"remote_branch": {
"safe_to_delete_remote": True,
"head_branch": (pr.get("head") or {}).get("ref"),
},
"local_worktree": {
"safe_to_remove_worktree": True,
"worktree_path": (
"/tmp/branches/fix-issue-844-exclude-epic-containers"
),
},
"planned_execution_order": (
mcp_server.merged_cleanup_reconcile.plan_cleanup_execution_order(
remote_assessment={"safe_to_delete_remote": True},
local_assessment={"safe_to_remove_worktree": True},
)
),
}
for pr in kwargs.get("closed_prs") or []
if pr.get("merged_at") or pr.get("merged")
],
"reviewer_scratch_entries": [],
"merged_pr_count": len(kwargs.get("closed_prs") or []),
},
).start()
res = gitea_reconcile_merged_cleanups(
dry_run=True,
pr_number=848,
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
)
self.assertTrue(res.get("success"))
self.assertFalse(res.get("performed"))
self.assertEqual(res.get("selection_mode"), "exact_pr")
self.assertEqual(res.get("selected_pr_number"), 848)
entries = res.get("entries") or []
self.assertEqual(len(entries), 1, entries)
self.assertEqual(entries[0].get("pr_number"), 848)
self.assertEqual(
entries[0].get("head_branch"),
"fix/issue-844-exclude-epic-containers",
)
# No other PR appears in plan.
self.assertEqual(list((res.get("planned_execution_orders") or {}).keys()), ["848"])
plan = (res.get("planned_execution_orders") or {}).get("848") or []
actions = [s.get("action") for s in plan]
self.assertEqual(
actions,
[
"remove_local_worktree",
"reassess_branch_ownership",
"delete_remote_branch",
],
)
# Prove we never scanned the multi-PR closed batch.
self.assertFalse(any("state=closed" in (u or "") for u, _ in batch_fetch_calls))
# closed_batch fixture must remain unused (sanity).
self.assertEqual(closed_batch[0]["number"], 852)
def test_exact_pr_execute_only_mutates_selected_pr(self):
"""Execute with pr_number must never touch #845/#846/#849/#852."""
from mcp_server import gitea_reconcile_merged_cleanups
pr_848 = self._merged_pr(848, "fix/issue-844-exclude-epic-containers")
worktree_path = "/tmp/branches/fix-issue-844-exclude-epic-containers"
remove_calls = []
delete_api_calls = []
ownership_branches = []
def fake_api(method, url, *args, **kwargs):
if method == "GET" and url.rstrip("/").endswith("/pulls/848"):
return dict(pr_848)
if method == "DELETE":
delete_api_calls.append(url)
# Forbid foreign PR branch deletion by URL content.
for forbidden in ("845", "846", "849", "852"):
self.assertNotIn(forbidden, url)
return {}
def fake_remove(project_root, branch, worktree_path=None):
remove_calls.append({"branch": branch, "worktree_path": worktree_path})
return {
"success": True,
"performed": True,
"message": f"removed {worktree_path}",
"worktree_path": worktree_path,
}
def fake_collect(**kwargs):
ownership_branches.append(kwargs.get("branch"))
return {"records": [], "inventory_error": False}
def fake_probe(h, o, r, auth, br):
return guard.classify_branch_readback_http_status(
404, not_found_scope=guard.NOT_FOUND_SCOPE_BRANCH
)
self.mock_api.side_effect = fake_api
self.mock_all.side_effect = lambda url, auth, limit=None: []
patch("mcp_server._remote_branch_exists", return_value=True).start()
patch(
"mcp_server.merged_cleanup_reconcile.build_reconciliation_report",
return_value={
"entries": [
{
"pr_number": 848,
"head_branch": "fix/issue-844-exclude-epic-containers",
"remote_branch": {"safe_to_delete_remote": True},
"local_worktree": {
"safe_to_remove_worktree": True,
"worktree_path": worktree_path,
},
"planned_execution_order": [
{"action": "remove_local_worktree", "phase": 1},
{"action": "reassess_branch_ownership", "phase": 2},
{"action": "delete_remote_branch", "phase": 3},
],
}
],
"reviewer_scratch_entries": [
# Foreign scratch must be filtered before report execute loop;
# if present here it would still be a test failure if acted on.
],
"merged_pr_count": 1,
},
).start()
patch(
"mcp_server.merged_cleanup_reconcile.remove_local_worktree",
side_effect=fake_remove,
).start()
patch(
"mcp_server._collect_branch_ownership_records",
side_effect=fake_collect,
).start()
patch("mcp_server._probe_remote_branch", side_effect=fake_probe).start()
res = gitea_reconcile_merged_cleanups(
dry_run=False,
execute_confirmed=True,
pr_number=848,
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
)
self.assertTrue(res.get("performed") or res.get("executed"))
self.assertEqual(res.get("selection_mode"), "exact_pr")
self.assertEqual(res.get("selected_pr_number"), 848)
actions = res.get("actions") or []
pr_numbers_touched = {
a.get("pr_number") for a in actions if a.get("pr_number") is not None
}
self.assertTrue(pr_numbers_touched.issubset({None, 848}) or not pr_numbers_touched)
removes = [a for a in actions if a.get("action") == "remove_local_worktree"]
deletes = [a for a in actions if a.get("action") == "delete_remote_branch"]
self.assertEqual(len(removes), 1)
self.assertEqual(remove_calls[0]["branch"], "fix/issue-844-exclude-epic-containers")
self.assertEqual(len(deletes), 1)
self.assertTrue(deletes[0].get("success"))
self.assertTrue(deletes[0].get("after_worktree_removal"))
self.assertEqual(len(delete_api_calls), 1)
self.assertEqual(
ownership_branches, ["fix/issue-844-exclude-epic-containers"]
)
def test_exact_pr_unknown_fails_closed_without_mutation(self):
from mcp_server import gitea_reconcile_merged_cleanups
def fake_api(method, url, *args, **kwargs):
if method == "GET" and "/pulls/99999" in url:
raise RuntimeError("HTTP 404 Not Found")
raise AssertionError(f"unexpected API call {method} {url}")
self.mock_api.side_effect = fake_api
res = gitea_reconcile_merged_cleanups(
dry_run=True,
pr_number=99999,
remote="prgs",
)
self.assertFalse(res.get("success"))
self.assertFalse(res.get("performed"))
self.assertEqual(res.get("blocker_kind"), "pr_unresolvable")
self.assertIn("99999", " ".join(res.get("reasons") or []))
def test_exact_pr_not_merged_fails_closed(self):
from mcp_server import gitea_reconcile_merged_cleanups
def fake_api(method, url, *args, **kwargs):
if method == "GET" and url.rstrip("/").endswith("/pulls/900"):
return {
"number": 900,
"merged": False,
"merged_at": None,
"state": "open",
"head": {"ref": "feat/x", "sha": "a" * 40},
}
raise AssertionError(f"unexpected {method} {url}")
self.mock_api.side_effect = fake_api
res = gitea_reconcile_merged_cleanups(
dry_run=False,
execute_confirmed=True,
pr_number=900,
remote="prgs",
)
self.assertFalse(res.get("success"))
self.assertFalse(res.get("performed"))
self.assertEqual(res.get("blocker_kind"), "pr_not_merged")
def test_exact_pr_invalid_number_fails_closed(self):
from mcp_server import gitea_reconcile_merged_cleanups
res = gitea_reconcile_merged_cleanups(
dry_run=True,
pr_number=0,
remote="prgs",
)
self.assertFalse(res.get("success"))
self.assertEqual(res.get("blocker_kind"), "invalid_pr_number")
self.mock_api.assert_not_called()
def test_batch_mode_still_works_without_pr_number(self):
"""Unfiltered batch path remains backward compatible."""
from mcp_server import gitea_reconcile_merged_cleanups
self.mock_all.side_effect = lambda url, auth, limit=None: []
self.mock_api.side_effect = lambda *a, **k: {}
patch(
"mcp_server.merged_cleanup_reconcile.build_reconciliation_report",
return_value={
"entries": [],
"reviewer_scratch_entries": [],
"merged_pr_count": 0,
},
).start()
res = gitea_reconcile_merged_cleanups(dry_run=True, remote="prgs", limit=10)
self.assertTrue(res.get("success"))
self.assertEqual(res.get("selection_mode"), "batch")
self.assertIsNone(res.get("selected_pr_number"))
if __name__ == "__main__":
@@ -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()
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,244 @@
"""#855 AC4: an expired reviewer lease must not indefinitely protect an
already-merged branch when no live claimant exists.
Two layers are covered:
* ``branch_cleanup_guard.assess_expired_reviewer_lease_reclaim`` — the pure,
fail-closed reclaim decision. Every condition must be provably satisfied or
the lease keeps protecting the branch.
* ``gitea_mcp_server._collect_branch_ownership_records`` — the wiring that
supplies authoritative evidence (PR merged state, owner-process liveness,
competing ownership) to that decision, and flips an expired reviewer lease
to reclaimable only under the full policy.
All inputs are fabricated; no real repository, lease, or credential is used.
"""
import importlib
import unittest
from unittest.mock import patch
import branch_cleanup_guard
mcp_server = importlib.import_module("gitea_mcp_server")
FAKE_AUTH = "token fake"
REMOTE = "prgs"
ORG = "Scaled-Tech-Consulting"
REPO = "Gitea-Tools"
HOST = "gitea.prgs.cc"
BRANCH = "feat/issue-638-webui-app-shell-phase1"
PR_NUMBER = 818
class TestAssessExpiredReviewerLeaseReclaim(unittest.TestCase):
"""Pure fail-closed reclaim decision (#855 AC4)."""
def _call(self, **overrides):
base = dict(
role="reviewer",
status="expired",
pr_merged=True,
owner_pid_alive=False,
competing_active_claimant=False,
)
base.update(overrides)
return branch_cleanup_guard.assess_expired_reviewer_lease_reclaim(**base)
def test_full_policy_satisfied_allows_reclaim(self):
out = self._call()
self.assertTrue(out["reclaim_allowed"])
self.assertEqual(out["reasons"], [])
self.assertEqual(out["decision"], "reclaim_expired_reviewer_lease")
def test_stale_dead_process_reviewer_also_reclaimable(self):
out = self._call(status="stale_dead_process")
self.assertTrue(out["reclaim_allowed"])
def test_non_reviewer_role_never_reclaims(self):
for role in ("author", "merger", "controller", "reconciler", "unknown"):
with self.subTest(role=role):
out = self._call(role=role)
self.assertFalse(out["reclaim_allowed"])
self.assertTrue(out["reasons"])
self.assertEqual(out["decision"], "keep_protecting")
def test_active_status_never_reclaims(self):
out = self._call(status="active")
self.assertFalse(out["reclaim_allowed"])
def test_pr_not_merged_blocks_reclaim(self):
out = self._call(pr_merged=False)
self.assertFalse(out["reclaim_allowed"])
def test_pr_merged_unknown_fails_closed(self):
out = self._call(pr_merged=None)
self.assertFalse(out["reclaim_allowed"])
def test_owner_process_alive_blocks_reclaim(self):
out = self._call(owner_pid_alive=True)
self.assertFalse(out["reclaim_allowed"])
def test_owner_liveness_unknown_fails_closed(self):
out = self._call(owner_pid_alive=None)
self.assertFalse(out["reclaim_allowed"])
def test_competing_active_claimant_blocks_reclaim(self):
out = self._call(competing_active_claimant=True)
self.assertFalse(out["reclaim_allowed"])
def test_competing_claimant_unknown_fails_closed(self):
out = self._call(competing_active_claimant=None)
self.assertFalse(out["reclaim_allowed"])
def test_reasons_never_leak_secrets(self):
out = self._call(role="author")
blob = " ".join(out["reasons"]).lower()
self.assertNotIn("token", blob)
self.assertNotIn("password", blob)
class _FakeLease(dict):
pass
class TestCollectorExpiredReviewerReclaimWiring(unittest.TestCase):
"""`_collect_branch_ownership_records` supplies authoritative evidence and
flips an expired reviewer lease to reclaimable only under the full policy."""
def _run(
self,
*,
lease_role="reviewer",
lease_freshness="stale_dead_process",
owner_pid_alive=False,
pr_merged=True,
extra_leases=None,
worktree_on_branch=False,
):
lease = _FakeLease(
role=lease_role,
work_kind="pr",
work_number=PR_NUMBER,
branch=BRANCH,
status="active",
owner_pid=999999,
remote=REMOTE,
org=ORG,
repo=REPO,
host=HOST,
freshness={
"freshness": lease_freshness,
"owner_pid": 999999,
"owner_pid_alive": owner_pid_alive,
"expired_by_time": lease_freshness == "expired",
},
)
leases = [lease] + list(extra_leases or [])
pr_payload = {
"number": PR_NUMBER,
"merged": pr_merged,
"merged_at": "2026-07-23T00:00:00Z" if pr_merged else None,
"head": {"ref": BRANCH},
}
def fake_api_request(method, url, *a, **k):
if method == "GET" and f"/pulls/{PR_NUMBER}" in url:
return pr_payload
raise AssertionError(f"unexpected api_request {method} {url}")
wt_entries = []
if worktree_on_branch:
wt_entries = [{"branch": BRANCH, "path": f"/x/branches/{BRANCH}"}]
with patch.object(
mcp_server.lease_lifecycle,
"list_active_leases",
return_value={"leases": leases},
), patch.object(
mcp_server.control_plane_db, "get_db", return_value=object(), create=True
), patch.object(
mcp_server.issue_lock_store, "iter_lock_files", return_value=[]
), patch.object(
mcp_server.worktree_cleanup_audit,
"list_worktrees",
return_value=wt_entries,
), patch.object(
mcp_server, "api_get_all", return_value=[]
), patch.object(
mcp_server, "api_request", side_effect=fake_api_request
):
return mcp_server._collect_branch_ownership_records(
remote=REMOTE,
host=HOST,
org=ORG,
repo=REPO,
branch=BRANCH,
pr_number=PR_NUMBER,
project_root="/x",
auth=FAKE_AUTH,
base_api="https://gitea.prgs.cc/api/v1/repos/x/y",
)
def _reviewer_records(self, bundle):
return [
rec
for rec in bundle["records"]
if rec.get("category")
== branch_cleanup_guard.OWNERSHIP_CATEGORY_REVIEWER_LEASE
]
def test_merged_dead_uncontested_reviewer_lease_is_reclaimable(self):
bundle = self._run()
self.assertFalse(bundle["inventory_error"])
recs = self._reviewer_records(bundle)
self.assertEqual(len(recs), 1)
self.assertTrue(recs[0]["reclaim_allowed"])
# And the guard consequently does not block deletion on it.
ownership = branch_cleanup_guard.assess_active_branch_ownership(
remote=REMOTE, org=ORG, repo=REPO, branch=BRANCH, host=HOST,
records=bundle["records"],
)
self.assertFalse(ownership["block"])
def test_unmerged_pr_keeps_reviewer_lease_protective(self):
bundle = self._run(pr_merged=False)
recs = self._reviewer_records(bundle)
self.assertEqual(len(recs), 1)
self.assertFalse(recs[0]["reclaim_allowed"])
ownership = branch_cleanup_guard.assess_active_branch_ownership(
remote=REMOTE, org=ORG, repo=REPO, branch=BRANCH, host=HOST,
records=bundle["records"],
)
self.assertTrue(ownership["block"])
def test_owner_process_alive_keeps_reviewer_lease_protective(self):
bundle = self._run(owner_pid_alive=True, lease_freshness="expired")
recs = self._reviewer_records(bundle)
self.assertFalse(recs[0]["reclaim_allowed"])
def test_competing_worktree_binding_keeps_reviewer_lease_protective(self):
bundle = self._run(worktree_on_branch=True)
recs = self._reviewer_records(bundle)
self.assertFalse(recs[0]["reclaim_allowed"])
ownership = branch_cleanup_guard.assess_active_branch_ownership(
remote=REMOTE, org=ORG, repo=REPO, branch=BRANCH, host=HOST,
records=bundle["records"],
)
self.assertTrue(ownership["block"])
def test_expired_author_lease_never_reclaimed_by_reviewer_policy(self):
bundle = self._run(lease_role="author")
author_recs = [
rec
for rec in bundle["records"]
if rec.get("category")
== branch_cleanup_guard.OWNERSHIP_CATEGORY_AUTHOR_LEASE
]
self.assertEqual(len(author_recs), 1)
self.assertFalse(author_recs[0]["reclaim_allowed"])
if __name__ == "__main__":
unittest.main()
@@ -0,0 +1,551 @@
"""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()
@@ -0,0 +1,630 @@
"""Durable linked-issue lock head refresh + merge-sync dead-session recovery (#871).
``gitea_update_pr_branch_by_merge`` advances a PR's *remote* head but historically
never advanced the linked durable issue lock's recorded head. After the owning
session died the drifted lock became unrecoverable and no further synchronization
was possible (PR #866 / issue #855).
Two halves are covered:
* the write-side refresh (``issue_lock_store.assess/apply_durable_lock_head_refresh``)
that records the new synced head under compare-and-swap with read-after-write; and
* the read-side recovery relation (``issue_lock_recovery`` +
``issue_lock_worktree.read_merge_sync_provenance``) that lets a dead-session lock
whose recorded head is a merge-sync *ancestor* of the live PR head be recovered —
and nothing else.
"""
from __future__ import annotations
import os
import subprocess
import sys
import tempfile
import unittest
from datetime import datetime, timedelta, timezone
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
import issue_lock_recovery # noqa: E402
import issue_lock_store # noqa: E402
import issue_lock_worktree # noqa: E402
ISSUE = 8710
PR_NUMBER = 8711
BRANCH = f"fix/issue-{ISSUE}-durable-lock-head-refresh"
IDENTITY = "example-user"
PROFILE = "example-author"
OLD = "a" * 40
NEW1 = "b" * 40
NEW2 = "c" * 40
BASE = "d" * 40
REMOTE = "prgs"
ORG = "ExampleOrg"
REPO = "ExampleRepo"
def dead_pid() -> int:
proc = subprocess.Popen([sys.executable, "-c", "pass"])
proc.wait()
return proc.pid
def future_ts(hours: int = 4) -> str:
return (
(datetime.now(timezone.utc) + timedelta(hours=hours))
.isoformat()
.replace("+00:00", "Z")
)
def _git(cwd, *args):
return subprocess.run(
["git", "-C", cwd, *args],
capture_output=True,
text=True,
check=True,
)
def _rev(cwd, ref="HEAD") -> str:
return _git(cwd, "rev-parse", ref).stdout.strip()
def build_merge_sync_repo(tmp: str) -> dict:
"""Build a repo where a feature branch was synced by merging master in.
Returns a dict with the prior (branch) head, the synced merge-commit head,
the master tip, plus a rebase-style linear descendant and an unrelated head.
"""
_git(tmp, "init", "-q", "-b", "master")
_git(tmp, "config", "user.email", "[email protected]")
_git(tmp, "config", "user.name", "T")
Path(tmp, "base.txt").write_text("base\n")
_git(tmp, "add", "-A")
_git(tmp, "commit", "-q", "-m", "root")
# Feature branch cut from root, one commit — this is the PRIOR/recorded head.
_git(tmp, "checkout", "-q", "-b", BRANCH)
Path(tmp, "feature.txt").write_text("feature\n")
_git(tmp, "add", "-A")
_git(tmp, "commit", "-q", "-m", "feature work")
prior = _rev(tmp)
# Master advances (the base the sync will merge in).
_git(tmp, "checkout", "-q", "master")
Path(tmp, "base.txt").write_text("base\nmore\n")
_git(tmp, "add", "-A")
_git(tmp, "commit", "-q", "-m", "master advance")
master_tip = _rev(tmp)
# Sync: merge master INTO the feature branch → merge commit, first parent = prior.
_git(tmp, "checkout", "-q", BRANCH)
_git(tmp, "merge", "-q", "--no-ff", "-m", "Merge master into feature", "master")
synced = _rev(tmp)
# A plain linear descendant of prior (NOT a merge) — a rebase/extra-commit shape.
_git(tmp, "checkout", "-q", "-b", "linear-branch", prior)
Path(tmp, "extra.txt").write_text("extra\n")
_git(tmp, "add", "-A")
_git(tmp, "commit", "-q", "-m", "extra linear commit")
linear = _rev(tmp)
# An unrelated root (force-push / rewritten history shape).
unrelated_dir = tempfile.mkdtemp()
_git(unrelated_dir, "init", "-q", "-b", "x")
_git(unrelated_dir, "config", "user.email", "[email protected]")
_git(unrelated_dir, "config", "user.name", "T")
Path(unrelated_dir, "z.txt").write_text("z\n")
_git(unrelated_dir, "add", "-A")
_git(unrelated_dir, "commit", "-q", "-m", "unrelated")
unrelated = _rev(unrelated_dir)
# Leave the worktree checked out on the feature branch at the PRIOR head, as
# a dead author session that never advanced would have left it.
_git(tmp, "checkout", "-q", BRANCH)
_git(tmp, "reset", "-q", "--hard", prior)
return {
"prior": prior,
"master_tip": master_tip,
"synced": synced,
"linear": linear,
"unrelated": unrelated,
}
# ─────────────────────────── write-side refresh ───────────────────────────
class TestDurableLockHeadRefresh(unittest.TestCase):
def setUp(self):
self.lock_dir = tempfile.mkdtemp()
self.wt = tempfile.mkdtemp()
lock_data = {
"issue_number": ISSUE,
"branch_name": BRANCH,
"worktree_path": self.wt,
"remote": REMOTE,
"org": ORG,
"repo": REPO,
"claimant": {"username": IDENTITY, "profile": PROFILE},
"work_lease": {
"operation_type": issue_lock_store.AUTHOR_ISSUE_WORK_LEASE,
"issue_number": ISSUE,
"branch": BRANCH,
"worktree_path": self.wt,
"claimant": {"username": IDENTITY, "profile": PROFILE},
"expires_at": future_ts(),
},
}
issue_lock_store.bind_session_lock(lock_data, lock_dir=self.lock_dir)
def _apply(self, **over):
kw = dict(
remote=REMOTE, org=ORG, repo=REPO, issue_number=ISSUE,
branch_name=BRANCH, worktree_path=self.wt, pr_number=PR_NUMBER,
identity=IDENTITY, profile=PROFILE, current_pid=os.getpid(),
expected_old_head=OLD, new_head=NEW1, synced_at=future_ts(0),
base_head=BASE, lock_dir=self.lock_dir,
)
kw.update(over)
return issue_lock_store.apply_durable_lock_head_refresh(**kw)
def _load(self):
return issue_lock_store.load_issue_lock(
remote=REMOTE, org=ORG, repo=REPO, issue_number=ISSUE,
lock_dir=self.lock_dir,
)
def test_first_sync_updates_recorded_head(self):
"""AC1: first base sync writes the resulting head to the durable lock."""
res = self._apply()
self.assertTrue(res["refreshed"], res["reasons"])
self.assertTrue(res["read_after_write_ok"])
self.assertEqual(self._load().get("synced_pr_head"), NEW1)
def test_second_sync_after_master_advance(self):
"""AC2: a later master advance permits a second sanctioned sync."""
self.assertTrue(self._apply()["refreshed"])
res2 = self._apply(expected_old_head=NEW1, new_head=NEW2)
self.assertTrue(res2["refreshed"], res2["reasons"])
self.assertEqual(self._load().get("synced_pr_head"), NEW2)
history = self._load().get("branch_sync_history")
self.assertEqual(len(history), 2)
self.assertEqual(history[0]["last_synced_pr_head"], NEW1)
self.assertEqual(history[1]["prior_pr_head"], NEW1)
def test_cas_detects_concurrent_head_change(self):
"""AC6: CAS refuses when the recorded synced head is not the old head."""
self.assertTrue(self._apply()["refreshed"]) # recorded head now NEW1
# A second sync claiming the old head is still OLD must fail closed.
res = self._apply(expected_old_head=OLD, new_head=NEW2)
self.assertFalse(res["refreshed"])
self.assertTrue(any("CAS" in r or "concurrent" in r for r in res["reasons"]))
self.assertEqual(self._load().get("synced_pr_head"), NEW1)
def test_wrong_issue_fails_closed(self):
res = self._apply(issue_number=999999)
self.assertFalse(res["refreshed"])
def test_wrong_branch_fails_closed(self):
res = self._apply(branch_name="fix/issue-8710-wrong")
self.assertFalse(res["refreshed"])
def test_wrong_repo_fails_closed(self):
res = self._apply(repo="OtherRepo")
self.assertFalse(res["refreshed"])
def test_wrong_identity_fails_closed(self):
res = self._apply(identity="intruder")
self.assertFalse(res["refreshed"])
def test_wrong_profile_fails_closed(self):
res = self._apply(profile="prgs-reviewer")
self.assertFalse(res["refreshed"])
def test_foreign_session_fails_closed(self):
"""A refresh is not a recovery: the current process must own the lock."""
path = issue_lock_store.lock_file_path(
remote=REMOTE, org=ORG, repo=REPO, issue_number=ISSUE,
lock_dir=self.lock_dir,
)
rec = issue_lock_store.read_lock_file(path)
rec["session_pid"] = dead_pid()
rec["pid"] = rec["session_pid"]
issue_lock_store.save_lock_file(path, rec)
res = self._apply()
self.assertFalse(res["refreshed"])
self.assertTrue(any("current session" in r or "live owner" in r for r in res["reasons"]))
def test_new_equals_old_fails_closed(self):
res = self._apply(expected_old_head=OLD, new_head=OLD)
self.assertFalse(res["refreshed"])
def test_non_full_sha_fails_closed(self):
self.assertFalse(self._apply(new_head="deadbeef")["refreshed"])
self.assertFalse(self._apply(expected_old_head="xyz")["refreshed"])
def test_no_lock_fails_closed(self):
assessment = issue_lock_store.assess_durable_lock_head_refresh(
None, remote=REMOTE, org=ORG, repo=REPO, issue_number=ISSUE,
branch_name=BRANCH, worktree_path=self.wt, pr_number=PR_NUMBER,
identity=IDENTITY, profile=PROFILE, current_pid=os.getpid(),
expected_old_head=OLD, new_head=NEW1,
)
self.assertFalse(assessment["allowed"])
# ─────────────────────── merge-sync provenance (real git) ───────────────────
class TestMergeSyncProvenanceObservation(unittest.TestCase):
def setUp(self):
self.tmp = tempfile.mkdtemp()
self.shas = build_merge_sync_repo(self.tmp)
def test_merge_sync_is_recognized(self):
obs = issue_lock_worktree.read_merge_sync_provenance(
self.tmp, prior_head_sha=self.shas["prior"],
synced_head_sha=self.shas["synced"],
)
self.assertTrue(obs["is_merge_sync"], obs["reasons"])
self.assertTrue(obs["prior_is_ancestor"])
self.assertTrue(obs["synced_is_merge"])
self.assertTrue(obs["first_parent_reaches_prior"])
def test_linear_descendant_is_not_a_merge_sync(self):
"""A plain non-merge descendant (rebase/extra commit) is not a sync."""
obs = issue_lock_worktree.read_merge_sync_provenance(
self.tmp, prior_head_sha=self.shas["prior"],
synced_head_sha=self.shas["linear"],
)
self.assertTrue(obs["probe_ok"])
self.assertFalse(obs["is_merge_sync"])
self.assertFalse(obs["synced_is_merge"])
def test_unrelated_history_fails_closed(self):
"""A rewritten/force-pushed head where prior is unreachable fails closed."""
obs = issue_lock_worktree.read_merge_sync_provenance(
self.tmp, prior_head_sha=self.shas["prior"],
synced_head_sha=self.shas["unrelated"],
)
self.assertFalse(obs["is_merge_sync"])
def test_missing_args_fail_closed(self):
obs = issue_lock_worktree.read_merge_sync_provenance(
self.tmp, prior_head_sha=None, synced_head_sha=self.shas["synced"],
)
self.assertFalse(obs["is_merge_sync"])
# ──────────────────── merge-sync dead-session recovery ──────────────────────
def make_dead_lock(worktree, **over):
pid = dead_pid()
lock = {
"issue_number": ISSUE,
"branch_name": BRANCH,
"worktree_path": worktree,
"remote": REMOTE,
"org": ORG,
"repo": REPO,
"session_pid": pid,
"pid": pid,
"claimant": {"username": IDENTITY, "profile": PROFILE},
"work_lease": {
"operation_type": issue_lock_store.AUTHOR_ISSUE_WORK_LEASE,
"issue_number": ISSUE,
"branch": BRANCH,
"worktree_path": worktree,
"claimant": {"username": IDENTITY, "profile": PROFILE},
"expires_at": future_ts(),
},
}
lock.update(over)
return lock
def sync_prov(prior, synced, **over):
d = {
"prior_head_sha": prior,
"synced_head_sha": synced,
"probe_ok": True,
"prior_present": True,
"synced_present": True,
"prior_is_ancestor": True,
"synced_is_merge": True,
"first_parent_reaches_prior": True,
"is_merge_sync": True,
"first_parent_sha": prior,
"parent_count": 2,
"proof": f"{synced} merged base into branch above {prior}",
"reasons": [],
}
d.update(over)
return d
class TestMergeSyncRecovery(unittest.TestCase):
def setUp(self):
self.tmp = tempfile.mkdtemp()
self.shas = build_merge_sync_repo(self.tmp)
self.prior = self.shas["prior"]
self.synced = self.shas["synced"]
def _assess(self, **over):
lock = over.pop("_lock", None) or make_dead_lock(self.tmp)
kw = dict(
issue_number=ISSUE, branch_name=BRANCH, worktree_path=self.tmp,
remote=REMOTE, org=ORG, repo=REPO, identity=IDENTITY, profile=PROFILE,
current_branch=BRANCH, porcelain_status="",
head_sha=self.prior, remote_head_sha=self.synced,
pr_head_sha=self.synced, pr_number=PR_NUMBER,
competing_live_locks=[], candidate_branches=[BRANCH],
current_pid=os.getpid(),
remote_branch_exists=True,
sync_provenance=sync_prov(self.prior, self.synced),
)
kw.update(over)
return issue_lock_recovery.assess_dead_session_lock_recovery(lock, **kw)
def test_merge_sync_drift_is_recoverable(self):
"""AC3/AC4: dead session, recorded head is a merge-sync ancestor of PR head."""
res = self._assess()
self.assertEqual(res["outcome"], issue_lock_recovery.RECOVERY_SANCTIONED, res["reasons"])
self.assertEqual(
res["evidence"]["head_relation"],
issue_lock_recovery.HEAD_RELATION_REMOTE_MERGE_SYNCED,
)
self.assertEqual(res["evidence"]["accepted_head"], self.synced)
def test_missing_provenance_fails_closed(self):
"""No server-derived provenance → cannot accept a remote ahead of local."""
res = self._assess(sync_provenance=None)
self.assertEqual(res["outcome"], issue_lock_recovery.REFUSED)
def test_non_ancestor_recorded_head_fails_closed(self):
"""AC7: provenance that does not prove ancestry is rejected."""
res = self._assess(
sync_provenance=sync_prov(
self.prior, self.synced, prior_is_ancestor=False, is_merge_sync=False,
reasons=["prior head is not an ancestor"],
)
)
self.assertEqual(res["outcome"], issue_lock_recovery.REFUSED)
def test_force_pushed_history_fails_closed(self):
"""AC8: a rewritten head (not a merge sync) stays protected."""
res = self._assess(
sync_provenance=sync_prov(
self.prior, self.synced, is_merge_sync=False, synced_is_merge=False,
reasons=["not a merge-based sync"],
)
)
self.assertEqual(res["outcome"], issue_lock_recovery.REFUSED)
def test_provenance_for_other_commits_fails_closed(self):
"""Provenance whose endpoints differ from the heads under assessment is rejected."""
res = self._assess(
sync_provenance=sync_prov("f" * 40, self.synced),
)
self.assertEqual(res["outcome"], issue_lock_recovery.REFUSED)
def test_dirty_worktree_fails_closed(self):
"""AC11: dirty worktrees remain protected."""
res = self._assess(porcelain_status=" M feature.txt\n")
self.assertEqual(res["outcome"], issue_lock_recovery.REFUSED)
def test_live_owner_fails_closed(self):
"""AC10: a live recorded owner is not a dead-session recovery."""
lock = make_dead_lock(self.tmp, session_pid=os.getpid(), pid=os.getpid())
res = self._assess(_lock=lock)
self.assertEqual(res["outcome"], issue_lock_recovery.REFUSED)
def test_competing_claimant_fails_closed(self):
"""AC13: a competing live lock blocks recovery."""
res = self._assess(
competing_live_locks=[{
"issue_number": ISSUE, "branch_name": BRANCH,
"worktree_path": "/some/other/wt", "pid": os.getpid(),
}]
)
self.assertEqual(res["outcome"], issue_lock_recovery.REFUSED)
def test_wrong_branch_fails_closed(self):
"""AC9: worktree on a different branch fails closed."""
res = self._assess(current_branch="fix/issue-8710-other")
self.assertEqual(res["outcome"], issue_lock_recovery.REFUSED)
def test_wrong_identity_fails_closed(self):
res = self._assess(identity="intruder")
self.assertEqual(res["outcome"], issue_lock_recovery.REFUSED)
def test_pr_head_mismatch_fails_closed(self):
"""The open PR must sit at the synced remote head."""
res = self._assess(pr_head_sha="e" * 40)
self.assertEqual(res["outcome"], issue_lock_recovery.REFUSED)
def test_owning_pr_evidence_for_merge_sync(self):
res = self._assess()
ev = issue_lock_recovery.owning_pr_recovery_evidence(res)
self.assertIsNotNone(ev)
self.assertEqual(ev["pr_number"], PR_NUMBER)
self.assertEqual(ev["head_sha"], self.synced)
self.assertEqual(
ev["head_relation"],
issue_lock_recovery.HEAD_RELATION_REMOTE_MERGE_SYNCED,
)
def test_recovered_owning_pr_from_persisted_record(self):
res = self._assess()
record = issue_lock_recovery.build_recovery_record(res, recovered_at=future_ts(0))
lock = {"issue_number": ISSUE, "branch_name": BRANCH,
"dead_session_recovery": record}
rebuilt = issue_lock_recovery.recovered_owning_pr_from_lock(lock)
self.assertIsNotNone(rebuilt)
self.assertEqual(rebuilt["head_sha"], self.synced)
self.assertEqual(
rebuilt["head_relation"],
issue_lock_recovery.HEAD_RELATION_REMOTE_MERGE_SYNCED,
)
class TestExistingRelationsUnchanged(unittest.TestCase):
"""AC14/AC15: equal-head recovery still works; merge-sync did not weaken it."""
def setUp(self):
self.tmp = tempfile.mkdtemp()
self.shas = build_merge_sync_repo(self.tmp)
def test_equal_head_recovery_still_sanctioned(self):
# Worktree at prior head; remote also at prior head → the #753 equal case.
prior = self.shas["prior"]
lock = make_dead_lock(self.tmp)
res = issue_lock_recovery.assess_dead_session_lock_recovery(
lock, issue_number=ISSUE, branch_name=BRANCH, worktree_path=self.tmp,
remote=REMOTE, org=ORG, repo=REPO, identity=IDENTITY, profile=PROFILE,
current_branch=BRANCH, porcelain_status="",
head_sha=prior, remote_head_sha=prior,
pr_head_sha=prior, pr_number=PR_NUMBER,
competing_live_locks=[], candidate_branches=[BRANCH],
current_pid=os.getpid(), remote_branch_exists=True,
)
self.assertEqual(res["outcome"], issue_lock_recovery.RECOVERY_SANCTIONED, res["reasons"])
self.assertEqual(
res["evidence"]["head_relation"], issue_lock_recovery.HEAD_RELATION_EQUAL,
)
class TestUpdatePrWrapperPartialFailure(unittest.TestCase):
"""AC5/AC16: the tool advances the remote head then refreshes the durable lock.
When the durable refresh fails after the remote advance, the tool must report a
partial lifecycle failure and NOT a fully successful synchronization. Exact PR-
head / base-head pinning is preserved (delegated to the real preflight, stubbed
here only to isolate the post-update lifecycle branch).
"""
def setUp(self):
import gitea_mcp_server as gms # noqa: E402
self.gms = gms
self._orig = {}
def _patch(name, value):
self._orig[name] = getattr(gms, name)
setattr(gms, name, value)
_patch("get_profile", lambda *a, **k: {
"allowed_operations": ["gitea.branch.push"],
"forbidden_operations": [],
"profile_name": "prgs-author",
})
_patch("_role_kind", lambda *a, **k: "author")
_patch("_profile_operation_gate", lambda *a, **k: None)
_patch("_permission_block_report", lambda *a, **k: {})
_patch("_resolve", lambda *a, **k: ("gitea.prgs.cc", ORG, REPO))
_patch("_verify_role_mutation_workspace", lambda *a, **k: None)
_patch("_get_workspace_porcelain", lambda *a, **k: "")
_patch("_canonical_local_git_root", lambda *a, **k: "/x")
_patch("_master_parity_block", lambda *a, **k: None)
_patch("_auth", lambda *a, **k: {"token": "x"})
_patch("repo_api_url", lambda *a, **k: "http://api")
_patch("_redact", lambda s: s)
_patch("_work_lease_claimant", lambda *a, **k: {
"username": IDENTITY, "profile": PROFILE,
})
_patch("_prove_author_ownership_for_pr", lambda *a, **k: {
"has_author_lock": True, "matched_issue": ISSUE,
"matched_via": "branch", "linked_issues": [ISSUE],
"recovered_owning_pr": None, "reasons": [],
})
# Real preflight is unit-tested elsewhere; stub it to isolate the
# post-update durable-lock lifecycle branch under test.
orig_pf = gms.pr_sync_status.assess_update_pr_branch_preflight
self._orig_pf = orig_pf
gms.pr_sync_status.assess_update_pr_branch_preflight = (
lambda *a, **k: {"mutation_allowed": True, "reasons": [], "performed": False}
)
# Sequence the two GET /pulls calls: OLD before update, NEW after.
self._pull_calls = {"n": 0}
def fake_api_request(method, url, auth, *a, **k):
m = method.upper()
if m == "GET" and url.endswith(f"/pulls/{PR_NUMBER}"):
self._pull_calls["n"] += 1
head = OLD if self._pull_calls["n"] == 1 else NEW1
return {
"state": "open",
"head": {"sha": head, "ref": BRANCH},
"base": {"sha": BASE, "ref": "master"},
"mergeable": True, "title": "t", "body": "b",
}
if m == "GET" and "/branches/" in url:
return {"commit": {"id": BASE}}
if m == "POST" and "/update" in url:
return {}
return {}
_patch("api_request", fake_api_request)
def tearDown(self):
for name, value in self._orig.items():
setattr(self.gms, name, value)
self.gms.pr_sync_status.assess_update_pr_branch_preflight = self._orig_pf
def _run(self):
return self.gms.gitea_update_pr_branch_by_merge(
pr_number=PR_NUMBER,
expected_pr_head_sha=OLD,
expected_base_head_sha=BASE,
remote=REMOTE,
worktree_path="/tmp/branches/wt-871",
)
def test_partial_failure_when_refresh_fails(self):
self._orig["apply_durable_lock_head_refresh"] = (
self.gms.issue_lock_store.apply_durable_lock_head_refresh
)
self.gms.issue_lock_store.apply_durable_lock_head_refresh = (
lambda **k: {"refreshed": False, "reasons": ["forced refresh failure"]}
)
try:
res = self._run()
finally:
self.gms.issue_lock_store.apply_durable_lock_head_refresh = (
self._orig["apply_durable_lock_head_refresh"]
)
self.assertTrue(res["performed"])
self.assertEqual(res["new_pr_head_sha"], NEW1)
self.assertFalse(res["success"])
self.assertTrue(res["partial_lifecycle_failure"])
self.assertFalse(res["durable_lock_refreshed"])
def test_full_success_when_refresh_succeeds(self):
self._orig["apply_durable_lock_head_refresh"] = (
self.gms.issue_lock_store.apply_durable_lock_head_refresh
)
self.gms.issue_lock_store.apply_durable_lock_head_refresh = (
lambda **k: {"refreshed": True, "read_after_write_ok": True,
"new_head": NEW1, "reasons": ["ok"]}
)
try:
res = self._run()
finally:
self.gms.issue_lock_store.apply_durable_lock_head_refresh = (
self._orig["apply_durable_lock_head_refresh"]
)
self.assertTrue(res["success"])
self.assertTrue(res["performed"])
self.assertTrue(res["durable_lock_refreshed"])
self.assertTrue(res["fully_synchronized"])
self.assertEqual(res["new_pr_head_sha"], NEW1)
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(
+107
View File
@@ -0,0 +1,107 @@
"""Documentation acceptance for the MCP restart governance ADR (#656).
Enforces issue #656 acceptance criteria:
* AC1 — policy document exists with an authorization matrix and the recorded
v1 decision (controller approval + automated safety gates).
* AC2 — restart is stated as a last resort with enumerated narrower recoveries.
* AC3 — a unilateral LLM full restart with affected sessions is forbidden.
* AC4 — break-glass conditions are listed.
* AC5 — the ADR is linked to #655, #652, #653, #630, #642, and is cross-linked
from the safety model and the web-console deployment boundary docs.
"""
from pathlib import Path
REPO_ROOT = Path(__file__).resolve().parent.parent
ADR = REPO_ROOT / "docs" / "architecture" / "mcp-restart-governance.md"
ADR_BASENAME = "mcp-restart-governance.md"
CROSS_LINK_DOCS = (
REPO_ROOT / "docs" / "safety-model.md",
REPO_ROOT / "docs" / "webui-deployment.md",
)
LINKED_ISSUES = ("#655", "#652", "#653", "#630", "#642")
POLICY_IDS = ("RG-01", "RG-02", "RG-03", "RG-04", "RG-05", "RG-06", "RG-07", "RG-08")
def _read(path: Path) -> str:
assert path.is_file(), f"missing {path.relative_to(REPO_ROOT)}"
return path.read_text(encoding="utf-8")
def test_ac1_adr_exists_with_matrix_and_v1_decision():
text = _read(ADR)
lower = text.lower()
assert text.lstrip().startswith("#"), "ADR lacks a title"
assert "#656" in text
assert "authorization matrix" in lower
# The matrix is a real table with the worker and privileged roles.
for role in ("author", "reviewer", "merger", "reconciler", "controller",
"operator", "admin"):
assert role in lower, f"authorization matrix missing role {role!r}"
# Recorded v1 decision.
assert "restart-governance/v1" in text
assert "controller approval" in lower and "automated safety gates" in lower
def test_ac2_restart_is_last_resort_with_narrower_recoveries():
text = _read(ADR)
lower = text.lower()
assert "last resort" in lower
# Enumerated narrower recoveries precede full restart on the ladder.
for rung in ("reconnect", "rebind", "scoped restart", "full restart",
"host"):
assert rung in lower, f"recovery ladder missing rung {rung!r}"
def test_ac3_forbids_unilateral_llm_full_restart_with_affected_sessions():
text = _read(ADR)
lower = text.lower()
assert "forbidden" in lower
assert "llm" in lower and "restart" in lower
assert "unilateral" in lower
# A worker role must not perform or authorize full/host restart.
assert "must not" in lower
def test_ac4_break_glass_conditions_listed():
text = _read(ADR)
lower = text.lower()
assert "break-glass" in lower
assert "incident" in lower
assert "audit" in lower
def test_ac5_adr_links_issue_lineage():
text = _read(ADR)
for issue in LINKED_ISSUES:
assert issue in text, f"ADR must link issue {issue}"
def test_ac5_safety_model_and_deployment_cross_link_adr():
for path in CROSS_LINK_DOCS:
text = _read(path)
assert ADR_BASENAME in text, (
f"{path.relative_to(REPO_ROOT)} must cross-link {ADR_BASENAME} "
f"(issue #656 acceptance criterion 5)"
)
def test_policy_ids_present_for_enforcement_code():
text = _read(ADR)
for pid in POLICY_IDS:
assert pid in text, f"policy id {pid} missing from ADR"
def test_failure_behavior_denies_on_ambiguity():
text = _read(ADR)
lower = text.lower()
assert "ambiguous" in lower and "deny" in lower
def test_cross_links_do_not_embed_secrets():
for path in (ADR,) + CROSS_LINK_DOCS:
text = _read(path)
for marker in ("ghp_", "BEGIN PRIVATE KEY", "Authorization: Bearer"):
assert marker not in text, f"{path} contains {marker!r}"
+146
View File
@@ -0,0 +1,146 @@
"""Tests for the MCP restart-path inventory and guards (#657).
Covers:
* the registry is well-formed and every path is classified;
* unknown restart attempts fail closed (AC "fail closed on unknown restart");
* the previously-unguarded full-restart primitives stay guarded/absent
against the real source tree (AC "tests for at least one previously
unguarded path");
* pkill of the daemon is still classified as contamination (#630, AC3);
* the inventory doc and module stay in lock-step.
"""
import os
import tempfile
import unittest
from pathlib import Path
import mcp_restart_paths as rp
import runtime_recovery_guard
REPO_ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
DOC_PATH = os.path.join(REPO_ROOT, "docs", "mcp-restart-path-inventory.md")
class TestRegistryWellformed(unittest.TestCase):
def test_registry_is_wellformed(self):
# Must not raise.
rp.assert_registry_wellformed()
def test_every_path_has_valid_classification(self):
for path in rp.iter_restart_paths():
self.assertIn(path.classification, rp.VALID_CLASSIFICATIONS)
self.assertTrue(path.guard.strip(), path.path_id)
self.assertTrue(path.references, path.path_id)
self.assertTrue(path.locations, path.path_id)
def test_ids_are_unique(self):
ids = [p.path_id for p in rp.iter_restart_paths()]
self.assertEqual(len(ids), len(set(ids)))
def test_covers_every_classification(self):
present = {p.classification for p in rp.iter_restart_paths()}
self.assertEqual(present, set(rp.VALID_CLASSIFICATIONS))
class TestUnknownAttemptFailsClosed(unittest.TestCase):
def test_unknown_path_raises(self):
with self.assertRaises(rp.UnknownRestartPathError):
rp.assert_restart_attempt_registered("totally_novel_restart_hack")
def test_get_unknown_raises(self):
with self.assertRaises(rp.UnknownRestartPathError):
rp.get_restart_path("nope")
def test_registered_attempt_returns_path(self):
path = rp.assert_restart_attempt_registered("manual_daemon_kill")
self.assertEqual(path.classification, rp.CLASS_FORBIDDEN)
class TestDaemonNeverSelfReplaces(unittest.TestCase):
"""Previously-unguarded full-restart primitive: daemon self-replacement."""
def test_no_self_replacement_in_source(self):
# The live daemon modules must contain no os.execv/os.kill/os._exit
# self-restart call. Must not raise.
rp.assert_no_daemon_self_replacement(REPO_ROOT)
def test_scanner_flags_injected_violation(self):
# Guard the guard: prove the scanner catches a real self-replace call.
with tempfile.TemporaryDirectory() as tmp:
bad = Path(tmp) / "gitea_mcp_server.py"
bad.write_text(
"import os\n"
"def restart():\n"
" os.execv('/usr/bin/python', ['python'])\n",
encoding="utf-8",
)
found = rp.scan_daemon_self_replacement(tmp)
self.assertTrue(found)
with self.assertRaises(AssertionError):
rp.assert_no_daemon_self_replacement(tmp)
def test_scanner_ignores_comment_and_docstring_mentions(self):
with tempfile.TemporaryDirectory() as tmp:
ok = Path(tmp) / "gitea_mcp_server.py"
ok.write_text(
"import os\n"
"# NOT os.execv() to re-point the interpreter here.\n"
'"""Never calls os._exit to restart."""\n'
"value = 1\n",
encoding="utf-8",
)
self.assertEqual(rp.scan_daemon_self_replacement(tmp), [])
class TestLegacyAutoRestartHelperRemoved(unittest.TestCase):
"""Previously-unguarded full-restart path: _trigger_mcp_auto_restart."""
def test_helper_absent_in_source(self):
# Must not raise: helper was removed in #685.
rp.assert_auto_restart_helper_absent(REPO_ROOT)
def test_scanner_flags_reintroduced_helper(self):
with tempfile.TemporaryDirectory() as tmp:
bad = Path(tmp) / "mcp_server.py"
bad.write_text(
"def _trigger_mcp_auto_restart():\n return True\n",
encoding="utf-8",
)
with self.assertRaises(AssertionError):
rp.assert_auto_restart_helper_absent(tmp)
class TestPkillStaysForbidden(unittest.TestCase):
"""AC3: pkill of the daemon remains forbidden/contaminating (#630)."""
def test_manual_daemon_kill_registered_as_forbidden(self):
path = rp.get_restart_path("manual_daemon_kill")
self.assertEqual(path.classification, rp.CLASS_FORBIDDEN)
def test_pkill_classified_as_contamination(self):
assessment = runtime_recovery_guard.assess_recovery_command(
"pkill -f mcp_server.py"
)
self.assertTrue(assessment["contaminated"])
def test_read_only_probe_not_contamination(self):
assessment = runtime_recovery_guard.assess_recovery_command(
"ps aux | grep mcp_server"
)
self.assertFalse(assessment["contaminated"])
class TestInventoryDocInSync(unittest.TestCase):
def test_doc_exists(self):
self.assertTrue(os.path.exists(DOC_PATH), DOC_PATH)
def test_doc_mentions_every_path_id(self):
with open(DOC_PATH, encoding="utf-8") as handle:
doc = handle.read()
for path in rp.iter_restart_paths():
self.assertIn(path.path_id, doc, f"doc missing {path.path_id}")
if __name__ == "__main__":
unittest.main()
+53
View File
@@ -12,6 +12,59 @@ import merged_cleanup_reconcile as mcr # noqa: E402
class TestMergedCleanupAssessment(unittest.TestCase):
def test_issue_851_plan_order_worktree_then_reassess_then_remote(self):
"""#851 dry-run plan: remove worktree, reassess ownership, then remote."""
plan = mcr.plan_cleanup_execution_order(
remote_assessment={"safe_to_delete_remote": True},
local_assessment={"safe_to_remove_worktree": True},
)
actions = [s["action"] for s in plan]
self.assertEqual(
actions,
[
"remove_local_worktree",
"reassess_branch_ownership",
"delete_remote_branch",
],
)
self.assertEqual(plan[0]["phase"], 1)
self.assertEqual(plan[-1]["phase"], 3)
self.assertIn("independently_safe", plan[0]["reason"])
self.assertIn("reassessment", plan[-1]["reason"])
def test_issue_851_plan_remote_only_when_worktree_not_safe(self):
plan = mcr.plan_cleanup_execution_order(
remote_assessment={"safe_to_delete_remote": True},
local_assessment={"safe_to_remove_worktree": False},
)
self.assertEqual([s["action"] for s in plan], ["delete_remote_branch"])
self.assertNotIn("reassess_branch_ownership", [s["action"] for s in plan])
def test_issue_851_plan_worktree_only_when_remote_not_safe(self):
plan = mcr.plan_cleanup_execution_order(
remote_assessment={"safe_to_delete_remote": False},
local_assessment={"safe_to_remove_worktree": True},
)
self.assertEqual([s["action"] for s in plan], ["remove_local_worktree"])
def test_issue_851_entry_includes_planned_execution_order(self):
entry = mcr.build_pr_cleanup_entry(
pr={
"number": 848,
"title": "Closes #844",
"body": "",
"merged_at": "2026-07-23T00:00:00Z",
"head": {"ref": "fix/issue-844-x", "sha": "a" * 40},
},
project_root="/tmp/not-a-real-root",
open_pr_heads=set(),
remote_branch_exists=True,
head_on_master=True,
delete_capability_allowed=True,
)
self.assertIn("planned_execution_order", entry)
self.assertIsInstance(entry["planned_execution_order"], list)
def test_extract_linked_issue_from_closes(self):
issue = mcr.extract_linked_issue(
"feat: cleanup (Closes #269)",
+16 -2
View File
@@ -37,6 +37,7 @@ def _live_lock(
"operation_type": issue_lock_store.AUTHOR_ISSUE_WORK_LEASE,
"acquired_at": now.isoformat(),
"expires_at": (now + timedelta(hours=2)).isoformat(),
"session_pid": os.getpid(),
"owner_pid": os.getpid(),
"status": "active",
}
@@ -177,11 +178,24 @@ class TestAuthorOwnershipIssuePrMismatch(unittest.TestCase):
self.assertFalse(result["proven"], result)
self.assertTrue(any("branch" in r for r in result["reasons"]))
def test_no_lock_fail_closed(self):
def test_pidless_durable_lock_rejected(self):
"""A lock without any PID identity must be classified as malformed/non-live and fail closed."""
lock = _live_lock(issue_number=727)
lock.pop("session_pid", None)
lock.pop("owner_pid", None)
lock.pop("pid", None)
path = issue_lock_store.lock_file_path(
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
issue_number=727,
lock_dir=self.lock_dir,
)
issue_lock_store.save_lock_file(path, lock)
result = mcp._prove_author_ownership_for_pr(
pr_number=728,
pr_title="feat: pr sync",
pr_body="Closes #727",
pr_body="Fixes #727",
source_branch="feat/issue-727-pr-sync-status",
remote="prgs",
host=None,
+153
View File
@@ -19,6 +19,7 @@ from pr_work_lease import ( # noqa: E402
assess_reviewer_mutation_blocked,
assess_reviewer_stale_head_final_report,
format_conflict_fix_lease_body,
find_active_conflict_fix_lease,
parse_conflict_fix_lease_comment,
parse_reviewer_lease_comment,
)
@@ -203,5 +204,157 @@ class TestFormatLease(unittest.TestCase):
self.assertEqual(parsed["pr_number"], 376)
class TestConflictFixLeaseLifecycle(unittest.TestCase):
def test_claim_followed_by_matching_release(self):
claim_body = _conflict_fix_body(phase="claimed", worktree="branches/fix-376")
expires = (NOW + timedelta(minutes=60)).isoformat().replace("+00:00", "Z")
release_body = "\n".join([
CONFLICT_FIX_LEASE_MARKER,
"pr: #376",
"branch: feat/fix-376",
"worktree: branches/fix-376",
"profile: prgs-author",
"phase: released",
f"head_before: {HEAD_A}",
f"head_after: {HEAD_B}",
f"expires_at: {expires}",
])
comments = [{"body": claim_body}, {"body": release_body}]
lease = find_active_conflict_fix_lease(comments, pr_number=376, now=NOW)
self.assertIsNone(lease)
def test_expired_claim_without_release(self):
past_expires = (NOW - timedelta(minutes=10)).isoformat().replace("+00:00", "Z")
claim_body = "\n".join([
CONFLICT_FIX_LEASE_MARKER,
"pr: #376",
"phase: claimed",
f"head_before: {HEAD_A}",
f"expires_at: {past_expires}",
"profile: prgs-author",
])
comments = [{"body": claim_body}]
lease = find_active_conflict_fix_lease(comments, pr_number=376, now=NOW)
self.assertIsNone(lease)
def test_mismatched_release_different_head(self):
claim_body = _conflict_fix_body(phase="claimed", worktree="branches/fix-376")
expires = (NOW + timedelta(minutes=60)).isoformat().replace("+00:00", "Z")
release_body = "\n".join([
CONFLICT_FIX_LEASE_MARKER,
"pr: #376",
"profile: prgs-author",
"phase: released",
f"head_before: {HEAD_B}",
f"expires_at: {expires}",
])
comments = [{"body": claim_body}, {"body": release_body}]
lease = find_active_conflict_fix_lease(comments, pr_number=376, now=NOW)
self.assertIsNotNone(lease)
self.assertEqual(lease["phase"], "claimed")
def test_mismatched_release_different_branch(self):
claim_body = "\n".join([
CONFLICT_FIX_LEASE_MARKER,
"pr: #376",
"branch: feat/branch-A",
"phase: claimed",
f"head_before: {HEAD_A}",
f"expires_at: {(NOW + timedelta(minutes=60)).isoformat().replace('+00:00', 'Z')}",
"profile: prgs-author",
])
release_body = "\n".join([
CONFLICT_FIX_LEASE_MARKER,
"pr: #376",
"branch: feat/branch-B",
"phase: released",
f"head_before: {HEAD_A}",
f"expires_at: {(NOW + timedelta(minutes=60)).isoformat().replace('+00:00', 'Z')}",
"profile: prgs-author",
])
comments = [{"body": claim_body}, {"body": release_body}]
lease = find_active_conflict_fix_lease(comments, pr_number=376, now=NOW)
self.assertIsNotNone(lease)
self.assertEqual(lease["phase"], "claimed")
def test_release_followed_by_newer_claim(self):
claim_1 = _conflict_fix_body(phase="claimed", worktree="branches/fix-376")
expires = (NOW + timedelta(minutes=60)).isoformat().replace("+00:00", "Z")
release_1 = "\n".join([
CONFLICT_FIX_LEASE_MARKER,
"pr: #376",
"profile: prgs-author",
"phase: released",
f"head_before: {HEAD_A}",
f"head_after: {HEAD_B}",
f"expires_at: {expires}",
])
claim_2 = "\n".join([
CONFLICT_FIX_LEASE_MARKER,
"pr: #376",
"profile: prgs-author",
"phase: claimed",
f"head_before: {HEAD_B}",
f"expires_at: {expires}",
])
comments = [{"body": claim_1}, {"body": release_1}, {"body": claim_2}]
lease = find_active_conflict_fix_lease(comments, pr_number=376, now=NOW)
self.assertIsNotNone(lease)
self.assertEqual(lease["head_before"], HEAD_B)
def test_malformed_or_ambiguous_markers(self):
malformed_release = "\n".join([
CONFLICT_FIX_LEASE_MARKER,
"pr: #376",
"phase: released",
# missing head_before and profile
])
claim_body = _conflict_fix_body(phase="claimed")
comments = [{"body": claim_body}, {"body": malformed_release}]
lease = find_active_conflict_fix_lease(comments, pr_number=376, now=NOW)
self.assertIsNotNone(lease)
def test_pr818_historical_sequence(self):
comment_14696 = "\n".join([
"<!-- mcp-conflict-fix-lease:v1 -->",
"pr: #818",
"branch: feat/issue-638-webui-app-shell-phase1",
"worktree: /Users/jasonwalker/Development/Gitea-Tools/branches/issue-638-webui-app-shell-phase1",
"profile: prgs-author",
"session_id: unknown",
"phase: claimed",
"head_before: 08061b7b8aebdd099a37d1abf5dafcf38e4fd3fb",
"expires_at: 2026-07-23T07:12:13Z",
"reviewer_active: no",
])
comment_14730 = "\n".join([
"<!-- mcp-conflict-fix-lease:v1 -->",
"pr: #818",
"branch: feat/issue-638-webui-app-shell-phase1",
"worktree: /Users/jasonwalker/Development/Gitea-Tools/branches/issue-638-webui-app-shell-phase1",
"profile: prgs-author",
"session_id: prgs-author-61241-e5129c60",
"phase: released",
"head_before: 08061b7b8aebdd099a37d1abf5dafcf38e4fd3fb",
"head_after: 64b6eb5d5402663098de5ded3b0617cc3b3df98f",
"expires_at: 2026-07-23T06:05:00Z",
"reviewer_active: no",
])
comments = [{"body": comment_14696}, {"body": comment_14730}]
check_now = datetime(2026, 7, 23, 6, 30, tzinfo=timezone.utc)
lease = find_active_conflict_fix_lease(comments, pr_number=818, now=check_now)
self.assertIsNone(lease)
reviewer_gate = assess_reviewer_mutation_blocked(
pr_number=818,
comments=comments,
reviewed_head_sha="64b6eb5d5402663098de5ded3b0617cc3b3df98f",
live_head_sha="64b6eb5d5402663098de5ded3b0617cc3b3df98f",
mutation="approve",
now=check_now,
)
self.assertTrue(reviewer_gate["mutation_allowed"])
if __name__ == "__main__":
unittest.main()
+340
View File
@@ -0,0 +1,340 @@
"""Tests for the MCP restart coordinator and impact analysis (#658).
Multi-session fixtures exercise every verdict branch: safe, unsafe (live work),
override, and the fail-closed deny on incomplete inventory. Also covers the
critical-section deny path and the new ``ControlPlaneDB.list_sessions``.
"""
from __future__ import annotations
import os
import tempfile
import unittest
from datetime import datetime, timedelta, timezone
import restart_coordinator as rc
from control_plane_db import ControlPlaneDB
NOW = datetime(2026, 7, 24, 6, 0, 0, tzinfo=timezone.utc)
def _ts(dt: datetime) -> str:
return dt.isoformat()
def _live_pid() -> int:
return os.getpid()
def _dead_pid() -> int:
# A pid that is essentially never alive. os.kill(0) on it raises
# ProcessLookupError → is_process_alive False.
return 2_000_000_000
def _session(session_id, *, pid, status="active", heartbeat=None, role="author"):
return {
"session_id": session_id,
"role": role,
"profile": "prgs-author",
"pid": pid,
"status": status,
"last_heartbeat_at": _ts(heartbeat or NOW),
}
def _lease(
lease_id,
*,
session_id,
freshness,
kind="issue",
number=658,
phase="allocated",
worktree=None,
role="author",
):
return {
"lease_id": lease_id,
"session_id": session_id,
"role": role,
"phase": phase,
"work_kind": kind,
"work_number": number,
"worktree_path": worktree,
"freshness": {"freshness": freshness},
}
class EvaluateRestartImpactTest(unittest.TestCase):
def test_incomplete_inventory_denies_fail_closed(self) -> None:
report = rc.evaluate_restart_impact(
{"inventory_complete": False, "incomplete_reasons": ["db down"]},
now=NOW,
)
self.assertEqual(report.verdict, rc.VERDICT_UNSAFE)
self.assertFalse(report.allow_restart)
self.assertFalse(report.restart_performed)
self.assertIn("db down", report.incomplete_reasons)
self.assertTrue(
any("fail closed" in reasoning for reasoning in report.reasons)
)
def test_missing_completeness_flag_denies(self) -> None:
# No inventory_complete key at all → treated as incomplete.
report = rc.evaluate_restart_impact({}, now=NOW)
self.assertEqual(report.verdict, rc.VERDICT_UNSAFE)
self.assertFalse(report.allow_restart)
def test_no_other_work_is_safe(self) -> None:
report = rc.evaluate_restart_impact(
{
"inventory_complete": True,
"sessions": [_session("requester", pid=_live_pid())],
"leases": [],
},
now=NOW,
requesting_session_id="requester",
)
self.assertEqual(report.verdict, rc.VERDICT_SAFE)
self.assertTrue(report.allow_restart)
self.assertEqual(report.blast_radius, rc.BLAST_NONE)
self.assertEqual(report.affected_issues, [])
def test_dead_foreign_session_and_lease_are_not_disruptive(self) -> None:
report = rc.evaluate_restart_impact(
{
"inventory_complete": True,
"sessions": [
_session("requester", pid=_live_pid()),
_session("dead", pid=_dead_pid()),
],
"leases": [
_lease("l-dead", session_id="dead", freshness="stale_dead_process")
],
},
now=NOW,
requesting_session_id="requester",
)
self.assertEqual(report.verdict, rc.VERDICT_SAFE)
self.assertTrue(report.allow_restart)
self.assertEqual(report.counts["leases_disruptive"], 0)
self.assertEqual(report.counts["sessions_live_other"], 0)
def test_live_foreign_lease_denies_without_override(self) -> None:
report = rc.evaluate_restart_impact(
{
"inventory_complete": True,
"sessions": [
_session("requester", pid=_live_pid()),
_session("worker", pid=_live_pid()),
],
"leases": [
_lease(
"l1",
session_id="worker",
freshness="active",
worktree="/tmp/wt-658",
phase="implementing",
)
],
},
now=NOW,
requesting_session_id="requester",
)
self.assertEqual(report.verdict, rc.VERDICT_UNSAFE)
self.assertFalse(report.allow_restart)
# Critical section detected: active lease with a live owner.
self.assertEqual(len(report.critical_sections), 1)
self.assertEqual(report.affected_issues, [658])
self.assertEqual(report.counts["mutations"], 1)
self.assertTrue(report.override_would_allow)
self.assertEqual(report.blast_radius, rc.BLAST_HIGH)
# Placeholder ack state for the affected session.
self.assertEqual(report.ack_state.get("worker"), "pending")
def test_operator_override_allows_despite_live_work(self) -> None:
inv = {
"inventory_complete": True,
"sessions": [
_session("requester", pid=_live_pid()),
_session("worker", pid=_live_pid()),
],
"leases": [_lease("l1", session_id="worker", freshness="active")],
}
report = rc.evaluate_restart_impact(
inv,
now=NOW,
requesting_session_id="requester",
operator_override=True,
)
self.assertEqual(report.verdict, rc.VERDICT_OVERRIDE)
self.assertTrue(report.allow_restart)
self.assertFalse(report.restart_performed)
def test_deny_when_critical_section_open(self) -> None:
# A single live author lease in a mutating phase is a critical section
# that must deny an un-overridden restart.
report = rc.evaluate_restart_impact(
{
"inventory_complete": True,
"sessions": [_session("worker", pid=_live_pid())],
"leases": [
_lease(
"l1",
session_id="worker",
freshness="active",
phase="merging",
kind="pr",
number=900,
)
],
},
now=NOW,
requesting_session_id="requester",
)
self.assertEqual(report.verdict, rc.VERDICT_UNSAFE)
self.assertFalse(report.allow_restart)
self.assertEqual(report.affected_prs, [900])
self.assertEqual(len(report.critical_sections), 1)
def test_terminal_lock_makes_restart_unsafe(self) -> None:
report = rc.evaluate_restart_impact(
{
"inventory_complete": True,
"sessions": [_session("requester", pid=_live_pid())],
"leases": [],
"terminal_lock": {"terminal_pr": 812},
},
now=NOW,
requesting_session_id="requester",
)
self.assertEqual(report.verdict, rc.VERDICT_UNSAFE)
self.assertFalse(report.allow_restart)
self.assertIsNotNone(report.terminal_lock)
self.assertTrue(
any("terminal" in reasoning for reasoning in report.reasons)
)
def test_other_live_session_without_lease_is_disruptive(self) -> None:
report = rc.evaluate_restart_impact(
{
"inventory_complete": True,
"sessions": [
_session("requester", pid=_live_pid()),
_session("idle-but-live", pid=_live_pid()),
],
"leases": [],
},
now=NOW,
requesting_session_id="requester",
)
self.assertEqual(report.verdict, rc.VERDICT_UNSAFE)
self.assertEqual(report.counts["sessions_live_other"], 1)
def test_stale_heartbeat_session_not_counted_live(self) -> None:
stale = NOW - timedelta(hours=2)
report = rc.evaluate_restart_impact(
{
"inventory_complete": True,
"sessions": [
_session("requester", pid=_live_pid()),
_session("stale", pid=_live_pid(), heartbeat=stale),
],
"leases": [],
},
now=NOW,
requesting_session_id="requester",
)
self.assertEqual(report.verdict, rc.VERDICT_SAFE)
self.assertEqual(report.counts["sessions_live_other"], 0)
def test_prior_recovery_attempts_echoed(self) -> None:
report = rc.evaluate_restart_impact(
{
"inventory_complete": True,
"sessions": [_session("requester", pid=_live_pid())],
"leases": [],
"prior_recovery_attempts": [
{"kind": "client_reconnect", "at": _ts(NOW)}
],
},
now=NOW,
requesting_session_id="requester",
)
self.assertEqual(len(report.prior_recovery_attempts), 1)
self.assertEqual(report.counts["prior_recovery_attempts"], 1)
def test_bare_string_freshness_accepted(self) -> None:
lease = _lease("l1", session_id="worker", freshness="active")
lease["freshness"] = "active" # bare string, not a dict
report = rc.evaluate_restart_impact(
{
"inventory_complete": True,
"sessions": [_session("worker", pid=_live_pid())],
"leases": [lease],
},
now=NOW,
requesting_session_id="requester",
)
self.assertEqual(report.counts["leases_disruptive"], 1)
def test_as_dict_is_serializable_dto(self) -> None:
import json
report = rc.evaluate_restart_impact(
{
"inventory_complete": True,
"sessions": [_session("requester", pid=_live_pid())],
"leases": [],
},
now=NOW,
requesting_session_id="requester",
)
payload = report.as_dict()
# Round-trips through JSON — safe for the console DTO.
encoded = json.dumps(payload)
decoded = json.loads(encoded)
self.assertEqual(decoded["verdict"], rc.VERDICT_SAFE)
self.assertIn("audit_record", decoded)
self.assertEqual(decoded["audit_record"]["event"], "restart_impact_evaluated")
self.assertFalse(decoded["restart_performed"])
self.assertIn("coordinator_version", decoded)
class ListSessionsTest(unittest.TestCase):
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory()
self.db = ControlPlaneDB(os.path.join(self._tmp.name, "cp.sqlite3"))
def tearDown(self) -> None:
self._tmp.cleanup()
def test_list_sessions_filters_by_status(self) -> None:
self.db.upsert_session(session_id="a", role="author", pid=1, status="active")
self.db.upsert_session(session_id="b", role="author", pid=2, status="ended")
active = self.db.list_sessions(statuses=("active",))
ids = {row["session_id"] for row in active}
self.assertEqual(ids, {"a"})
every = self.db.list_sessions()
self.assertEqual({row["session_id"] for row in every}, {"a", "b"})
def test_list_sessions_feeds_coordinator(self) -> None:
self.db.upsert_session(
session_id="requester", role="author", pid=os.getpid(), status="active"
)
report = rc.evaluate_restart_impact(
{
"inventory_complete": True,
"sessions": self.db.list_sessions(statuses=("active",)),
"leases": [],
},
now=NOW,
requesting_session_id="requester",
)
self.assertEqual(report.counts["sessions_total"], 1)
if __name__ == "__main__": # pragma: no cover
unittest.main()
+135
View File
@@ -0,0 +1,135 @@
"""Tests for the Phase 1 operator console application shell (#638)."""
import sys
import unittest
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from starlette.routing import Route
from starlette.testclient import TestClient
from webui import layout
from webui.app import create_app
from webui.nav import NAV_GROUPS, STUB_PAGES, nav_hrefs
class TestShellNav(unittest.TestCase):
def setUp(self):
self.client = TestClient(create_app())
def test_nav_group_labels_present(self):
text = self.client.get("/").text
for group in NAV_GROUPS:
with self.subTest(group=group.label):
self.assertIn(f">{group.label}<", text)
def test_phase1_group_labels_cover_expected_ia(self):
labels = {group.label for group in NAV_GROUPS}
for expected in (
"Health",
"Traffic",
"Runtime/Sessions",
"Projects",
"Inventory",
"Timeline",
"Policy",
"Insights",
):
with self.subTest(label=expected):
self.assertIn(expected, labels)
def test_every_nav_href_resolves_to_a_get_route(self):
app = create_app()
get_paths = {
route.path
for route in app.routes
if isinstance(route, Route) and "GET" in route.methods
}
for href in nav_hrefs():
with self.subTest(href=href):
self.assertIn(href, get_paths, f"nav href {href} has no GET route")
def test_legacy_hrefs_still_navigable(self):
text = self.client.get("/").text
for href in ("/queue", "/projects", "/prompts", "/runtime",
"/audit", "/worktrees", "/leases", "/actions"):
with self.subTest(href=href):
self.assertIn(f'href="{href}"', text)
class TestShellBadges(unittest.TestCase):
def setUp(self):
self.client = TestClient(create_app())
def test_mode_badge_present(self):
self.assertIn("mode: read-only", self.client.get("/").text)
def test_environment_badge_present(self):
self.assertIn("env:", self.client.get("/").text)
def test_default_environment_is_local(self):
self.assertEqual(layout.environment_label(), "local")
def test_remote_bind_reports_remote_environment(self):
import os
prior = os.environ.get("WEBUI_HOST")
os.environ["WEBUI_HOST"] = "10.0.0.5"
try:
self.assertEqual(layout.environment_label(), "remote")
finally:
if prior is None:
os.environ.pop("WEBUI_HOST", None)
else:
os.environ["WEBUI_HOST"] = prior
def test_docs_link_present(self):
text = self.client.get("/").text
self.assertIn(layout.DOCS_URL, text)
self.assertIn(">Docs<", text)
class TestShellStubs(unittest.TestCase):
def setUp(self):
self.client = TestClient(create_app())
def test_stub_routes_render_200(self):
for path, (title, _desc) in STUB_PAGES.items():
with self.subTest(path=path):
response = self.client.get(path)
self.assertEqual(response.status_code, 200, path)
self.assertIn(title, response.text)
self.assertIn("placeholder", response.text)
def test_stub_routes_are_read_only(self):
for path in STUB_PAGES:
with self.subTest(path=path):
response = self.client.post(path)
self.assertEqual(response.status_code, 405)
self.assertEqual(response.json()["error"], "read-only-mvp")
def test_stub_pages_carry_nav_and_badges(self):
response = self.client.get("/inventory")
self.assertIn("mode: read-only", response.text)
self.assertIn('href="/queue"', response.text)
class TestShellHome(unittest.TestCase):
def setUp(self):
self.client = TestClient(create_app())
def test_home_summarizes_console(self):
text = self.client.get("/").text
self.assertIn("Operator console", text)
self.assertIn("Phase 1", text)
def test_home_links_legacy_pages(self):
text = self.client.get("/").text
self.assertIn("MVP legacy pages", text)
for href in ("/queue", "/audit", "/leases"):
with self.subTest(href=href):
self.assertIn(f'href="{href}"', text)
if __name__ == "__main__":
unittest.main()
+345
View File
@@ -0,0 +1,345 @@
"""Tests for the system-health dashboard view (#639).
Covers the acceptance criteria directly: the page renders the health DTO
fields (AC1), degraded dependencies are visible (AC2), stale runtime is warned
prominently and never rendered as mutation-safe (AC3), healthy and degraded
fixtures both render (AC4), and the shell carries a nav entry (AC5).
"""
import sys
import unittest
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from starlette.testclient import TestClient
from webui.app import create_app
from webui.deployment_boundary import scan_text_for_client_secrets
from webui.layout import render_page
from webui.nav import iter_nav_items
from webui.system_health import (
STATUS_DEGRADED,
STATUS_DOWN,
STATUS_OK,
STATUS_SKIPPED,
STATUS_UNPROVEN,
DependencyProbe,
StaleRuntime,
SystemHealthSnapshot,
VersionInfo,
)
from webui.system_health_views import render_system_health_page
DASHBOARD_PATH = "/system-health"
def _version(*, known: bool = True) -> VersionInfo:
return VersionInfo(
git_sha="1c455b6ec0f9cb761fe6248de68c17e061fb5ecd" if known else None,
git_describe="v0.4.1-12-g1c455b6" if known else None,
control_plane_schema_version=4 if known else None,
python_version="3.13.1",
known=known,
)
def _parity(*, stale: bool = False, determinable: bool = True) -> StaleRuntime:
if stale:
return StaleRuntime(
daemon_head="aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa",
checkout_head="bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb",
remote_head="cccccccccccccccccccccccccccccccccccccccc",
stale=True,
determinable=True,
mutation_safe=False,
reasons=("runtime, checkout, and remote commits disagree",),
)
if not determinable:
return StaleRuntime(
daemon_head=None,
checkout_head=None,
remote_head=None,
stale=False,
determinable=False,
mutation_safe=False,
reasons=("local checkout HEAD could not be read",),
)
return StaleRuntime(
daemon_head="1c455b6ec0f9cb761fe6248de68c17e061fb5ecd",
checkout_head="1c455b6ec0f9cb761fe6248de68c17e061fb5ecd",
remote_head="1c455b6ec0f9cb761fe6248de68c17e061fb5ecd",
stale=False,
determinable=True,
mutation_safe=True,
reasons=(),
)
def _snapshot(
*,
status: str = STATUS_OK,
ready: bool = True,
readiness_complete: bool = True,
readiness_reasons: tuple[str, ...] = (),
dependencies: tuple[DependencyProbe, ...] | None = None,
parity: StaleRuntime | None = None,
namespaces: tuple[dict, ...] = (),
probe_errors: tuple[str, ...] = (),
version_known: bool = True,
) -> SystemHealthSnapshot:
if dependencies is None:
dependencies = (
DependencyProbe(
name="control_plane_db",
kind="sqlite",
status=STATUS_OK,
detail="schema version 4",
required=True,
latency_ms=1.25,
metadata={"schema_version": 4},
),
)
return SystemHealthSnapshot(
status=status,
ready=ready,
readiness_complete=readiness_complete,
readiness_reasons=readiness_reasons,
service="mcp-control-plane-webui",
mode="read-only",
version=_version(known=version_known),
started_at="2026-07-23T19:50:47+00:00",
uptime_seconds=3661.5,
timestamp="2026-07-23T20:51:48+00:00",
deep_probes_requested=False,
dependencies=dependencies,
mcp_namespaces=namespaces,
stale_runtime=parity if parity is not None else _parity(),
probe_errors=probe_errors,
)
class TestHealthyRender(unittest.TestCase):
"""AC1 / AC4 — every health DTO field reaches the page."""
def setUp(self):
self.html = render_system_health_page(_snapshot())
def test_readiness_fields_render(self):
self.assertIn("System health", self.html)
self.assertIn("Ready", self.html)
self.assertIn("mcp-control-plane-webui", self.html)
self.assertIn("read-only", self.html)
self.assertIn("2026-07-23T20:51:48+00:00", self.html)
def test_version_and_uptime_render(self):
self.assertIn("1c455b6ec0f9cb761fe6248de68c17e061fb5ecd", self.html)
self.assertIn("v0.4.1-12-g1c455b6", self.html)
self.assertIn("3.13.1", self.html)
self.assertIn("3661.500s", self.html)
self.assertIn("1.02h", self.html)
def test_dependency_row_renders_with_latency(self):
self.assertIn("control_plane_db", self.html)
self.assertIn("sqlite", self.html)
self.assertIn("schema version 4", self.html)
self.assertIn("1.2 ms", self.html)
def test_healthy_page_shows_no_stale_warning(self):
self.assertNotIn("Stale runtime:", self.html)
self.assertNotIn("Staleness", self.html)
def test_unknown_version_is_labelled_not_faked(self):
html = render_system_health_page(_snapshot(version_known=False))
self.assertIn("unknown", html)
self.assertIn("unresolved", html)
class TestDegradedRender(unittest.TestCase):
"""AC2 — a degraded or unrun dependency is visible, not swallowed."""
def setUp(self):
self.deps = (
DependencyProbe(
name="control_plane_db",
kind="sqlite",
status=STATUS_OK,
detail="schema version 4",
required=True,
latency_ms=0.9,
),
DependencyProbe(
name="repository",
kind="git",
status=STATUS_DOWN,
detail="repository root is not a git checkout",
required=True,
latency_ms=4.0,
),
DependencyProbe(
name="gitea",
kind="http",
status=STATUS_SKIPPED,
detail="deep probe not requested",
required=False,
),
)
self.html = render_system_health_page(
_snapshot(
status=STATUS_DEGRADED,
ready=False,
readiness_complete=False,
readiness_reasons=("required dependency 'repository' is down",),
dependencies=self.deps,
)
)
def test_degraded_banner_names_the_dependency(self):
self.assertIn("Degraded dependencies:", self.html)
self.assertIn("repository", self.html)
def test_not_run_probe_is_reported_separately(self):
self.assertIn("Not probed:", self.html)
self.assertIn("gitea", self.html)
self.assertIn("not counted", self.html)
def test_not_ready_headline_and_reason(self):
self.assertIn("Not ready", self.html)
self.assertIn("required dependency &#x27;repository&#x27; is down", self.html)
def test_degraded_status_badge_present(self):
self.assertIn("badge-health-degraded", self.html)
self.assertIn("badge-health-down", self.html)
def test_ready_but_incomplete_is_not_shown_as_plain_ready(self):
html = render_system_health_page(
_snapshot(ready=True, readiness_complete=False)
)
self.assertIn("Ready (incomplete evidence)", html)
class TestStaleRuntimeWarning(unittest.TestCase):
"""AC3 — staleness is prominent and never claims mutation safety."""
def test_stale_runtime_warns_and_denies_mutation_safety(self):
html = render_system_health_page(_snapshot(parity=_parity(stale=True)))
self.assertIn("Stale runtime:", html)
self.assertIn("do not treat this runtime as mutation-safe", html)
self.assertIn("<tr><th>Mutation safe</th><td>False</td></tr>", html)
def test_indeterminate_parity_is_not_reported_safe(self):
html = render_system_health_page(
_snapshot(parity=_parity(determinable=False))
)
self.assertIn("Staleness", html)
self.assertIn("<tr><th>Mutation safe</th><td>False</td></tr>", html)
self.assertIn("<tr><th>Determinable</th><td>False</td></tr>", html)
def test_healthy_parity_reports_mutation_safe_true(self):
html = render_system_health_page(_snapshot())
self.assertIn("<tr><th>Mutation safe</th><td>True</td></tr>", html)
class TestNamespacesAndErrors(unittest.TestCase):
def test_unproven_namespace_rows_render(self):
html = render_system_health_page(
_snapshot(
namespaces=(
{
"namespace": "gitea-author",
"required_tool": "gitea_lock_issue",
"status": STATUS_UNPROVEN,
"ide_namespace_proven": False,
"reason": "the web console cannot invoke the IDE-managed MCP client",
},
)
)
)
self.assertIn("gitea-author", html)
self.assertIn("gitea_lock_issue", html)
self.assertIn("badge-health-unproven", html)
def test_no_namespaces_degrades_gracefully(self):
html = render_system_health_page(_snapshot(namespaces=()))
self.assertIn("No MCP namespaces are declared.", html)
def test_probe_errors_render_when_present(self):
html = render_system_health_page(
_snapshot(probe_errors=("probe raised: disk offline",))
)
self.assertIn("Probe errors", html)
self.assertIn("disk offline", html)
def test_probe_error_card_absent_when_clean(self):
self.assertNotIn("Probe errors", render_system_health_page(_snapshot()))
class TestReadOnlyAndRedaction(unittest.TestCase):
def test_no_restart_or_kill_controls(self):
html = render_system_health_page(_snapshot())
self.assertNotIn("<button", html)
self.assertNotIn("<form", html)
self.assertNotIn("pkill", html)
self.assertIn("read-only", html)
def test_recovery_points_at_sanctioned_path(self):
html = render_system_health_page(_snapshot())
self.assertIn("Reconnect the MCP client", html)
self.assertIn("Never kill the daemon process manually", html)
def test_secret_shaped_detail_is_redacted(self):
leaky = DependencyProbe(
name="gitea",
kind="http",
status=STATUS_DOWN,
detail="auth failed for token=ghp_ABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789",
required=False,
latency_ms=12.0,
)
html = render_system_health_page(_snapshot(dependencies=(leaky,)))
self.assertNotIn("ghp_ABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789", html)
def test_html_in_detail_is_escaped(self):
hostile = DependencyProbe(
name="repository",
kind="git",
status=STATUS_DOWN,
detail="<script>alert(1)</script>",
required=True,
)
html = render_system_health_page(_snapshot(dependencies=(hostile,)))
self.assertNotIn("<script>", html)
self.assertIn("&lt;script&gt;", html)
class TestNavAndRoute(unittest.TestCase):
"""AC5 — the shell links the dashboard, and the route serves it."""
def setUp(self):
self.client = TestClient(create_app())
def test_nav_contains_system_health(self):
self.assertIn(
(DASHBOARD_PATH, "System health"),
[(item.href, item.label) for item in iter_nav_items()],
)
def test_rendered_shell_links_dashboard(self):
page = render_page(title="Home", body_html="<p>x</p>")
self.assertIn(f'href="{DASHBOARD_PATH}"', page)
def test_route_renders_dashboard(self):
response = self.client.get(DASHBOARD_PATH)
self.assertEqual(response.status_code, 200)
self.assertIn("System health", response.text)
self.assertIn("Stale-runtime parity", response.text)
def test_route_is_read_only(self):
self.assertEqual(self.client.post(DASHBOARD_PATH).status_code, 405)
def test_live_page_leaks_no_client_secret(self):
findings = scan_text_for_client_secrets(self.client.get(DASHBOARD_PATH).text)
self.assertEqual(findings, [])
if __name__ == "__main__": # pragma: no cover
unittest.main()
File diff suppressed because it is too large Load Diff
+24 -2
View File
@@ -134,13 +134,35 @@ class TestClassification(unittest.TestCase):
self.assertEqual(cls, wca.CLASS_ACTIVE_OPEN_PR)
self.assertFalse(wca.is_removable(cls))
def test_stale_clean_issue_worktree_removable(self):
# Scenario 5: clean issue worktree, TTL expired, no lock -> removable.
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.
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))