Compare commits

..
Author SHA1 Message Date
sysadmin ba3ea3012c fix(author): harden dirty-session rebind inventory and journal identity (#868)
Complete dirty-inventory revalidation immediately before and after
bind_session_lock so added, removed, or renamed paths fail closed.
Persist and validate full recovery-journal operation identity (remote,
org, repo, claimant identity, claimant profile) on execute, resume,
retry, and already_rebound. Add focused regression coverage and the
reconciler success-path integration test.

Closes #868.
2026-07-24 01:13:50 -04:00
6 changed files with 915 additions and 963 deletions
-84
View File
@@ -525,90 +525,6 @@ def assess_ownership_record_activity(record: dict[str, Any]) -> dict[str, Any]:
}
# Reviewer-lease reclaim is only reachable from a non-live (expired/stale) lease.
_RECLAIMABLE_REVIEWER_STATUSES = _EXPIRED_STATUSES | _STALE_STATUSES
def is_active_ownership_status(status: str | None) -> bool:
"""True when *status* denotes live/active ownership of a branch (#855).
Used to decide whether a *competing* active claimant still uses a branch
when weighing an expired reviewer lease for reclaim. Expired, stale,
released, and terminal statuses are not active.
"""
return _norm_str(status).lower() in _ACTIVE_OWNERSHIP_STATUSES
def assess_expired_reviewer_lease_reclaim(
*,
role: str,
status: str,
pr_merged: bool | None,
owner_pid_alive: bool | None,
competing_active_claimant: bool | None,
) -> dict[str, Any]:
"""Decide, explicitly and fail-closed, whether an expired reviewer lease
may stop protecting an already-merged branch (#855 AC4).
An expired reviewer lease should not protect a merged branch forever once
its work is done and no live claimant remains. Reclaim is permitted only
when **every** condition below is provably satisfied; any unknown
(``None``) or contrary value keeps the lease protective:
- the lease is a ``reviewer`` lease (author/merger/controller/reconciler
leases are out of scope and always keep protecting);
- its status is expired or stale (never an active/live lease);
- the PR is proven merged (``pr_merged is True``);
- the lease owner process is proven dead (``owner_pid_alive is False``);
- no competing active claimant uses the branch
(``competing_active_claimant is False``).
Returns a decision dict with ``reclaim_allowed`` and, when refused, the
fail-closed ``reasons``. The reasons never contain secrets — only the
role, the status, and which condition was unproven.
"""
reasons: list[str] = []
normalized_role = _norm_str(role).lower()
normalized_status = _norm_str(status).lower()
if normalized_role != "reviewer":
reasons.append(
f"lease role '{normalized_role or 'unknown'}' is not a reviewer "
"lease; expired-reviewer reclaim does not apply"
)
if normalized_status not in _RECLAIMABLE_REVIEWER_STATUSES:
reasons.append(
f"lease status '{normalized_status or 'unknown'}' is not expired "
"or stale; only a non-live reviewer lease may be reclaimed"
)
if pr_merged is not True:
reasons.append(
"PR merged state is not proven true; reclaim requires an "
"already-merged PR (fail closed)"
)
if owner_pid_alive is not False:
reasons.append(
"lease owner process liveness is not proven dead; a live owner "
"still protects the branch (fail closed)"
)
if competing_active_claimant is not False:
reasons.append(
"a competing active claimant may still use the branch; reclaim "
"requires no other active ownership (fail closed)"
)
allowed = not reasons
return {
"reclaim_allowed": allowed,
"role": normalized_role,
"status": normalized_status,
"decision": (
"reclaim_expired_reviewer_lease" if allowed else "keep_protecting"
),
"reasons": [] if allowed else reasons,
}
def assess_active_branch_ownership(
*,
remote: str,
+401 -68
View File
@@ -1,4 +1,4 @@
"""Dirty-preserving same-claimant author-session rebind (#864).
"""Dirty-preserving same-claimant author-session rebind (#864 / #868).
A registered issue worktree can be dirty while its durable lock owner PID is
provably dead. Ordinary ``gitea_lock_issue`` refuses dirty trees, and dead-session
@@ -14,6 +14,14 @@ This operation:
heartbeat)
* preserves every tracked/untracked byte
* does NOT sync remote, create recovery worktrees, clean, reset, or change heads
#868 hardens:
* complete dirty-inventory revalidation (full path set + fingerprints)
immediately before and after ``bind_session_lock``
* durable recovery-journal operation identity (remote, org, repo, claimant
identity, claimant profile) validated on execute / resume / retry /
already_rebound
"""
from __future__ import annotations
@@ -66,6 +74,16 @@ REQUIRED_LOCK_FIELDS = (
"repo",
)
# Durable journal operation identity (#868 F2). All five must be persisted on
# JOURNAL_PHASE_ASSESSED and re-validated on resume / retry / already_rebound.
REQUIRED_JOURNAL_IDENTITY_FIELDS = (
"remote",
"org",
"repo",
"claimant_identity",
"claimant_profile",
)
def _utc_now_iso() -> str:
return datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
@@ -196,6 +214,174 @@ def collect_dirty_inventory(worktree_path: str) -> dict[str, Any]:
}
def revalidate_complete_dirty_inventory(
worktree_path: str,
*,
expected_dirty_paths: Sequence[str] | None,
expected_fingerprints: Mapping[str, str] | None,
phase: str = "inventory",
) -> dict[str, Any]:
"""Collect the full dirty inventory and require exact pin equality (#868 F1).
Unlike fingerprint-only checks over the expected path list, this recollects
the authoritative tracked+untracked inventory and refuses added, removed,
or renamed paths as well as fingerprint movement.
"""
reasons: list[str] = []
inv = collect_dirty_inventory(worktree_path)
if inv.get("ok") is False:
reasons.extend(list(inv.get("reasons") or []) or [f"{phase}: dirty inventory collection failed"])
observed_paths = sorted(
{_text(p) for p in (inv.get("dirty_paths") or []) if _text(p)}
)
pin_paths = sorted(
{_text(p) for p in (expected_dirty_paths or []) if _text(p)}
)
if not pin_paths:
reasons.append(
f"{phase}: expected_dirty_paths pin is empty; complete inventory "
"revalidation requires a non-empty pin (fail closed)"
)
if set(observed_paths) != set(pin_paths):
extra = sorted(set(observed_paths) - set(pin_paths))
missing = sorted(set(pin_paths) - set(observed_paths))
if extra:
reasons.append(
f"{phase}: complete dirty inventory path-set disagreement: "
f"unexpected paths {extra}"
)
if missing:
reasons.append(
f"{phase}: complete dirty inventory path-set disagreement: "
f"missing expected paths {missing}"
)
obs_fps = {
_text(k): _text(v)
for k, v in dict(inv.get("fingerprints") or {}).items()
if _text(k)
}
pin_fps = {
_text(k): _text(v)
for k, v in dict(expected_fingerprints or {}).items()
if _text(k)
}
if not pin_fps:
reasons.append(
f"{phase}: expected_fingerprints pin is empty; byte-level pins "
"are required (fail closed)"
)
else:
for rel, expected_hash in pin_fps.items():
if rel not in set(pin_paths):
reasons.append(
f"{phase}: expected_fingerprints contains '{rel}' which is "
"not in expected_dirty_paths"
)
continue
actual_hash = obs_fps.get(rel)
if not actual_hash:
reasons.append(
f"{phase}: fingerprint missing for dirty path '{rel}'"
)
elif actual_hash != expected_hash:
reasons.append(
f"{phase}: fingerprint disagreement for '{rel}': "
f"observed {actual_hash}, expected {expected_hash}"
)
for rel in observed_paths:
if rel not in pin_fps:
reasons.append(
f"{phase}: observed dirty path '{rel}' has no fingerprint pin"
)
return {
"ok": not reasons,
"reasons": reasons,
"inventory": inv,
"observed_dirty_paths": observed_paths,
"expected_dirty_paths": pin_paths,
"observed_fingerprints": obs_fps,
"expected_fingerprints": pin_fps,
"phase": phase,
}
def build_journal_operation_identity(
*,
remote: str,
org: str,
repo: str,
claimant_identity: str | None,
claimant_profile: str | None,
) -> dict[str, str]:
"""Return the five-field durable operation identity for the recovery journal."""
return {
"remote": _text(remote),
"org": _text(org),
"repo": _text(repo),
"claimant_identity": _text(claimant_identity),
"claimant_profile": _text(claimant_profile),
}
def validate_journal_operation_identity(
journal: Mapping[str, Any] | None,
*,
remote: str,
org: str,
repo: str,
claimant_identity: str | None,
claimant_profile: str | None,
require_present: bool = True,
) -> list[str]:
"""Validate durable journal identity fields (#868 F2).
Rejects missing, mismatched, stale, cross-repository, or cross-claimant
journal state. When *require_present* is True, incomplete legacy journals
(any of the five fields absent/empty) fail closed.
"""
reasons: list[str] = []
if not isinstance(journal, Mapping):
if require_present:
reasons.append(
"recovery journal is missing or unreadable; complete operation "
"identity cannot be proven (fail closed)"
)
return reasons
expected = build_journal_operation_identity(
remote=remote,
org=org,
repo=repo,
claimant_identity=claimant_identity,
claimant_profile=claimant_profile,
)
for field in REQUIRED_JOURNAL_IDENTITY_FIELDS:
observed = _text(journal.get(field))
want = expected[field]
if not observed:
reasons.append(
f"recovery journal omits operation identity field '{field}' "
"(incomplete legacy or malformed journal identity; fail closed)"
)
continue
if not want:
reasons.append(
f"caller pin for journal identity field '{field}' is empty "
"(fail closed)"
)
continue
if observed != want:
reasons.append(
f"recovery journal identity mismatch for '{field}': "
f"journal={observed!r}, expected={want!r} "
"(cross-repository / cross-claimant / replay refused)"
)
return reasons
def journal_path(lock_dir: str, issue_number: int) -> str:
root = (lock_dir or "").strip()
return os.path.join(root, f".rebind-journal-{int(issue_number)}.json")
@@ -789,10 +975,22 @@ def _already_rebound(
existing_lock: Mapping[str, Any],
current_pid: int,
worktree_path: str,
expected_dirty_paths: Sequence[str] | None,
expected_fingerprints: Mapping[str, str],
worktree_for_fps: str,
remote: str,
org: str,
repo: str,
claimant_identity: str | None,
claimant_profile: str | None,
journal: Mapping[str, Any] | None = None,
) -> tuple[bool, list[str]]:
"""Return (True, notes) when lock is already rebound to this session."""
"""Return (True, notes) when lock is already rebound to this session.
#868: require complete matching operation identity (remote/org/repo/
claimant) and complete dirty-inventory revalidation, not fingerprint-only
checks. Incomplete or mismatched journal identity fails closed.
"""
notes: list[str] = []
pid = _recorded_pid(existing_lock)
try:
@@ -803,19 +1001,61 @@ def _already_rebound(
return False, []
if not _same_realpath(_text(existing_lock.get("worktree_path")), worktree_path):
return False, []
# Fingerprints must still match pins (byte preservation).
for rel, expected in (expected_fingerprints or {}).items():
abs_path = os.path.join(worktree_for_fps, rel)
try:
actual = content_fingerprint(abs_path)
except OSError as exc:
notes.append(f"could not re-fingerprint '{rel}' for already_rebound: {exc}")
return False, notes
if actual != _text(expected):
# Durable lock repo binding must still match the caller's target.
for field, expected in (("remote", remote), ("org", org), ("repo", repo)):
observed = _text(existing_lock.get(field))
want = _text(expected)
if observed and want and observed != want:
notes.append(
f"fingerprint drift on already-rebound check for '{rel}'"
f"already_rebound refused: lock {field}={observed!r} does not "
f"match expected {want!r} (cross-repository replay)"
)
return False, notes
lock_claimant = _lock_claimant(existing_lock)
locked_identity = _text(lock_claimant.get("username"))
locked_profile = _text(lock_claimant.get("profile"))
pin_identity = _text(claimant_identity)
pin_profile = _text(claimant_profile)
if pin_identity and locked_identity and pin_identity != locked_identity:
notes.append(
f"already_rebound refused: lock claimant '{locked_identity}' does "
f"not match pin '{pin_identity}' (cross-claimant replay)"
)
return False, notes
if pin_profile and locked_profile and pin_profile != locked_profile:
notes.append(
f"already_rebound refused: lock profile '{locked_profile}' does "
f"not match pin '{pin_profile}' (cross-claimant replay)"
)
return False, notes
# When a durable journal is present, require complete matching identity.
if isinstance(journal, Mapping) and journal:
id_reasons = validate_journal_operation_identity(
journal,
remote=remote,
org=org,
repo=repo,
claimant_identity=claimant_identity,
claimant_profile=claimant_profile,
require_present=True,
)
if id_reasons:
notes.extend(id_reasons)
return False, notes
inv_check = revalidate_complete_dirty_inventory(
worktree_for_fps,
expected_dirty_paths=expected_dirty_paths,
expected_fingerprints=expected_fingerprints,
phase="already_rebound",
)
if not inv_check["ok"]:
notes.extend(list(inv_check["reasons"] or []))
return False, notes
gen = lock_generation(existing_lock)
if gen < 1:
# A never-written generation is suspicious for a completed rebind, but
@@ -917,6 +1157,53 @@ def apply_dirty_same_claimant_session_rebind(
"remote_head": remote_head,
}
root_for_journal = (lock_dir or "").strip() or None
if root_for_journal is None and isinstance(existing_lock, Mapping):
root_for_journal = os.path.dirname(
_text(existing_lock.get("lock_file_path"))
or lock_file_path(
remote=remote, org=org, repo=repo, issue_number=issue_number
)
)
jpath_probe = (
journal_path(root_for_journal, issue_number) if root_for_journal else ""
)
existing_journal = _read_json(jpath_probe) if jpath_probe else None
# Resume / retry: reject incomplete, mismatched, or cross-repo journal
# identity before treating any prior journal as authoritative (#868 F2).
if isinstance(existing_journal, Mapping) and existing_journal:
journal_id_reasons = validate_journal_operation_identity(
existing_journal,
remote=remote,
org=org,
repo=repo,
claimant_identity=claimant_identity,
claimant_profile=claimant_profile,
require_present=True,
)
# Incomplete legacy journals from pre-#868 apply paths must fail closed
# when any identity field is missing — even if the rest of the payload
# looks familiar. Only a complete matching identity may proceed.
phase = _text(existing_journal.get("phase"))
if journal_id_reasons and phase not in ("", JOURNAL_PHASE_COMPLETE):
# Allow a completed journal with missing legacy identity only when
# already_rebound path will re-validate lock + inventory; for
# mid-flight incomplete journals, refuse.
if phase in (
JOURNAL_PHASE_ASSESSED,
JOURNAL_PHASE_PRE_BIND,
JOURNAL_PHASE_BOUND,
"bind_failed",
):
return {
**base_result,
"success": False,
"reasons": journal_id_reasons,
"journal_path": jpath_probe,
"journal_phase": phase or None,
}
# Retry-safe: if already rebound to this session, succeed even when assess
# refuses because old_pid no longer matches the (updated) lock.
if (
@@ -928,8 +1215,15 @@ def apply_dirty_same_claimant_session_rebind(
existing_lock=existing_lock,
current_pid=pid_now,
worktree_path=worktree_path,
expected_dirty_paths=expected_dirty_paths,
expected_fingerprints=expected_fingerprints,
worktree_for_fps=wt,
remote=remote,
org=org,
repo=repo,
claimant_identity=claimant_identity,
claimant_profile=claimant_profile,
journal=existing_journal,
)
if done:
lock_path = _text(existing_lock.get("lock_file_path")) or lock_file_path(
@@ -956,6 +1250,20 @@ def apply_dirty_same_claimant_session_rebind(
"generation_after": lock_generation(existing_lock),
"journal_phase": JOURNAL_PHASE_ALREADY_REBOUND,
}
# Same-pid candidate that failed complete identity/inventory checks
# must not fall through into a fresh bind that would re-mint authority.
if _recorded_pid(existing_lock) is not None:
try:
if int(_recorded_pid(existing_lock)) == int(pid_now) and notes:
return {
**base_result,
"success": False,
"already_rebound": False,
"reasons": notes,
"journal_path": jpath_probe or None,
}
except (TypeError, ValueError):
pass
if not assessment["rebind_sanctioned"]:
return base_result
@@ -982,6 +1290,25 @@ def apply_dirty_same_claimant_session_rebind(
issue_number,
)
op_identity = build_journal_operation_identity(
remote=remote,
org=org,
repo=repo,
claimant_identity=claimant_identity,
claimant_profile=claimant_profile,
)
# Refuse incomplete caller identity before any durable write.
for field, value in op_identity.items():
if not value:
return {
**base_result,
"success": False,
"reasons": [
f"cannot write recovery journal: operation identity field "
f"'{field}' is empty (fail closed)"
],
}
journal = {
"phase": JOURNAL_PHASE_ASSESSED,
"issue_number": issue_number,
@@ -996,36 +1323,40 @@ def apply_dirty_same_claimant_session_rebind(
"remote_head": remote_head,
"started_at": _utc_now_iso(),
"source": SOURCE,
# #868 F2 — complete durable operation identity
**op_identity,
}
_atomic_write_json(jpath, journal)
# Immediate pre-bind fingerprint re-verification.
pre_fps: dict[str, str] = {}
for rel in expected_dirty_paths or []:
abs_path = os.path.join(wt, rel)
try:
pre_fps[rel] = content_fingerprint(abs_path)
except OSError as exc:
return {
**base_result,
"success": False,
"reasons": [f"pre-bind fingerprint failed for '{rel}': {exc}"],
"journal_path": jpath,
}
for rel, expected in (expected_fingerprints or {}).items():
if pre_fps.get(rel) != _text(expected):
return {
**base_result,
"success": False,
"reasons": [
f"pre-bind fingerprint drift for '{rel}': "
f"observed {pre_fps.get(rel)}, expected {expected}"
],
"journal_path": jpath,
}
# #868 F1 — complete dirty-inventory revalidation immediately before mutation.
# Fail closed with no bind so failures cannot leave a newly authoritative
# live session.
pre_inv = revalidate_complete_dirty_inventory(
wt,
expected_dirty_paths=expected_dirty_paths,
expected_fingerprints=expected_fingerprints,
phase="pre-bind",
)
if not pre_inv["ok"]:
journal["phase"] = "pre_bind_inventory_failed"
journal["pre_bind_inventory"] = {
"observed_dirty_paths": pre_inv.get("observed_dirty_paths"),
"reasons": pre_inv.get("reasons"),
}
_atomic_write_json(jpath, journal)
return {
**base_result,
"success": False,
"reasons": list(pre_inv["reasons"] or []),
"journal_path": jpath,
"journal_phase": "pre_bind_inventory_failed",
}
pre_fps = dict(pre_inv.get("observed_fingerprints") or {})
journal["phase"] = JOURNAL_PHASE_PRE_BIND
journal["pre_bind_fingerprints"] = pre_fps
journal["pre_bind_dirty_paths"] = list(pre_inv.get("observed_dirty_paths") or [])
_atomic_write_json(jpath, journal)
now = _utc_now_iso()
@@ -1095,38 +1426,37 @@ def apply_dirty_same_claimant_session_rebind(
journal["lock_path"] = lock_path
_atomic_write_json(jpath, journal)
# Post-bind fingerprint verification — every byte unchanged.
post_fps: dict[str, str] = {}
for rel in expected_dirty_paths or []:
abs_path = os.path.join(wt, rel)
try:
post_fps[rel] = content_fingerprint(abs_path)
except OSError as exc:
return {
**base_result,
"success": False,
"reasons": [
f"post-bind fingerprint failed for '{rel}': {exc}; "
"lock may be rebound but content verification failed"
],
"lock_path": lock_path,
"journal_path": jpath,
"generation_before": gen_before,
}
for rel, expected in (expected_fingerprints or {}).items():
if post_fps.get(rel) != _text(expected):
return {
**base_result,
"success": False,
"reasons": [
f"post-bind fingerprint drift for '{rel}': "
f"observed {post_fps.get(rel)}, expected {expected}"
],
"lock_path": lock_path,
"journal_path": jpath,
"generation_before": gen_before,
"fingerprints_after": post_fps,
}
# #868 F1 — complete dirty-inventory revalidation immediately after mutation.
# Path set must remain exactly equal; fingerprints must be unchanged.
post_inv = revalidate_complete_dirty_inventory(
wt,
expected_dirty_paths=expected_dirty_paths,
expected_fingerprints=expected_fingerprints,
phase="post-bind",
)
post_fps = dict(post_inv.get("observed_fingerprints") or {})
if not post_inv["ok"]:
journal["phase"] = "post_bind_inventory_failed"
journal["post_bind_inventory"] = {
"observed_dirty_paths": post_inv.get("observed_dirty_paths"),
"reasons": post_inv.get("reasons"),
}
journal["post_bind_fingerprints"] = post_fps
_atomic_write_json(jpath, journal)
return {
**base_result,
"success": False,
"reasons": list(post_inv["reasons"] or []) + [
"post-bind complete inventory revalidation failed after "
"bind_session_lock; lock may be rebound but content/path "
"verification failed (fail closed)"
],
"lock_path": lock_path,
"journal_path": jpath,
"journal_phase": "post_bind_inventory_failed",
"generation_before": gen_before,
"fingerprints_after": post_fps,
}
# Remove stale session pointer for old_pid when it points at this lock.
removed_old_pointer = False
@@ -1154,6 +1484,9 @@ def apply_dirty_same_claimant_session_rebind(
journal["generation_after"] = gen_after
journal["removed_old_session_pointer"] = removed_old_pointer
journal["post_bind_fingerprints"] = post_fps
journal["post_bind_dirty_paths"] = list(
post_inv.get("observed_dirty_paths") or []
)
_atomic_write_json(jpath, journal)
return {
+11 -213
View File
@@ -11116,9 +11116,6 @@ def _collect_branch_ownership_records(
"""
records: list[dict] = []
inventory_error = False
# #855 AC4: expired/stale reviewer-lease records eligible for an explicit
# reclaim decision, evaluated after the full ownership inventory is built.
reviewer_reclaim_candidates: list[tuple[dict, bool | None]] = []
target_branch = (branch or "").strip()
if not target_branch:
return {"records": records, "inventory_error": False}
@@ -11267,28 +11264,15 @@ def _collect_branch_ownership_records(
else:
status = freshness_status
reclaim_allowed = False
rec = _base_rec(
category=category,
status=status,
reclaim_allowed=reclaim_allowed,
role=role,
host=lease_host or host_n or host,
records.append(
_base_rec(
category=category,
status=status,
reclaim_allowed=reclaim_allowed,
role=role,
host=lease_host or host_n or host,
)
)
records.append(rec)
# #855 AC4: a reviewer lease that is expired/stale (its owner
# gone) becomes a candidate for an explicit, fail-closed
# reclaim decision made once the full inventory is known.
if (
role == "reviewer"
and status
in branch_cleanup_guard._RECLAIMABLE_REVIEWER_STATUSES
):
owner_alive = (
fr.get("owner_pid_alive") if isinstance(fr, dict) else None
)
reviewer_reclaim_candidates.append(
(rec, owner_alive if isinstance(owner_alive, bool) else None)
)
except Exception:
# O1: fail closed on control-plane inventory errors.
inventory_error = True
@@ -11359,44 +11343,6 @@ def _collect_branch_ownership_records(
)
)
# #855 AC4: decide, explicitly and fail-closed, whether any expired/stale
# reviewer lease may stop protecting an already-merged branch. This runs
# only after the full ownership inventory is built, so a competing active
# claimant (an active lease, author session, worktree binding, or active
# reviewer comment lease) is visible. An inventory failure keeps every
# reclaim candidate protective (reclaim_allowed stays False).
if reviewer_reclaim_candidates and not inventory_error:
pr_merged_state: bool | None = None
if pr_number is not None and auth and base_api:
try:
pr_live = api_request(
"GET", f"{base_api}/pulls/{int(pr_number)}", auth
)
if isinstance(pr_live, dict) and pr_live:
pr_merged_state = bool(
pr_live.get("merged") or pr_live.get("merged_at")
)
except Exception:
# Unknown merged state fails closed (candidate stays protective).
pr_merged_state = None
for cand_rec, owner_alive in reviewer_reclaim_candidates:
competing = any(
other is not cand_rec
and branch_cleanup_guard.is_active_ownership_status(
other.get("status")
)
for other in records
)
decision = branch_cleanup_guard.assess_expired_reviewer_lease_reclaim(
role=str(cand_rec.get("role")),
status=str(cand_rec.get("status")),
pr_merged=pr_merged_state,
owner_pid_alive=owner_alive,
competing_active_claimant=competing,
)
cand_rec["reclaim_allowed"] = decision["reclaim_allowed"]
cand_rec["reclaim_decision"] = decision["decision"]
return {"records": records, "inventory_error": inventory_error}
@@ -11459,7 +11405,6 @@ def gitea_reconcile_merged_cleanups(
dry_run: bool = True,
execute_confirmed: bool = False,
limit: int = 50,
pr_number: int | None = None,
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
@@ -11470,11 +11415,7 @@ def gitea_reconcile_merged_cleanups(
Args:
dry_run: Defaults to True. When True, only builds the reconciliation report.
execute_confirmed: Must be True when dry_run=False.
limit: Max number of closed PRs to inspect (batch mode only; ignored when
``pr_number`` is set).
pr_number: Optional exact merged PR selector (#855). When set, only that
PR is assessed/acted on (fail closed if missing, unmerged, or
ambiguous). When omitted, existing batch behaviour is preserved.
limit: Max number of closed PRs to inspect.
remote: Known Gitea instance ('dadeschools' or 'prgs').
host: Override the Gitea host.
org: Override the owner/organization.
@@ -11509,120 +11450,11 @@ def gitea_reconcile_merged_cleanups(
"audit_phase": audit_reconciliation_mode.current_phase(),
}
# #855: optional exact PR pin. Fail closed before any inventory mutation.
exact_pr: int | None = None
if pr_number is not None:
try:
exact_pr = int(pr_number)
except (TypeError, ValueError):
return {
"success": False,
"performed": False,
"executed": False,
"dry_run": bool(dry_run),
"selection_mode": "exact_pr",
"selected_pr_number": pr_number,
"reasons": [
f"pr_number={pr_number!r} is not a valid integer "
"(fail closed; no mutation)"
],
"blocker_kind": "invalid_pr_number",
}
if exact_pr <= 0:
return {
"success": False,
"performed": False,
"executed": False,
"dry_run": bool(dry_run),
"selection_mode": "exact_pr",
"selected_pr_number": exact_pr,
"reasons": [
f"pr_number={exact_pr} must be a positive integer "
"(fail closed; no mutation)"
],
"blocker_kind": "invalid_pr_number",
}
h, o, r = _resolve(remote, host, org, repo)
auth = _auth(h)
base = repo_api_url(h, o, r)
selection_mode = "batch"
closed_prs: list[dict] = []
open_prs: list[dict] = []
if exact_pr is not None:
selection_mode = "exact_pr"
try:
pr_live = api_request("GET", f"{base}/pulls/{exact_pr}", auth)
except Exception as exc:
return {
"success": False,
"performed": False,
"executed": False,
"dry_run": bool(dry_run),
"selection_mode": selection_mode,
"selected_pr_number": exact_pr,
"reasons": [
f"PR #{exact_pr} could not be uniquely resolved "
f"(fail closed; no mutation): {_redact(str(exc))}"
],
"blocker_kind": "pr_unresolvable",
}
if not isinstance(pr_live, dict) or not pr_live:
return {
"success": False,
"performed": False,
"executed": False,
"dry_run": bool(dry_run),
"selection_mode": selection_mode,
"selected_pr_number": exact_pr,
"reasons": [
f"PR #{exact_pr} could not be uniquely resolved "
"(empty response; fail closed; no mutation)"
],
"blocker_kind": "pr_unresolvable",
}
live_number = pr_live.get("number")
try:
live_number_int = int(live_number) if live_number is not None else None
except (TypeError, ValueError):
live_number_int = None
if live_number_int != exact_pr:
return {
"success": False,
"performed": False,
"executed": False,
"dry_run": bool(dry_run),
"selection_mode": selection_mode,
"selected_pr_number": exact_pr,
"reasons": [
f"PR #{exact_pr} resolution is ambiguous or mismatched "
f"(live number={live_number!r}; fail closed; no mutation)"
],
"blocker_kind": "pr_ambiguous",
}
if not (pr_live.get("merged") or pr_live.get("merged_at")):
return {
"success": False,
"performed": False,
"executed": False,
"dry_run": bool(dry_run),
"selection_mode": selection_mode,
"selected_pr_number": exact_pr,
"reasons": [
f"PR #{exact_pr} is not merged "
"(exact-target cleanup requires a merged PR; "
"fail closed; no mutation)"
],
"blocker_kind": "pr_not_merged",
}
closed_prs = [pr_live]
# Exact mode still needs open heads for remote-delete safety gates.
open_prs = api_get_all(f"{base}/pulls?state=open", auth)
else:
# Preserve historical call order (closed then open) for batch callers/tests.
closed_prs = api_get_all(f"{base}/pulls?state=closed", auth, limit=limit)
open_prs = api_get_all(f"{base}/pulls?state=open", auth)
closed_prs = api_get_all(f"{base}/pulls?state=closed", auth, limit=limit)
open_prs = api_get_all(f"{base}/pulls?state=open", auth)
merged_closed: list[dict] = []
remote_branch_exists: dict[str, bool] = {}
@@ -11649,13 +11481,6 @@ def gitea_reconcile_merged_cleanups(
scratch_candidates = merged_cleanup_reconcile.discover_reviewer_scratch_worktrees(
_canonical_local_git_root()
)
# #855: exact-target never inventories or mutates foreign PR scratch trees.
if exact_pr is not None:
scratch_candidates = [
s
for s in scratch_candidates
if int(s.get("pr_number") or 0) == int(exact_pr)
]
active_reviewer_leases: dict[int, bool] = {}
pr_states: dict[int, dict] = {}
for scratch in scratch_candidates:
@@ -11690,33 +11515,6 @@ def gitea_reconcile_merged_cleanups(
active_reviewer_leases=active_reviewer_leases,
pr_states=pr_states,
)
report["selection_mode"] = selection_mode
if exact_pr is not None:
report["selected_pr_number"] = exact_pr
# Fail closed if exact pin somehow produced other or zero entries.
entries = list(report.get("entries") or [])
entry_numbers = []
for entry in entries:
try:
entry_numbers.append(int(entry.get("pr_number")))
except (TypeError, ValueError):
entry_numbers.append(entry.get("pr_number"))
if entry_numbers != [exact_pr]:
return {
"success": False,
"performed": False,
"executed": False,
"dry_run": bool(dry_run),
"selection_mode": selection_mode,
"selected_pr_number": exact_pr,
"reasons": [
f"exact PR #{exact_pr} selection produced unexpected "
f"candidate set {entry_numbers!r} "
"(fail closed; no mutation)"
],
"blocker_kind": "exact_selection_mismatch",
"entries": entries,
}
if dry_run:
report["dry_run"] = True
-352
View File
@@ -1639,358 +1639,6 @@ class TestSecondRemediationIntegration(unittest.TestCase):
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__":
unittest.main()
@@ -1,7 +1,10 @@
"""Integration tests for dirty same-claimant author-session rebind (#864).
"""Integration tests for dirty same-claimant author-session rebind (#864 / #868).
Uses real temp git repos/worktrees and a temp GITEA_ISSUE_LOCK_DIR. Does not
mutate any real #860/#864 worktree on disk.
mutate any real #860/#864/#868 worktree on disk.
#868 adds complete dirty-inventory revalidation around bind_session_lock and
complete recovery-journal operation identity (remote/org/repo/claimant).
"""
from __future__ import annotations
@@ -371,6 +374,7 @@ def test_retry_after_journal_mid_state(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
jpath = rebind.journal_path(lock_dir, ISSUE)
# Complete operation identity required for mid-flight resume (#868 F2).
rebind._atomic_write_json(
jpath,
{
@@ -379,6 +383,12 @@ def test_retry_after_journal_mid_state(dirty_repo, lock_dir):
"old_pid": old,
"new_pid": os.getpid(),
"expected_generation": 1,
"remote": REMOTE,
"org": ORG,
"repo": REPO,
"claimant_identity": IDENTITY,
"claimant_profile": PROFILE,
"source": rebind.SOURCE,
},
)
result = rebind.apply_dirty_same_claimant_session_rebind(
@@ -843,3 +853,494 @@ def test_content_fingerprint_stable(tmp_path):
b = rebind.content_fingerprint(str(p))
assert a == b
assert len(a) == 64
# ── #868 F1 — Complete dirty-inventory revalidation ─────────────────────────
def test_extra_tracked_dirty_path_before_bind_refused(dirty_repo, lock_dir):
"""Extra tracked dirty path appearing immediately before binding fails closed."""
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
wt = dirty_repo["worktree"]
# Seed a tracked file, then dirty it without including it in the pin set.
tracked_extra = "tracked_extra_before_bind.txt"
path = Path(wt) / tracked_extra
path.write_text("seed tracked extra\n", encoding="utf-8")
_git(wt, "add", tracked_extra)
_git(wt, "commit", "-q", "-m", "seed extra tracked")
# Heads moved — re-pin heads so only inventory disagreement is tested.
head = _git(wt, "rev-parse", "HEAD").stdout.strip()
_git(wt, "push", "-q", "origin", BRANCH)
remote_head = _git(wt, "rev-parse", f"refs/remotes/origin/{BRANCH}").stdout.strip()
path.write_text("dirty tracked extra\n", encoding="utf-8")
# Pins still describe the original inventory (without tracked_extra).
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(
dirty_repo,
lock,
lock_dir,
old_pid=old,
expected_local_head=head,
expected_remote_head=remote_head,
local_head=head,
remote_head=remote_head,
dirty_inventory=None, # force live recollect in apply
)
)
assert not result["success"]
joined = " ".join(result["reasons"])
assert "unexpected paths" in joined or "path-set disagreement" in joined
# Must not leave a newly authoritative live session for this pid.
rebound = ils.read_lock_file(lock["lock_file_path"])
assert rebound is not None
assert int(rebound.get("session_pid") or 0) == old
def test_extra_untracked_path_before_bind_refused(dirty_repo, lock_dir):
"""Extra untracked path appearing immediately before binding fails closed."""
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
wt = dirty_repo["worktree"]
extra = Path(wt) / "surprise_untracked_before_bind.txt"
extra.write_text("sneaky\n", encoding="utf-8")
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(
dirty_repo,
lock,
lock_dir,
old_pid=old,
dirty_inventory=None,
)
)
assert not result["success"]
joined = " ".join(result["reasons"])
assert "unexpected paths" in joined or "path-set disagreement" in joined
rebound = ils.read_lock_file(lock["lock_file_path"])
assert int(rebound.get("session_pid") or 0) == old
def test_path_added_during_mutation_window_refused(dirty_repo, lock_dir, monkeypatch):
"""Path added during the mutation window is detected by post-bind inventory."""
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
wt = dirty_repo["worktree"]
real_bind = ils.bind_session_lock
def _bind_then_add_path(lock_payload, **kwargs):
path = real_bind(lock_payload, **kwargs)
surprise = Path(wt) / "added_during_bind.txt"
surprise.write_text("during bind\n", encoding="utf-8")
return path
monkeypatch.setattr(rebind, "bind_session_lock", _bind_then_add_path)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(dirty_repo, lock, lock_dir, old_pid=old, dirty_inventory=None)
)
assert not result["success"]
joined = " ".join(result["reasons"])
assert "post-bind" in joined
assert "unexpected paths" in joined or "path-set disagreement" in joined
assert result.get("journal_phase") == "post_bind_inventory_failed"
def test_path_removed_during_mutation_window_refused(dirty_repo, lock_dir, monkeypatch):
"""Path removed during the mutation window is detected by post-bind inventory."""
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
wt = dirty_repo["worktree"]
victim = dirty_repo["dirty_paths"][-1] # prefer untracked for easy remove
real_bind = ils.bind_session_lock
def _bind_then_remove_path(lock_payload, **kwargs):
path = real_bind(lock_payload, **kwargs)
target = Path(wt) / victim
if target.exists():
target.unlink()
return path
monkeypatch.setattr(rebind, "bind_session_lock", _bind_then_remove_path)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(dirty_repo, lock, lock_dir, old_pid=old, dirty_inventory=None)
)
assert not result["success"]
joined = " ".join(result["reasons"])
assert "post-bind" in joined
assert "missing expected" in joined or "path-set disagreement" in joined
def test_fingerprint_movement_unchanged_path_set_refused(dirty_repo, lock_dir, monkeypatch):
"""Fingerprint movement with unchanged path set fails pre- or post-bind check."""
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
wt = dirty_repo["worktree"]
victim = dirty_repo["dirty_paths"][0]
real_bind = ils.bind_session_lock
def _bind_then_mutate_bytes(lock_payload, **kwargs):
path = real_bind(lock_payload, **kwargs)
target = Path(wt) / victim
target.write_text(target.read_text(encoding="utf-8") + "mutated\n", encoding="utf-8")
return path
monkeypatch.setattr(rebind, "bind_session_lock", _bind_then_mutate_bytes)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(dirty_repo, lock, lock_dir, old_pid=old, dirty_inventory=None)
)
assert not result["success"]
joined = " ".join(result["reasons"])
assert "fingerprint" in joined
assert "post-bind" in joined
# ── #868 F2 — Complete recovery-journal identity ────────────────────────────
def _complete_journal(**overrides):
base = {
"phase": rebind.JOURNAL_PHASE_ASSESSED,
"issue_number": ISSUE,
"branch_name": BRANCH,
"worktree_path": "/tmp/wt",
"old_pid": 1,
"new_pid": os.getpid(),
"remote": REMOTE,
"org": ORG,
"repo": REPO,
"claimant_identity": IDENTITY,
"claimant_profile": PROFILE,
"source": rebind.SOURCE,
}
base.update(overrides)
return base
def test_journal_remote_mismatch_refused(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
jpath = rebind.journal_path(lock_dir, ISSUE)
rebind._atomic_write_json(
jpath,
_complete_journal(
phase=rebind.JOURNAL_PHASE_PRE_BIND,
old_pid=old,
remote="dadeschools",
),
)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(dirty_repo, lock, lock_dir, old_pid=old)
)
assert not result["success"]
assert any("remote" in r and "mismatch" in r for r in result["reasons"])
def test_journal_org_mismatch_refused(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
jpath = rebind.journal_path(lock_dir, ISSUE)
rebind._atomic_write_json(
jpath,
_complete_journal(
phase=rebind.JOURNAL_PHASE_PRE_BIND,
old_pid=old,
org="Other-Org",
),
)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(dirty_repo, lock, lock_dir, old_pid=old)
)
assert not result["success"]
assert any("org" in r and "mismatch" in r for r in result["reasons"])
def test_journal_repo_mismatch_refused(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
jpath = rebind.journal_path(lock_dir, ISSUE)
rebind._atomic_write_json(
jpath,
_complete_journal(
phase=rebind.JOURNAL_PHASE_PRE_BIND,
old_pid=old,
repo="Other-Repo",
),
)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(dirty_repo, lock, lock_dir, old_pid=old)
)
assert not result["success"]
assert any("repo" in r and "mismatch" in r for r in result["reasons"])
def test_journal_claimant_identity_mismatch_refused(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
jpath = rebind.journal_path(lock_dir, ISSUE)
rebind._atomic_write_json(
jpath,
_complete_journal(
phase=rebind.JOURNAL_PHASE_PRE_BIND,
old_pid=old,
claimant_identity="intruder",
),
)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(dirty_repo, lock, lock_dir, old_pid=old)
)
assert not result["success"]
assert any("claimant_identity" in r for r in result["reasons"])
def test_journal_claimant_profile_mismatch_refused(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
jpath = rebind.journal_path(lock_dir, ISSUE)
rebind._atomic_write_json(
jpath,
_complete_journal(
phase=rebind.JOURNAL_PHASE_PRE_BIND,
old_pid=old,
claimant_profile="prgs-reviewer",
),
)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(dirty_repo, lock, lock_dir, old_pid=old)
)
assert not result["success"]
assert any("claimant_profile" in r for r in result["reasons"])
def test_incomplete_legacy_journal_identity_refused(dirty_repo, lock_dir):
"""Pre-#868 journals missing the five identity fields fail closed mid-flight."""
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
jpath = rebind.journal_path(lock_dir, ISSUE)
rebind._atomic_write_json(
jpath,
{
"phase": rebind.JOURNAL_PHASE_PRE_BIND,
"issue_number": ISSUE,
"old_pid": old,
"new_pid": os.getpid(),
# deliberately omit remote/org/repo/claimant_*
},
)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(dirty_repo, lock, lock_dir, old_pid=old)
)
assert not result["success"]
assert any("omits operation identity" in r for r in result["reasons"])
def test_journal_replay_cross_repository_refused(dirty_repo, lock_dir):
"""Replaying a journal from another repository is refused."""
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
jpath = rebind.journal_path(lock_dir, ISSUE)
rebind._atomic_write_json(
jpath,
_complete_journal(
phase=rebind.JOURNAL_PHASE_ASSESSED,
old_pid=old,
remote="dadeschools",
org="Other-Org",
repo="Other-Repo",
),
)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(dirty_repo, lock, lock_dir, old_pid=old)
)
assert not result["success"]
assert any("replay" in r or "mismatch" in r for r in result["reasons"])
def test_journal_replay_cross_claimant_refused(dirty_repo, lock_dir):
"""Replaying a journal from another claimant is refused."""
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
jpath = rebind.journal_path(lock_dir, ISSUE)
rebind._atomic_write_json(
jpath,
_complete_journal(
phase=rebind.JOURNAL_PHASE_ASSESSED,
old_pid=old,
claimant_identity="other-user",
claimant_profile="other-profile",
),
)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(dirty_repo, lock, lock_dir, old_pid=old)
)
assert not result["success"]
assert any("claimant" in r for r in result["reasons"])
def test_successful_exact_retry_with_complete_identity(dirty_repo, lock_dir):
"""Exact retry after success is already_rebound with complete matching identity."""
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
r1 = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(dirty_repo, lock, lock_dir, old_pid=old)
)
assert r1["success"], r1
# Journal must persist the five identity fields.
jpath = rebind.journal_path(lock_dir, ISSUE)
journal = rebind._read_json(jpath)
assert journal is not None
assert journal["phase"] == rebind.JOURNAL_PHASE_COMPLETE
for field in rebind.REQUIRED_JOURNAL_IDENTITY_FIELDS:
assert journal.get(field), field
assert journal["remote"] == REMOTE
assert journal["org"] == ORG
assert journal["repo"] == REPO
assert journal["claimant_identity"] == IDENTITY
assert journal["claimant_profile"] == PROFILE
rebound = ils.read_lock_file(r1["lock_path"])
r2 = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(
dirty_repo,
rebound,
lock_dir,
old_pid=old,
existing_lock=rebound,
)
)
assert r2["success"], r2
assert r2["already_rebound"] is True
def test_already_rebound_requires_complete_matching_identity(dirty_repo, lock_dir):
"""already_rebound with mismatched journal identity fails closed."""
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
r1 = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(dirty_repo, lock, lock_dir, old_pid=old)
)
assert r1["success"], r1
# Corrupt journal identity after success.
jpath = rebind.journal_path(lock_dir, ISSUE)
journal = rebind._read_json(jpath)
assert journal is not None
journal["claimant_identity"] = "not-the-owner"
rebind._atomic_write_json(jpath, journal)
rebound = ils.read_lock_file(r1["lock_path"])
r2 = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(
dirty_repo,
rebound,
lock_dir,
old_pid=old,
existing_lock=rebound,
)
)
# Same pid on lock would look like already_rebound, but identity must match.
assert not r2["success"]
assert any("claimant_identity" in r or "mismatch" in r for r in r2["reasons"])
def test_ordinary_dirty_worktree_refusal_preserved(dirty_repo):
"""#868 must not weaken ordinary dirty-worktree refusal on lock_issue path."""
porcelain = dirty_repo["inventory"]["porcelain_status"]
assessment = issue_lock_worktree.assess_issue_lock_worktree(
worktree_path=dirty_repo["worktree"],
current_branch=BRANCH,
porcelain_status=porcelain,
base_equivalent=False,
)
assert assessment["block"] is True
# ── #868 — Reconciler success path ──────────────────────────────────────────
def test_reconciler_success_path_tightly_pinned(dirty_repo, lock_dir):
"""Reconciler with authorize_reconciler_execute=True may execute rebind.
Reconciler execution grants no commit/push/publication/review/merge
capability — only the tightly pinned session rebind.
"""
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(
dirty_repo,
lock,
lock_dir,
old_pid=old,
role_kind="reconciler",
authorize_reconciler_execute=True,
# Reconciler may act for the recorded claimant without being that
# identity in the active session (still pin-checked against lock).
current_identity="sysadmin",
current_profile="prgs-reconciler",
)
)
assert result["success"], result
assert result["outcome"] == rebind.REBIND_SANCTIONED
rebound = ils.read_lock_file(result["lock_path"])
assert int(rebound["session_pid"]) == os.getpid()
# Provenance records the rebind tool; no publication authority is granted.
assert (
rebound.get("lock_provenance", {}).get("source")
== issue_lock_provenance.SOURCE_DIRTY_SAME_CLAIMANT_REBIND
)
# Reconciler rebind does not stamp commit/push/review/merge capabilities.
prov = rebound.get("lock_provenance") or {}
blob = json.dumps(prov)
for forbidden in (
"gitea.repo.commit",
"gitea.branch.push",
"gitea.pr.approve",
"gitea.pr.merge",
"gitea.pr.create",
):
assert forbidden not in blob
def test_revalidate_complete_dirty_inventory_helper(dirty_repo):
ok = rebind.revalidate_complete_dirty_inventory(
dirty_repo["worktree"],
expected_dirty_paths=dirty_repo["dirty_paths"],
expected_fingerprints=dirty_repo["fingerprints"],
phase="unit",
)
assert ok["ok"] is True
bad = rebind.revalidate_complete_dirty_inventory(
dirty_repo["worktree"],
expected_dirty_paths=dirty_repo["dirty_paths"][:-1],
expected_fingerprints={
p: dirty_repo["fingerprints"][p] for p in dirty_repo["dirty_paths"][:-1]
},
phase="unit",
)
assert bad["ok"] is False
assert any("unexpected paths" in r for r in bad["reasons"])
def test_validate_journal_operation_identity_helper():
complete = _complete_journal()
assert (
rebind.validate_journal_operation_identity(
complete,
remote=REMOTE,
org=ORG,
repo=REPO,
claimant_identity=IDENTITY,
claimant_profile=PROFILE,
)
== []
)
incomplete = {"phase": "assessed", "remote": REMOTE}
reasons = rebind.validate_journal_operation_identity(
incomplete,
remote=REMOTE,
org=ORG,
repo=REPO,
claimant_identity=IDENTITY,
claimant_profile=PROFILE,
)
assert any("omits operation identity" in r for r in reasons)
assert any("org" in r for r in reasons)
@@ -1,244 +0,0 @@
"""#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()