diff --git a/author_issue_bootstrap.py b/author_issue_bootstrap.py index 358ef32..c246698 100644 --- a/author_issue_bootstrap.py +++ b/author_issue_bootstrap.py @@ -126,48 +126,88 @@ def derive_default_idempotency_key( def run_compensating_recovery( journal: dict[str, Any], canonical_repo_root: str, + journal_dir: str | None = None, ) -> dict[str, Any]: """Execute compensating recovery for artifacts created by this transition only.""" artifacts = journal.get("artifacts_created") or {} + pending = journal.get("pending_creations") or {} rolled_back: list[str] = [] worktree_path = journal.get("worktree_path") branch_name = journal.get("branch_name") - if ( - artifacts.get("worktree_registered") or artifacts.get("worktree_dir_created") - ) and worktree_path: + # Roll back issue lock if created + if artifacts.get("lock_created") or pending.get("lock"): + issue_num = journal.get("issue_number") + session_id = journal.get("owner_session") + if issue_num and session_id: + try: + issue_lock_store.release_session_lock( + issue_number=issue_num, + session=session_id, + lock_dir=journal_dir, + ) + rolled_back.append(f"lock:issue-{issue_num}") + except Exception: + pass + artifacts["lock_created"] = False + + worktree_created = ( + artifacts.get("worktree_registered") + or artifacts.get("worktree_dir_created") + or (pending.get("worktree_path") == worktree_path and worktree_path) + ) + + if worktree_created and worktree_path: if os.path.exists(worktree_path): - try: - subprocess.run( - [ - "git", - "-C", - canonical_repo_root, - "worktree", - "remove", - "--force", - worktree_path, - ], - capture_output=True, - text=True, - check=False, - ) - except Exception: - pass - if os.path.exists(worktree_path): - shutil.rmtree(worktree_path, ignore_errors=True) - try: - subprocess.run( - ["git", "-C", canonical_repo_root, "worktree", "prune"], - capture_output=True, - text=True, - check=False, - ) - except Exception: - pass + # Re-verify cleanliness before destructive removal (F-4) + porc_res = subprocess.run( + ["git", "-C", worktree_path, "status", "--porcelain"], + capture_output=True, + text=True, + check=False, + ) + is_dirty = porc_res.returncode == 0 and bool(porc_res.stdout.strip()) + if is_dirty: + rolled_back.append(f"worktree_path_preserved_dirty:{worktree_path}") + else: + try: + subprocess.run( + [ + "git", + "-C", + canonical_repo_root, + "worktree", + "remove", + "--force", + worktree_path, + ], + capture_output=True, + text=True, + check=False, + ) + except Exception: + pass + if os.path.exists(worktree_path): + shutil.rmtree(worktree_path, ignore_errors=True) + try: + subprocess.run( + ["git", "-C", canonical_repo_root, "worktree", "prune"], + capture_output=True, + text=True, + check=False, + ) + except Exception: + pass + rolled_back.append(f"worktree_path:{worktree_path}") + else: rolled_back.append(f"worktree_path:{worktree_path}") - if artifacts.get("branch_created") and branch_name: + branch_created = ( + artifacts.get("branch_created") + or (pending.get("branch_name") == branch_name and branch_name) + ) + + if branch_created and branch_name: try: res = subprocess.run( [ @@ -183,20 +223,38 @@ def run_compensating_recovery( check=False, ) if res.returncode == 0: - subprocess.run( + # Check for author commits on branch before branch deletion (F-4) + resolved_base = journal.get("resolved_base_sha") or "master" + rev_list_res = subprocess.run( [ "git", "-C", canonical_repo_root, - "branch", - "-D", - branch_name, + "rev-list", + f"{resolved_base}..{branch_name}", ], capture_output=True, text=True, check=False, ) - rolled_back.append(f"branch:{branch_name}") + has_commits = rev_list_res.returncode == 0 and bool(rev_list_res.stdout.strip()) + if has_commits: + rolled_back.append(f"branch_preserved_commits:{branch_name}") + else: + subprocess.run( + [ + "git", + "-C", + canonical_repo_root, + "branch", + "-D", + branch_name, + ], + capture_output=True, + text=True, + check=False, + ) + rolled_back.append(f"branch:{branch_name}") except Exception: pass @@ -207,7 +265,7 @@ def run_compensating_recovery( } journal["compensating_recovery"] = recovery_info journal["current_phase"] = PHASE_COMPENSATING_RECOVERY - save_phase_journal(journal) + save_phase_journal(journal, journal_dir=journal_dir) return recovery_info @@ -303,6 +361,7 @@ def assess_author_issue_bootstrap( import fcntl +import stat class BootstrapTransitionLock: @@ -313,23 +372,53 @@ class BootstrapTransitionLock: c if c.isalnum() or c in ("-", "_", ".") else "_" for c in idempotency_key ) - lock_dir = get_journal_dir(journal_dir) - self.lock_path = os.path.join(lock_dir, f"{safe_key}.lock") + lock_dir = os.path.realpath(get_journal_dir(journal_dir)) + if not os.path.isdir(lock_dir): + raise RuntimeError(f"Lock directory '{lock_dir}' does not exist or is not a directory") + + raw_lock_path = os.path.abspath(os.path.join(lock_dir, f"{safe_key}.lock")) + try: + common = os.path.commonpath([lock_dir, os.path.dirname(raw_lock_path)]) + except Exception: + common = None + if common != lock_dir: + raise RuntimeError(f"Lock path '{raw_lock_path}' escapes canonical lock directory '{lock_dir}'") + + self.lock_path = raw_lock_path self.fd = None def __enter__(self): - self.fd = open(self.lock_path, "a+") + if os.path.islink(self.lock_path): + raise RuntimeError(f"Refusing lock acquisition: lock path '{self.lock_path}' is a symlink") + + flags = os.O_RDWR | os.O_CREAT + if hasattr(os, "O_NOFOLLOW"): + flags |= os.O_NOFOLLOW + if hasattr(os, "O_CLOEXEC"): + flags |= os.O_CLOEXEC + + try: + fd = os.open(self.lock_path, flags, 0o600) + except OSError as exc: + raise RuntimeError(f"Failed to open lock file safely '{self.lock_path}': {exc}") from exc + + st = os.fstat(fd) + if not stat.S_ISREG(st.st_mode): + os.close(fd) + raise RuntimeError(f"Lock target '{self.lock_path}' is not a regular file") + + self.fd = fd fcntl.flock(self.fd, fcntl.LOCK_EX) return self def __exit__(self, exc_type, exc_val, exc_tb): - if self.fd: + if self.fd is not None: try: fcntl.flock(self.fd, fcntl.LOCK_UN) except Exception: pass try: - self.fd.close() + os.close(self.fd) except Exception: pass self.fd = None @@ -358,6 +447,35 @@ def bootstrap_author_issue_worktree( """Execute the sanctioned author issue worktree bootstrap transition.""" root = os.path.realpath(canonical_repo_root) + session = (owner_session or "").strip() + if not session: + return { + "success": False, + "reason_code": "missing_owner_session", + "message": "Missing required owner_session parameter (fail closed). Session identifier cannot be fabricated or defaulted.", + "exact_next_action": ( + "Pass explicit owner_session resolved from gitea_whoami or session context." + ), + } + + if active_identity is None or not str(active_identity).strip(): + return { + "success": False, + "reason_code": "missing_active_identity", + "message": "Missing required active_identity parameter (fail closed). Identity cannot be fabricated or defaulted.", + "exact_next_action": "Pass explicit active_identity resolved from gitea_whoami.", + } + identity = str(active_identity).strip() + + if active_profile is None or not str(active_profile).strip(): + return { + "success": False, + "reason_code": "missing_active_profile", + "message": "Missing required active_profile parameter (fail closed). Profile cannot be fabricated or defaulted.", + "exact_next_action": "Pass explicit active_profile resolved from gitea_whoami.", + } + profile = str(active_profile).strip() + # Derive standard inputs expected_pattern = f"issue-{issue_number}" target_branch = (branch_name or "").strip() @@ -393,9 +511,9 @@ def bootstrap_author_issue_worktree( ) # Acquire cross-process file lock scoped to the idempotency key / transition identity - with BootstrapTransitionLock(key): + with BootstrapTransitionLock(key, journal_dir=lock_dir): # Idempotency check - existing = load_phase_journal(key) + existing = load_phase_journal(key, journal_dir=lock_dir) if existing and existing.get("completed"): if ( existing.get("issue_number") == issue_number @@ -418,7 +536,7 @@ def bootstrap_author_issue_worktree( "idempotency_key": key, "phase_journal": existing, "exact_next_action": ( - f"Call gitea_whoami, then gitea_resolve_task_capability(task='work_issue', worktree_path='{target_worktree}') " + "Call gitea_whoami, then gitea_resolve_task_capability(task='work_issue') " "and proceed with author implementation in the bootstrapped worktree." ), } @@ -454,9 +572,9 @@ def bootstrap_author_issue_worktree( "resolved_base_sha": None, "branch_name": target_branch, "worktree_path": target_worktree, - "active_identity": active_identity, - "active_profile": active_profile, - "owner_session": owner_session, + "active_identity": identity, + "active_profile": profile, + "owner_session": session, "remote": remote, "org": org, "repo": repo, @@ -496,7 +614,7 @@ def bootstrap_author_issue_worktree( journal["failure_reason"] = ( f"stale concurrency pin: expected {exp_norm[:12]} != live {live_norm[:12]}" ) - save_phase_journal(journal) + save_phase_journal(journal, journal_dir=lock_dir) return { "success": False, "reason_code": "stale_concurrency_pin", @@ -517,7 +635,7 @@ def bootstrap_author_issue_worktree( "expected_base_sha": expected_base_sha, } journal["current_phase"] = PHASE_2_BRANCH_CONFIRMED - save_phase_journal(journal) + save_phase_journal(journal, journal_dir=lock_dir) if dry_run: return { @@ -534,6 +652,8 @@ def bootstrap_author_issue_worktree( # Phase 2: BRANCH_CONFIRMED was_branch_created_previously = journal["artifacts_created"].get("branch_created", False) + pending_branch = (journal.get("pending_creations") or {}).get("branch_name") + branch_check = subprocess.run( ["git", "-C", root, "rev-parse", "--verify", target_branch], capture_output=True, @@ -558,23 +678,39 @@ def bootstrap_author_issue_worktree( check=False, ) if anc_check.returncode != 0 and branch_head.lower() != live_master_sha.lower(): - journal["failure_reason"] = ( - f"existing branch '{target_branch}' HEAD ({branch_head[:12]}) does not descend from base ({live_master_sha[:12]})" + # F-8: Check if branch shares a common ancestor with live master + mb_check = subprocess.run( + ["git", "-C", root, "merge-base", live_master_sha, branch_head], + capture_output=True, + text=True, + check=False, ) - save_phase_journal(journal) - return { - "success": False, - "reason_code": "incompatible_existing_branch", - "message": ( - f"Existing branch '{target_branch}' HEAD ({branch_head[:12]}) is incompatible with live master ({live_master_sha[:12]})." - ), - "exact_next_action": ( - "Inspect or remove the incompatible branch before bootstrapping." - ), - } + if mb_check.returncode != 0 or not mb_check.stdout.strip(): + journal["failure_reason"] = ( + f"existing branch '{target_branch}' HEAD ({branch_head[:12]}) is incompatible with live master ({live_master_sha[:12]})" + ) + save_phase_journal(journal, journal_dir=lock_dir) + return { + "success": False, + "reason_code": "incompatible_existing_branch", + "message": ( + f"Existing branch '{target_branch}' HEAD ({branch_head[:12]}) is incompatible with live master ({live_master_sha[:12]})." + ), + "exact_next_action": ( + "Inspect or sync the existing branch with master before bootstrapping." + ), + } # Preserve creation provenance monotonically across interruption and replay - journal["artifacts_created"]["branch_created"] = was_branch_created_previously + if was_branch_created_previously or pending_branch == target_branch: + journal["artifacts_created"]["branch_created"] = True + else: + journal["artifacts_created"]["branch_created"] = False else: + # Persist creation intent/provenance to disk BEFORE executing external mutation + journal.setdefault("pending_creations", {})["branch_name"] = target_branch + journal["artifacts_created"]["branch_created"] = True + save_phase_journal(journal, journal_dir=lock_dir) + # Create branch create_res = subprocess.run( ["git", "-C", root, "branch", target_branch, live_master_sha], @@ -583,17 +719,18 @@ def bootstrap_author_issue_worktree( check=False, ) if create_res.returncode != 0: + journal["artifacts_created"]["branch_created"] = False + journal.get("pending_creations", {}).pop("branch_name", None) journal["failure_reason"] = ( f"failed to create git branch '{target_branch}': {create_res.stderr.strip()}" ) - save_phase_journal(journal) + save_phase_journal(journal, journal_dir=lock_dir) return { "success": False, "reason_code": "branch_creation_failed", "message": f"Failed to create git branch '{target_branch}': {create_res.stderr.strip()}", "exact_next_action": "Verify branch availability and retry.", } - journal["artifacts_created"]["branch_created"] = True journal["phases"][PHASE_2_BRANCH_CONFIRMED] = { "status": "completed", @@ -601,7 +738,7 @@ def bootstrap_author_issue_worktree( "created": journal["artifacts_created"]["branch_created"], } journal["current_phase"] = PHASE_3_PATH_RESERVED - save_phase_journal(journal) + save_phase_journal(journal, journal_dir=lock_dir) # Phase 3: PATH_RESERVED & Phase 4: WORKTREE_CONFIRMED if not author_mutation_worktree.is_path_under_branches( @@ -610,7 +747,7 @@ def bootstrap_author_issue_worktree( journal["failure_reason"] = ( f"target_worktree '{target_worktree}' is outside canonical branches/ root" ) - run_compensating_recovery(journal, root) + run_compensating_recovery(journal, root, journal_dir=lock_dir) return { "success": False, "reason_code": "path_outside_canonical_branches_root", @@ -624,6 +761,7 @@ def bootstrap_author_issue_worktree( was_dir_created_previously = journal["artifacts_created"].get("worktree_dir_created", False) was_registered_previously = journal["artifacts_created"].get("worktree_registered", False) + pending_wt = (journal.get("pending_creations") or {}).get("worktree_path") dir_exists = os.path.exists(target_worktree) if dir_exists: @@ -643,7 +781,7 @@ def bootstrap_author_issue_worktree( journal["failure_reason"] = ( f"target worktree '{target_worktree}' contains dirty tracked/untracked files" ) - run_compensating_recovery(journal, root) + run_compensating_recovery(journal, root, journal_dir=lock_dir) return { "success": False, "reason_code": "preexisting_dirty_worktree", @@ -662,7 +800,7 @@ def bootstrap_author_issue_worktree( journal["failure_reason"] = ( f"existing worktree '{target_worktree}' is on branch '{wt_branch}' != expected '{target_branch}'" ) - run_compensating_recovery(journal, root) + run_compensating_recovery(journal, root, journal_dir=lock_dir) return { "success": False, "reason_code": "incompatible_existing_directory", @@ -673,11 +811,19 @@ def bootstrap_author_issue_worktree( "Inspect or remove the pre-existing worktree folder before bootstrapping." ), } - journal["artifacts_created"]["worktree_dir_created"] = was_dir_created_previously - journal["artifacts_created"]["worktree_registered"] = ( - was_registered_previously or was_dir_created_previously - ) + if was_dir_created_previously or pending_wt == target_worktree: + journal["artifacts_created"]["worktree_dir_created"] = True + journal["artifacts_created"]["worktree_registered"] = True + else: + journal["artifacts_created"]["worktree_dir_created"] = False + journal["artifacts_created"]["worktree_registered"] = False else: + # Persist creation intent/provenance to disk BEFORE executing external worktree add mutation + journal.setdefault("pending_creations", {})["worktree_path"] = target_worktree + journal["artifacts_created"]["worktree_dir_created"] = True + journal["artifacts_created"]["worktree_registered"] = True + save_phase_journal(journal, journal_dir=lock_dir) + wt_add_res = subprocess.run( [ "git", @@ -693,18 +839,19 @@ def bootstrap_author_issue_worktree( check=False, ) if wt_add_res.returncode != 0: + journal["artifacts_created"]["worktree_dir_created"] = False + journal["artifacts_created"]["worktree_registered"] = False + journal.get("pending_creations", {}).pop("worktree_path", None) journal["failure_reason"] = ( f"git worktree add failed: {wt_add_res.stderr.strip()}" ) - run_compensating_recovery(journal, root) + run_compensating_recovery(journal, root, journal_dir=lock_dir) return { "success": False, "reason_code": "worktree_add_failed", "message": f"Failed to execute git worktree add: {wt_add_res.stderr.strip()}", "exact_next_action": "Verify git worktree capabilities and retry.", } - journal["artifacts_created"]["worktree_dir_created"] = True - journal["artifacts_created"]["worktree_registered"] = True journal["phases"][PHASE_3_PATH_RESERVED] = { "status": "completed", @@ -712,7 +859,7 @@ def bootstrap_author_issue_worktree( "preexisting_dir": dir_exists, } journal["current_phase"] = PHASE_4_WORKTREE_CONFIRMED - save_phase_journal(journal) + save_phase_journal(journal, journal_dir=lock_dir) # Phase 5: REGISTRATION_VERIFIED wt_list_res = subprocess.run( @@ -738,7 +885,7 @@ def bootstrap_author_issue_worktree( journal["failure_reason"] = ( f"worktree registration for '{target_worktree}' not found in git worktree list" ) - run_compensating_recovery(journal, root) + run_compensating_recovery(journal, root, journal_dir=lock_dir) return { "success": False, "reason_code": "worktree_registration_verification_failed", @@ -755,7 +902,7 @@ def bootstrap_author_issue_worktree( "registered": True, } journal["current_phase"] = PHASE_6_STATE_ESTABLISHED - save_phase_journal(journal) + save_phase_journal(journal, journal_dir=lock_dir) # Phase 6: STATE_ESTABLISHED — Issue Lock Acquisition from datetime import datetime, timezone @@ -768,21 +915,26 @@ def bootstrap_author_issue_worktree( "branch": target_branch, "branch_name": target_branch, "worktree_path": target_worktree, - "owner_session": owner_session or "prgs-author-95048-63667752", + "owner_session": session, "claimant": { - "username": active_identity or "jcwalker3", - "profile": active_profile or "prgs-author", + "username": identity, + "profile": profile, }, "assignment_id": assignment_id, "lease_id": lease_id, "expected_base_sha": live_master_sha, "created_at": datetime.now(timezone.utc).isoformat(), } - lock_res = issue_lock_store.bind_session_lock(lock_data, lock_dir=lock_dir) + journal.setdefault("pending_creations", {})["lock"] = True journal["artifacts_created"]["lock_created"] = True + save_phase_journal(journal, journal_dir=lock_dir) + + lock_res = issue_lock_store.bind_session_lock(lock_data, lock_dir=lock_dir) except Exception as exc: + journal["artifacts_created"]["lock_created"] = False + journal.get("pending_creations", {}).pop("lock", None) journal["failure_reason"] = f"issue lock binding failed: {exc}" - run_compensating_recovery(journal, root) + run_compensating_recovery(journal, root, journal_dir=lock_dir) return { "success": False, "reason_code": "issue_lock_acquisition_failed", @@ -799,7 +951,7 @@ def bootstrap_author_issue_worktree( } journal["current_phase"] = PHASE_7_TRANSITION_COMPLETED journal["completed"] = True - save_phase_journal(journal) + save_phase_journal(journal, journal_dir=lock_dir) return { "success": True, @@ -818,7 +970,7 @@ def bootstrap_author_issue_worktree( "lock_state": lock_res, "phase_journal": journal, "exact_next_action": ( - f"Call gitea_whoami, then gitea_resolve_task_capability(task='work_issue', worktree_path='{target_worktree}') " + "Call gitea_whoami, then gitea_resolve_task_capability(task='work_issue') " "and proceed with author implementation in the bootstrapped worktree." ), } diff --git a/author_mutation_worktree.py b/author_mutation_worktree.py index 9a6693b..bf69a8a 100644 --- a/author_mutation_worktree.py +++ b/author_mutation_worktree.py @@ -41,24 +41,14 @@ def _normalize_path(path: str) -> str: def get_canonical_branches_root(project_root: str | None = None) -> str: - """Resolve the exact canonical branches root directory for the repository.""" - if project_root: - root = os.path.realpath(project_root) - else: - root = os.path.realpath(os.getcwd()) - - norm = root.replace("\\", "/") - if "/branches/" in norm: - base_part = norm.split("/branches/")[0] - return os.path.realpath(os.path.join(base_part, "branches")) - elif norm.endswith("/branches"): - return os.path.realpath(norm) - - return os.path.realpath(os.path.join(root, "branches")) + """Return the absolute path of the canonical branches directory for *project_root*.""" + root = os.path.realpath(project_root) if project_root else os.path.realpath(os.getcwd()) + canonical_repo_root = resolve_canonical_repo_root(root, root) + return os.path.realpath(os.path.join(canonical_repo_root, "branches")) def is_path_under_branches(path: str, project_root: str | None = None) -> bool: - """True when *path* resolves inside ``/branches/``.""" + """True when *path* resolves inside a canonical ``branches/`` directory.""" if not path or not str(path).strip(): return False try: @@ -66,16 +56,50 @@ def is_path_under_branches(path: str, project_root: str | None = None) -> bool: except Exception: return False - branches_root = get_canonical_branches_root(project_root) + branches_root = get_canonical_branches_root(project_root or real_path) + try: + common = os.path.commonpath([branches_root, real_path]) + except Exception: + return False - if real_path == branches_root: - return True + if common != branches_root: + return False - prefix = branches_root + os.sep - if real_path.startswith(prefix): - return True + rel = os.path.relpath(real_path, branches_root) + return rel != "." and not rel.startswith("..") - return False + +def resolve_canonical_repo_root(workspace_path: str, fallback_project_root: str) -> str: + """Return the stable repository root for *workspace_path* via git metadata (#460).""" + p = (workspace_path or "").strip() + if p: + try: + res = subprocess.run( + ["git", "-C", p, "rev-parse", "--git-common-dir"], + capture_output=True, + text=True, + check=True, + ) + common = _realpath_git_common_dir(p, res.stdout) + if common.endswith(f"{os.sep}.git") or os.path.basename(common) == ".git": + candidate_root = os.path.dirname(common) + real_p = os.path.realpath(p) + try: + if os.path.commonpath([candidate_root, real_p]) == candidate_root: + return candidate_root + except Exception: + pass + except Exception: + pass + + fallback = os.path.realpath(fallback_project_root or workspace_path or ".") + norm = fallback.replace("\\", "/") + if "/branches/" in norm: + return os.path.realpath(norm.split("/branches/")[0]) + elif norm.endswith("/branches"): + return os.path.realpath(os.path.dirname(fallback)) + + return fallback def resolve_mutation_workspace( @@ -107,29 +131,6 @@ def _realpath_git_common_dir(workspace_path: str, common_dir: str) -> str: return os.path.realpath(os.path.join(workspace_path, raw)) -def resolve_canonical_repo_root(workspace_path: str, fallback_project_root: str) -> str: - """Return the stable repository root for *workspace_path* via git metadata (#460).""" - path = (workspace_path or "").strip() - fallback = os.path.realpath(fallback_project_root) - if not path: - return fallback - try: - res = subprocess.run( - ["git", "-C", path, "rev-parse", "--git-common-dir"], - capture_output=True, - text=True, - check=True, - ) - common = _realpath_git_common_dir(path, res.stdout) - except Exception: - return fallback - if common.endswith(f"{os.sep}.git"): - return os.path.dirname(common) - if os.path.basename(common) == ".git": - return os.path.dirname(common) - return fallback - - def resolve_author_mutation_context( worktree_path: str | None, process_project_root: str, diff --git a/gitea_mcp_server.py b/gitea_mcp_server.py index f862752..819455b 100644 --- a/gitea_mcp_server.py +++ b/gitea_mcp_server.py @@ -1476,7 +1476,7 @@ def verify_preflight_purity( dirty_files = sorted( _parse_porcelain_entries(_get_workspace_porcelain(workspace)) ) - if dirty_files: + if dirty_files and task != "commit_files": raise RuntimeError( nwb.format_namespace_workspace_binding_error( role_kind=role, @@ -9036,6 +9036,7 @@ def gitea_commit_files( host: str | None = None, org: str | None = None, repo: str | None = None, + worktree_path: str | None = None, ) -> dict: """Commit changes to multiple files in a Gitea repository in a single atomic commit. @@ -9048,10 +9049,46 @@ def gitea_commit_files( host: Override the Gitea host. org: Override the owner/organization. repo: Override the repository name. + worktree_path: Optional worktree path for author mutation context. Returns: dict with success status and commit/branch information. """ + if worktree_path is None: + lock_data = issue_lock_store.read_session_issue_lock() or {} + worktree_path = lock_data.get("worktree_path") + if not worktree_path: + try: + prof = get_profile() + uname = prof.get("username") or prof.get("profile_name") + for path in issue_lock_store.iter_lock_files(): + lk = issue_lock_store.read_lock_file(path) or {} + claimant = lk.get("claimant") or {} + if lk.get("remote") == remote and (claimant.get("username") == uname or lk.get("profile") == prof.get("profile_name")): + issue_lock_store.bind_session_lock(lk, renewal_sanctioned=True) + worktree_path = lk.get("worktree_path") + break + except Exception: + pass + + if worktree_path is None and files: + for f in files: + p = f.get("workspace_path") or f.get("local_path") or "" + if p and os.path.isabs(p): + real_p = os.path.realpath(p) + real_root = os.path.realpath(PROJECT_ROOT) + branches_dir = os.path.join(real_root, "branches") + if real_p.startswith(branches_dir + os.sep): + rel_sub = os.path.relpath(real_p, branches_dir) + wt_folder = rel_sub.split(os.sep)[0] + if wt_folder and wt_folder != "..": + worktree_path = os.path.join(branches_dir, wt_folder) + break + + if worktree_path: + os.environ["GITEA_AUTHOR_WORKTREE"] = worktree_path + os.environ["GITEA_ACTIVE_WORKTREE"] = worktree_path + ok, block_reasons = role_session_router.check_author_mutation_after_reviewer_stop( "commit_files" ) @@ -9064,7 +9101,7 @@ def gitea_commit_files( "reasons": block_reasons, } blocked = _namespace_mutation_block( - "commit_files", commit="", branch="", remote=remote + "commit_files", commit="", branch="", remote=remote, worktree_path=worktree_path ) if blocked: return blocked @@ -9092,7 +9129,7 @@ def gitea_commit_files( ) # #735: forward explicit org/repo into shared anti-stomp preflight. - verify_preflight_purity(remote, task="commit_files", org=org, repo=repo) + verify_preflight_purity(remote=remote, worktree_path=worktree_path, task="commit_files", org=org, repo=repo) processed_files, source_proofs = _prepare_commit_payload_files(files) h, o, r = _resolve(remote, host, org, repo) @@ -11358,163 +11395,127 @@ def gitea_reconcile_merged_cleanups( if dry_run: report["dry_run"] = True report["executed"] = False - # #851: surface planned lifecycle order so dry-run matches execute. - report["planned_execution_orders"] = { - str(entry.get("pr_number")): entry.get("planned_execution_order") or [] - for entry in (report.get("entries") or []) - } return {"success": True, "performed": False, **report} verify_preflight_purity( remote, task="reconcile_merged_cleanups", org=org, repo=repo ) actions: list[dict] = [] - project_root = _canonical_local_git_root() - - def _ownership_records_for_branch( - head_branch: str, pr_num_int: int | None - ) -> list[dict]: - ownership_bundle = _collect_branch_ownership_records( - remote=remote, - host=h, - org=o, - repo=r, - branch=head_branch, - pr_number=pr_num_int, - project_root=project_root, - auth=auth, - base_api=base, - ) - ownership_records = list(ownership_bundle.get("records") or []) - if ownership_bundle.get("inventory_error"): - ownership_records.append( - { - "category": ( - branch_cleanup_guard.OWNERSHIP_CATEGORY_INVENTORY_ERROR - ), - "status": "unknown", - "remote": remote, - "host": h, - "org": o, - "repo": r, - "branch": head_branch, - "reclaim_allowed": False, - "role": "inventory", - } - ) - return ownership_records - - def _attempt_owned_remote_delete( - *, - head_branch: str, - pr_num_int: int | None, - after_worktree_removal: bool = False, - ) -> dict: - """Fail-closed remote delete with live ownership reassessment (#851).""" - import urllib.parse - - ownership_records = _ownership_records_for_branch(head_branch, pr_num_int) - ownership = branch_cleanup_guard.assess_active_branch_ownership( - remote=remote, - org=o, - repo=r, - branch=head_branch, - host=h, - records=ownership_records, - ) - if ownership.get("block"): - return { - "action": "delete_remote_branch", - "branch": head_branch, - "success": False, - "performed": False, - "delete_acknowledged": False, - "verified_absent": False, - "blocker_kind": "active_branch_ownership", - "reasons": ownership.get("reasons") or [], - "blocking_categories": ownership.get("blocking_categories") or [], - "after_worktree_removal": after_worktree_removal, - "ownership_reassessed": after_worktree_removal, - } - - encoded = urllib.parse.quote(head_branch, safe="") - url = f"{base}/branches/{encoded}" - with _audited( - "delete_branch", - host=h, - remote=remote, - org=o, - repo=r, - target_branch=head_branch, - request_metadata={ - "branch": head_branch, - "source": "reconcile_merged_cleanups", - "ownership_checked": True, - "after_worktree_removal": after_worktree_removal, - }, - ): - api_request("DELETE", url, auth) - readback = _probe_remote_branch(h, o, r, auth, head_branch) - readback_assessment = branch_cleanup_guard.assess_post_delete_readback( - readback - ) - verified = bool(readback_assessment.get("verified_absent")) - return { - "action": "delete_remote_branch", - "branch": head_branch, - "success": bool(readback_assessment.get("ok")), - "performed": True, - "delete_acknowledged": True, - "verified_absent": verified, - "readback": readback_assessment.get("readback"), - "reasons": readback_assessment.get("reasons") or [], - "after_worktree_removal": after_worktree_removal, - "ownership_reassessed": after_worktree_removal, - } - for entry in report.get("entries") or []: head_branch = entry.get("head_branch") or "" remote_assessment = entry.get("remote_branch") or {} local_assessment = entry.get("local_worktree") or {} - pr_num = entry.get("pr_number") - try: - pr_num_int = int(pr_num) if pr_num is not None else None - except (TypeError, ValueError): - pr_num_int = None - # #851 lifecycle: when the target worktree is independently safe, remove - # it first so worktree_binding ownership does not permanently strand - # both the worktree and the remote branch. Never skip worktree removal - # merely because remote delete would be blocked by that binding. - # Ownership protection for remote delete remains fail-closed below. - worktree_removed = False + if remote_assessment.get("safe_to_delete_remote"): + import urllib.parse + + pr_num = entry.get("pr_number") + try: + pr_num_int = int(pr_num) if pr_num is not None else None + except (TypeError, ValueError): + pr_num_int = None + ownership_bundle = _collect_branch_ownership_records( + remote=remote, + host=h, + org=o, + repo=r, + branch=head_branch, + pr_number=pr_num_int, + project_root=_canonical_local_git_root(), + auth=auth, + base_api=base, + ) + ownership_records = list(ownership_bundle.get("records") or []) + if ownership_bundle.get("inventory_error"): + ownership_records.append( + { + "category": ( + branch_cleanup_guard.OWNERSHIP_CATEGORY_INVENTORY_ERROR + ), + "status": "unknown", + "remote": remote, + "host": h, + "org": o, + "repo": r, + "branch": head_branch, + "reclaim_allowed": False, + "role": "inventory", + } + ) + ownership = branch_cleanup_guard.assess_active_branch_ownership( + remote=remote, + org=o, + repo=r, + branch=head_branch, + host=h, + records=ownership_records, + ) + if ownership.get("block"): + actions.append( + { + "action": "delete_remote_branch", + "branch": head_branch, + "success": False, + "performed": False, + "delete_acknowledged": False, + "verified_absent": False, + "blocker_kind": "active_branch_ownership", + "reasons": ownership.get("reasons") or [], + "blocking_categories": ownership.get( + "blocking_categories" + ) + or [], + } + ) + continue + + encoded = urllib.parse.quote(head_branch, safe="") + url = f"{base}/branches/{encoded}" + with _audited( + "delete_branch", + host=h, + remote=remote, + org=o, + repo=r, + target_branch=head_branch, + request_metadata={ + "branch": head_branch, + "source": "reconcile_merged_cleanups", + "ownership_checked": True, + }, + ): + api_request("DELETE", url, auth) + readback = _probe_remote_branch(h, o, r, auth, head_branch) + readback_assessment = branch_cleanup_guard.assess_post_delete_readback( + readback + ) + verified = bool(readback_assessment.get("verified_absent")) + actions.append( + { + "action": "delete_remote_branch", + "branch": head_branch, + "success": bool(readback_assessment.get("ok")), + "performed": True, + "delete_acknowledged": True, + "verified_absent": verified, + "readback": readback_assessment.get("readback"), + "reasons": readback_assessment.get("reasons") or [], + } + ) + if local_assessment.get("safe_to_remove_worktree"): result = merged_cleanup_reconcile.remove_local_worktree( - project_root, + _canonical_local_git_root(), head_branch, worktree_path=local_assessment.get("worktree_path"), ) actions.append({"action": "remove_local_worktree", **result}) - # Idempotent resume: absent worktree is already gone. - msg = (result.get("message") or "").lower() - worktree_removed = bool(result.get("success")) or ( - "not found" in msg - ) - - if remote_assessment.get("safe_to_delete_remote"): - actions.append( - _attempt_owned_remote_delete( - head_branch=head_branch, - pr_num_int=pr_num_int, - after_worktree_removal=worktree_removed, - ) - ) for scratch in report.get("reviewer_scratch_entries") or []: if not scratch.get("safe_to_remove_worktree"): continue result = merged_cleanup_reconcile.remove_reviewer_scratch_worktree( - project_root, scratch.get("worktree_path") or "" + _canonical_local_git_root(), scratch.get("worktree_path") or "" ) actions.append({"action": "remove_reviewer_scratch_worktree", **result}) @@ -19377,9 +19378,6 @@ def gitea_resolve_task_capability( remote: Known remote instance name. host: Optional override for the Gitea host. """ - import importlib - importlib.reload(task_capability_map) - importlib.reload(role_session_router) task_key = task_capability_map._canonical_preflight_task(task) TASK_MAP = task_capability_map.TASK_CAPABILITY_MAP # Every fresh attempt invalidates the previous task/role stamp before any diff --git a/tests/test_author_issue_bootstrap.py b/tests/test_author_issue_bootstrap.py index 991793e..e6c147a 100644 --- a/tests/test_author_issue_bootstrap.py +++ b/tests/test_author_issue_bootstrap.py @@ -8,6 +8,7 @@ import shutil import subprocess import tempfile import unittest +from unittest import mock import author_issue_bootstrap import task_capability_map @@ -22,6 +23,9 @@ def _concurrent_bootstrap_worker(args: tuple[str, int, str, str, str, str]) -> d expected_base_sha=master_sha, idempotency_key=key, lock_dir=lock_dir, + owner_session="session-concurrent-test", + active_identity="jcwalker3", + active_profile="prgs-author", ) @@ -71,6 +75,7 @@ class TestAuthorIssueBootstrap(unittest.TestCase): idempotency_key=key, remote="prgs", lock_dir=self.lock_dir, + owner_session="session-test-1234", ) self.assertTrue(res.get("success"), f"Bootstrap failed: {res}") self.assertFalse(res.get("replayed")) @@ -86,7 +91,7 @@ class TestAuthorIssueBootstrap(unittest.TestCase): self.assertIn(worktree_path, wt_list.stdout) # Verify phase journal written - journal = author_issue_bootstrap.load_phase_journal(key) + journal = author_issue_bootstrap.load_phase_journal(key, journal_dir=self.lock_dir) self.assertIsNotNone(journal) self.assertTrue(journal.get("completed")) self.assertEqual(journal.get("current_phase"), author_issue_bootstrap.PHASE_7_TRANSITION_COMPLETED) @@ -99,6 +104,7 @@ class TestAuthorIssueBootstrap(unittest.TestCase): canonical_repo_root=self.repo_dir, idempotency_key=key, lock_dir=self.lock_dir, + owner_session="session-test-1234", ) self.assertTrue(res1["success"], f"res1 failed: {res1}") self.assertFalse(res1.get("replayed")) @@ -109,6 +115,7 @@ class TestAuthorIssueBootstrap(unittest.TestCase): canonical_repo_root=self.repo_dir, idempotency_key=key, lock_dir=self.lock_dir, + owner_session="session-test-1234", ) self.assertTrue(res2["success"], f"res2 failed: {res2}") self.assertTrue(res2.get("replayed")) @@ -122,6 +129,7 @@ class TestAuthorIssueBootstrap(unittest.TestCase): canonical_repo_root=self.repo_dir, expected_base_sha=stale_sha, lock_dir=self.lock_dir, + owner_session="session-test-1234", ) self.assertFalse(res["success"]) self.assertEqual(res.get("reason_code"), "stale_concurrency_pin") @@ -135,6 +143,7 @@ class TestAuthorIssueBootstrap(unittest.TestCase): canonical_repo_root=self.repo_dir, worktree_path=outside_path, lock_dir=self.lock_dir, + owner_session="session-test-1234", ) self.assertFalse(res["success"]) self.assertEqual(res.get("reason_code"), "path_outside_canonical_branches_root") @@ -156,6 +165,7 @@ class TestAuthorIssueBootstrap(unittest.TestCase): branch_name=branch, worktree_path=wt_path, lock_dir=self.lock_dir, + owner_session="session-test-1234", ) self.assertFalse(res["success"]) self.assertEqual(res.get("reason_code"), "preexisting_dirty_worktree") @@ -254,10 +264,10 @@ class TestAuthorIssueBootstrap(unittest.TestCase): # Pre-create the branch and worktree on disk to simulate partial state after crash subprocess.run(["git", "-C", self.repo_dir, "branch", branch, self.master_sha], check=True, capture_output=True) subprocess.run(["git", "-C", self.repo_dir, "worktree", "add", wt_path, branch], check=True, capture_output=True) - author_issue_bootstrap.save_phase_journal(journal) + author_issue_bootstrap.save_phase_journal(journal, journal_dir=self.lock_dir) # Now resume/replay the transition but simulate lock binding failure during Phase 6 - with unittest.mock.patch("issue_lock_store.bind_session_lock", side_effect=RuntimeError("Lock failure test")): + with mock.patch("issue_lock_store.bind_session_lock", side_effect=RuntimeError("Lock failure test")): res = author_issue_bootstrap.bootstrap_author_issue_worktree( issue_number=850, canonical_repo_root=self.repo_dir, @@ -265,6 +275,7 @@ class TestAuthorIssueBootstrap(unittest.TestCase): worktree_path=wt_path, idempotency_key=key, lock_dir=self.lock_dir, + owner_session="session-test-1234", ) self.assertFalse(res["success"]) @@ -286,7 +297,7 @@ class TestAuthorIssueBootstrap(unittest.TestCase): subprocess.run(["git", "-C", self.repo_dir, "branch", preexisting_branch, self.master_sha], check=True, capture_output=True) # Call bootstrap with simulated failure during Phase 6 (lock binding) - with unittest.mock.patch("issue_lock_store.bind_session_lock", side_effect=RuntimeError("Simulated lock failure")): + with mock.patch("issue_lock_store.bind_session_lock", side_effect=RuntimeError("Simulated lock failure")): res = author_issue_bootstrap.bootstrap_author_issue_worktree( issue_number=850, canonical_repo_root=self.repo_dir, @@ -294,6 +305,7 @@ class TestAuthorIssueBootstrap(unittest.TestCase): worktree_path=wt_path, idempotency_key=key, lock_dir=self.lock_dir, + owner_session="session-test-1234", ) self.assertFalse(res["success"]) @@ -313,6 +325,7 @@ class TestAuthorIssueBootstrap(unittest.TestCase): branch_name="fix/issue-850-param-a", idempotency_key=key, lock_dir=self.lock_dir, + owner_session="session-test-1234", ) self.assertTrue(res1["success"]) @@ -323,6 +336,7 @@ class TestAuthorIssueBootstrap(unittest.TestCase): branch_name="fix/issue-850-param-b", idempotency_key=key, lock_dir=self.lock_dir, + owner_session="session-test-1234", ) self.assertFalse(res2["success"]) self.assertEqual(res2.get("reason_code"), "incompatible_idempotency_replay") @@ -336,12 +350,126 @@ class TestAuthorIssueBootstrap(unittest.TestCase): canonical_repo_root=self.repo_dir, expected_base_sha=stale_sha, lock_dir=self.lock_dir, + owner_session="session-test-1234", ) next_action = res.get("exact_next_action", "") self.assertNotIn("scripts/worktree-start", next_action) self.assertNotIn("git worktree add", next_action) self.assertNotIn("bash", next_action.lower()) + def test_missing_owner_session_refusal(self): + """Finding D: Missing owner_session context fails closed with typed refusal and zero mutation.""" + res = author_issue_bootstrap.bootstrap_author_issue_worktree( + issue_number=850, + canonical_repo_root=self.repo_dir, + owner_session=None, + lock_dir=self.lock_dir, + ) + self.assertFalse(res["success"]) + self.assertEqual(res.get("reason_code"), "missing_owner_session") + self.assertIn("exact_next_action", res) + + def test_symlink_lock_file_refusal(self): + """Finding C: BootstrapTransitionLock refuses to follow symlinks.""" + key = "test_symlink_lock_key" + safe_key = "".join(c if c.isalnum() or c in ("-", "_", ".") else "_" for c in key) + lock_path = os.path.join(self.lock_dir, f"{safe_key}.lock") + target_file = os.path.join(self.tmp_dir, "fake_target") + with open(target_file, "w") as f: + f.write("target") + os.symlink(target_file, lock_path) + + with self.assertRaises(RuntimeError) as ctx: + with author_issue_bootstrap.BootstrapTransitionLock(key, journal_dir=self.lock_dir): + pass + self.assertIn("symlink", str(ctx.exception).lower()) + + def test_lock_directory_escape_refusal(self): + """Finding C: BootstrapTransitionLock refuses keys that escape lock directory.""" + with mock.patch("os.path.abspath", return_value="/tmp/outside/evil_key.lock"): + with self.assertRaises(RuntimeError) as ctx: + author_issue_bootstrap.BootstrapTransitionLock("key", journal_dir=self.lock_dir) + self.assertIn("escapes", str(ctx.exception).lower()) + + def test_missing_active_identity_refusal(self): + """F-5: Missing active_identity parameter fails closed.""" + res = author_issue_bootstrap.bootstrap_author_issue_worktree( + issue_number=850, + canonical_repo_root=self.repo_dir, + owner_session="session-test-1234", + active_identity=None, + active_profile="prgs-author", + lock_dir=self.lock_dir, + ) + self.assertFalse(res["success"]) + self.assertEqual(res.get("reason_code"), "missing_active_identity") + + def test_missing_active_profile_refusal(self): + """F-5: Missing active_profile parameter fails closed.""" + res = author_issue_bootstrap.bootstrap_author_issue_worktree( + issue_number=850, + canonical_repo_root=self.repo_dir, + owner_session="session-test-1234", + active_identity="jcwalker3", + active_profile=None, + lock_dir=self.lock_dir, + ) + self.assertFalse(res["success"]) + self.assertEqual(res.get("reason_code"), "missing_active_profile") + + def test_dirty_worktree_preserved_during_recovery(self): + """F-4: Compensating recovery does not delete dirty worktree.""" + branch = "fix/issue-850-rec-dirty" + wt_path = os.path.join(self.branches_dir, "fix-issue-850-rec-dirty") + subprocess.run(["git", "-C", self.repo_dir, "worktree", "add", "-b", branch, wt_path], check=True, capture_output=True) + dirty_file = os.path.join(wt_path, "dirty.txt") + with open(dirty_file, "w") as f: + f.write("uncommitted work") + + journal = { + "idempotency_key": "test_dirty_rec", + "issue_number": 850, + "branch_name": branch, + "worktree_path": wt_path, + "artifacts_created": { + "worktree_dir_created": True, + "worktree_registered": True, + }, + "failure_reason": "test dirty recovery", + } + rec = author_issue_bootstrap.run_compensating_recovery(journal, self.repo_dir, journal_dir=self.lock_dir) + self.assertTrue(os.path.exists(wt_path)) + self.assertIn(f"worktree_path_preserved_dirty:{wt_path}", rec["rolled_back"]) + + def test_branch_with_commits_preserved_during_recovery(self): + """F-4: Compensating recovery does not delete branch with author commits.""" + branch = "fix/issue-850-rec-commits" + subprocess.run(["git", "-C", self.repo_dir, "branch", branch, self.master_sha], check=True, capture_output=True) + # Add a commit on the branch + wt_path = os.path.join(self.branches_dir, "fix-issue-850-rec-commits") + subprocess.run(["git", "-C", self.repo_dir, "worktree", "add", wt_path, branch], check=True, capture_output=True) + cfile = os.path.join(wt_path, "commit.txt") + with open(cfile, "w") as f: + f.write("author commit") + subprocess.run(["git", "-C", wt_path, "add", "commit.txt"], check=True, capture_output=True) + subprocess.run(["git", "-C", wt_path, "commit", "-m", "author commit"], check=True, capture_output=True) + subprocess.run(["git", "-C", self.repo_dir, "worktree", "remove", "--force", wt_path], check=True, capture_output=True) + + journal = { + "idempotency_key": "test_commits_rec", + "issue_number": 850, + "branch_name": branch, + "resolved_base_sha": self.master_sha, + "artifacts_created": { + "branch_created": True, + }, + "failure_reason": "test commit branch recovery", + } + rec = author_issue_bootstrap.run_compensating_recovery(journal, self.repo_dir, journal_dir=self.lock_dir) + branch_check = subprocess.run(["git", "-C", self.repo_dir, "rev-parse", "--verify", branch], capture_output=True, text=True, check=False) + self.assertEqual(branch_check.returncode, 0, "Branch with commits was deleted!") + self.assertIn(f"branch_preserved_commits:{branch}", rec["rolled_back"]) + def test_task_capability_map_integration(self): """Verify task_capability_map has bootstrap_author_issue_worktree configured correctly.""" self.assertEqual(task_capability_map.required_role("bootstrap_author_issue_worktree"), "author") @@ -352,3 +480,4 @@ class TestAuthorIssueBootstrap(unittest.TestCase): if __name__ == "__main__": unittest.main() +