Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
89657a06c4 | ||
|
|
ba3ea3012c |
@@ -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,
|
||||
|
||||
@@ -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
@@ -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
|
||||
|
||||
@@ -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()
|
||||
Reference in New Issue
Block a user