"""Scoped on-demand identifier registrar for the ADR-007 interim. The production fleet sweep is intentionally disabled. This command provides one explicit, repository-scoped way to run the existing State Hub registration adapter without teaching ordinary agent sessions to set ``STATEHUB_REGISTRAR`` themselves. """ from __future__ import annotations import fcntl import os import shutil import subprocess import uuid from dataclasses import dataclass from pathlib import Path from typing import Any import httpx from repo_manager.gitops import GitError, commit_paths, push_ff from repo_manager.parse.record import iter_record_files, parse_record_file from repo_manager.parse.workplan import parse_workplan_file LOCK_PATH = Path("/tmp/repo-manager-identifier-registrar.lock") @dataclass class RegistrarResult: status: str evidence: dict[str, Any] error: dict[str, Any] | None = None correlation_id: str = "" def to_dict(self) -> dict[str, Any]: result: dict[str, Any] = { "command": "repo.identifiers.registrar_reconcile", "status": self.status, "correlation_id": self.correlation_id, "evidence": self.evidence, } if self.error: result["error"] = self.error return result def _git(repo: Path, *args: str) -> subprocess.CompletedProcess[str]: return subprocess.run( ["git", *args], cwd=repo, capture_output=True, text=True, check=False, ) def _missing_identifiers(repo: Path) -> dict[str, list[str]]: workplans: list[str] = [] tasks: list[str] = [] intakes: list[str] = [] decisions: list[str] = [] workplans_dir = repo / "workplans" if workplans_dir.is_dir(): # Closed archives are frozen under ADR-007 and are not registration work. for path in sorted(workplans_dir.glob("*.md")): parsed = parse_workplan_file(path, repo_root=repo) if parsed.frontmatter.get("type") != "workplan" or not parsed.id: continue # Closed records are frozen provenance under ADR-007. State Hub also # deliberately refuses to create missing task rows for them. if (parsed.status or "").strip().lower() in {"finished", "archived"}: continue if not parsed.state_hub_workstream_id: workplans.append(parsed.id) for task in parsed.tasks: if task.id and not task.state_hub_task_id: tasks.append(task.id) for path in iter_record_files(repo): for record in parse_record_file(path, repo_root=repo): if record.uuid: continue target = intakes if record.kind == "intake" else decisions target.append(record.id) return { "workplans": workplans, "tasks": tasks, "intakes": intakes, "decisions": decisions, } def _check_git(repo: Path) -> tuple[dict[str, Any], str | None]: status = _git(repo, "status", "--porcelain") if status.returncode != 0: return {}, status.stderr.strip() or "not a Git repository" if status.stdout.strip(): return {}, "worktree must be clean before registrar reconciliation" upstream = _git(repo, "rev-parse", "--abbrev-ref", "--symbolic-full-name", "@{u}") if upstream.returncode != 0: return {}, "current branch must have an upstream" counts = _git(repo, "rev-list", "--left-right", "--count", "@{u}...HEAD") if counts.returncode != 0: return {}, counts.stderr.strip() or "could not compare HEAD with upstream" try: behind, ahead = (int(value) for value in counts.stdout.split()) except (TypeError, ValueError): return {}, "could not parse Git ahead/behind state" if behind or ahead: return ( {"upstream": upstream.stdout.strip(), "behind": behind, "ahead": ahead}, "branch must match its upstream; pull or push before registrar reconciliation", ) remote = _git(repo, "remote", "get-url", "origin") if remote.returncode != 0: return {}, "origin remote is required" remote_url = remote.stdout.strip() stale_markers = ("gitea-remote", "gitea.coulomb.social", "92.205.130.254") if any(marker in remote_url for marker in stale_markers): return {"origin": remote_url}, "origin still targets the retired Gitea lineage" return { "upstream": upstream.stdout.strip(), "behind": behind, "ahead": ahead, "origin": remote_url, }, None def _check_primary(api_base: str) -> tuple[dict[str, Any], str | None]: try: response = httpx.get(f"{api_base.rstrip('/')}/state/health", timeout=10.0) response.raise_for_status() payload = response.json() except (httpx.HTTPError, ValueError) as exc: return {}, f"primary State Hub health check failed: {exc}" if payload.get("status") != "ok" or payload.get("db") != "connected": return payload, "State Hub is not healthy and database-connected" return payload, None def _run_statehub(command: list[str], *, env: dict[str, str]) -> subprocess.CompletedProcess[str]: return subprocess.run( command, capture_output=True, text=True, check=False, timeout=300, env=env, ) def registrar_reconcile( path: Path, *, api_base: str = "http://127.0.0.1:8000", statehub_bin: str | None = None, confirm_primary: bool = False, push: bool = False, ) -> RegistrarResult: """Register missing workplan/task UUIDs through one scoped child process.""" cid = str(uuid.uuid4()) repo = path.expanduser().resolve() before = _missing_identifiers(repo) evidence: dict[str, Any] = { "repo_path": str(repo), "repo_slug": repo.name, "api_base": api_base.rstrip("/"), "missing_before": before, } if not confirm_primary: return RegistrarResult( "rejected", evidence, { "code": "confirmation_required", "message": "pass --confirm-primary after verifying this is the authoritative hub", }, cid, ) if not any(before.values()): evidence["missing_after"] = before return RegistrarResult("noop", evidence, None, cid) git_evidence, git_error = _check_git(repo) evidence["git"] = git_evidence if git_error: return RegistrarResult( "rejected", evidence, {"code": "git_precondition_failed", "message": git_error}, cid ) health, health_error = _check_primary(api_base) evidence["state_hub_health"] = health if health_error: return RegistrarResult( "rejected", evidence, {"code": "primary_unavailable", "message": health_error}, cid ) executable = statehub_bin or shutil.which("statehub") if not executable: return RegistrarResult( "rejected", evidence, {"code": "statehub_missing", "message": "statehub CLI is not installed"}, cid, ) LOCK_PATH.parent.mkdir(parents=True, exist_ok=True) with LOCK_PATH.open("a+", encoding="utf-8") as lock: try: fcntl.flock(lock.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB) except BlockingIOError: return RegistrarResult( "rejected", evidence, {"code": "registrar_busy", "message": "another registrar reconciliation is active"}, cid, ) child_env = {**os.environ, "STATEHUB_REGISTRAR": "1"} command = [ executable, "fix-consistency", "--path", str(repo), "--api-base", api_base.rstrip("/"), ] completed = _run_statehub(command, env=child_env) evidence["statehub_exit_code"] = completed.returncode evidence["statehub_stdout_tail"] = completed.stdout[-4000:] evidence["statehub_stderr_tail"] = completed.stderr[-2000:] after = _missing_identifiers(repo) evidence["missing_after"] = after if completed.returncode not in {0, 2} or any(after.values()): return RegistrarResult( "failed", evidence, { "code": "registration_incomplete", "message": "State Hub reconciliation did not assign every requested identifier", }, cid, ) changed = _git(repo, "status", "--porcelain") paths = [line[3:] for line in changed.stdout.splitlines() if len(line) > 3] if paths: try: evidence["git_sha"] = commit_paths( repo, paths, "chore(registrar): assign State Hub identifiers", ) except GitError as exc: return RegistrarResult( "failed", evidence, {"code": "commit_failed", "message": str(exc)}, cid, ) if push: pushed, message = push_ff(repo) evidence["push"] = {"ok": pushed, "message": message} if not pushed: return RegistrarResult( "failed", evidence, {"code": "push_failed", "message": message}, cid, ) return RegistrarResult("applied", evidence, None, cid)