"""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 recovery (#753) also requires cleanliness. This module is the *only* sanctioned path that rebinds session/lock provenance onto the *same* worktree without touching tracked or untracked content. This is SEPARATE from #860 dirty-orphan recovery (PID-less + remote sync). This operation: * acts only on an already-registered dirty worktree * updates only stale lock/session provenance (PID, session pointer, generation, 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 import hashlib import json import os import subprocess import tempfile from datetime import datetime, timezone from typing import Any, Mapping, Sequence from author_mutation_worktree import is_path_under_branches from issue_lock_provenance import ( SOURCE_DIRTY_SAME_CLAIMANT_REBIND, build_sanctioned_lock_provenance, ) from issue_lock_store import ( AUTHOR_ISSUE_WORK_LEASE, bind_session_lock, is_process_alive, lock_file_path, lock_generation, read_lock_file, ) from reviewer_worktree import parse_dirty_tracked_files # Outcomes REBIND_SANCTIONED = "REBIND_SANCTIONED" REFUSED = "REFUSED" NO_CANDIDATE = "NO_CANDIDATE" # Provenance / tool identity SOURCE_TOOL = SOURCE_DIRTY_SAME_CLAIMANT_REBIND SOURCE = SOURCE_DIRTY_SAME_CLAIMANT_REBIND # Journal phases (crash-safe apply) JOURNAL_PHASE_ASSESSED = "assessed" JOURNAL_PHASE_PRE_BIND = "pre_bind" JOURNAL_PHASE_BOUND = "bound" JOURNAL_PHASE_COMPLETE = "complete" JOURNAL_PHASE_ALREADY_REBOUND = "already_rebound" REQUIRED_LOCK_FIELDS = ( "issue_number", "branch_name", "worktree_path", "remote", "org", "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") def _text(value: Any) -> str: return str(value or "").strip() def _same_realpath(left: str | None, right: str | None) -> bool: if not left or not right: return False try: return os.path.realpath(left) == os.path.realpath(right) except OSError: return left == right def _lock_claimant(lock: Mapping[str, Any]) -> dict[str, Any]: claimant = lock.get("claimant") if not isinstance(claimant, Mapping): lease = lock.get("work_lease") claimant = lease.get("claimant") if isinstance(lease, Mapping) else None return dict(claimant) if isinstance(claimant, Mapping) else {} def _recorded_pid(lock: Mapping[str, Any]) -> Any: pid = lock.get("session_pid") if pid is None: pid = lock.get("pid") return pid def content_fingerprint(path: str) -> str: """Return sha256 hex digest of file bytes at *path*. Missing or unreadable files raise ``OSError`` / ``FileNotFoundError`` so callers fail closed rather than inventing an empty hash. """ digest = hashlib.sha256() with open(path, "rb") as handle: while True: chunk = handle.read(1024 * 1024) if not chunk: break digest.update(chunk) return digest.hexdigest() def parse_dirty_paths(porcelain: str) -> list[str]: """Tracked + untracked paths from ``git status --porcelain -uall`` output.""" paths: list[str] = [] seen: set[str] = set() for line in (porcelain or "").splitlines(): if not line or len(line) < 4: continue if line.startswith("??"): path = line[3:].strip() else: path = line[3:].strip() if " -> " in path: path = path.split(" -> ", 1)[1].strip() if not path or path in seen: continue seen.add(path) paths.append(path) return paths def collect_dirty_inventory(worktree_path: str) -> dict[str, Any]: """Observe dirty tracked + untracked paths and content fingerprints. Uses ``git status --porcelain -uall`` so every untracked file is listed individually (not collapsed into a directory). """ path = (worktree_path or "").strip() if not path: return { "worktree_path": path, "porcelain_status": "", "dirty_paths": [], "fingerprints": {}, "ok": False, "reasons": ["worktree path is empty"], } status_res = subprocess.run( ["git", "-C", path, "status", "--porcelain", "-uall"], capture_output=True, text=True, check=False, ) if status_res.returncode != 0: err = (status_res.stderr or status_res.stdout or "").strip() return { "worktree_path": path, "porcelain_status": "", "dirty_paths": [], "fingerprints": {}, "ok": False, "reasons": [f"git status failed in '{path}': {err or 'unknown error'}"], } porcelain = status_res.stdout or "" dirty_paths = parse_dirty_paths(porcelain) fingerprints: dict[str, str] = {} reasons: list[str] = [] for rel in dirty_paths: abs_path = os.path.join(path, rel) if os.path.isdir(abs_path) and not os.path.islink(abs_path): # Directories appear only if git reported them; fingerprinting a # directory is not defined — fail closed. reasons.append(f"dirty path '{rel}' is a directory; cannot fingerprint") continue try: fingerprints[rel] = content_fingerprint(abs_path) except OSError as exc: reasons.append(f"could not fingerprint '{rel}': {exc}") return { "worktree_path": os.path.realpath(path), "porcelain_status": porcelain, "dirty_paths": dirty_paths, "fingerprints": fingerprints, "ok": not reasons, "reasons": reasons, "tracked_dirty": parse_dirty_tracked_files(porcelain), } 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") def _atomic_write_json(path: str, data: dict[str, Any]) -> None: parent = os.path.dirname(path) or "." os.makedirs(parent, mode=0o700, exist_ok=True) payload = json.dumps(data, indent=2, sort_keys=True) + "\n" fd, temp_path = tempfile.mkstemp(prefix=".rebind-j-", suffix=".json", dir=parent) try: with os.fdopen(fd, "w", encoding="utf-8") as handle: handle.write(payload) handle.flush() os.fsync(handle.fileno()) os.replace(temp_path, path) finally: if os.path.exists(temp_path): try: os.remove(temp_path) except OSError: pass def _read_json(path: str) -> dict[str, Any] | None: if not path or not os.path.exists(path): return None try: with open(path, encoding="utf-8") as handle: data = json.load(handle) except (OSError, json.JSONDecodeError): return None return data if isinstance(data, dict) else None def _malformed_lock_reasons(lock: Mapping[str, Any]) -> list[str]: missing: list[str] = [] for field in REQUIRED_LOCK_FIELDS: if not _text(lock.get(field)): missing.append(field) pid = _recorded_pid(lock) if pid is None or _text(pid) == "": missing.append("session_pid/pid") else: try: if int(pid) <= 0: missing.append("session_pid/pid") except (TypeError, ValueError): missing.append("session_pid/pid") return missing def _canonical_under_branches(worktree_path: str, repo_root: str | None) -> tuple[bool, list[str]]: """Prove worktree is a realpath under ``/branches/`` with no symlink escape.""" reasons: list[str] = [] path = (worktree_path or "").strip() if not path: return False, ["worktree path is empty"] try: real = os.path.realpath(path) except OSError as exc: return False, [f"worktree path could not be realpath-resolved: {exc}"] if not os.path.isdir(real): reasons.append(f"worktree path '{path}' is not an existing directory") root = (repo_root or "").strip() if root: try: root_real = os.path.realpath(root) except OSError as exc: return False, [f"repo root could not be realpath-resolved: {exc}"] if not is_path_under_branches(real, root_real): reasons.append( f"worktree '{real}' is not under branches/ of repo root '{root_real}' " "(unregistered/noncanonical worktree; fail closed)" ) # Symlink escape: the declared path must not resolve outside branches/. declared_abs = os.path.abspath(path) if os.path.islink(path) or declared_abs != real: if not is_path_under_branches(real, root_real): reasons.append( f"worktree path '{path}' escapes branches/ via symlink/realpath " f"(resolves to '{real}')" ) else: # Without an explicit repo root, still require a /branches/ segment. if not is_path_under_branches(real, None): reasons.append( f"worktree '{real}' is not under a branches/ directory " "(unregistered/noncanonical worktree; fail closed)" ) return not reasons, reasons def assess_dirty_same_claimant_session_rebind( *, remote: str, org: str, repo: str, issue_number: int, branch_name: str, worktree_path: str, claimant_identity: str | None, claimant_profile: str | None, old_pid: int | None, expected_local_head: str | None, expected_remote_head: str | None, expected_dirty_paths: Sequence[str] | None, expected_fingerprints: Mapping[str, str] | None, existing_lock: Mapping[str, Any] | None, current_identity: str | None, current_profile: str | None, role_kind: str | None, current_pid: int | None, current_branch: str | None, local_head: str | None, remote_head: str | None, porcelain_status: str | None = None, dirty_inventory: Mapping[str, Any] | None = None, competing_live_locks: Sequence[Mapping[str, Any]] | None = None, competing_sessions: Sequence[Mapping[str, Any]] | None = None, workflow_lease_active: bool = False, authorize_reconciler_execute: bool = False, permission_allowed: bool = False, repo_root: str | None = None, ) -> dict[str, Any]: """Pure assessment: may this dirty same-claimant lock be session-rebound? Every pin must agree. ``permission_allowed=True`` alone is never ownership proof. Fail closed on live old PID, foreign identity/profile, pin mismatch, unregistered/noncanonical worktree, head movement, dirty path/fingerprint disagreement, competing ownership, malformed lock, empty PID, wrong role. """ reasons: list[str] = [] evidence: dict[str, Any] = { "issue_number": issue_number, "branch_name": branch_name, "worktree_path": worktree_path, "remote": remote, "org": org, "repo": repo, "old_pid": old_pid, "current_pid": current_pid if current_pid is not None else os.getpid(), "role_kind": _text(role_kind).lower() or None, "permission_allowed": bool(permission_allowed), } if not existing_lock: return _assessment_result( NO_CANDIDATE, False, ["no existing durable lock for this issue; not a rebind candidate"], evidence, ) lock = dict(existing_lock) if lock.get("issue_number") != issue_number: return _assessment_result( NO_CANDIDATE, False, [ f"existing lock targets issue #{lock.get('issue_number')}, " f"not #{issue_number}; not a rebind candidate" ], evidence, ) missing = _malformed_lock_reasons(lock) if missing: return _assessment_result( REFUSED, False, [ "durable lock record is incomplete and cannot prove ownership " f"(missing/unusable: {', '.join(missing)})" ], evidence, ) recorded_pid = _recorded_pid(lock) evidence["recorded_pid"] = recorded_pid evidence["lock_generation"] = lock_generation(lock) # ── Role gate ─────────────────────────────────────────────────────────── role = _text(role_kind).lower() if role in {"reviewer", "merger"}: reasons.append( f"role '{role}' cannot rebind dirty same-claimant author sessions " "(fail closed)" ) elif role == "reconciler": if not authorize_reconciler_execute: reasons.append( "reconciler role requires authorize_reconciler_execute=True " "to execute dirty same-claimant rebind (fail closed)" ) elif role == "author": pass elif role: reasons.append(f"role '{role}' is not authorized for dirty same-claimant rebind") else: reasons.append("role_kind is unknown; dirty same-claimant rebind refused") # permission_allowed is explicitly NOT ownership proof evidence["note_permission_not_ownership"] = ( "permission_allowed is not treated as ownership proof" ) # ── Repository / issue / branch / worktree pins ───────────────────────── for field, expected in (("remote", remote), ("org", org), ("repo", repo)): actual = _text(lock.get(field)) if actual != _text(expected): reasons.append( f"lock {field} '{actual}' does not match requested '{_text(expected)}'" ) locked_branch = _text(lock.get("branch_name")) if locked_branch != _text(branch_name): reasons.append( f"lock branch '{locked_branch}' does not match requested " f"'{_text(branch_name)}'" ) checked_out = _text(current_branch) if not checked_out: reasons.append( "worktree is not on a named branch (detached HEAD); locked-branch " "occupancy could not be proven" ) elif checked_out != locked_branch: reasons.append( f"worktree is on branch '{checked_out}', not the locked branch " f"'{locked_branch}'" ) locked_worktree = _text(lock.get("worktree_path")) if not _same_realpath(locked_worktree, worktree_path): reasons.append( f"lock worktree '{locked_worktree}' does not match declared " f"'{_text(worktree_path)}'" ) evidence["locked_worktree_path"] = locked_worktree under_ok, under_reasons = _canonical_under_branches(worktree_path, repo_root) if not under_ok: reasons.extend(under_reasons) # ── old_pid pin + liveness ────────────────────────────────────────────── if old_pid is None or _text(old_pid) == "": reasons.append("old_pid pin is empty; rebind refused (fail closed)") else: try: old_pid_i = int(old_pid) except (TypeError, ValueError): reasons.append(f"old_pid '{old_pid}' is not a valid PID") old_pid_i = None if old_pid_i is not None: if old_pid_i <= 0: reasons.append("old_pid must be a positive integer (fail closed)") try: recorded_i = int(recorded_pid) except (TypeError, ValueError): recorded_i = None if recorded_i is None or recorded_i != old_pid_i: reasons.append( f"old_pid {old_pid_i} does not match lock session_pid/pid " f"{recorded_pid}" ) if is_process_alive(old_pid_i): reasons.append( f"old_pid {old_pid_i} is still alive; dirty same-claimant " "rebind requires a provably dead owner (fail closed)" ) evidence["old_pid_alive"] = is_process_alive(old_pid_i) if current_pid is not None: try: if int(current_pid) == old_pid_i: reasons.append( "old_pid is the current session PID; nothing to rebind" ) except (TypeError, ValueError): pass # ── Claimant identity / profile ───────────────────────────────────────── lock_claimant = _lock_claimant(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) active_identity = _text(current_identity) active_profile = _text(current_profile) evidence["locked_identity"] = locked_identity or None evidence["locked_profile"] = locked_profile or None if not locked_identity or not locked_profile: reasons.append( "durable lock does not record a claimant identity/profile; " "ownership could not be proven" ) if not pin_identity or not pin_profile: reasons.append( "claimant_identity/claimant_profile pins are required (fail closed)" ) if locked_identity and pin_identity and locked_identity != pin_identity: reasons.append( f"claimant_identity pin '{pin_identity}' does not match lock " f"claimant '{locked_identity}'" ) if locked_profile and pin_profile and locked_profile != pin_profile: reasons.append( f"claimant_profile pin '{pin_profile}' does not match lock profile " f"'{locked_profile}'" ) # Author path: active session must be the same claimant. Reconciler execute # may rebind for the recorded claimant when explicitly authorized. if role == "author": if not active_identity or not active_profile: reasons.append( "active session identity/profile is unknown; author ownership " "could not be proven" ) if locked_identity and active_identity and locked_identity != active_identity: reasons.append( f"lock claimant '{locked_identity}' does not match active " f"identity '{active_identity}' (foreign claimant refused)" ) if locked_profile and active_profile and locked_profile != active_profile: reasons.append( f"lock profile '{locked_profile}' does not match active profile " f"'{active_profile}' (profile mismatch refused)" ) if pin_identity and active_identity and pin_identity != active_identity: reasons.append( f"claimant_identity pin '{pin_identity}' does not match active " f"identity '{active_identity}'" ) if pin_profile and active_profile and pin_profile != active_profile: reasons.append( f"claimant_profile pin '{pin_profile}' does not match active " f"profile '{active_profile}'" ) # ── Heads (must match pins and each other for this rebind class) ──────── obs_local = _text(local_head) obs_remote = _text(remote_head) pin_local = _text(expected_local_head) pin_remote = _text(expected_remote_head) evidence["local_head"] = obs_local or None evidence["remote_head"] = obs_remote or None evidence["expected_local_head"] = pin_local or None evidence["expected_remote_head"] = pin_remote or None if not pin_local or not pin_remote: reasons.append( "expected_local_head and expected_remote_head pins are required " "(fail closed)" ) if not obs_local: reasons.append("local head SHA could not be determined") if not obs_remote: reasons.append("remote head SHA could not be determined") if pin_local and obs_local and pin_local != obs_local: reasons.append( f"local head moved or mismatched pin: observed {obs_local}, " f"expected {pin_local}" ) if pin_remote and obs_remote and pin_remote != obs_remote: reasons.append( f"remote head moved or mismatched pin: observed {obs_remote}, " f"expected {pin_remote}" ) if obs_local and obs_remote and obs_local != obs_remote: # Dirty rebind does not allow unpublished head movement; heads must agree. reasons.append( f"local head {obs_local} does not match remote head {obs_remote}; " "dirty same-claimant rebind requires matching heads (fail closed)" ) # ── Dirty inventory + fingerprint pins ────────────────────────────────── inv: dict[str, Any] if isinstance(dirty_inventory, Mapping) and dirty_inventory.get("dirty_paths") is not None: inv = dict(dirty_inventory) if not inv.get("fingerprints") and porcelain_status is not None: # Allow fingerprints-only refresh via recompute if needed. pass elif porcelain_status is not None: # Porcelain alone proves path set, not bytes. Fingerprints must come from # dirty_inventory (or apply()'s collect_dirty_inventory) — never from the # caller's expected_fingerprints pin (that would make the pin tautological). dirty_paths_obs = parse_dirty_paths(porcelain_status) inv = { "porcelain_status": porcelain_status, "dirty_paths": dirty_paths_obs, "fingerprints": {}, "ok": True, "reasons": [], } else: reasons.append( "neither dirty_inventory nor porcelain_status was provided; " "dirty state could not be proven" ) inv = {"dirty_paths": [], "fingerprints": {}, "ok": False} if inv.get("ok") is False and inv.get("reasons"): reasons.extend(list(inv.get("reasons") or [])) 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)}) evidence["observed_dirty_paths"] = observed_paths evidence["expected_dirty_paths"] = pin_paths if not pin_paths: reasons.append( "expected_dirty_paths pin is empty; dirty same-claimant rebind " "requires a non-empty dirty inventory pin (fail closed)" ) if set(observed_paths) != set(pin_paths): extra = sorted(set(observed_paths) - set(pin_paths)) missing_p = sorted(set(pin_paths) - set(observed_paths)) if extra: reasons.append( f"dirty path set disagreement: unexpected paths {extra}" ) if missing_p: reasons.append( f"dirty path set disagreement: missing expected paths {missing_p}" ) 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) } evidence["observed_fingerprints"] = obs_fps evidence["expected_fingerprints"] = pin_fps if not pin_fps: reasons.append( "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"expected_fingerprints contains '{rel}' which is not in " "expected_dirty_paths" ) actual_hash = obs_fps.get(rel) if not actual_hash: reasons.append( f"fingerprint missing for dirty path '{rel}'" ) elif actual_hash != expected_hash: reasons.append( f"fingerprint disagreement for '{rel}': observed " f"{actual_hash}, expected {expected_hash}" ) for rel in obs_fps: if rel in set(pin_paths) and rel not in pin_fps: reasons.append( f"expected_fingerprints missing pin for observed dirty path '{rel}'" ) # ── Competing ownership ───────────────────────────────────────────────── competing: list[dict[str, Any]] = [] for entry in competing_live_locks or (): if not isinstance(entry, Mapping): continue same_issue = entry.get("issue_number") == issue_number same_branch = _text(entry.get("branch_name")) == locked_branch if not (same_issue or same_branch): continue if ( same_issue and same_branch and _same_realpath(_text(entry.get("worktree_path")), worktree_path) ): # The lock we are rebinding is not competition with itself, but a # *live* competing owner on the same worktree is still a problem. entry_pid = entry.get("pid") or entry.get("session_pid") try: entry_pid_i = int(entry_pid) if entry_pid is not None else None except (TypeError, ValueError): entry_pid_i = None if entry_pid_i is not None and is_process_alive(entry_pid_i): if old_pid is None or entry_pid_i != int(old_pid): competing.append( { "issue_number": entry.get("issue_number"), "branch_name": entry.get("branch_name"), "worktree_path": entry.get("worktree_path"), "pid": entry_pid_i, } ) continue competing.append( { "issue_number": entry.get("issue_number"), "branch_name": entry.get("branch_name"), "worktree_path": entry.get("worktree_path"), "pid": entry.get("pid") or entry.get("session_pid"), } ) if competing: described = ", ".join( f"issue #{c['issue_number']} branch '{c['branch_name']}' pid={c.get('pid')}" for c in competing ) reasons.append(f"competing live lock exists ({described})") evidence["competing_live_locks"] = competing competing_sess: list[dict[str, Any]] = [] for entry in competing_sessions or (): if not isinstance(entry, Mapping): continue sess_pid = entry.get("pid") or entry.get("session_pid") try: sess_pid_i = int(sess_pid) if sess_pid is not None else None except (TypeError, ValueError): sess_pid_i = None if sess_pid_i is None: continue if current_pid is not None and sess_pid_i == int(current_pid): continue if old_pid is not None: try: if sess_pid_i == int(old_pid) and not is_process_alive(sess_pid_i): continue except (TypeError, ValueError): pass if is_process_alive(sess_pid_i) or entry.get("live") is True: competing_sess.append( { "pid": sess_pid_i, "lock_file_path": entry.get("lock_file_path"), } ) if competing_sess: reasons.append( "competing live session pointer(s) claim this lock: " + ", ".join(str(s["pid"]) for s in competing_sess) ) evidence["competing_sessions"] = competing_sess if workflow_lease_active: reasons.append( "workflow lease is active for this scope; dirty same-claimant " "rebind refused (fail closed)" ) evidence["workflow_lease_active"] = bool(workflow_lease_active) if reasons: return _assessment_result(REFUSED, False, reasons, evidence) proof = [ f"registered dirty worktree for issue #{issue_number} on branch " f"'{locked_branch}' matches claimant '{locked_identity}' / profile " f"'{locked_profile}'; old_pid {recorded_pid} is dead; heads " f"{obs_local} match; {len(pin_paths)} dirty paths fingerprint-pinned; " "provenance-only rebind sanctioned" ] return _assessment_result(REBIND_SANCTIONED, True, proof, evidence) def _assessment_result( outcome: str, sanctioned: bool, reasons: list[str], evidence: dict[str, Any], ) -> dict[str, Any]: return { "outcome": outcome, "rebind_sanctioned": sanctioned, "is_candidate": outcome != NO_CANDIDATE, "reasons": reasons, "evidence": evidence, "expected_generation": evidence.get("lock_generation"), } 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. #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: pid_i = int(pid) if pid is not None else None except (TypeError, ValueError): return False, [] if pid_i != int(current_pid): return False, [] if not _same_realpath(_text(existing_lock.get("worktree_path")), worktree_path): return False, [] # 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"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 # a same-pid lock with matching fingerprints is still "ours". notes.append("lock generation is 0; treating same-pid match as rebound") return True, notes or ["lock already bound to current session PID"] def apply_dirty_same_claimant_session_rebind( *, remote: str, org: str, repo: str, issue_number: int, branch_name: str, worktree_path: str, claimant_identity: str | None, claimant_profile: str | None, old_pid: int | None, expected_local_head: str | None, expected_remote_head: str | None, expected_dirty_paths: Sequence[str] | None, expected_fingerprints: Mapping[str, str] | None, existing_lock: Mapping[str, Any] | None, current_identity: str | None, current_profile: str | None, role_kind: str | None, current_pid: int | None = None, current_branch: str | None, local_head: str | None, remote_head: str | None, porcelain_status: str | None = None, dirty_inventory: Mapping[str, Any] | None = None, competing_live_locks: Sequence[Mapping[str, Any]] | None = None, competing_sessions: Sequence[Mapping[str, Any]] | None = None, workflow_lease_active: bool = False, authorize_reconciler_execute: bool = False, permission_allowed: bool = False, repo_root: str | None = None, dry_run: bool = False, lock_dir: str | None = None, ) -> dict[str, Any]: """Assess and (unless dry_run) apply a dirty same-claimant session rebind.""" pid_now = int(current_pid) if current_pid is not None else os.getpid() wt = os.path.realpath((worktree_path or "").strip()) if worktree_path else "" # Prefer a live inventory when applying so fingerprints are re-observed. inv = dict(dirty_inventory) if isinstance(dirty_inventory, Mapping) else None if inv is None and wt: inv = collect_dirty_inventory(wt) assessment = assess_dirty_same_claimant_session_rebind( remote=remote, org=org, repo=repo, issue_number=issue_number, branch_name=branch_name, worktree_path=worktree_path, claimant_identity=claimant_identity, claimant_profile=claimant_profile, old_pid=old_pid, expected_local_head=expected_local_head, expected_remote_head=expected_remote_head, expected_dirty_paths=expected_dirty_paths, expected_fingerprints=expected_fingerprints, existing_lock=existing_lock, current_identity=current_identity, current_profile=current_profile, role_kind=role_kind, current_pid=pid_now, current_branch=current_branch, local_head=local_head, remote_head=remote_head, porcelain_status=porcelain_status if porcelain_status is not None else (inv or {}).get("porcelain_status"), dirty_inventory=inv, competing_live_locks=competing_live_locks, competing_sessions=competing_sessions, workflow_lease_active=workflow_lease_active, authorize_reconciler_execute=authorize_reconciler_execute, permission_allowed=permission_allowed, repo_root=repo_root, ) base_result: dict[str, Any] = { "success": False, "dry_run": bool(dry_run), "outcome": assessment["outcome"], "rebind_sanctioned": assessment["rebind_sanctioned"], "reasons": list(assessment.get("reasons") or []), "evidence": assessment.get("evidence") or {}, "old_pid": old_pid, "new_pid": pid_now, "already_rebound": False, "dirty_paths": list(expected_dirty_paths or []), "fingerprints": dict(expected_fingerprints or {}), "local_head": local_head, "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 ( isinstance(existing_lock, Mapping) and expected_fingerprints and wt ): done, notes = _already_rebound( 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( remote=remote, org=org, repo=repo, issue_number=issue_number, lock_dir=lock_dir, ) session_ptr = os.path.join( (lock_dir or os.path.dirname(lock_path) or "."), f"session-{pid_now}.json", ) return { **base_result, "success": True, "outcome": REBIND_SANCTIONED, "rebind_sanctioned": True, "already_rebound": True, "reasons": notes, "lock_path": lock_path, "session_pointer": session_ptr, "generation_before": lock_generation(existing_lock), "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 if dry_run: return { **base_result, "success": True, "message": "dry_run: rebind sanctioned; no lock/session writes performed", "generation_before": assessment.get("expected_generation"), "generation_after": assessment.get("expected_generation"), } lock = dict(existing_lock or {}) gen_before = lock_generation(lock) root = (lock_dir or "").strip() or None jpath = journal_path( root or os.path.dirname( _text(lock.get("lock_file_path")) or lock_file_path( remote=remote, org=org, repo=repo, issue_number=issue_number ) ), 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, "branch_name": branch_name, "worktree_path": wt, "old_pid": old_pid, "new_pid": pid_now, "expected_generation": gen_before, "expected_fingerprints": dict(expected_fingerprints or {}), "expected_dirty_paths": list(expected_dirty_paths or []), "local_head": local_head, "remote_head": remote_head, "started_at": _utc_now_iso(), "source": SOURCE, # #868 F2 — complete durable operation identity **op_identity, } _atomic_write_json(jpath, journal) # #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() new_lock = dict(lock) new_lock["session_pid"] = pid_now new_lock["pid"] = pid_now new_lock["last_heartbeat_at"] = now new_lock["remote"] = remote new_lock["org"] = org new_lock["repo"] = repo new_lock["issue_number"] = issue_number new_lock["branch_name"] = branch_name new_lock["worktree_path"] = _text(lock.get("worktree_path")) or wt # Preserve work_lease (including expires_at); refresh heartbeat only. lease = new_lock.get("work_lease") if isinstance(lease, dict): lease = dict(lease) lease["last_heartbeat_at"] = now if not lease.get("operation_type"): lease["operation_type"] = AUTHOR_ISSUE_WORK_LEASE new_lock["work_lease"] = lease claimant = _lock_claimant(lock) new_lock["lock_provenance"] = build_sanctioned_lock_provenance( tool=SOURCE_TOOL, source=SOURCE, claimant=claimant or { "username": claimant_identity, "profile": claimant_profile, }, ) new_lock["rebind_record"] = { "source": SOURCE, "old_pid": old_pid, "new_pid": pid_now, "rebound_at": now, "local_head": local_head, "remote_head": remote_head, "dirty_path_count": len(list(expected_dirty_paths or [])), "generation_before": gen_before, "evidence": { "fingerprints": dict(expected_fingerprints or {}), "dirty_paths": list(expected_dirty_paths or []), }, } try: lock_path = bind_session_lock( new_lock, lock_dir=root, expected_generation=gen_before, ) except Exception as exc: journal["phase"] = "bind_failed" journal["error"] = str(exc) _atomic_write_json(jpath, journal) return { **base_result, "success": False, "reasons": [f"bind_session_lock failed: {exc}"], "journal_path": jpath, "generation_before": gen_before, } journal["phase"] = JOURNAL_PHASE_BOUND journal["lock_path"] = lock_path _atomic_write_json(jpath, journal) # #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 if old_pid is not None and root: old_ptr = os.path.join(root, f"session-{int(old_pid)}.json") if os.path.exists(old_ptr): ptr = _read_json(old_ptr) or {} ptr_lock = _text(ptr.get("lock_file_path")) if not ptr_lock or os.path.realpath(ptr_lock) == os.path.realpath(lock_path): try: os.remove(old_ptr) removed_old_pointer = True except OSError: pass bound = read_lock_file(lock_path) or new_lock gen_after = lock_generation(bound) session_ptr = os.path.join( root or os.path.dirname(lock_path), f"session-{pid_now}.json", ) journal["phase"] = JOURNAL_PHASE_COMPLETE journal["completed_at"] = _utc_now_iso() 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 { **base_result, "success": True, "message": ( f"Rebound dirty same-claimant author session for issue #{issue_number} " f"from dead pid {old_pid} to pid {pid_now}; dirty bytes preserved" ), "lock_path": lock_path, "session_pointer": session_ptr, "generation_before": gen_before, "generation_after": gen_after, "fingerprints": post_fps, "fingerprints_before": pre_fps, "removed_old_session_pointer": removed_old_pointer, "journal_path": jpath, "journal_phase": JOURNAL_PHASE_COMPLETE, "rebind_record": bound.get("rebind_record") or new_lock.get("rebind_record"), "lock_provenance": bound.get("lock_provenance"), } def build_issue_860_regression_fixture_spec() -> dict[str, Any]: """Data-only fixture describing the #860 class scenario (no real mutation). Claimant jcwalker3 / prgs-author, dead PID, no live session pointer, seven dirty paths with fingerprint pins, matching local/remote heads. """ dirty_paths = [ "dirty_same_claimant_session_rebind.py", "issue_lock_provenance.py", "task_capability_map.py", "gitea_mcp_server.py", "tests/test_dirty_same_claimant_session_rebind.py", "docs/runbook-dirty-rebind.md", "scratch/notes-untracked.txt", ] # Stable placeholder digests — tests replace with real fingerprints when # constructing on-disk fixtures. These exist so the spec is self-describing. fingerprints = { path: hashlib.sha256(f"issue-860-fixture:{path}".encode()).hexdigest() for path in dirty_paths } head = "a" * 40 dead = 424860 return { "issue_class": "issue-860-dirty-orphan-class-fixture", "description": ( "Registered dirty worktree, same claimant, dead owner PID, no live " "session pointer, seven fingerprint-pinned dirty paths, matching heads. " "Data only — does not mutate any real worktree." ), "remote": "prgs", "org": "Scaled-Tech-Consulting", "repo": "Gitea-Tools", "issue_number": 860, "branch_name": "fix/issue-860-dirty-orphan-recovery", "worktree_path": "/scratch/branches/fix-issue-860-dirty-orphan-recovery", "claimant_identity": "jcwalker3", "claimant_profile": "prgs-author", "old_pid": dead, "old_pid_alive": False, "live_session_pointer": None, "expected_local_head": head, "expected_remote_head": head, "expected_dirty_paths": dirty_paths, "expected_fingerprints": fingerprints, "dirty_path_count": 7, "role_kind": "author", "notes": [ "Distinct from #864 apply path: this fixture documents the #860 class " "inputs (dead PID + dirty inventory) without remote sync or recovery " "worktree creation.", ], }