Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
99fda93bcc | ||
|
|
9eb0f29cef | ||
|
|
0b29404031 | ||
|
|
5463f58933 | ||
|
|
5032965e3a | ||
|
|
57a52b1a99 | ||
|
|
e33b8d3712 | ||
|
|
344dc41ce2 | ||
|
|
620ed6e9a9 | ||
|
|
0f19773076 | ||
|
|
324b4b3e93 |
@@ -87,6 +87,7 @@ MUTATION_TASKS = frozenset({
|
||||
"edit_pr",
|
||||
"commit_files",
|
||||
"gitea_commit_files",
|
||||
"publish_unpublished_branch",
|
||||
"delete_branch",
|
||||
"cleanup_merged_pr_branch",
|
||||
"cleanup_stale_claims",
|
||||
|
||||
@@ -0,0 +1,591 @@
|
||||
"""Publish an unpublished local commit on a registered issue worktree (#812 AC20).
|
||||
|
||||
Entry point B of #812 is the state where an author's work has already advanced
|
||||
to a local commit: the worktree is registered, clean, on the issue branch, and
|
||||
carries the only copy of the implementation, but the branch has never been
|
||||
published. That state deadlocks, because two individually correct predicates
|
||||
close a cycle:
|
||||
|
||||
* ``issue_lock_renewal.assess_exact_owner_lease_renewal`` refuses to renew an
|
||||
expired lease without an observable remote head — an unpublished branch has
|
||||
none.
|
||||
* Every publication path (``gitea_commit_files``, ``gitea_create_pr``) derives
|
||||
its workspace from the author issue lock under #618, so nothing can create
|
||||
that remote head without first holding the lock.
|
||||
|
||||
This module supplies the missing operation: it publishes an *already committed*
|
||||
local head to the remote branch, so exact-owner renewal has the evidence it
|
||||
requires. It deliberately does **not** renew, reclaim, rebind, or clear any
|
||||
lock. Publication is the whole of its authority.
|
||||
|
||||
Why this is not a lock bypass
|
||||
-----------------------------
|
||||
The operation can only publish a branch whose **durable issue-lock record
|
||||
already names the caller as claimant**. Ownership is read from the lock file on
|
||||
disk (``issue_lock_store``), never from a caller-supplied flag, so the tool
|
||||
cannot manufacture a claim it does not already hold. Nothing here weakens the
|
||||
#510/#618/#713 guards: a dirty tree, an unregistered worktree, a foreign
|
||||
claimant, a changed HEAD, or a divergent remote head each refuse, exactly as
|
||||
they do today. The only thing this adds is the ability to make an existing,
|
||||
owned, committed, clean branch observable on the remote.
|
||||
|
||||
Separation of records (#812 AC23)
|
||||
---------------------------------
|
||||
The durable **issue-lock file** and the control-plane **workflow lease** are
|
||||
distinct records. This module reads the former as ownership evidence and writes
|
||||
neither. Publishing changes remote git state only; no lock is renewed,
|
||||
abandoned, reclaimed, or generation-bumped here.
|
||||
|
||||
Process evidence (#812 AC24)
|
||||
----------------------------
|
||||
Liveness of the lock's recorded pid is **not consulted**. That is deliberate:
|
||||
the recorded pid routinely belongs to the long-running MCP daemon rather than to
|
||||
an active author client, and the existing reclaim predicate
|
||||
(``assess_expired_lock_reclaim``) can never be satisfied while that daemon runs.
|
||||
Publication does not require the recording process to be dead, so this module
|
||||
never asserts, infers, or depends on a process being dead. Ownership is proven
|
||||
by identity and profile match against the recorded claimant instead.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
import os
|
||||
import re
|
||||
import subprocess
|
||||
|
||||
from reviewer_worktree import parse_dirty_tracked_files
|
||||
from stable_branch_push_guard import is_stable_ref, redact_command
|
||||
|
||||
# Assessment outcomes.
|
||||
PUBLISH_SANCTIONED = "publish_sanctioned"
|
||||
ALREADY_PUBLISHED = "already_published"
|
||||
REFUSED = "refused"
|
||||
|
||||
#: Implementation branches must stay traceable to their issue (#713 lineage).
|
||||
ISSUE_BRANCH_RE = re.compile(r"^(fix|feat|docs|chore)/issue-(\d+)-.+$")
|
||||
|
||||
_SHA_RE = re.compile(r"^[0-9a-f]{40}$")
|
||||
|
||||
|
||||
def _text(value: object) -> str:
|
||||
return value.strip() if isinstance(value, str) else ""
|
||||
|
||||
|
||||
def _realpath(value: str | None) -> str | None:
|
||||
path = _text(value)
|
||||
return os.path.realpath(path) if path else None
|
||||
|
||||
|
||||
def parse_untracked_files(porcelain_status: str) -> list[str]:
|
||||
"""Return untracked paths from ``git status --porcelain`` output.
|
||||
|
||||
``reviewer_worktree.parse_dirty_tracked_files`` deliberately skips ``??``
|
||||
entries. Publication needs both halves: an untracked file in the worktree is
|
||||
unpublished content that the commit does not carry, so publishing would
|
||||
silently leave it behind.
|
||||
"""
|
||||
untracked: list[str] = []
|
||||
for line in (porcelain_status or "").splitlines():
|
||||
if not line.startswith("??"):
|
||||
continue
|
||||
path = line[2:].strip()
|
||||
if path:
|
||||
untracked.append(path)
|
||||
return untracked
|
||||
|
||||
|
||||
def hash_worktree_files(worktree_path: str, paths) -> dict[str, str | None]:
|
||||
"""SHA-256 each path under *worktree_path*; ``None`` when unreadable."""
|
||||
root = _text(worktree_path)
|
||||
hashes: dict[str, str | None] = {}
|
||||
for rel in paths or ():
|
||||
rel_text = _text(rel)
|
||||
if not rel_text:
|
||||
continue
|
||||
full = os.path.join(root, rel_text)
|
||||
try:
|
||||
with open(full, "rb") as handle:
|
||||
digest = hashlib.sha256()
|
||||
for chunk in iter(lambda: handle.read(65536), b""):
|
||||
digest.update(chunk)
|
||||
hashes[rel_text] = digest.hexdigest()
|
||||
except OSError:
|
||||
hashes[rel_text] = None
|
||||
return hashes
|
||||
|
||||
|
||||
def read_remote_branch_head(
|
||||
worktree_path: str, remote_name: str, branch_name: str
|
||||
) -> dict:
|
||||
"""Observe the remote head for *branch_name*, read-only.
|
||||
|
||||
``probe_ok`` False means git could not answer at all. That is kept distinct
|
||||
from "the branch does not exist": an unobservable remote must fail closed
|
||||
rather than be mistaken for an absent branch, because the two lead to
|
||||
opposite dispositions.
|
||||
"""
|
||||
path = _text(worktree_path)
|
||||
remote = _text(remote_name)
|
||||
branch = _text(branch_name)
|
||||
result: dict = {
|
||||
"probe_ok": False,
|
||||
"remote_branch_exists": False,
|
||||
"remote_head_sha": None,
|
||||
"reasons": [],
|
||||
}
|
||||
if not (path and remote and branch):
|
||||
result["reasons"].append(
|
||||
"remote head probe requires a worktree path, remote name, and branch"
|
||||
)
|
||||
return result
|
||||
try:
|
||||
res = subprocess.run(
|
||||
["git", "-C", path, "ls-remote", remote, f"refs/heads/{branch}"],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
check=False,
|
||||
)
|
||||
except OSError as exc: # git unavailable — fail closed, never assume absent
|
||||
result["reasons"].append(f"remote head probe could not run: {exc}")
|
||||
return result
|
||||
if res.returncode != 0:
|
||||
result["reasons"].append(
|
||||
f"remote head probe failed for '{branch}' on remote '{remote}'"
|
||||
)
|
||||
return result
|
||||
|
||||
result["probe_ok"] = True
|
||||
for line in (res.stdout or "").splitlines():
|
||||
parts = line.split()
|
||||
if len(parts) >= 2 and parts[1] == f"refs/heads/{branch}":
|
||||
result["remote_branch_exists"] = True
|
||||
result["remote_head_sha"] = parts[0].strip()
|
||||
break
|
||||
return result
|
||||
|
||||
|
||||
def read_is_ancestor(
|
||||
worktree_path: str, ancestor_sha: str, descendant_sha: str
|
||||
) -> dict:
|
||||
"""Observe whether *ancestor_sha* is an ancestor of *descendant_sha*."""
|
||||
path = _text(worktree_path)
|
||||
ancestor = _text(ancestor_sha)
|
||||
descendant = _text(descendant_sha)
|
||||
result: dict = {"probe_ok": False, "is_ancestor": False, "reasons": []}
|
||||
if not (path and ancestor and descendant):
|
||||
result["reasons"].append(
|
||||
"ancestry probe requires a worktree path and both commit SHAs"
|
||||
)
|
||||
return result
|
||||
try:
|
||||
present = subprocess.run(
|
||||
["git", "-C", path, "rev-parse", "--verify", "--quiet",
|
||||
f"{ancestor}^{{commit}}"],
|
||||
capture_output=True, text=True, check=False,
|
||||
)
|
||||
if present.returncode != 0:
|
||||
result["reasons"].append(
|
||||
f"remote head {ancestor} is not present locally, so it cannot be "
|
||||
"proven an ancestor of the commit being published"
|
||||
)
|
||||
return result
|
||||
res = subprocess.run(
|
||||
["git", "-C", path, "merge-base", "--is-ancestor", ancestor, descendant],
|
||||
capture_output=True, text=True, check=False,
|
||||
)
|
||||
except OSError as exc:
|
||||
result["reasons"].append(f"ancestry probe could not run: {exc}")
|
||||
return result
|
||||
result["probe_ok"] = res.returncode in (0, 1)
|
||||
result["is_ancestor"] = res.returncode == 0
|
||||
return result
|
||||
|
||||
|
||||
def assess_unpublished_commit_publication(
|
||||
existing_lock,
|
||||
*,
|
||||
issue_number: int,
|
||||
branch_name: str,
|
||||
worktree_path: str,
|
||||
expected_head: str,
|
||||
remote: str,
|
||||
org: str,
|
||||
repo: str,
|
||||
identity: str | None,
|
||||
profile: str | None,
|
||||
worktree_state,
|
||||
worktree_registered: bool | None = None,
|
||||
remote_probe=None,
|
||||
ancestry=None,
|
||||
competing_open_prs=(),
|
||||
expected_file_hashes=None,
|
||||
observed_file_hashes=None,
|
||||
) -> dict:
|
||||
"""Decide whether an unpublished local commit may be published (#812 AC20).
|
||||
|
||||
Pure predicate. Every input is either a caller-declared expectation that
|
||||
must be *matched* against observation, or a server-side observation. No
|
||||
caller-supplied boolean is accepted as proof of ownership, liveness, or
|
||||
eligibility: ``existing_lock`` comes from the durable lock file and the
|
||||
git/PR state is observed by the server.
|
||||
|
||||
The single mutating disposition it can return is "publish this exact commit
|
||||
to this exact branch". It never sanctions renewal, reclamation, force
|
||||
updates, history rewriting, or publication of uncommitted content.
|
||||
"""
|
||||
reasons: list[str] = []
|
||||
branch = _text(branch_name)
|
||||
head = _text(expected_head).lower()
|
||||
workspace = _realpath(worktree_path)
|
||||
state = worktree_state if isinstance(worktree_state, dict) else {}
|
||||
lock = existing_lock if isinstance(existing_lock, dict) else None
|
||||
|
||||
evidence: dict = {
|
||||
"issue_number": issue_number,
|
||||
"branch_name": branch or None,
|
||||
"worktree_path": workspace,
|
||||
"expected_head": head or None,
|
||||
"remote": _text(remote) or None,
|
||||
"org": _text(org) or None,
|
||||
"repo": _text(repo) or None,
|
||||
"identity": _text(identity) or None,
|
||||
"profile": _text(profile) or None,
|
||||
"lock_record_present": lock is not None,
|
||||
"recorded_claimant": None,
|
||||
"recorded_branch": None,
|
||||
"recorded_worktree": None,
|
||||
"lock_generation": None,
|
||||
"local_head_sha": _text(state.get("head_sha")) or None,
|
||||
"current_branch": _text(state.get("current_branch")) or None,
|
||||
"dirty_tracked_files": [],
|
||||
"untracked_files": [],
|
||||
"worktree_registered": worktree_registered,
|
||||
"remote_branch_exists": None,
|
||||
"remote_head_sha": None,
|
||||
"fast_forward_from_remote": None,
|
||||
"competing_open_prs": [],
|
||||
"file_hashes_verified": None,
|
||||
"hash_mismatches": [],
|
||||
# Recorded explicitly so no reader mistakes silence for a liveness
|
||||
# claim, and so the audit shows which records were left alone (AC23/AC24).
|
||||
"owner_pid_liveness_consulted": False,
|
||||
"workflow_lease_touched": False,
|
||||
"issue_lock_record_mutated": False,
|
||||
}
|
||||
|
||||
# ── declared shape ────────────────────────────────────────────────────
|
||||
if not branch:
|
||||
reasons.append("branch name not declared; fail closed")
|
||||
if not head:
|
||||
reasons.append(
|
||||
"expected_head not declared; publication must name the exact commit"
|
||||
)
|
||||
elif not _SHA_RE.match(head):
|
||||
reasons.append(
|
||||
f"expected_head '{head}' is not a full 40-character commit SHA; "
|
||||
"abbreviated or symbolic revisions are refused"
|
||||
)
|
||||
if not workspace:
|
||||
reasons.append("worktree path not declared; fail closed")
|
||||
|
||||
if branch:
|
||||
match = ISSUE_BRANCH_RE.match(branch)
|
||||
if not match:
|
||||
reasons.append(
|
||||
f"branch '{branch}' is not an issue-linked implementation branch "
|
||||
"((fix|feat|docs|chore)/issue-<number>-<description>); fail closed"
|
||||
)
|
||||
elif int(match.group(2)) != int(issue_number):
|
||||
reasons.append(
|
||||
f"branch '{branch}' does not carry issue number {issue_number}; "
|
||||
"fail closed"
|
||||
)
|
||||
if is_stable_ref(branch):
|
||||
reasons.append(
|
||||
f"refusing to publish stable branch '{branch}'; this operation "
|
||||
"publishes issue branches only"
|
||||
)
|
||||
|
||||
# ── ownership: durable issue-lock record only (AC8, AC20, AC24) ───────
|
||||
if lock is None:
|
||||
reasons.append(
|
||||
"no durable issue-lock record for this issue; publication requires an "
|
||||
"existing recorded claim naming the caller, so this operation cannot "
|
||||
"be used to bypass the author lock"
|
||||
)
|
||||
else:
|
||||
lease = lock.get("work_lease")
|
||||
lease = lease if isinstance(lease, dict) else {}
|
||||
claimant = lease.get("claimant")
|
||||
claimant = claimant if isinstance(claimant, dict) else {}
|
||||
recorded_user = _text(claimant.get("username"))
|
||||
recorded_profile = _text(claimant.get("profile"))
|
||||
recorded_branch = _text(lock.get("branch_name")) or _text(lease.get("branch"))
|
||||
recorded_worktree = _realpath(
|
||||
_text(lock.get("worktree_path")) or _text(lease.get("worktree_path"))
|
||||
)
|
||||
evidence["recorded_claimant"] = {
|
||||
"username": recorded_user or None,
|
||||
"profile": recorded_profile or None,
|
||||
}
|
||||
evidence["recorded_branch"] = recorded_branch or None
|
||||
evidence["recorded_worktree"] = recorded_worktree
|
||||
try:
|
||||
evidence["lock_generation"] = int(lock.get("lock_generation") or 0)
|
||||
except (TypeError, ValueError):
|
||||
evidence["lock_generation"] = 0
|
||||
|
||||
try:
|
||||
recorded_issue = int(lock.get("issue_number") or 0)
|
||||
except (TypeError, ValueError):
|
||||
recorded_issue = 0
|
||||
if recorded_issue != int(issue_number):
|
||||
reasons.append(
|
||||
f"durable lock records issue {lock.get('issue_number')}, not "
|
||||
f"{issue_number}; ambiguous ownership, fail closed"
|
||||
)
|
||||
for field, declared in (
|
||||
("remote", _text(remote)),
|
||||
("org", _text(org)),
|
||||
("repo", _text(repo)),
|
||||
):
|
||||
recorded = _text(lock.get(field))
|
||||
if recorded and declared and recorded != declared:
|
||||
reasons.append(
|
||||
f"durable lock records {field} '{recorded}' but the request "
|
||||
f"declares '{declared}'; repository mismatch, fail closed"
|
||||
)
|
||||
if recorded_branch and branch and recorded_branch != branch:
|
||||
reasons.append(
|
||||
f"durable lock records branch '{recorded_branch}' but the request "
|
||||
f"declares '{branch}'; fail closed"
|
||||
)
|
||||
if recorded_worktree and workspace and recorded_worktree != workspace:
|
||||
reasons.append(
|
||||
f"durable lock records worktree '{recorded_worktree}' but the "
|
||||
f"request declares '{workspace}'; fail closed"
|
||||
)
|
||||
if not recorded_user or not recorded_profile:
|
||||
reasons.append(
|
||||
"durable lock does not record a claimant username and profile; "
|
||||
"ownership cannot be proven, fail closed"
|
||||
)
|
||||
else:
|
||||
if recorded_user != _text(identity):
|
||||
reasons.append(
|
||||
f"durable lock claimant '{recorded_user}' is not the acting "
|
||||
f"identity '{_text(identity) or '(unknown)'}'; foreign claim, "
|
||||
"fail closed"
|
||||
)
|
||||
if recorded_profile != _text(profile):
|
||||
reasons.append(
|
||||
f"durable lock claimant profile '{recorded_profile}' is not "
|
||||
f"the active profile '{_text(profile) or '(unknown)'}'; "
|
||||
"fail closed"
|
||||
)
|
||||
|
||||
# ── worktree: registered, on-branch, clean, at the expected commit ────
|
||||
if worktree_registered is False:
|
||||
reasons.append(
|
||||
f"worktree '{workspace}' is not listed in git worktree list; #713 "
|
||||
"requires a genuinely registered worktree, fail closed"
|
||||
)
|
||||
|
||||
current_branch = _text(state.get("current_branch"))
|
||||
if not current_branch:
|
||||
reasons.append("worktree branch could not be observed; fail closed")
|
||||
elif branch and current_branch != branch:
|
||||
reasons.append(
|
||||
f"worktree is on branch '{current_branch}', not '{branch}'; fail closed"
|
||||
)
|
||||
|
||||
porcelain = state.get("porcelain_status") or ""
|
||||
dirty_tracked = parse_dirty_tracked_files(porcelain)
|
||||
untracked = parse_untracked_files(porcelain)
|
||||
evidence["dirty_tracked_files"] = dirty_tracked
|
||||
evidence["untracked_files"] = untracked
|
||||
if dirty_tracked:
|
||||
reasons.append(
|
||||
"worktree has dirty tracked files, so the commit is not the whole of "
|
||||
f"the work: {', '.join(dirty_tracked)}. This operation publishes an "
|
||||
"existing clean commit only; uncommitted content is out of scope"
|
||||
)
|
||||
if untracked:
|
||||
reasons.append(
|
||||
"worktree has untracked files that the commit does not carry: "
|
||||
f"{', '.join(untracked)}. Publishing would silently leave them "
|
||||
"behind; fail closed"
|
||||
)
|
||||
|
||||
local_head = _text(state.get("head_sha")).lower()
|
||||
if not local_head:
|
||||
reasons.append("local HEAD could not be observed; fail closed")
|
||||
elif head and local_head != head:
|
||||
reasons.append(
|
||||
f"worktree HEAD is {local_head} but the request declares {head}; the "
|
||||
"local commit changed since it was recorded, fail closed"
|
||||
)
|
||||
|
||||
# ── remote state ──────────────────────────────────────────────────────
|
||||
probe = remote_probe if isinstance(remote_probe, dict) else {}
|
||||
already_published = False
|
||||
if not probe.get("probe_ok"):
|
||||
reasons.append(
|
||||
"remote branch head could not be observed; publication must not "
|
||||
"proceed against an unknown remote state, fail closed"
|
||||
)
|
||||
reasons.extend(probe.get("reasons") or [])
|
||||
else:
|
||||
remote_exists = bool(probe.get("remote_branch_exists"))
|
||||
remote_head = _text(probe.get("remote_head_sha")).lower() or None
|
||||
evidence["remote_branch_exists"] = remote_exists
|
||||
evidence["remote_head_sha"] = remote_head
|
||||
if remote_exists and remote_head and head:
|
||||
if remote_head == head:
|
||||
already_published = True
|
||||
evidence["fast_forward_from_remote"] = True
|
||||
else:
|
||||
anc = ancestry if isinstance(ancestry, dict) else {}
|
||||
is_anc = bool(anc.get("probe_ok")) and bool(anc.get("is_ancestor"))
|
||||
evidence["fast_forward_from_remote"] = is_anc
|
||||
if not is_anc:
|
||||
reasons.append(
|
||||
f"remote branch '{branch}' already exists at {remote_head}, "
|
||||
f"which is not an ancestor of {head}; publishing would "
|
||||
"discard or rewrite published history, fail closed"
|
||||
)
|
||||
reasons.extend(anc.get("reasons") or [])
|
||||
elif remote_exists and not remote_head:
|
||||
reasons.append(
|
||||
f"remote branch '{branch}' exists but its head could not be read; "
|
||||
"fail closed"
|
||||
)
|
||||
|
||||
# ── competing claims ──────────────────────────────────────────────────
|
||||
competing = [p for p in (competing_open_prs or ()) if p]
|
||||
evidence["competing_open_prs"] = list(competing)
|
||||
if competing:
|
||||
reasons.append(
|
||||
f"open pull request(s) {competing} already claim issue {issue_number} "
|
||||
"or this branch; ambiguous ownership, fail closed"
|
||||
)
|
||||
|
||||
# ── content verification before publication ───────────────────────────
|
||||
if expected_file_hashes:
|
||||
observed = (
|
||||
observed_file_hashes if isinstance(observed_file_hashes, dict) else {}
|
||||
)
|
||||
mismatches: list[str] = []
|
||||
for path, expected_digest in dict(expected_file_hashes).items():
|
||||
actual = observed.get(path)
|
||||
if actual is None:
|
||||
mismatches.append(f"{path}: missing or unreadable in the worktree")
|
||||
elif _text(actual).lower() != _text(expected_digest).lower():
|
||||
mismatches.append(
|
||||
f"{path}: expected {expected_digest}, observed {actual}"
|
||||
)
|
||||
evidence["hash_mismatches"] = mismatches
|
||||
evidence["file_hashes_verified"] = not mismatches
|
||||
if mismatches:
|
||||
reasons.append(
|
||||
"declared content hashes do not match the worktree: "
|
||||
+ "; ".join(mismatches)
|
||||
+ ". Refusing to publish content that is not what was recorded"
|
||||
)
|
||||
|
||||
if reasons:
|
||||
return {
|
||||
"outcome": REFUSED,
|
||||
"publish_sanctioned": False,
|
||||
"already_published": False,
|
||||
"reasons": reasons,
|
||||
"evidence": evidence,
|
||||
}
|
||||
return {
|
||||
"outcome": ALREADY_PUBLISHED if already_published else PUBLISH_SANCTIONED,
|
||||
# Idempotent retry: a remote head that already equals the assessed commit
|
||||
# needs no second push, so the caller verifies instead of acting.
|
||||
"publish_sanctioned": not already_published,
|
||||
"already_published": already_published,
|
||||
"reasons": [],
|
||||
"evidence": evidence,
|
||||
}
|
||||
|
||||
|
||||
def publish_commit_to_remote_branch(
|
||||
*,
|
||||
worktree_path: str,
|
||||
remote_name: str,
|
||||
branch_name: str,
|
||||
expected_head: str,
|
||||
) -> dict:
|
||||
"""Send exactly *expected_head* to ``refs/heads/<branch_name>``.
|
||||
|
||||
The refspec names the commit SHA explicitly rather than ``HEAD`` or the
|
||||
local branch, so what lands is the commit that was assessed and nothing
|
||||
else. No force, no lease, no ``+`` prefix: a non-fast-forward is rejected by
|
||||
git itself, the last of several independent guards against overwriting
|
||||
published history.
|
||||
"""
|
||||
path = _text(worktree_path)
|
||||
remote = _text(remote_name)
|
||||
branch = _text(branch_name)
|
||||
head = _text(expected_head)
|
||||
result: dict = {
|
||||
"success": False,
|
||||
"pushed_ref": f"refs/heads/{branch}" if branch else None,
|
||||
"pushed_sha": head or None,
|
||||
"stderr": None,
|
||||
"reasons": [],
|
||||
}
|
||||
if not (path and remote and branch and head):
|
||||
result["reasons"].append(
|
||||
"publication requires a worktree path, remote, branch, and commit SHA"
|
||||
)
|
||||
return result
|
||||
|
||||
refspec = f"{head}:refs/heads/{branch}"
|
||||
try:
|
||||
res = subprocess.run(
|
||||
["git", "-C", path, "push", remote, refspec],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
check=False,
|
||||
)
|
||||
except OSError as exc:
|
||||
result["reasons"].append(f"publication could not run: {exc}")
|
||||
return result
|
||||
|
||||
if res.returncode != 0:
|
||||
# Redact before surfacing: failures can echo credentialed remote URLs.
|
||||
result["stderr"] = redact_command(res.stderr or "")
|
||||
result["reasons"].append(
|
||||
f"publication of {head} to '{branch}' on remote '{remote}' failed"
|
||||
)
|
||||
return result
|
||||
|
||||
result["success"] = True
|
||||
return result
|
||||
|
||||
|
||||
def verify_published_head(
|
||||
*, worktree_path: str, remote_name: str, branch_name: str, expected_head: str
|
||||
) -> dict:
|
||||
"""Read-after-write: confirm the remote head equals *expected_head* (AC20)."""
|
||||
probe = read_remote_branch_head(worktree_path, remote_name, branch_name)
|
||||
head = _text(expected_head).lower()
|
||||
observed = _text(probe.get("remote_head_sha")).lower() or None
|
||||
verified = bool(head) and bool(probe.get("probe_ok")) and observed == head
|
||||
reasons: list[str] = list(probe.get("reasons") or [])
|
||||
if probe.get("probe_ok") and not verified:
|
||||
reasons.append(
|
||||
"read-after-write verification failed: remote head is "
|
||||
f"{observed or '(absent)'}, expected {head}"
|
||||
)
|
||||
return {
|
||||
"verified": verified,
|
||||
"remote_head_sha": observed,
|
||||
"expected_head": head or None,
|
||||
"reasons": reasons,
|
||||
}
|
||||
@@ -0,0 +1,201 @@
|
||||
# ADR: MCP Control Plane Web Console architecture and information architecture
|
||||
|
||||
- **Status:** Proposed (documentation only; blocks no code, gates every #631 child)
|
||||
- **Date:** 2026-07-22
|
||||
- **Tracking issue:** [#632](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/632) — architecture and information architecture (Phase 1)
|
||||
- **Parent epic:** [#631](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/631) — MCP Control Plane Web Console
|
||||
- **Foundation (closed, extend — do not recreate):** #425 tracker and children #426 skeleton, #427 projects, #428 prompts, #429 queue, #430 runtime, #431 audit paste, #432 worktrees, #433 leases, #434 gated actions, #435 auth/deployment boundary, #436 tests/CI
|
||||
- **Related:** `mcp-allocator-control-plane-observability-adr.md`, `mcp-stable-control-runtime-policy-adr.md`, `control-plane-db-substrate.md`, `../safety-model.md`, `../tool-boundaries.md`, `../credential-isolation.md`, `../webui-local-dev.md`, `../webui-deployment.md`
|
||||
|
||||
## 1. Context
|
||||
|
||||
The MVP web UI shipped under `webui/` as a read-only Starlette application with ten operator routes and a JSON export beside most of them. It is a working foundation, not the console product described by epic #631, and it carries no durable architecture record: no layer contract, no authority boundary, no API versioning rule, no page map, and no statement of which phase may open a write path.
|
||||
|
||||
Twenty children (#632–#651) hang off #631. Without one architecture document each implementer re-derives boundaries, and the most likely failure is not a bad view — it is a privileged action wired into the browser before the authorization and audit model of #633 exists.
|
||||
|
||||
This ADR is the single retrievable design source for the console. It decides structure only. It implements no UI, no API, and no change to deployment topology.
|
||||
|
||||
## 2. Decision summary (core)
|
||||
|
||||
| Layer | Owns | Must not |
|
||||
|-------|------|----------|
|
||||
| **Browser UI** | Rendering, navigation, operator affordances | Hold tokens, call Gitea/providers directly, or execute an action the server did not gate |
|
||||
| **HTTP route layer** (`webui/app.py`) | Versioned routing, authentication, authorization, redaction boundary, audit emission | Contain domain logic or reach past a loader to a raw credential |
|
||||
| **Domain loaders** (`webui/*_loader.py`, `*_scanner.py`, `runtime_health.py`, `project_registry.py`) | Assembling read models from authoritative sources | Mutate anything, or emit unredacted secrets across the boundary |
|
||||
| **Gitea** | Durable work record: issues, PRs, comments, reviews, labels, merges | Be the concurrency lock under multi-session load |
|
||||
| **Control-plane DB** | Sessions, assignment, leases, heartbeats, events | Replace Gitea history |
|
||||
| **MCP tools / capability gates** | Mutation authorization | Be re-implemented, mirrored, or bypassed by console code |
|
||||
| **External providers** (Sentry/GlitchTip, AI providers) | Incident and usage data | Assign work or mutate Gitea outside the #612 bridge |
|
||||
|
||||
**One-liner:** **Gitea records. The DB coordinates. MCP tools authorize. The console projects state and executes only capability-checked, audited actions. Providers observe.**
|
||||
|
||||
## 3. Console surface today versus target
|
||||
|
||||
`webui/app.py` currently registers these routes (see `../webui-local-dev.md` for the operator-facing table): `/`, `/health`, `/queue`, `/projects`, `/projects/{id}`, `/prompts`, `/runtime`, `/audit`, `/worktrees`, `/leases`, `/actions`, and the unversioned exports `/api/queue`, `/api/projects`, `/api/prompts`, `/api/runtime`, `/api/audit`, `/api/worktrees`, `/api/leases`, `/api/actions`, `/api/actions/{id}/preview`, `/api/actions/{id}/attempt`.
|
||||
|
||||
Every one of these is **retained and evolved**. No child issue may recreate a route from scratch; each states in its PR which MVP surface it extends and what it changes.
|
||||
|
||||
## 4. Authority boundaries
|
||||
|
||||
### 4.1 Gitea (durable record)
|
||||
|
||||
Authoritative for issue and PR identity and state, comments, reviews and verdicts, labels, merges, and branch refs. When the console and Gitea disagree about durable state, Gitea wins and the console view is refreshed — never the reverse.
|
||||
|
||||
### 4.2 Control-plane DB (coordination)
|
||||
|
||||
Authoritative for live coordination: which session holds which assignment or lease, heartbeat freshness, expiry, and the allocation event log. The console reads it; only allocator and lease tools write it.
|
||||
|
||||
### 4.3 MCP capability gates (authorization)
|
||||
|
||||
`task_capability_map.py` and `gitea_resolve_task_capability` remain the only authority that decides whether a mutation may run. The console asks; it never answers. A console action that cannot name the MCP tool it delegates to is not an action — it is a defect.
|
||||
|
||||
### 4.4 Filesystem and git (local state)
|
||||
|
||||
Issue lock files, `branches/` worktrees, and registered git worktrees are read through existing scanners. The console never deletes, rebinds, or force-clears local state outside a Phase 2 gated action.
|
||||
|
||||
### 4.5 Providers (observe only)
|
||||
|
||||
Sentry/GlitchTip and AI providers are read surfaces. The #612 incident bridge is the only path that turns an observation into Gitea work.
|
||||
|
||||
## 5. Request flow and the redaction boundary
|
||||
|
||||
```text
|
||||
browser ──HTTP──> route layer ──> domain loader ──> Gitea REST
|
||||
│ ├──> control-plane DB
|
||||
│ ├──> filesystem / git
|
||||
│ └──> providers
|
||||
│
|
||||
[redaction boundary]
|
||||
│
|
||||
audit event
|
||||
```
|
||||
|
||||
| Stage | May hold credentials | Emits |
|
||||
|-------|----------------------|-------|
|
||||
| Loader → route layer | yes (server-side, via `gitea_auth`) | domain objects |
|
||||
| Route layer → browser | **no** | redacted DTOs, HTML |
|
||||
|
||||
Two invariants govern the boundary and are non-negotiable for every child:
|
||||
|
||||
1. **No secrets to the browser.** Tokens, keychain identifiers, Authorization headers, raw provider endpoints, and credential-bearing URLs are redacted by default, consistent with `../safety-model.md` §3 and `../credential-isolation.md`. Serializers redact; templates do not sanitize after the fact.
|
||||
2. **No ungated mutations.** A write reaches an authoritative system only by delegating to an MCP tool that passed its own capability gate. HTML forms and JSON endpoints are transport, never authority.
|
||||
|
||||
## 6. API naming and versioning
|
||||
|
||||
**Decision:** all console APIs added from Phase 1 onward are served under `/api/v1/...`.
|
||||
|
||||
- Nouns are plural and hierarchical: `/api/v1/inventory/leases`, `/api/v1/system/health`.
|
||||
- Read endpoints are `GET` and side-effect free.
|
||||
- Phase 2 action endpoints are `POST /api/v1/actions/{action_id}/preview` and `POST /api/v1/actions/{action_id}/execute`; `preview` stays side-effect free and returns a mutation ledger.
|
||||
- The existing unversioned MVP exports remain as **compatibility aliases** for the whole of Phase 1 so the current operator flow never breaks. They may be retired no earlier than Phase 2, and only after the replacing `v1` route ships and `../webui-local-dev.md` records the swap.
|
||||
- A breaking change to a `v1` payload requires `/api/v2/...`, not an in-place edit.
|
||||
- Every JSON payload carries enough provenance for an auditor to tell where the data came from — at minimum the source system and whether the inventory was complete, matching the pagination-proof habit the MVP queue export already established.
|
||||
|
||||
## 7. Page map
|
||||
|
||||
| Page | Purpose | Owning child | Evolves |
|
||||
|------|---------|--------------|---------|
|
||||
| `/` | Console shell, navigation, next-safe-action summary | #638 | MVP `/` (#426) |
|
||||
| `/system` | System-health dashboard | #639 | new, backed by #634 |
|
||||
| `/traffic` | Workflow traffic control, queues, blockers | #640 | MVP `/queue` (#429) |
|
||||
| `/runtime` | Runtime and session view | #641 | MVP `/runtime` (#430) |
|
||||
| `/projects`, `/projects/{id}` | Project registry and onboarding | #635 | MVP `/projects` (#427) |
|
||||
| `/inventory` | Sessions, leases, locks, worktrees in one surface | #636 | MVP `/leases` (#433) + `/worktrees` (#432) |
|
||||
| `/timeline` | Workflow events and conversation timeline | #637 | new |
|
||||
| `/actions` | Gated action registry, preview, execution | #642, #643, #644 | MVP `/actions` (#434) |
|
||||
| `/gitea` | Issue and PR linkage console | #645 | new |
|
||||
| `/policy` | Guardrail visibility, then versioned editing | #646, #647 | new |
|
||||
| `/notifications` | Human-attention routing | #648 | new |
|
||||
| `/observability` | Sentry/GlitchTip correlation and durable issue creation | #649 | new |
|
||||
| `/providers` | AI-provider connections and insights | #650 | new |
|
||||
| `/analytics` | Usage, token cost, latency, workflow performance | #651 | new |
|
||||
| `/audit` | Final-report validator preview and audit log | #431 foundation, extended by #633 | MVP `/audit` (#431) |
|
||||
| `/prompts`, `/prompts/{id}` | Canonical prompt library | #638 | MVP `/prompts` (#428) |
|
||||
| `/health` | Liveness and deployment metadata | #634 | MVP `/health` (#435) |
|
||||
|
||||
## 8. Component ownership for every epic child
|
||||
|
||||
Each #631 child maps to at least one architectural component defined above.
|
||||
|
||||
| Child | Capability area | Primary component | Phase |
|
||||
|-------|-----------------|-------------------|-------|
|
||||
| #632 | Architecture and information architecture | this ADR | 1 |
|
||||
| #633 | Authorization, RBAC, secret redaction, audit and retention | route layer + redaction boundary (§5) | 1 |
|
||||
| #634 | Read-only system-health API | `/api/v1/system/health` + health loader | 1 |
|
||||
| #635 | Project registry API evolution | `/api/v1/projects` + `project_registry.py` | 1 |
|
||||
| #636 | Session, lease, lock, worktree inventory API | `/api/v1/inventory/*` + `lease_loader.py`, `worktree_scanner.py` | 1 |
|
||||
| #637 | Workflow-event and conversation timeline model | `/api/v1/events` + control-plane DB event log | 1 |
|
||||
| #638 | Application shell evolution | browser UI layer + `layout.py` | 1 |
|
||||
| #639 | System-health dashboard | `/system` page over #634 | 1 |
|
||||
| #640 | Workflow traffic-control view | `/traffic` page over the queue loader | 1 |
|
||||
| #641 | Runtime and session view | `/runtime` page over `runtime_health.py` | 1 |
|
||||
| #642 | Sanctioned restart and graceful reload controls | gated action framework, restart class | 2 |
|
||||
| #643 | Requests, intent preview, authorization, workflow initiation | `/api/v1/actions/*` execute path | 2 |
|
||||
| #644 | Stale-runtime recovery, worktree rebinding, reconciliation controls | gated actions over filesystem/git authority | 2 |
|
||||
| #645 | Gitea issue and PR linkage console | `/gitea` page over Gitea authority | 3 |
|
||||
| #646 | Workflow policy and guardrail visibility | `/policy` read view over the capability map | 3 |
|
||||
| #647 | Versioned policy editing, validation, simulation, approval, rollback | `/policy` write path, gated | 3 |
|
||||
| #648 | Notifications and human-attention routing | notification component over the event model | 3 |
|
||||
| #649 | Sentry/GlitchTip connections, correlation, durable issue creation | provider layer + #612 incident bridge | 4 |
|
||||
| #650 | AI-provider connections and operational insights | provider layer | 4 |
|
||||
| #651 | Model usage, token cost, latency, workflow analytics | analytics component over the event model | 4 |
|
||||
|
||||
Related but **outside** this epic: #667 (restart status, impact preview, and approval controls) belongs to the #655 restart-governance umbrella and must reuse the #642 action class rather than adding a second restart surface.
|
||||
|
||||
## 9. Phase gates
|
||||
|
||||
| Phase | May ship | Entry condition |
|
||||
|-------|----------|-----------------|
|
||||
| **1 — read-only visibility** | `GET` pages and `GET /api/v1/...` | this ADR accepted |
|
||||
| **2 — controlled actions** | gated `POST` action execution | #633 authorization, RBAC, and audit model landed |
|
||||
| **3 — orchestration and policy** | linkage, policy visibility, versioned policy editing | Phase 1 inventory plus the Phase 2 action framework |
|
||||
| **4 — insights** | provider correlation, analytics | evidence-backed sources from Phases 1–3 |
|
||||
|
||||
Phase 1 must not open a mutation endpoint, and the read-only guard that returns `405 read-only-mvp` stays in force until the Phase 2 entry condition is met. A phase is not entered by exception; if a control is urgent, the entry condition is what gets prioritized.
|
||||
|
||||
## 10. Security and workflow safety
|
||||
|
||||
- **Fail closed** on unknown authentication, missing RBAC mapping, or ambiguous lease ownership. An unknown state renders as blocked, never as permitted.
|
||||
- **Redact by default**, per §5.
|
||||
- **Every privileged action** requires a resolved capability, an explicit operator confirmation, and a durable audit event naming actor, action, target, and outcome.
|
||||
- **Contamination surfaces.** Session contamination — including a manually killed MCP daemon (#630) — must be shown and must block clean claims rather than being silently repaired.
|
||||
- **Deployment boundary unchanged.** Loopback by default, with the existing refusal of public binds (#435). This ADR documents that target; it does not widen it.
|
||||
|
||||
## 11. Forbidden paths
|
||||
|
||||
These are rejected designs, not preferences:
|
||||
|
||||
1. **Raw provider incidents as work.** The allocator never receives an unclassified Sentry/GlitchTip incident; only the #612 bridge turns an observation into a Gitea issue.
|
||||
2. **Browser-held tokens.** No credential, keychain identifier, or Authorization header is ever sent to the browser or embedded in a client bundle.
|
||||
3. **Process-kill recovery.** The console must not expose `pkill`, process-identifier termination, or any host process kill as a recovery affordance (#630). Restart is the sanctioned, operator-owned path of #642 and the #655 umbrella.
|
||||
4. **Ungated browser mutations.** No review, approval, merge, close, or comment may originate from the browser without passing an MCP capability gate.
|
||||
5. **Policy invented in the console.** The console projects policy from the capability map and canonical workflows; it never encodes a second copy.
|
||||
6. **Recreating MVP scope.** Re-implementing a #426–#436 surface without an explicit evolve-or-extend statement is out of bounds.
|
||||
|
||||
## 12. Approval checklist (readable without chat history)
|
||||
|
||||
A controller can accept or reject this ADR against these six points alone:
|
||||
|
||||
1. Layers and their owners are defined (§2) and each authority is named (§4).
|
||||
2. The redaction boundary and the two invariants are stated (§5).
|
||||
3. API versioning is decided, including what happens to the existing unversioned routes (§6).
|
||||
4. A page map exists and names an owning child for every page (§7).
|
||||
5. Every #631 child maps to at least one component and one phase (§8).
|
||||
6. Phase gates and forbidden paths are explicit (§9, §11).
|
||||
|
||||
## 13. Open questions and follow-ups
|
||||
|
||||
Unresolved choices are recorded here rather than settled by implication. Each needs its own durable issue before the phase that depends on it:
|
||||
|
||||
- **Authentication mechanism.** Whether the console authenticates via an access proxy (Cloudflare Access or equivalent) or an application-level session is deferred to #633. This ADR requires only that it fail closed.
|
||||
- **Event model substrate.** Whether the #637 timeline reads the control-plane event log directly or through a projection is deferred to #637.
|
||||
- **CI path filter coverage.** `webui/ci_paths.py` triggers the web UI suite on `webui/`, `tests/test_webui_*`, and `docs/webui*`. This ADR lives under `docs/architecture/`, so editing it alone does not trigger that gate; the accompanying `tests/test_webui_architecture_docs.py` does run in the full suite. Widening the filter is a small follow-up, deliberately not bundled into a documentation-only change.
|
||||
- **Retention.** Audit-event retention duration is owned by #633.
|
||||
|
||||
## 14. Acceptance
|
||||
|
||||
Accepting this ADR means:
|
||||
|
||||
- Phase 1 children may proceed against the layers, page map, and API rules above.
|
||||
- Phase 2 children may not open a write path until #633 lands.
|
||||
- Any deviation is recorded as an amendment to this file with its own issue reference, not as an undocumented divergence in code.
|
||||
@@ -120,6 +120,7 @@ that gates each call, not which tools exist.
|
||||
- `gitea_observability_list_projects`
|
||||
- `gitea_observability_reconcile_incident`
|
||||
- `gitea_post_heartbeat`
|
||||
- `gitea_publish_unpublished_issue_branch`
|
||||
- `gitea_quarantine_contaminated_review`
|
||||
- `gitea_reclaim_expired_workflow_lease`
|
||||
- `gitea_reconcile_already_landed_pr`
|
||||
|
||||
@@ -37,6 +37,12 @@ Optional environment variables:
|
||||
See [webui-deployment.md](webui-deployment.md) for internal-only serving,
|
||||
Cloudflare Access/WARP/VPN guidance, and unsafe bind overrides (#435).
|
||||
|
||||
See
|
||||
[architecture/webui-control-plane-console-architecture-adr.md](architecture/webui-control-plane-console-architecture-adr.md)
|
||||
for the console architecture: layer and authority boundaries, the redaction
|
||||
boundary, `/api/v1/...` versioning, the target page map, and the phase gates
|
||||
that govern when a write path may open (#632, epic #631).
|
||||
|
||||
## Routes (MVP)
|
||||
|
||||
| Path | Description |
|
||||
|
||||
+296
-9
@@ -2024,6 +2024,7 @@ import review_quarantine # noqa: E402 # #695 contaminated formal-review quaran
|
||||
import mcp_daemon_guard # noqa: E402 # #695 native transport provenance
|
||||
import already_landed_reconcile # noqa: E402
|
||||
import author_mutation_worktree # noqa: E402
|
||||
import branch_publish # noqa: E402 # #812 AC20 unpublished-commit publication
|
||||
import root_checkout_guard # noqa: E402
|
||||
import workflow_scope_guard # noqa: E402 # #683 production scope / force-on guards
|
||||
import stable_branch_push_guard # noqa: E402
|
||||
@@ -9082,6 +9083,237 @@ def gitea_commit_files(
|
||||
}
|
||||
|
||||
|
||||
def _publication_block(reasons: list[str], **extra) -> dict:
|
||||
"""Uniform fail-closed shape for publication refusals (#812 AC20)."""
|
||||
payload = {
|
||||
"success": False,
|
||||
"performed": False,
|
||||
"published": False,
|
||||
"verified": False,
|
||||
"outcome": branch_publish.REFUSED,
|
||||
"reasons": reasons,
|
||||
# State every record this operation left alone, so a refusal can never
|
||||
# be misread as a lock mutation (#812 AC23).
|
||||
"issue_lock_record_mutated": False,
|
||||
"workflow_lease_touched": False,
|
||||
}
|
||||
payload.update(extra)
|
||||
return payload
|
||||
|
||||
|
||||
@mcp.tool()
|
||||
def gitea_publish_unpublished_issue_branch(
|
||||
issue_number: int,
|
||||
branch_name: str,
|
||||
worktree_path: str,
|
||||
expected_head: str,
|
||||
remote: str = "dadeschools",
|
||||
host: str | None = None,
|
||||
org: str | None = None,
|
||||
repo: str | None = None,
|
||||
git_remote_name: str | None = None,
|
||||
expected_file_hashes: dict | None = None,
|
||||
dry_run: bool = False,
|
||||
) -> dict:
|
||||
"""Publish an already-committed, unpublished issue branch (#812 AC20).
|
||||
|
||||
Creates the remote head for a branch whose work is *already* a local commit
|
||||
on a registered, clean worktree, so exact-owner lease renewal
|
||||
(``issue_lock_renewal``) has the published head its evidence model requires.
|
||||
This is the one step of the entry point B deadlock that no existing tool can
|
||||
perform: publication is otherwise lock-derived under #618, and the lock
|
||||
itself is withheld until a remote head exists.
|
||||
|
||||
Not a lock bypass. The branch's **durable issue-lock record must already
|
||||
name the caller as claimant** — ownership is read from the lock file, never
|
||||
asserted by the caller — so this can only publish work the caller already
|
||||
owns. It refuses a dirty or untracked-carrying worktree, an unregistered
|
||||
worktree, a changed local HEAD, a remote head that is not an ancestor of the
|
||||
commit, a competing open PR on another branch for the same issue, and any
|
||||
declared-hash mismatch. It renews, reclaims, and clears nothing: the durable
|
||||
issue-lock file and the control-plane workflow lease are both left untouched
|
||||
(#812 AC23), and the recorded owner pid's liveness is never consulted or
|
||||
asserted (#812 AC24).
|
||||
|
||||
Args:
|
||||
issue_number: The issue whose recorded claim authorizes publication.
|
||||
branch_name: Issue branch to publish, ``(fix|feat|docs|chore)/issue-N-…``.
|
||||
worktree_path: Registered worktree holding the commit.
|
||||
expected_head: Full 40-character SHA of the commit to publish. Required:
|
||||
publication names the exact commit, and a mismatch fails closed.
|
||||
remote: Known instance — 'dadeschools' or 'prgs'.
|
||||
host: Override the Gitea host.
|
||||
org: Override the owner/organization.
|
||||
repo: Override the repository name.
|
||||
git_remote_name: Git remote to publish to; defaults to *remote*.
|
||||
expected_file_hashes: Optional ``{path: sha256}`` verified against the
|
||||
worktree before publication. Any mismatch or missing file refuses.
|
||||
dry_run: Report the decision and evidence, mutate nothing.
|
||||
|
||||
Returns:
|
||||
dict with 'success', 'performed', 'published', 'verified', 'outcome',
|
||||
'remote_head_sha', 'reasons', and 'evidence'.
|
||||
"""
|
||||
task = "publish_unpublished_branch"
|
||||
ok, block_reasons = role_session_router.check_author_mutation_after_reviewer_stop(
|
||||
task
|
||||
)
|
||||
if not ok:
|
||||
return _publication_block(block_reasons)
|
||||
|
||||
blocked = _namespace_mutation_block(task, remote=remote)
|
||||
if blocked:
|
||||
return blocked
|
||||
blocked = _profile_permission_block(
|
||||
task_capability_map.required_permission(task),
|
||||
remote=remote, host=host, org=org, repo=repo,
|
||||
org_explicit=org is not None,
|
||||
repo_explicit=repo is not None,
|
||||
)
|
||||
if blocked:
|
||||
return blocked
|
||||
|
||||
verify_preflight_purity(remote, task=task, org=org, repo=repo)
|
||||
|
||||
h, o, r = _resolve(remote, host, org, repo)
|
||||
git_remote = (git_remote_name or remote or "").strip()
|
||||
workspace = os.path.realpath(os.path.abspath((worktree_path or "").strip() or "."))
|
||||
|
||||
existing_lock = issue_lock_store.load_issue_lock(
|
||||
remote=remote, org=o, repo=r, issue_number=int(issue_number)
|
||||
)
|
||||
git_state = issue_lock_worktree.read_worktree_git_state(workspace)
|
||||
registered = author_mutation_worktree.path_in_git_worktree_list(
|
||||
workspace, PROJECT_ROOT
|
||||
)
|
||||
|
||||
remote_probe = branch_publish.read_remote_branch_head(
|
||||
workspace, git_remote, branch_name
|
||||
)
|
||||
ancestry = None
|
||||
probe_head = (remote_probe.get("remote_head_sha") or "").strip()
|
||||
if probe_head and probe_head.lower() != (expected_head or "").strip().lower():
|
||||
ancestry = branch_publish.read_is_ancestor(workspace, probe_head, expected_head)
|
||||
|
||||
# A PR on this very branch is this work's own PR, not a rival claim. Only an
|
||||
# open PR for the same issue on a *different* branch is a competing claim.
|
||||
competing: list = []
|
||||
for pull in _list_open_pulls(h, o, r, _auth(h)):
|
||||
ref = str((pull.get("head") or {}).get("ref") or "")
|
||||
if not ref or ref == branch_name:
|
||||
continue
|
||||
if issue_lock_adoption.branch_carries_issue_marker(ref, int(issue_number)):
|
||||
competing.append(pull.get("number"))
|
||||
|
||||
observed_hashes = None
|
||||
if expected_file_hashes:
|
||||
observed_hashes = branch_publish.hash_worktree_files(
|
||||
workspace, list(expected_file_hashes.keys())
|
||||
)
|
||||
|
||||
claimant = _work_lease_claimant(h)
|
||||
assessment = branch_publish.assess_unpublished_commit_publication(
|
||||
existing_lock,
|
||||
issue_number=int(issue_number),
|
||||
branch_name=branch_name,
|
||||
worktree_path=workspace,
|
||||
expected_head=expected_head,
|
||||
remote=remote,
|
||||
org=o,
|
||||
repo=r,
|
||||
identity=claimant.get("username"),
|
||||
profile=claimant.get("profile"),
|
||||
worktree_state=git_state,
|
||||
worktree_registered=registered,
|
||||
remote_probe=remote_probe,
|
||||
ancestry=ancestry,
|
||||
competing_open_prs=competing,
|
||||
expected_file_hashes=expected_file_hashes,
|
||||
observed_file_hashes=observed_hashes,
|
||||
)
|
||||
|
||||
if assessment["outcome"] == branch_publish.REFUSED:
|
||||
return _publication_block(
|
||||
assessment["reasons"], evidence=assessment["evidence"]
|
||||
)
|
||||
|
||||
if dry_run:
|
||||
return {
|
||||
"success": True,
|
||||
"performed": False,
|
||||
"published": False,
|
||||
"verified": False,
|
||||
"dry_run": True,
|
||||
"outcome": assessment["outcome"],
|
||||
"would_publish": assessment["publish_sanctioned"],
|
||||
"remote_head_sha": assessment["evidence"].get("remote_head_sha"),
|
||||
"reasons": [],
|
||||
"evidence": assessment["evidence"],
|
||||
"issue_lock_record_mutated": False,
|
||||
"workflow_lease_touched": False,
|
||||
}
|
||||
|
||||
performed = False
|
||||
if assessment["publish_sanctioned"]:
|
||||
with _audited(
|
||||
task, host=h, remote=remote, org=o, repo=r,
|
||||
target_branch=branch_name,
|
||||
request_metadata={
|
||||
"issue_number": int(issue_number),
|
||||
"expected_head": expected_head,
|
||||
"git_remote": git_remote,
|
||||
},
|
||||
):
|
||||
push = branch_publish.publish_commit_to_remote_branch(
|
||||
worktree_path=workspace,
|
||||
remote_name=git_remote,
|
||||
branch_name=branch_name,
|
||||
expected_head=expected_head,
|
||||
)
|
||||
if not push.get("success"):
|
||||
return _publication_block(
|
||||
push.get("reasons") or ["publication failed"],
|
||||
evidence=assessment["evidence"],
|
||||
stderr=push.get("stderr"),
|
||||
)
|
||||
performed = True
|
||||
|
||||
verification = branch_publish.verify_published_head(
|
||||
worktree_path=workspace,
|
||||
remote_name=git_remote,
|
||||
branch_name=branch_name,
|
||||
expected_head=expected_head,
|
||||
)
|
||||
if not verification.get("verified"):
|
||||
return _publication_block(
|
||||
verification.get("reasons") or ["read-after-write verification failed"],
|
||||
evidence=assessment["evidence"],
|
||||
performed=performed,
|
||||
published=performed,
|
||||
remote_head_sha=verification.get("remote_head_sha"),
|
||||
)
|
||||
|
||||
return {
|
||||
"success": True,
|
||||
"performed": performed,
|
||||
"published": True,
|
||||
"verified": True,
|
||||
"outcome": assessment["outcome"],
|
||||
"remote_head_sha": verification.get("remote_head_sha"),
|
||||
"reasons": [],
|
||||
"evidence": assessment["evidence"],
|
||||
# Publication is the whole of this operation's authority (#812 AC23).
|
||||
"issue_lock_record_mutated": False,
|
||||
"workflow_lease_touched": False,
|
||||
"exact_next_action": (
|
||||
f"Remote head for '{branch_name}' is now observable. Call "
|
||||
f"gitea_lock_issue(issue_number={int(issue_number)}, "
|
||||
f"branch_name='{branch_name}', worktree_path='{workspace}') to renew "
|
||||
"the exact-owner lease, then gitea_create_pr."
|
||||
),
|
||||
}
|
||||
|
||||
|
||||
# Merge methods supported by the Gitea merge API.
|
||||
_MERGE_METHODS = ("merge", "squash", "rebase")
|
||||
|
||||
@@ -12536,10 +12768,39 @@ def _try_auto_switch_for_operation(op: str, host: str | None = None) -> bool:
|
||||
return False
|
||||
|
||||
|
||||
def _git_default_remote_name(root: str) -> str:
|
||||
"""First configured git remote name for *root*, defaulting to 'origin'.
|
||||
|
||||
Used to resolve the live remote master target for parity (#610). Best
|
||||
effort: any failure falls back to 'origin' so callers never raise.
|
||||
"""
|
||||
try:
|
||||
res = subprocess.run(
|
||||
["git", "-C", root, "remote"],
|
||||
capture_output=True, text=True, check=False,
|
||||
)
|
||||
except Exception:
|
||||
return "origin"
|
||||
if res.returncode != 0:
|
||||
return "origin"
|
||||
names = [n.strip() for n in (res.stdout or "").splitlines() if n.strip()]
|
||||
return names[0] if names else "origin"
|
||||
|
||||
|
||||
def _current_master_parity() -> dict:
|
||||
"""Assess this process's code against the on-disk master HEAD (#420)."""
|
||||
"""Assess this process's code against local and live remote master (#420/#610).
|
||||
|
||||
Compares the daemon's startup commit, the on-disk checkout HEAD, and the
|
||||
live remote master target. A stale daemon relative to live master fails
|
||||
closed for mutations even when the local checkout HEAD still matches the
|
||||
startup commit. The live-remote read is best effort: an unresolved live
|
||||
head leaves read-only diagnostics unblocked but is never mutation-safe.
|
||||
"""
|
||||
current_head = master_parity_gate.read_git_head(PROJECT_ROOT)
|
||||
return master_parity_gate.assess_master_parity(_STARTUP_PARITY, current_head)
|
||||
live_head = master_parity_gate.read_remote_master_head(
|
||||
PROJECT_ROOT, remote=_git_default_remote_name(PROJECT_ROOT))
|
||||
return master_parity_gate.assess_master_parity(
|
||||
_STARTUP_PARITY, current_head, live_remote_head=live_head)
|
||||
|
||||
|
||||
def _current_runtime_mode_report(refresh: bool = False) -> dict:
|
||||
@@ -16480,15 +16741,34 @@ def gitea_get_runtime_context(
|
||||
"restart_required": parity["restart_required"],
|
||||
"startup_head": parity["startup_head"],
|
||||
"current_head": parity["current_head"],
|
||||
# #610 distinguished mutation-safety signals:
|
||||
"daemon_start_head": parity["daemon_start_head"],
|
||||
"local_head": parity["local_head"],
|
||||
"live_remote_head": parity["live_remote_head"],
|
||||
"live_known": parity["live_known"],
|
||||
"live_stale": parity["live_stale"],
|
||||
"mutation_safe": parity["mutation_safe"],
|
||||
"summary": master_parity_gate.format_parity(parity),
|
||||
"mutation_gate_enforced": not master_parity_gate.gate_disabled(),
|
||||
# #610: the capability resolver is authoritative for mutation safety;
|
||||
# local parity alone must never authorize a mutation.
|
||||
"resolver_authoritative_for_mutation_safety": True,
|
||||
}
|
||||
if parity["stale"] and not master_parity_gate.gate_disabled():
|
||||
safe_next_action = (
|
||||
"Server code is stale relative to master; restart the Gitea MCP "
|
||||
"server to load current capability gates before mutating. "
|
||||
f"({master_parity_gate.format_parity(parity)})"
|
||||
)
|
||||
if parity["restart_required"] and not master_parity_gate.gate_disabled():
|
||||
if parity["live_stale"]:
|
||||
safe_next_action = (
|
||||
"Daemon is stale relative to LIVE remote master "
|
||||
f"(started {parity['startup_head'][:12] if parity['startup_head'] else 'unknown'}, "
|
||||
f"live master {parity['live_remote_head'][:12] if parity['live_remote_head'] else 'unknown'}); "
|
||||
"restart/reconnect the Gitea MCP server before mutating. The "
|
||||
"capability resolver is authoritative for mutation safety."
|
||||
)
|
||||
else:
|
||||
safe_next_action = (
|
||||
"Server code is stale relative to master; restart the Gitea MCP "
|
||||
"server to load current capability gates before mutating. "
|
||||
f"({master_parity_gate.format_parity(parity)})"
|
||||
)
|
||||
result["safe_next_action"] = safe_next_action
|
||||
|
||||
if reveal and h:
|
||||
@@ -16537,6 +16817,13 @@ def gitea_assess_master_parity(
|
||||
"determinable": parity["determinable"],
|
||||
"startup_head": parity["startup_head"],
|
||||
"current_head": parity["current_head"],
|
||||
# #610 distinguished mutation-safety signals:
|
||||
"daemon_start_head": parity["daemon_start_head"],
|
||||
"local_head": parity["local_head"],
|
||||
"live_remote_head": parity["live_remote_head"],
|
||||
"live_known": parity["live_known"],
|
||||
"live_stale": parity["live_stale"],
|
||||
"mutation_safe": parity["mutation_safe"],
|
||||
"mutation_gate_enforced": enforced,
|
||||
"summary": master_parity_gate.format_parity(parity),
|
||||
"reasons": parity["reasons"],
|
||||
@@ -16553,7 +16840,7 @@ def gitea_assess_master_parity(
|
||||
source=canonical_source,
|
||||
),
|
||||
}
|
||||
if parity["stale"] and enforced:
|
||||
if parity["restart_required"] and enforced:
|
||||
out["report"] = master_parity_gate.parity_report(parity)
|
||||
return out
|
||||
|
||||
|
||||
+213
-17
@@ -24,12 +24,50 @@ from __future__ import annotations
|
||||
|
||||
import os
|
||||
import subprocess
|
||||
import time
|
||||
|
||||
# Live-remote head cache: the parity gate runs on every mutation and every
|
||||
# runtime-context read, so the ``git ls-remote`` result is cached briefly to
|
||||
# avoid a network round-trip per call (#610). Keyed by (root, remote, branch).
|
||||
_REMOTE_HEAD_CACHE: dict[tuple[str, str, str], tuple[float, str | None]] = {}
|
||||
_REMOTE_HEAD_TTL = 60.0
|
||||
|
||||
# When True, ``read_remote_master_head`` never performs ``git ls-remote`` unless
|
||||
# ``GITEA_TEST_LIVE_REMOTE_HEAD`` is set. Conftest enables this suite-wide so
|
||||
# feature worktrees (whose HEAD differs from live master) cannot flip legacy
|
||||
# runtime-context assertions to live_stale, and so unit tests never depend on
|
||||
# a live network (PR #788 F1/F2 / issue #610). Module-level (not env-only) so
|
||||
# ``patch.dict(os.environ, …, clear=True)`` cannot re-enable the probe.
|
||||
_HERMETIC_TEST_MODE: bool = False
|
||||
|
||||
|
||||
def _clear_remote_head_cache() -> None:
|
||||
"""Reset the live-remote head cache (test isolation / forced refresh)."""
|
||||
_REMOTE_HEAD_CACHE.clear()
|
||||
|
||||
|
||||
def set_hermetic_test_mode(enabled: bool) -> None:
|
||||
"""Enable or disable suite-wide hermetic live-remote reads (tests only)."""
|
||||
global _HERMETIC_TEST_MODE
|
||||
_HERMETIC_TEST_MODE = bool(enabled)
|
||||
_clear_remote_head_cache()
|
||||
|
||||
|
||||
def hermetic_test_mode() -> bool:
|
||||
"""Return whether hermetic live-remote reads are active."""
|
||||
return bool(_HERMETIC_TEST_MODE)
|
||||
|
||||
|
||||
# Environment escape hatches (ops + tests):
|
||||
# GITEA_MCP_DISABLE_PARITY_GATE -> disable enforcement entirely (fail open).
|
||||
# GITEA_TEST_CURRENT_HEAD -> force the "current" HEAD read, for tests.
|
||||
ENV_DISABLE = "GITEA_MCP_DISABLE_PARITY_GATE"
|
||||
ENV_TEST_CURRENT_HEAD = "GITEA_TEST_CURRENT_HEAD"
|
||||
# GITEA_TEST_LIVE_REMOTE_HEAD -> force the live remote master read, for tests.
|
||||
ENV_TEST_LIVE_REMOTE_HEAD = "GITEA_TEST_LIVE_REMOTE_HEAD"
|
||||
# GITEA_TEST_ALLOW_LIVE_REMOTE_PROBE -> opt a single test into a real ls-remote
|
||||
# even when hermetic mode is on (rare; prefer ENV_TEST_LIVE_REMOTE_HEAD).
|
||||
ENV_TEST_ALLOW_LIVE_REMOTE_PROBE = "GITEA_TEST_ALLOW_LIVE_REMOTE_PROBE"
|
||||
|
||||
|
||||
def read_git_head(root: str) -> str | None:
|
||||
@@ -58,6 +96,75 @@ def read_git_head(root: str) -> str | None:
|
||||
return (res.stdout or "").strip() or None
|
||||
|
||||
|
||||
def read_remote_master_head(
|
||||
root: str,
|
||||
remote: str = "origin",
|
||||
branch: str = "master",
|
||||
ttl: float = _REMOTE_HEAD_TTL,
|
||||
) -> str | None:
|
||||
"""Return the live remote ``branch`` commit SHA, or ``None`` (#610).
|
||||
|
||||
Resolves the *live* target commit via ``git ls-remote`` so parity can tell
|
||||
a daemon that is behind the live remote master apart from one whose local
|
||||
checkout simply hasn't been pulled. ``None`` means the live head could not
|
||||
be resolved (offline, no such remote, git unavailable, error) -- callers
|
||||
must treat unknown live state as *not mutation-safe* while never blocking
|
||||
read-only diagnostics. A ``GITEA_TEST_LIVE_REMOTE_HEAD`` override takes
|
||||
precedence so the wiring can be exercised deterministically and offline.
|
||||
|
||||
The result is cached for *ttl* seconds per (root, remote, branch) so the
|
||||
gate does not run a network probe on every mutation/read (``ttl=0`` forces
|
||||
a live probe). Both hits and ``None`` misses are cached to bound offline
|
||||
latency; the env override bypasses the cache and the subprocess entirely.
|
||||
|
||||
Under suite hermetic mode (``set_hermetic_test_mode(True)``, set by
|
||||
conftest) a missing override returns ``None`` without network I/O so
|
||||
feature-worktree test runs cannot observe live_stale against real master
|
||||
(PR #788 F1) and unit tests stay offline (F2). Opt out with an explicit
|
||||
``GITEA_TEST_LIVE_REMOTE_HEAD`` pin or ``GITEA_TEST_ALLOW_LIVE_REMOTE_PROBE``.
|
||||
"""
|
||||
forced = os.environ.get(ENV_TEST_LIVE_REMOTE_HEAD)
|
||||
if forced is not None:
|
||||
return forced.strip() or None
|
||||
if _HERMETIC_TEST_MODE and not (
|
||||
os.environ.get(ENV_TEST_ALLOW_LIVE_REMOTE_PROBE) or ""
|
||||
).strip():
|
||||
# Hermetic default: live head unknown. live_stale stays False;
|
||||
# mutation_safe is False when live is unknown (documented #610 note).
|
||||
return None
|
||||
# Defense in depth: even without the module flag, never probe while pytest
|
||||
# is running unless the test opted into a real probe or set an override.
|
||||
if (os.environ.get("PYTEST_CURRENT_TEST") or "").strip() and not (
|
||||
os.environ.get(ENV_TEST_ALLOW_LIVE_REMOTE_PROBE) or ""
|
||||
).strip():
|
||||
return None
|
||||
if not root:
|
||||
return None
|
||||
key = (root, remote, branch)
|
||||
now = time.monotonic()
|
||||
if ttl > 0:
|
||||
cached = _REMOTE_HEAD_CACHE.get(key)
|
||||
if cached is not None and (now - cached[0]) < ttl:
|
||||
return cached[1]
|
||||
sha: str | None = None
|
||||
try:
|
||||
res = subprocess.run(
|
||||
["git", "-C", root, "ls-remote", remote, f"refs/heads/{branch}"],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
check=False,
|
||||
timeout=5,
|
||||
)
|
||||
if res.returncode == 0:
|
||||
lines = (res.stdout or "").strip().splitlines()
|
||||
if lines:
|
||||
sha = lines[0].split("\t", 1)[0].split()[0].strip() or None
|
||||
except Exception:
|
||||
sha = None
|
||||
_REMOTE_HEAD_CACHE[key] = (now, sha)
|
||||
return sha
|
||||
|
||||
|
||||
def capture_startup_parity(root: str, head: str | None = None) -> dict:
|
||||
"""Capture the process source-tree baseline once at server startup.
|
||||
|
||||
@@ -72,18 +179,38 @@ def _short(sha: str | None) -> str:
|
||||
return sha[:12] if sha else "unknown"
|
||||
|
||||
|
||||
def assess_master_parity(startup: dict | None, current_head: str | None) -> dict:
|
||||
def assess_master_parity(
|
||||
startup: dict | None,
|
||||
current_head: str | None,
|
||||
live_remote_head: str | None = None,
|
||||
) -> dict:
|
||||
"""Compare the startup baseline against the current on-disk ``HEAD``.
|
||||
|
||||
Pure: both HEADs are supplied by the caller. Returns a structured result:
|
||||
Pure: all HEADs are supplied by the caller. Returns a structured result:
|
||||
|
||||
- ``in_parity`` -- server code matches the on-disk master (or parity
|
||||
could not be determined, which is not treated as stale).
|
||||
- ``stale`` -- the on-disk master has definitively advanced past the
|
||||
running process.
|
||||
- ``restart_required`` -- alias of ``stale``; the recovery action.
|
||||
- ``determinable`` -- whether both HEADs were known well enough to compare.
|
||||
- ``restart_required`` -- ``stale`` or ``live_stale``; the recovery action.
|
||||
- ``determinable`` -- whether both local HEADs were known well enough to
|
||||
compare.
|
||||
- ``startup_head`` / ``current_head`` / ``reasons``.
|
||||
|
||||
#610 adds live-remote awareness so a daemon that is stale relative to the
|
||||
*live* remote master cannot report a mutation-safe result even when the
|
||||
local checkout HEAD still matches the daemon's startup commit:
|
||||
|
||||
- ``daemon_start_head`` -- the commit the running process started at
|
||||
(alias of ``startup_head``, named for clarity in reports).
|
||||
- ``local_head`` -- the on-disk checkout HEAD (alias of ``current_head``).
|
||||
- ``live_remote_head`` -- the live remote target commit, or ``None`` when it
|
||||
could not be fetched.
|
||||
- ``live_known`` -- whether the live remote target was resolved.
|
||||
- ``live_stale`` -- the live remote master has advanced past the running
|
||||
process (daemon is behind live master) even if local parity is green.
|
||||
- ``mutation_safe`` -- the daemon code, local checkout, and live remote
|
||||
target all agree; the only state in which a mutation may rely on parity.
|
||||
"""
|
||||
startup_head = (startup or {}).get("startup_head")
|
||||
reasons: list[str] = []
|
||||
@@ -91,32 +218,56 @@ def assess_master_parity(startup: dict | None, current_head: str | None) -> dict
|
||||
if startup_head is None:
|
||||
reasons.append(
|
||||
"startup commit was not captured; code parity cannot be enforced")
|
||||
return _result(True, False, False, startup_head, current_head, reasons)
|
||||
return _result(True, False, False, startup_head, current_head,
|
||||
live_remote_head, False, reasons)
|
||||
|
||||
if current_head is None:
|
||||
reasons.append(
|
||||
"current workspace HEAD could not be read; code parity cannot be "
|
||||
"enforced")
|
||||
return _result(True, False, False, startup_head, current_head, reasons)
|
||||
return _result(True, False, False, startup_head, current_head,
|
||||
live_remote_head, False, reasons)
|
||||
|
||||
if startup_head == current_head:
|
||||
return _result(True, False, True, startup_head, current_head, reasons)
|
||||
local_in_parity = startup_head == current_head
|
||||
local_stale = not local_in_parity
|
||||
if local_stale:
|
||||
reasons.append(
|
||||
f"MCP server started at commit {_short(startup_head)} but the "
|
||||
f"workspace master is now {_short(current_head)}; restart the "
|
||||
f"server to load the current capability gates")
|
||||
|
||||
reasons.append(
|
||||
f"MCP server started at commit {_short(startup_head)} but the workspace "
|
||||
f"master is now {_short(current_head)}; restart the server to load the "
|
||||
f"current capability gates")
|
||||
return _result(False, True, True, startup_head, current_head, reasons)
|
||||
live_known = live_remote_head is not None
|
||||
live_stale = live_known and live_remote_head != startup_head
|
||||
if live_stale:
|
||||
reasons.append(
|
||||
f"live remote master is {_short(live_remote_head)} but the MCP "
|
||||
f"server started at {_short(startup_head)}; the daemon is stale "
|
||||
f"relative to live master -- restart/reconnect before mutating")
|
||||
|
||||
return _result(
|
||||
local_in_parity, local_stale, True, startup_head, current_head,
|
||||
live_remote_head, live_stale, reasons)
|
||||
|
||||
|
||||
def _result(in_parity, stale, determinable, startup_head, current_head, reasons):
|
||||
def _result(in_parity, stale, determinable, startup_head, current_head,
|
||||
live_remote_head, live_stale, reasons):
|
||||
live_known = live_remote_head is not None
|
||||
mutation_safe = (
|
||||
determinable and in_parity and live_known and not live_stale)
|
||||
return {
|
||||
"in_parity": in_parity,
|
||||
"stale": stale,
|
||||
"restart_required": stale,
|
||||
"restart_required": stale or live_stale,
|
||||
"determinable": determinable,
|
||||
"startup_head": startup_head,
|
||||
"current_head": current_head,
|
||||
# #610 distinguished signals:
|
||||
"daemon_start_head": startup_head,
|
||||
"local_head": current_head,
|
||||
"live_remote_head": live_remote_head,
|
||||
"live_known": live_known,
|
||||
"live_stale": live_stale,
|
||||
"mutation_safe": mutation_safe,
|
||||
"reasons": list(reasons),
|
||||
}
|
||||
|
||||
@@ -130,11 +281,13 @@ def parity_block_reasons(assessment: dict) -> list[str]:
|
||||
"""Block reasons for a mutation gate (empty when the mutation may proceed).
|
||||
|
||||
A disabled gate or an in-parity / non-determinable assessment yields no
|
||||
reasons; only a definitively stale server blocks.
|
||||
reasons. A definitively stale server blocks, and (#610) a daemon that is
|
||||
stale relative to the *live* remote master blocks even when the local
|
||||
checkout HEAD still matches the daemon's startup commit.
|
||||
"""
|
||||
if gate_disabled():
|
||||
return []
|
||||
if assessment.get("stale"):
|
||||
if assessment.get("stale") or assessment.get("live_stale"):
|
||||
return list(assessment.get("reasons") or
|
||||
["server code is stale relative to master (fail closed)"])
|
||||
return []
|
||||
@@ -147,6 +300,10 @@ def parity_report(assessment: dict) -> dict:
|
||||
"restart_required": True,
|
||||
"startup_head": assessment.get("startup_head"),
|
||||
"current_head": assessment.get("current_head"),
|
||||
# #610: name the live remote target so the report distinguishes a
|
||||
# local-code stale from a daemon-behind-live-master stale.
|
||||
"live_remote_head": assessment.get("live_remote_head"),
|
||||
"live_stale": bool(assessment.get("live_stale")),
|
||||
"reasons": list(assessment.get("reasons") or []),
|
||||
"recovery": [
|
||||
"The running MCP server is executing code older than the current "
|
||||
@@ -157,6 +314,45 @@ def parity_report(assessment: dict) -> dict:
|
||||
}
|
||||
|
||||
|
||||
def parity_resolver_disagreement(
|
||||
assessment: dict,
|
||||
resolver_restart_required: bool,
|
||||
) -> dict | None:
|
||||
"""Typed blocker when the resolver requires restart but parity looks green.
|
||||
|
||||
The capability resolver (``gitea_resolve_task_capability``) detects stale
|
||||
runtime authoritatively for mutation safety (#610). When it requires a
|
||||
restart, local-only parity must never override it: this returns a typed,
|
||||
fail-closed blocker that names the resolver as authoritative. Returns
|
||||
``None`` when the resolver does not require a restart.
|
||||
"""
|
||||
if not resolver_restart_required:
|
||||
return None
|
||||
parity_optimistic = bool(assessment.get("in_parity")) and not (
|
||||
assessment.get("stale") or assessment.get("live_stale"))
|
||||
return {
|
||||
"kind": "parity_resolver_disagreement",
|
||||
"restart_required": True,
|
||||
"resolver_authoritative": True,
|
||||
"parity_optimistic": parity_optimistic,
|
||||
"daemon_start_head": assessment.get("daemon_start_head"),
|
||||
"local_head": assessment.get("local_head"),
|
||||
"live_remote_head": assessment.get("live_remote_head"),
|
||||
"reasons": [
|
||||
"The capability resolver requires a restart/reconnect (stale "
|
||||
"runtime) but master-parity reported local code as in-parity. "
|
||||
"The resolver is authoritative for mutation safety; do not mutate "
|
||||
"on local parity alone. Restart/reconnect the Gitea MCP server "
|
||||
"and re-verify before mutating.",
|
||||
],
|
||||
"recovery": [
|
||||
"Trust the resolver: treat this session as stale.",
|
||||
"Restart or /mcp reconnect the Gitea MCP namespace so it reloads "
|
||||
"current master and live target state, then re-run preflight.",
|
||||
],
|
||||
}
|
||||
|
||||
|
||||
def format_parity(assessment: dict) -> str:
|
||||
"""One-line human summary for logs / runtime context."""
|
||||
if assessment.get("stale"):
|
||||
|
||||
@@ -62,6 +62,14 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
|
||||
"permission": "gitea.branch.push",
|
||||
"role": "author",
|
||||
},
|
||||
# #812 AC20: publish an already-committed, unpublished local head so
|
||||
# exact-owner lease renewal has an observable remote head to reason about.
|
||||
# Same authority as any other author push — deliberately not a new
|
||||
# operation name, so it cannot widen an already-configured author profile.
|
||||
"publish_unpublished_branch": {
|
||||
"permission": "gitea.branch.push",
|
||||
"role": "author",
|
||||
},
|
||||
"create_pr": {
|
||||
"permission": "gitea.pr.create",
|
||||
"role": "author",
|
||||
@@ -500,6 +508,7 @@ ROLE_EXCLUSIVE_TASKS: frozenset[str] = frozenset(
|
||||
"gitea_release_merger_pr_lease",
|
||||
"create_branch",
|
||||
"push_branch",
|
||||
"publish_unpublished_branch",
|
||||
"create_pr",
|
||||
"commit_files",
|
||||
"gitea_commit_files",
|
||||
|
||||
@@ -167,6 +167,35 @@ def _reset_mutation_authority(monkeypatch):
|
||||
import pytest
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def _hermetic_live_remote_master_head():
|
||||
"""#610 / PR #788 F1/F2: keep live-remote parity reads offline in tests.
|
||||
|
||||
``read_remote_master_head`` would otherwise ``git ls-remote`` whenever
|
||||
``GITEA_TEST_LIVE_REMOTE_HEAD`` is unset. Feature worktrees under
|
||||
``branches/`` always differ from live master, so legacy suites that assert
|
||||
runtime-context ``safe_next_action`` flip to live_stale. Module-level
|
||||
hermetic mode survives ``patch.dict(os.environ, …, clear=True)``.
|
||||
Tests that exercise the real probe path call
|
||||
``master_parity_gate.set_hermetic_test_mode(False)`` and/or set
|
||||
``GITEA_TEST_ALLOW_LIVE_REMOTE_PROBE``.
|
||||
"""
|
||||
try:
|
||||
import master_parity_gate as _mpg
|
||||
|
||||
_mpg.set_hermetic_test_mode(True)
|
||||
except Exception:
|
||||
_mpg = None
|
||||
try:
|
||||
yield
|
||||
finally:
|
||||
if _mpg is not None:
|
||||
try:
|
||||
_mpg.set_hermetic_test_mode(False)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def _deterministic_workspace_remotes():
|
||||
try:
|
||||
|
||||
@@ -0,0 +1,650 @@
|
||||
"""Publication of an unpublished local commit (#812 AC20).
|
||||
|
||||
Entry point B of #812: a registered worktree, clean, on its issue branch,
|
||||
holding a local commit that has never been published. Exact-owner lease renewal
|
||||
refuses such a claim for want of an observable remote head, and every existing
|
||||
publication path is lock-derived, so the two predicates close a cycle around
|
||||
work that is otherwise complete.
|
||||
|
||||
These tests exercise the disposition through its *evidence*, never through any
|
||||
particular issue number: every case uses an arbitrary issue number against a
|
||||
synthetic repository, and the same assertions hold for any other. Nothing here
|
||||
reads, writes, or references the live protected worktree named in #812 AC17 —
|
||||
that content is preserved evidence for the duration of this work, so the
|
||||
fixtures below build their own repositories from scratch.
|
||||
|
||||
The remote is a local bare repository, so publication and read-after-write
|
||||
verification are genuinely executed rather than mocked.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
import subprocess
|
||||
import sys
|
||||
import tempfile
|
||||
import unittest
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from unittest.mock import patch
|
||||
|
||||
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
|
||||
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
|
||||
|
||||
import branch_publish # noqa: E402
|
||||
import issue_lock_provenance # noqa: E402
|
||||
import issue_lock_renewal # noqa: E402
|
||||
import issue_lock_store # noqa: E402
|
||||
import mcp_server # noqa: E402
|
||||
from mutation_profile_fixture import shared_mutation_env # noqa: E402
|
||||
|
||||
ISSUE = 9812
|
||||
BRANCH = f"feat/issue-{ISSUE}-publish-fixture"
|
||||
IDENTITY = "example-user"
|
||||
PROFILE = "test-author-prgs"
|
||||
ORG = "Scaled-Tech-Consulting"
|
||||
REPO = "Gitea-Tools"
|
||||
GIT_REMOTE = "prgs"
|
||||
|
||||
|
||||
def _ts(hours: int) -> str:
|
||||
return (
|
||||
(datetime.now(timezone.utc) + timedelta(hours=hours))
|
||||
.isoformat()
|
||||
.replace("+00:00", "Z")
|
||||
)
|
||||
|
||||
|
||||
class _PublishBase(unittest.TestCase):
|
||||
"""Real git repo + real bare remote + durable lock naming the caller.
|
||||
|
||||
The recorded owner pid is deliberately **this live process**. That mirrors
|
||||
the production shape #812 documents, where the pid belongs to a long-running
|
||||
MCP daemon rather than to a dead author client, and it proves publication
|
||||
never depends on a dead process (#812 AC24).
|
||||
"""
|
||||
|
||||
def setUp(self):
|
||||
self.lock_dir = tempfile.TemporaryDirectory()
|
||||
self.addCleanup(self.lock_dir.cleanup)
|
||||
self.origin = tempfile.mkdtemp(prefix="issue812-origin-")
|
||||
self.repo = tempfile.mkdtemp(prefix="issue812-work-")
|
||||
for path in (self.origin, self.repo):
|
||||
self.addCleanup(
|
||||
lambda p=path: subprocess.run(["rm", "-rf", p], check=False)
|
||||
)
|
||||
self._init_repos()
|
||||
self.remotes = patch.dict(
|
||||
mcp_server.REMOTES,
|
||||
{"prgs": {"host": "gitea.prgs.cc", "org": ORG, "repo": REPO}},
|
||||
)
|
||||
self.remotes.start()
|
||||
self.addCleanup(patch.stopall)
|
||||
mcp_server._IDENTITY_CACHE.clear()
|
||||
|
||||
# ── fixture construction ─────────────────────────────────────────────
|
||||
def _git(self, *args, cwd=None):
|
||||
return subprocess.run(
|
||||
["git", "-C", cwd or self.repo, *args],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
check=True,
|
||||
)
|
||||
|
||||
def _init_repos(self):
|
||||
subprocess.run(
|
||||
["git", "init", "-q", "--bare", "-b", "master", self.origin], check=True
|
||||
)
|
||||
self._git("init", "-q", "-b", "master")
|
||||
self._git("config", "user.email", "[email protected]")
|
||||
self._git("config", "user.name", "Test")
|
||||
self._git("remote", "add", GIT_REMOTE, self.origin)
|
||||
|
||||
with open(os.path.join(self.repo, "seed.txt"), "w") as fh:
|
||||
fh.write("seed\n")
|
||||
self._git("add", "seed.txt")
|
||||
self._git("commit", "-q", "-m", "seed")
|
||||
self.base_sha = self._git("rev-parse", "HEAD").stdout.strip()
|
||||
self._git("push", "-q", GIT_REMOTE, "master")
|
||||
|
||||
self._git("checkout", "-q", "-b", BRANCH)
|
||||
with open(os.path.join(self.repo, "work.txt"), "w") as fh:
|
||||
fh.write("unpublished implementation\n")
|
||||
self._git("add", "work.txt")
|
||||
self._git("commit", "-q", "-m", "unpublished implementation")
|
||||
self.head_sha = self._git("rev-parse", "HEAD").stdout.strip()
|
||||
self.worktree = os.path.realpath(self.repo)
|
||||
|
||||
def lock_path(self):
|
||||
return issue_lock_store.lock_file_path(
|
||||
remote="prgs", org=ORG, repo=REPO, issue_number=ISSUE,
|
||||
lock_dir=self.lock_dir.name,
|
||||
)
|
||||
|
||||
def write_lock(self, **overrides):
|
||||
path = self.lock_path()
|
||||
claimant = overrides.pop(
|
||||
"claimant", {"username": IDENTITY, "profile": PROFILE}
|
||||
)
|
||||
pid = overrides.pop("session_pid", os.getpid())
|
||||
lease = {
|
||||
"operation_type": issue_lock_store.AUTHOR_ISSUE_WORK_LEASE,
|
||||
"issue_number": ISSUE,
|
||||
"pr_number": None,
|
||||
"branch": overrides.get("branch_name", BRANCH),
|
||||
"worktree_path": overrides.get("worktree_path", self.worktree),
|
||||
"claimant": claimant,
|
||||
"created_at": _ts(-2),
|
||||
"last_heartbeat_at": _ts(-2),
|
||||
# Expired: entry point B's lease has lapsed, which is precisely why
|
||||
# renewal — and therefore a published head — is needed.
|
||||
"expires_at": _ts(-1),
|
||||
}
|
||||
lease.update(overrides.pop("work_lease", {}))
|
||||
data = {
|
||||
"issue_number": ISSUE,
|
||||
"branch_name": BRANCH,
|
||||
"remote": "prgs",
|
||||
"org": ORG,
|
||||
"repo": REPO,
|
||||
"worktree_path": self.worktree,
|
||||
"session_pid": pid,
|
||||
"pid": pid,
|
||||
"lock_generation": 1,
|
||||
"work_lease": lease,
|
||||
"lock_provenance": issue_lock_provenance.build_sanctioned_lock_provenance(
|
||||
tool="gitea_lock_issue", claimant=claimant
|
||||
),
|
||||
}
|
||||
data.update(overrides)
|
||||
data["lock_file_path"] = path
|
||||
issue_lock_store.save_lock_file(path, data)
|
||||
return path
|
||||
|
||||
def _tool_env(self):
|
||||
env = shared_mutation_env(
|
||||
PROFILE, include_example_repo=True,
|
||||
GITEA_ISSUE_LOCK_DIR=self.lock_dir.name,
|
||||
)
|
||||
env["GITEA_ISSUE_LOCK_DIR"] = self.lock_dir.name
|
||||
# These tests repoint PROJECT_ROOT at a synthetic repository so the
|
||||
# registered-worktree proof runs for real. Pin the parity gate to the
|
||||
# server's own startup head so that repointing does not read as a stale
|
||||
# daemon; the gate itself stays live and enforced.
|
||||
startup_head = mcp_server._STARTUP_PARITY.get("startup_head") or ""
|
||||
env["GITEA_TEST_CURRENT_HEAD"] = startup_head
|
||||
env["GITEA_TEST_LIVE_REMOTE_HEAD"] = startup_head
|
||||
return env
|
||||
|
||||
# ── tool driver ──────────────────────────────────────────────────────
|
||||
def run_publish(self, *, open_prs=None, expected_head=None, **kwargs):
|
||||
"""Drive the public publication tool against the synthetic fixture."""
|
||||
env = self._tool_env()
|
||||
with patch(
|
||||
"mcp_server._list_open_pulls", return_value=list(open_prs or [])
|
||||
), patch(
|
||||
"mcp_server._auth", return_value="token x"
|
||||
), patch(
|
||||
"mcp_server.get_auth_header", return_value="token x"
|
||||
), patch(
|
||||
"mcp_server._work_lease_claimant",
|
||||
return_value={"username": IDENTITY, "profile": PROFILE},
|
||||
), patch.object(
|
||||
mcp_server, "PROJECT_ROOT", self.repo
|
||||
), patch.dict(os.environ, env, clear=True):
|
||||
os.environ["GITEA_ISSUE_LOCK_DIR"] = self.lock_dir.name
|
||||
return mcp_server.gitea_publish_unpublished_issue_branch(
|
||||
issue_number=kwargs.pop("issue_number", ISSUE),
|
||||
branch_name=kwargs.pop("branch_name", BRANCH),
|
||||
worktree_path=kwargs.pop("worktree_path", self.worktree),
|
||||
expected_head=expected_head or self.head_sha,
|
||||
remote="prgs",
|
||||
git_remote_name=kwargs.pop("git_remote_name", GIT_REMOTE),
|
||||
**kwargs,
|
||||
)
|
||||
|
||||
def remote_head(self, branch=BRANCH):
|
||||
res = subprocess.run(
|
||||
["git", "-C", self.origin, "rev-parse", "--verify", "--quiet", branch],
|
||||
capture_output=True, text=True, check=False,
|
||||
)
|
||||
return (res.stdout or "").strip() or None
|
||||
|
||||
|
||||
class TestSuccessfulPublication(_PublishBase):
|
||||
"""AC20 — the branch becomes observable and is verified after the write."""
|
||||
|
||||
def test_publishes_clean_unpublished_commit(self):
|
||||
self.write_lock()
|
||||
self.assertIsNone(self.remote_head(), "fixture must start unpublished")
|
||||
|
||||
result = self.run_publish()
|
||||
|
||||
self.assertTrue(result["success"], result.get("reasons"))
|
||||
self.assertTrue(result["performed"])
|
||||
self.assertTrue(result["published"])
|
||||
self.assertTrue(result["verified"], "read-after-write must be proven")
|
||||
self.assertEqual(result["remote_head_sha"], self.head_sha)
|
||||
self.assertEqual(self.remote_head(), self.head_sha)
|
||||
|
||||
def test_publication_does_not_rewrite_the_commit(self):
|
||||
self.write_lock()
|
||||
self.run_publish()
|
||||
# The published object is the same commit, not a copy or a rewrite.
|
||||
self.assertEqual(self.remote_head(), self.head_sha)
|
||||
self.assertEqual(
|
||||
self._git("rev-parse", "HEAD").stdout.strip(), self.head_sha
|
||||
)
|
||||
|
||||
def test_exact_next_action_names_the_lock_call(self):
|
||||
self.write_lock()
|
||||
result = self.run_publish()
|
||||
self.assertIn("gitea_lock_issue", result["exact_next_action"])
|
||||
|
||||
|
||||
class TestFailsClosed(_PublishBase):
|
||||
"""AC20/AC9 — each refusal reason, exercised independently."""
|
||||
|
||||
def test_changed_local_head_refuses(self):
|
||||
self.write_lock()
|
||||
stale = self.base_sha # a real commit, but not the declared head
|
||||
result = self.run_publish(expected_head=stale)
|
||||
self.assertFalse(result["success"])
|
||||
self.assertTrue(
|
||||
any("local commit changed" in r for r in result["reasons"]),
|
||||
result["reasons"],
|
||||
)
|
||||
self.assertIsNone(self.remote_head(), "refusal must not publish")
|
||||
|
||||
def test_abbreviated_sha_refuses(self):
|
||||
self.write_lock()
|
||||
result = self.run_publish(expected_head=self.head_sha[:8])
|
||||
self.assertFalse(result["success"])
|
||||
self.assertTrue(
|
||||
any("40-character" in r for r in result["reasons"]), result["reasons"]
|
||||
)
|
||||
|
||||
def test_dirty_tracked_worktree_refuses(self):
|
||||
self.write_lock()
|
||||
with open(os.path.join(self.repo, "work.txt"), "a") as fh:
|
||||
fh.write("uncommitted edit\n")
|
||||
|
||||
result = self.run_publish()
|
||||
|
||||
self.assertFalse(result["success"])
|
||||
self.assertTrue(
|
||||
any("dirty tracked files" in r for r in result["reasons"]),
|
||||
result["reasons"],
|
||||
)
|
||||
self.assertIn("work.txt", result["evidence"]["dirty_tracked_files"])
|
||||
self.assertIsNone(self.remote_head())
|
||||
|
||||
def test_untracked_file_refuses(self):
|
||||
self.write_lock()
|
||||
with open(os.path.join(self.repo, "stray.txt"), "w") as fh:
|
||||
fh.write("not committed\n")
|
||||
|
||||
result = self.run_publish()
|
||||
|
||||
self.assertFalse(result["success"])
|
||||
self.assertTrue(
|
||||
any("untracked files" in r for r in result["reasons"]), result["reasons"]
|
||||
)
|
||||
self.assertIn("stray.txt", result["evidence"]["untracked_files"])
|
||||
self.assertIsNone(self.remote_head())
|
||||
|
||||
def test_unexpected_remote_head_refuses(self):
|
||||
"""A remote head that is not an ancestor must never be overwritten."""
|
||||
self.write_lock()
|
||||
# Publish a divergent commit to the branch from a separate line.
|
||||
self._git("checkout", "-q", "-b", "divergent", self.base_sha)
|
||||
with open(os.path.join(self.repo, "other.txt"), "w") as fh:
|
||||
fh.write("someone else's work\n")
|
||||
self._git("add", "other.txt")
|
||||
self._git("commit", "-q", "-m", "divergent")
|
||||
divergent = self._git("rev-parse", "HEAD").stdout.strip()
|
||||
self._git("push", "-q", GIT_REMOTE, f"{divergent}:refs/heads/{BRANCH}")
|
||||
self._git("checkout", "-q", BRANCH)
|
||||
|
||||
result = self.run_publish()
|
||||
|
||||
self.assertFalse(result["success"])
|
||||
self.assertTrue(
|
||||
any("not an ancestor" in r for r in result["reasons"]), result["reasons"]
|
||||
)
|
||||
self.assertEqual(
|
||||
self.remote_head(), divergent, "the other head must survive intact"
|
||||
)
|
||||
|
||||
def test_fast_forward_remote_head_is_allowed(self):
|
||||
"""An ancestor head is an honest fast-forward, not a conflict."""
|
||||
self.write_lock()
|
||||
self._git("push", "-q", GIT_REMOTE, f"{self.base_sha}:refs/heads/{BRANCH}")
|
||||
|
||||
result = self.run_publish()
|
||||
|
||||
self.assertTrue(result["success"], result.get("reasons"))
|
||||
self.assertTrue(result["evidence"]["fast_forward_from_remote"])
|
||||
self.assertEqual(self.remote_head(), self.head_sha)
|
||||
|
||||
def test_content_hash_mismatch_refuses(self):
|
||||
self.write_lock()
|
||||
wrong = {"work.txt": "0" * 64}
|
||||
|
||||
result = self.run_publish(expected_file_hashes=wrong)
|
||||
|
||||
self.assertFalse(result["success"])
|
||||
self.assertTrue(
|
||||
any("declared content hashes" in r for r in result["reasons"]),
|
||||
result["reasons"],
|
||||
)
|
||||
self.assertFalse(result["evidence"]["file_hashes_verified"])
|
||||
self.assertIsNone(self.remote_head())
|
||||
|
||||
def test_matching_content_hashes_publish(self):
|
||||
self.write_lock()
|
||||
digests = branch_publish.hash_worktree_files(self.worktree, ["work.txt"])
|
||||
|
||||
result = self.run_publish(expected_file_hashes=digests)
|
||||
|
||||
self.assertTrue(result["success"], result.get("reasons"))
|
||||
self.assertTrue(result["evidence"]["file_hashes_verified"])
|
||||
|
||||
def test_missing_declared_file_refuses(self):
|
||||
self.write_lock()
|
||||
result = self.run_publish(expected_file_hashes={"absent.txt": "0" * 64})
|
||||
self.assertFalse(result["success"])
|
||||
self.assertTrue(
|
||||
any("missing or unreadable" in r for r in result["reasons"]),
|
||||
result["reasons"],
|
||||
)
|
||||
|
||||
def test_foreign_claimant_refuses(self):
|
||||
"""Ownership comes from the durable record, not from the caller."""
|
||||
self.write_lock(claimant={"username": "someone-else", "profile": PROFILE})
|
||||
|
||||
result = self.run_publish()
|
||||
|
||||
self.assertFalse(result["success"])
|
||||
self.assertTrue(
|
||||
any("foreign claim" in r for r in result["reasons"]), result["reasons"]
|
||||
)
|
||||
self.assertIsNone(self.remote_head())
|
||||
|
||||
def test_foreign_profile_refuses(self):
|
||||
self.write_lock(
|
||||
claimant={"username": IDENTITY, "profile": "test-reviewer-prgs"}
|
||||
)
|
||||
result = self.run_publish()
|
||||
self.assertFalse(result["success"])
|
||||
self.assertTrue(
|
||||
any("claimant profile" in r for r in result["reasons"]), result["reasons"]
|
||||
)
|
||||
|
||||
def test_absent_lock_record_refuses(self):
|
||||
"""No recorded claim means this cannot be used to bypass the lock."""
|
||||
result = self.run_publish() # no write_lock()
|
||||
|
||||
self.assertFalse(result["success"])
|
||||
self.assertTrue(
|
||||
any("no durable issue-lock record" in r for r in result["reasons"]),
|
||||
result["reasons"],
|
||||
)
|
||||
self.assertIsNone(self.remote_head())
|
||||
|
||||
def test_branch_mismatch_against_lock_refuses(self):
|
||||
self.write_lock(branch_name=f"feat/issue-{ISSUE}-different")
|
||||
result = self.run_publish()
|
||||
self.assertFalse(result["success"])
|
||||
self.assertTrue(
|
||||
any("records branch" in r for r in result["reasons"]), result["reasons"]
|
||||
)
|
||||
|
||||
def test_worktree_mismatch_against_lock_refuses(self):
|
||||
self.write_lock(worktree_path="/tmp/some/other/worktree")
|
||||
result = self.run_publish()
|
||||
self.assertFalse(result["success"])
|
||||
self.assertTrue(
|
||||
any("records worktree" in r for r in result["reasons"]), result["reasons"]
|
||||
)
|
||||
|
||||
def test_competing_open_pr_on_another_branch_refuses(self):
|
||||
self.write_lock()
|
||||
competing = [{"number": 4242, "head": {"ref": f"fix/issue-{ISSUE}-rival"}}]
|
||||
|
||||
result = self.run_publish(open_prs=competing)
|
||||
|
||||
self.assertFalse(result["success"])
|
||||
self.assertTrue(
|
||||
any("already claim issue" in r for r in result["reasons"]),
|
||||
result["reasons"],
|
||||
)
|
||||
self.assertIsNone(self.remote_head())
|
||||
|
||||
def test_open_pr_on_the_same_branch_is_not_competing(self):
|
||||
"""This branch's own PR is not a rival claim against itself."""
|
||||
self.write_lock()
|
||||
own = [{"number": 77, "head": {"ref": BRANCH}}]
|
||||
|
||||
result = self.run_publish(open_prs=own)
|
||||
|
||||
self.assertTrue(result["success"], result.get("reasons"))
|
||||
|
||||
|
||||
class TestGuardStrictnessPreserved(_PublishBase):
|
||||
"""AC15 — publication is an operation, never a weakening of the guards."""
|
||||
|
||||
def test_non_issue_branch_refuses(self):
|
||||
self._git("checkout", "-q", "-b", "scratch/not-issue-linked")
|
||||
self.write_lock(branch_name="scratch/not-issue-linked")
|
||||
|
||||
result = self.run_publish(branch_name="scratch/not-issue-linked")
|
||||
|
||||
self.assertFalse(result["success"])
|
||||
self.assertTrue(
|
||||
any("issue-linked" in r for r in result["reasons"]), result["reasons"]
|
||||
)
|
||||
|
||||
def test_stable_branch_refuses(self):
|
||||
self.write_lock(branch_name="master")
|
||||
result = self.run_publish(branch_name="master")
|
||||
self.assertFalse(result["success"])
|
||||
self.assertTrue(
|
||||
any("issue-linked" in r or "stable branch" in r for r in result["reasons"]),
|
||||
result["reasons"],
|
||||
)
|
||||
|
||||
def test_branch_number_must_match_the_issue(self):
|
||||
other = "feat/issue-7777-mismatched"
|
||||
self._git("checkout", "-q", "-b", other)
|
||||
self.write_lock(branch_name=other)
|
||||
result = self.run_publish(branch_name=other)
|
||||
self.assertFalse(result["success"])
|
||||
self.assertTrue(
|
||||
any("does not carry issue number" in r for r in result["reasons"]),
|
||||
result["reasons"],
|
||||
)
|
||||
|
||||
def test_unregistered_worktree_refuses(self):
|
||||
"""#713 — an improvised directory is not a registered worktree."""
|
||||
path = self.write_lock()
|
||||
assessment = branch_publish.assess_unpublished_commit_publication(
|
||||
issue_lock_store.read_lock_file(path),
|
||||
issue_number=ISSUE, branch_name=BRANCH, worktree_path=self.worktree,
|
||||
expected_head=self.head_sha, remote="prgs", org=ORG, repo=REPO,
|
||||
identity=IDENTITY, profile=PROFILE,
|
||||
worktree_state={
|
||||
"current_branch": BRANCH, "porcelain_status": "",
|
||||
"head_sha": self.head_sha,
|
||||
},
|
||||
worktree_registered=False,
|
||||
remote_probe={"probe_ok": True, "remote_branch_exists": False},
|
||||
)
|
||||
self.assertEqual(assessment["outcome"], branch_publish.REFUSED)
|
||||
self.assertTrue(
|
||||
any("not listed in git worktree list" in r
|
||||
for r in assessment["reasons"]),
|
||||
assessment["reasons"],
|
||||
)
|
||||
|
||||
def test_unobservable_remote_refuses(self):
|
||||
"""An unknown remote state must not be mistaken for an absent branch."""
|
||||
self.write_lock()
|
||||
result = self.run_publish(git_remote_name="no-such-remote")
|
||||
self.assertFalse(result["success"])
|
||||
self.assertTrue(
|
||||
any("could not be observed" in r for r in result["reasons"]),
|
||||
result["reasons"],
|
||||
)
|
||||
|
||||
|
||||
class TestRecordSeparation(_PublishBase):
|
||||
"""AC23 — the durable issue lock and the workflow lease are distinct."""
|
||||
|
||||
def test_publication_leaves_the_issue_lock_byte_identical(self):
|
||||
path = self.write_lock()
|
||||
with open(path, "rb") as fh:
|
||||
before = fh.read()
|
||||
|
||||
result = self.run_publish()
|
||||
|
||||
self.assertTrue(result["success"], result.get("reasons"))
|
||||
with open(path, "rb") as fh:
|
||||
after = fh.read()
|
||||
self.assertEqual(before, after, "publication must not mutate the lock record")
|
||||
self.assertFalse(result["issue_lock_record_mutated"])
|
||||
self.assertFalse(result["workflow_lease_touched"])
|
||||
|
||||
def test_refusal_also_reports_untouched_records(self):
|
||||
result = self.run_publish() # refuses: no lock record
|
||||
self.assertFalse(result["issue_lock_record_mutated"])
|
||||
self.assertFalse(result["workflow_lease_touched"])
|
||||
|
||||
def test_lock_generation_is_not_advanced(self):
|
||||
path = self.write_lock()
|
||||
self.run_publish()
|
||||
lock = issue_lock_store.read_lock_file(path)
|
||||
self.assertEqual(lock["lock_generation"], 1)
|
||||
|
||||
|
||||
class TestTruthfulProcessEvidence(_PublishBase):
|
||||
"""AC24 — a live daemon pid is never represented as a dead process."""
|
||||
|
||||
def test_live_recorded_pid_does_not_block_publication(self):
|
||||
# The recorded pid is this live process, standing in for the live MCP
|
||||
# daemon. Reclaim would refuse here; publication legitimately does not.
|
||||
path = self.write_lock(session_pid=os.getpid())
|
||||
lock = issue_lock_store.read_lock_file(path)
|
||||
self.assertEqual(lock["pid"], os.getpid())
|
||||
|
||||
result = self.run_publish()
|
||||
|
||||
self.assertTrue(result["success"], result.get("reasons"))
|
||||
self.assertEqual(self.remote_head(), self.head_sha)
|
||||
|
||||
def test_liveness_is_not_consulted_as_evidence(self):
|
||||
self.write_lock(session_pid=os.getpid())
|
||||
result = self.run_publish()
|
||||
self.assertFalse(result["evidence"]["owner_pid_liveness_consulted"])
|
||||
|
||||
def test_reclaim_still_refuses_for_the_same_live_pid(self):
|
||||
"""Publication does not soften the reclaim predicate it routes around."""
|
||||
path = self.write_lock(session_pid=os.getpid())
|
||||
lock = issue_lock_store.read_lock_file(path)
|
||||
reclaim = issue_lock_store.assess_expired_lock_reclaim(lock)
|
||||
self.assertFalse(reclaim["reclaim_allowed"])
|
||||
|
||||
|
||||
class TestIdempotentRetry(_PublishBase):
|
||||
"""AC20 — retry is safe and read-after-write is proven every time."""
|
||||
|
||||
def test_second_publication_reports_already_published(self):
|
||||
self.write_lock()
|
||||
first = self.run_publish()
|
||||
self.assertTrue(first["performed"])
|
||||
|
||||
second = self.run_publish()
|
||||
|
||||
self.assertTrue(second["success"], second.get("reasons"))
|
||||
self.assertFalse(second["performed"], "no second push is needed")
|
||||
self.assertTrue(second["published"])
|
||||
self.assertTrue(second["verified"])
|
||||
self.assertEqual(second["outcome"], branch_publish.ALREADY_PUBLISHED)
|
||||
self.assertEqual(self.remote_head(), self.head_sha)
|
||||
|
||||
|
||||
class TestDryRun(_PublishBase):
|
||||
"""AC12 — dry run reports the decision and mutates nothing."""
|
||||
|
||||
def test_dry_run_reports_intent_without_publishing(self):
|
||||
self.write_lock()
|
||||
|
||||
result = self.run_publish(dry_run=True)
|
||||
|
||||
self.assertTrue(result["success"])
|
||||
self.assertTrue(result["dry_run"])
|
||||
self.assertTrue(result["would_publish"])
|
||||
self.assertFalse(result["performed"])
|
||||
self.assertIsNone(self.remote_head(), "dry run must not publish")
|
||||
|
||||
def test_dry_run_and_apply_agree_on_a_refusal(self):
|
||||
"""AC11 — the reported decision does not depend on which mode ran."""
|
||||
self.write_lock(claimant={"username": "someone-else", "profile": PROFILE})
|
||||
|
||||
dry = self.run_publish(dry_run=True)
|
||||
applied = self.run_publish()
|
||||
|
||||
self.assertFalse(dry["success"])
|
||||
self.assertFalse(applied["success"])
|
||||
self.assertEqual(dry["reasons"], applied["reasons"])
|
||||
|
||||
|
||||
class TestRenewalUnblocked(_PublishBase):
|
||||
"""AC20/AC21 — renewal is permitted only after verified publication."""
|
||||
|
||||
def _renewal(self, remote_head):
|
||||
return issue_lock_renewal.assess_exact_owner_lease_renewal(
|
||||
issue_lock_store.read_lock_file(self.lock_path()),
|
||||
issue_number=ISSUE, branch_name=BRANCH, worktree_path=self.worktree,
|
||||
remote="prgs", org=ORG, repo=REPO,
|
||||
identity=IDENTITY, profile=PROFILE,
|
||||
current_branch=BRANCH, porcelain_status="", worktree_exists=True,
|
||||
head_sha=self.head_sha, remote_head_sha=remote_head,
|
||||
)
|
||||
|
||||
def test_renewal_refuses_before_publication(self):
|
||||
self.write_lock()
|
||||
decision = self._renewal(None)
|
||||
self.assertFalse(decision["renewal_sanctioned"])
|
||||
self.assertTrue(
|
||||
any("unpublished branch" in r for r in decision["reasons"]),
|
||||
decision["reasons"],
|
||||
)
|
||||
|
||||
def test_renewal_is_sanctioned_after_publication(self):
|
||||
self.write_lock()
|
||||
result = self.run_publish()
|
||||
self.assertTrue(result["verified"], result.get("reasons"))
|
||||
|
||||
decision = self._renewal(self.remote_head())
|
||||
|
||||
self.assertTrue(decision["renewal_sanctioned"], decision["reasons"])
|
||||
|
||||
|
||||
class TestProtectedAssetUntouched(unittest.TestCase):
|
||||
"""AC17 — no test or fixture may reference the protected worktree."""
|
||||
|
||||
def test_no_reference_to_the_protected_worktree(self):
|
||||
here = os.path.dirname(os.path.abspath(__file__))
|
||||
root = os.path.dirname(here)
|
||||
needle = "issue-635-project-registry" + "-api"
|
||||
for path in (
|
||||
os.path.join(here, "test_issue_812_publish_unpublished_commit.py"),
|
||||
os.path.join(root, "branch_publish.py"),
|
||||
):
|
||||
with open(path, "r", encoding="utf-8") as fh:
|
||||
body = fh.read()
|
||||
self.assertNotIn(needle, body)
|
||||
|
||||
|
||||
if __name__ == "__main__": # pragma: no cover
|
||||
unittest.main()
|
||||
@@ -78,6 +78,95 @@ class TestBlockReasonsAndReport(unittest.TestCase):
|
||||
self.assertTrue(report["recovery"])
|
||||
|
||||
|
||||
class TestLiveRemoteParity(unittest.TestCase):
|
||||
"""#610: parity must account for the live remote master, not just local.
|
||||
|
||||
The daemon can be stale relative to the live remote target while the local
|
||||
checkout HEAD still matches the daemon's startup commit, so local parity
|
||||
reports green even though a mutation would run against outdated code.
|
||||
"""
|
||||
|
||||
SHA_C = "c" * 40
|
||||
|
||||
def test_distinguishes_three_shas(self):
|
||||
res = mp.assess_master_parity(
|
||||
{"startup_head": SHA_A}, SHA_A, live_remote_head=SHA_B)
|
||||
self.assertEqual(res["daemon_start_head"], SHA_A)
|
||||
self.assertEqual(res["local_head"], SHA_A)
|
||||
self.assertEqual(res["live_remote_head"], SHA_B)
|
||||
|
||||
def test_mutation_safe_only_when_all_three_match(self):
|
||||
res = mp.assess_master_parity(
|
||||
{"startup_head": SHA_A}, SHA_A, live_remote_head=SHA_A)
|
||||
self.assertTrue(res["mutation_safe"])
|
||||
self.assertTrue(res["live_known"])
|
||||
self.assertFalse(res["live_stale"])
|
||||
|
||||
def test_live_stale_when_remote_advanced_past_daemon(self):
|
||||
# Local checkout still matches the daemon start (local parity green),
|
||||
# but the live remote master has advanced -> daemon is live-stale.
|
||||
res = mp.assess_master_parity(
|
||||
{"startup_head": SHA_A}, SHA_A, live_remote_head=SHA_B)
|
||||
self.assertTrue(res["in_parity"]) # local parity still green
|
||||
self.assertTrue(res["live_stale"])
|
||||
self.assertFalse(res["mutation_safe"])
|
||||
self.assertTrue(any("live" in r.lower() for r in res["reasons"]))
|
||||
|
||||
def test_live_unknown_is_not_mutation_safe_but_not_stale(self):
|
||||
# Non-goal: unfetchable live remote must not be treated as stale for
|
||||
# read-only, but a mutation-safe claim fails closed.
|
||||
res = mp.assess_master_parity(
|
||||
{"startup_head": SHA_A}, SHA_A, live_remote_head=None)
|
||||
self.assertFalse(res["live_known"])
|
||||
self.assertFalse(res["mutation_safe"])
|
||||
self.assertFalse(res["live_stale"])
|
||||
self.assertTrue(res["in_parity"])
|
||||
|
||||
def test_default_live_remote_preserves_legacy_shape(self):
|
||||
# Callers that do not supply a live head keep the pre-#610 behavior:
|
||||
# in-parity, not live-stale, no live-derived block.
|
||||
res = mp.assess_master_parity({"startup_head": SHA_A}, SHA_A)
|
||||
self.assertFalse(res["live_stale"])
|
||||
self.assertEqual(mp.parity_block_reasons(res), [])
|
||||
|
||||
|
||||
class TestLiveStaleBlockAndReport(unittest.TestCase):
|
||||
"""#610: live-staleness must block mutations and surface a typed blocker."""
|
||||
|
||||
def test_live_stale_produces_block_reasons(self):
|
||||
res = mp.assess_master_parity(
|
||||
{"startup_head": SHA_A}, SHA_A, live_remote_head=SHA_B)
|
||||
self.assertTrue(mp.parity_block_reasons(res))
|
||||
|
||||
def test_disable_env_suppresses_live_stale_block(self):
|
||||
res = mp.assess_master_parity(
|
||||
{"startup_head": SHA_A}, SHA_A, live_remote_head=SHA_B)
|
||||
with patch.dict(os.environ, {mp.ENV_DISABLE: "1"}):
|
||||
self.assertEqual(mp.parity_block_reasons(res), [])
|
||||
|
||||
def test_resolver_disagreement_returns_typed_blocker(self):
|
||||
# Parity says local-green, resolver says restart required -> disagreement
|
||||
# is a typed, fail-closed blocker naming the resolver as authoritative.
|
||||
res = mp.assess_master_parity({"startup_head": SHA_A}, SHA_A)
|
||||
blocker = mp.parity_resolver_disagreement(res, resolver_restart_required=True)
|
||||
self.assertIsNotNone(blocker)
|
||||
self.assertEqual(blocker["kind"], "parity_resolver_disagreement")
|
||||
self.assertTrue(blocker["restart_required"])
|
||||
self.assertTrue(blocker["resolver_authoritative"])
|
||||
|
||||
def test_no_disagreement_when_resolver_agrees(self):
|
||||
res = mp.assess_master_parity({"startup_head": SHA_A}, SHA_A)
|
||||
self.assertIsNone(
|
||||
mp.parity_resolver_disagreement(res, resolver_restart_required=False))
|
||||
|
||||
def test_live_stale_report_names_live_remote(self):
|
||||
res = mp.assess_master_parity(
|
||||
{"startup_head": SHA_A}, SHA_A, live_remote_head=SHA_B)
|
||||
report = mp.parity_report(res)
|
||||
self.assertEqual(report["live_remote_head"], SHA_B)
|
||||
self.assertTrue(report["restart_required"])
|
||||
|
||||
|
||||
class TestReadGitHead(unittest.TestCase):
|
||||
def test_test_override_takes_precedence(self):
|
||||
with patch.dict(os.environ, {mp.ENV_TEST_CURRENT_HEAD: SHA_B}):
|
||||
@@ -95,6 +184,149 @@ class TestReadGitHead(unittest.TestCase):
|
||||
self.assertIsNone(mp.read_git_head(""))
|
||||
|
||||
|
||||
class TestReadRemoteMasterHead(unittest.TestCase):
|
||||
"""#610: live remote master head reader (env-overridable, fails to None)."""
|
||||
|
||||
def test_test_override_takes_precedence(self):
|
||||
with patch.dict(os.environ, {mp.ENV_TEST_LIVE_REMOTE_HEAD: SHA_B}):
|
||||
self.assertEqual(mp.read_remote_master_head("/nonexistent"), SHA_B)
|
||||
|
||||
def test_blank_override_is_none(self):
|
||||
with patch.dict(os.environ, {mp.ENV_TEST_LIVE_REMOTE_HEAD: " "}):
|
||||
self.assertIsNone(mp.read_remote_master_head("/nonexistent"))
|
||||
|
||||
def test_unfetchable_remote_is_none(self):
|
||||
# No override; a bogus root/remote must fail closed to None, never raise.
|
||||
env = {k: v for k, v in os.environ.items()
|
||||
if k != mp.ENV_TEST_LIVE_REMOTE_HEAD}
|
||||
with patch.dict(os.environ, env, clear=True):
|
||||
self.assertIsNone(
|
||||
mp.read_remote_master_head("/nonexistent", remote="nope"))
|
||||
|
||||
|
||||
class TestRemoteHeadCache(unittest.TestCase):
|
||||
"""#610: live remote reads are cached with a TTL to stay off the network.
|
||||
|
||||
The parity gate runs on every mutation and every runtime-context read, so an
|
||||
unbounded ``git ls-remote`` per call would be a latency/flakiness regression.
|
||||
"""
|
||||
|
||||
def setUp(self):
|
||||
# These cases intentionally exercise the subprocess/cache path, so they
|
||||
# opt out of suite-wide hermetic mode (PR #788 F1).
|
||||
self._saved_hermetic = mp.hermetic_test_mode()
|
||||
mp.set_hermetic_test_mode(False)
|
||||
mp._clear_remote_head_cache()
|
||||
env = {
|
||||
k: v for k, v in os.environ.items()
|
||||
if k not in (mp.ENV_TEST_LIVE_REMOTE_HEAD,
|
||||
mp.ENV_TEST_ALLOW_LIVE_REMOTE_PROBE,
|
||||
"PYTEST_CURRENT_TEST")
|
||||
}
|
||||
# Allow the probe path under hermetic defenses while still mocking
|
||||
# subprocess so no real network call runs.
|
||||
env[mp.ENV_TEST_ALLOW_LIVE_REMOTE_PROBE] = "1"
|
||||
self._env = patch.dict(os.environ, env, clear=True)
|
||||
self._env.start()
|
||||
self.addCleanup(self._env.stop)
|
||||
self.addCleanup(mp._clear_remote_head_cache)
|
||||
self.addCleanup(
|
||||
lambda: mp.set_hermetic_test_mode(self._saved_hermetic)
|
||||
)
|
||||
|
||||
def _fake_run(self, sha):
|
||||
class _R:
|
||||
returncode = 0
|
||||
stdout = f"{sha}\trefs/heads/master\n"
|
||||
calls = {"n": 0}
|
||||
|
||||
def run(*args, **kwargs):
|
||||
calls["n"] += 1
|
||||
return _R()
|
||||
return run, calls
|
||||
|
||||
def test_second_call_within_ttl_uses_cache(self):
|
||||
run, calls = self._fake_run(SHA_B)
|
||||
with patch.object(mp.subprocess, "run", run):
|
||||
a = mp.read_remote_master_head("/repo", remote="prgs", ttl=100)
|
||||
b = mp.read_remote_master_head("/repo", remote="prgs", ttl=100)
|
||||
self.assertEqual(a, SHA_B)
|
||||
self.assertEqual(b, SHA_B)
|
||||
self.assertEqual(calls["n"], 1)
|
||||
|
||||
def test_zero_ttl_bypasses_cache(self):
|
||||
run, calls = self._fake_run(SHA_B)
|
||||
with patch.object(mp.subprocess, "run", run):
|
||||
mp.read_remote_master_head("/repo", remote="prgs", ttl=0)
|
||||
mp.read_remote_master_head("/repo", remote="prgs", ttl=0)
|
||||
self.assertEqual(calls["n"], 2)
|
||||
|
||||
def test_env_override_never_touches_subprocess(self):
|
||||
run, calls = self._fake_run(SHA_B)
|
||||
with patch.dict(os.environ, {mp.ENV_TEST_LIVE_REMOTE_HEAD: SHA_A}):
|
||||
with patch.object(mp.subprocess, "run", run):
|
||||
self.assertEqual(
|
||||
mp.read_remote_master_head("/repo", remote="prgs"), SHA_A)
|
||||
self.assertEqual(calls["n"], 0)
|
||||
|
||||
|
||||
class TestHermeticLiveRemoteReads(unittest.TestCase):
|
||||
"""#610 / PR #788 F1/F2: suite hermetic mode never hits the network."""
|
||||
|
||||
def setUp(self):
|
||||
self._saved = mp.hermetic_test_mode()
|
||||
mp.set_hermetic_test_mode(True)
|
||||
mp._clear_remote_head_cache()
|
||||
self.addCleanup(lambda: mp.set_hermetic_test_mode(self._saved))
|
||||
self.addCleanup(mp._clear_remote_head_cache)
|
||||
|
||||
def test_hermetic_mode_returns_none_without_subprocess(self):
|
||||
run_calls = {"n": 0}
|
||||
|
||||
def boom(*args, **kwargs):
|
||||
run_calls["n"] += 1
|
||||
raise AssertionError("ls-remote must not run under hermetic mode")
|
||||
|
||||
env = {
|
||||
k: v for k, v in os.environ.items()
|
||||
if k not in (mp.ENV_TEST_LIVE_REMOTE_HEAD,
|
||||
mp.ENV_TEST_ALLOW_LIVE_REMOTE_PROBE)
|
||||
}
|
||||
with patch.dict(os.environ, env, clear=True):
|
||||
with patch.object(mp.subprocess, "run", boom):
|
||||
self.assertIsNone(
|
||||
mp.read_remote_master_head("/repo", remote="prgs")
|
||||
)
|
||||
self.assertEqual(run_calls["n"], 0)
|
||||
|
||||
def test_hermetic_mode_survives_clear_true_env(self):
|
||||
"""Module flag, not env pin: clear=True cannot re-enable the probe."""
|
||||
run_calls = {"n": 0}
|
||||
|
||||
def boom(*args, **kwargs):
|
||||
run_calls["n"] += 1
|
||||
raise AssertionError("ls-remote must not run after clear=True")
|
||||
|
||||
with patch.dict(os.environ, {}, clear=True):
|
||||
with patch.object(mp.subprocess, "run", boom):
|
||||
self.assertIsNone(mp.read_remote_master_head("/repo"))
|
||||
self.assertEqual(run_calls["n"], 0)
|
||||
|
||||
def test_explicit_override_still_wins_under_hermetic(self):
|
||||
run_calls = {"n": 0}
|
||||
|
||||
def boom(*args, **kwargs):
|
||||
run_calls["n"] += 1
|
||||
raise AssertionError("override must bypass subprocess")
|
||||
|
||||
with patch.dict(os.environ, {mp.ENV_TEST_LIVE_REMOTE_HEAD: SHA_B}):
|
||||
with patch.object(mp.subprocess, "run", boom):
|
||||
self.assertEqual(
|
||||
mp.read_remote_master_head("/repo"), SHA_B
|
||||
)
|
||||
self.assertEqual(run_calls["n"], 0)
|
||||
|
||||
|
||||
class TestServerWiring(unittest.TestCase):
|
||||
"""Integration with the gate choke point in the server namespace."""
|
||||
|
||||
@@ -105,6 +337,13 @@ class TestServerWiring(unittest.TestCase):
|
||||
self._saved = self.srv._STARTUP_PARITY
|
||||
self.srv._STARTUP_PARITY = {"root": self.srv.PROJECT_ROOT,
|
||||
"startup_head": SHA_A}
|
||||
# Keep the live-remote read hermetic (no real ls-remote network call):
|
||||
# default the live master to the daemon start so parity is fully green
|
||||
# unless a test overrides the live head explicitly (#610).
|
||||
self._live_patch = patch.dict(
|
||||
os.environ, {mp.ENV_TEST_LIVE_REMOTE_HEAD: SHA_A})
|
||||
self._live_patch.start()
|
||||
self.addCleanup(self._live_patch.stop)
|
||||
|
||||
def tearDown(self):
|
||||
self.srv._STARTUP_PARITY = self._saved
|
||||
@@ -147,6 +386,36 @@ class TestServerWiring(unittest.TestCase):
|
||||
self.assertTrue(out["in_parity"])
|
||||
self.assertNotIn("report", out)
|
||||
|
||||
# --- #610: live-remote wiring -------------------------------------------
|
||||
|
||||
def test_live_stale_blocks_mutation_though_local_green(self):
|
||||
# Local checkout matches the daemon start (local parity green) but the
|
||||
# live remote master has advanced -> mutations must fail closed.
|
||||
with patch.dict(os.environ, {mp.ENV_TEST_CURRENT_HEAD: SHA_A,
|
||||
mp.ENV_TEST_LIVE_REMOTE_HEAD: SHA_B}):
|
||||
self.assertEqual(self.srv._master_parity_block("gitea.read"), [])
|
||||
self.assertTrue(
|
||||
self.srv._master_parity_block("gitea.pr.create"))
|
||||
|
||||
def test_assess_tool_exposes_three_distinct_shas(self):
|
||||
with patch.dict(os.environ, {mp.ENV_TEST_CURRENT_HEAD: SHA_A,
|
||||
mp.ENV_TEST_LIVE_REMOTE_HEAD: SHA_B}):
|
||||
out = self.srv.gitea_assess_master_parity(remote="prgs")
|
||||
self.assertEqual(out["daemon_start_head"], SHA_A)
|
||||
self.assertEqual(out["local_head"], SHA_A)
|
||||
self.assertEqual(out["live_remote_head"], SHA_B)
|
||||
self.assertTrue(out["live_stale"])
|
||||
self.assertFalse(out["mutation_safe"])
|
||||
self.assertIn("report", out)
|
||||
|
||||
def test_assess_tool_mutation_safe_when_all_three_match(self):
|
||||
with patch.dict(os.environ, {mp.ENV_TEST_CURRENT_HEAD: SHA_A,
|
||||
mp.ENV_TEST_LIVE_REMOTE_HEAD: SHA_A}):
|
||||
out = self.srv.gitea_assess_master_parity(remote="prgs")
|
||||
self.assertTrue(out["mutation_safe"])
|
||||
self.assertFalse(out["live_stale"])
|
||||
self.assertNotIn("report", out)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
|
||||
@@ -139,6 +139,9 @@ EXPECTED_ROLE_EXCLUSIVE_TASKS = frozenset(
|
||||
"gitea_release_merger_pr_lease",
|
||||
"create_branch",
|
||||
"push_branch",
|
||||
# #812 AC20: publishing an unpublished local head is author-only for the
|
||||
# same reason every other push is — it writes a branch to the remote.
|
||||
"publish_unpublished_branch",
|
||||
"create_pr",
|
||||
"commit_files",
|
||||
"gitea_commit_files",
|
||||
|
||||
@@ -0,0 +1,149 @@
|
||||
"""Documentation acceptance for the web console architecture ADR (#632 / epic #631).
|
||||
|
||||
Enforces the acceptance criteria of issue #632:
|
||||
|
||||
* AC1 — the ADR exists and covers layers, authority, phases, API versioning,
|
||||
and a page map.
|
||||
* AC2 — every #631 child (#632–#651) maps to at least one architectural
|
||||
component.
|
||||
* AC3 — the closed MVP (#425–#436) is stated as foundation, not recreated.
|
||||
* AC4 — forbidden paths are explicit: raw provider incidents as work,
|
||||
browser-held tokens, process-kill recovery.
|
||||
* AC5 — a controller can approve the document without reading chat history.
|
||||
|
||||
Plus the linkage requirement: ``docs/webui-local-dev.md`` cross-links the ADR.
|
||||
"""
|
||||
from pathlib import Path
|
||||
|
||||
REPO_ROOT = Path(__file__).resolve().parent.parent
|
||||
ADR = (
|
||||
REPO_ROOT
|
||||
/ "docs"
|
||||
/ "architecture"
|
||||
/ "webui-control-plane-console-architecture-adr.md"
|
||||
)
|
||||
ADR_BASENAME = "webui-control-plane-console-architecture-adr.md"
|
||||
LOCAL_DEV = REPO_ROOT / "docs" / "webui-local-dev.md"
|
||||
|
||||
# Epic #631 children, phases 1-4 (twenty capability areas).
|
||||
EPIC_CHILDREN = tuple(f"#{number}" for number in range(632, 652))
|
||||
|
||||
|
||||
def _read(path: Path) -> str:
|
||||
assert path.is_file(), f"missing {path.relative_to(REPO_ROOT)}"
|
||||
return path.read_text(encoding="utf-8")
|
||||
|
||||
|
||||
def test_ac1_adr_exists_with_required_sections():
|
||||
text = _read(ADR)
|
||||
lower = text.lower()
|
||||
assert text.lstrip().startswith("#"), "ADR lacks a title"
|
||||
assert "#631" in text and "#632" in text
|
||||
for heading in (
|
||||
"## 2. Decision summary",
|
||||
"## 4. Authority boundaries",
|
||||
"## 5. Request flow and the redaction boundary",
|
||||
"## 6. API naming and versioning",
|
||||
"## 7. Page map",
|
||||
"## 8. Component ownership",
|
||||
"## 9. Phase gates",
|
||||
"## 11. Forbidden paths",
|
||||
):
|
||||
assert heading in text, f"ADR must contain section {heading!r}"
|
||||
assert "browser ui" in lower and "domain loader" in lower
|
||||
assert "control-plane db" in lower and "capability gate" in lower
|
||||
|
||||
|
||||
def test_ac1_api_versioning_is_decided_including_legacy_routes():
|
||||
text = _read(ADR)
|
||||
assert "/api/v1/" in text, "ADR must decide the versioned API prefix"
|
||||
assert "/api/v2/" in text, "ADR must state how breaking changes are handled"
|
||||
lower = text.lower()
|
||||
assert "compatibility alias" in lower, (
|
||||
"ADR must say what happens to the existing unversioned MVP exports"
|
||||
)
|
||||
|
||||
|
||||
def test_ac1_page_map_covers_mvp_routes():
|
||||
text = _read(ADR)
|
||||
for route in ("`/`", "`/health`", "`/projects`", "`/prompts`", "`/runtime`",
|
||||
"`/audit`", "`/actions`"):
|
||||
assert route in text, f"page map must account for MVP route {route}"
|
||||
|
||||
|
||||
def test_ac2_every_epic_child_maps_to_a_component():
|
||||
text = _read(ADR)
|
||||
ownership = text.split("## 8. Component ownership", 1)[-1].split("## 9.", 1)[0]
|
||||
missing = [child for child in EPIC_CHILDREN if child not in ownership]
|
||||
assert not missing, (
|
||||
f"epic #631 children without an architectural component: {missing}"
|
||||
)
|
||||
|
||||
|
||||
def test_ac2_every_child_row_declares_a_phase():
|
||||
text = _read(ADR)
|
||||
ownership = text.split("## 8. Component ownership", 1)[-1].split("## 9.", 1)[0]
|
||||
for child in EPIC_CHILDREN:
|
||||
row = next(
|
||||
(line for line in ownership.splitlines() if line.startswith(f"| {child} ")),
|
||||
None,
|
||||
)
|
||||
assert row is not None, f"no ownership row for {child}"
|
||||
assert row.rstrip().endswith(("| 1 |", "| 2 |", "| 3 |", "| 4 |")), (
|
||||
f"ownership row for {child} must end with its phase: {row!r}"
|
||||
)
|
||||
|
||||
|
||||
def test_ac3_mvp_is_foundation_not_recreated():
|
||||
text = _read(ADR)
|
||||
assert "#425" in text and "#436" in text
|
||||
lower = text.lower()
|
||||
assert "do not recreate" in lower or "recreating mvp scope" in lower
|
||||
assert "retained and evolved" in lower
|
||||
|
||||
|
||||
def test_ac4_forbidden_paths_are_explicit():
|
||||
text = _read(ADR)
|
||||
forbidden = text.split("## 11. Forbidden paths", 1)[-1].split("## 12.", 1)[0]
|
||||
lower = forbidden.lower()
|
||||
assert "raw provider incidents" in lower and "#612" in forbidden
|
||||
assert "browser-held tokens" in lower
|
||||
assert "process-kill recovery" in lower and "#630" in forbidden
|
||||
assert "ungated browser mutations" in lower
|
||||
|
||||
|
||||
def test_ac5_approval_checklist_is_self_contained():
|
||||
text = _read(ADR)
|
||||
assert "## 12. Approval checklist" in text
|
||||
checklist = text.split("## 12. Approval checklist", 1)[-1].split("## 13.", 1)[0]
|
||||
for marker in ("1.", "2.", "3.", "4.", "5.", "6."):
|
||||
assert marker in checklist, f"approval checklist missing item {marker}"
|
||||
|
||||
|
||||
def test_adr_states_the_two_boundary_invariants():
|
||||
text = _read(ADR)
|
||||
lower = text.lower()
|
||||
assert "no secrets to the browser" in lower
|
||||
assert "no ungated mutations" in lower
|
||||
|
||||
|
||||
def test_open_questions_are_recorded_not_implied():
|
||||
text = _read(ADR)
|
||||
assert "## 13. Open questions and follow-ups" in text
|
||||
section = text.split("## 13. Open questions and follow-ups", 1)[-1]
|
||||
assert "#633" in section, "deferred authorization work must name its issue"
|
||||
|
||||
|
||||
def test_local_dev_doc_cross_links_the_adr():
|
||||
text = _read(LOCAL_DEV)
|
||||
assert ADR_BASENAME in text, (
|
||||
"docs/webui-local-dev.md must cross-link the console architecture ADR "
|
||||
"(issue #632 scope)"
|
||||
)
|
||||
|
||||
|
||||
def test_docs_do_not_embed_secrets():
|
||||
for path in (ADR, LOCAL_DEV):
|
||||
text = _read(path)
|
||||
for marker in ("ghp_", "BEGIN PRIVATE KEY", "Authorization: Bearer"):
|
||||
assert marker not in text, f"{path.name} contains {marker!r}"
|
||||
@@ -0,0 +1,458 @@
|
||||
"""Tests for the worker registry and configuration schema (#798, epic #797)."""
|
||||
import json
|
||||
import sys
|
||||
import tempfile
|
||||
import unittest
|
||||
from pathlib import Path
|
||||
|
||||
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
|
||||
|
||||
from webui.worker_registry import (
|
||||
ALLOWED_ROLES,
|
||||
SCHEMA_VERSION,
|
||||
RegistryValidationError,
|
||||
WorkerRegistry,
|
||||
default_registry_path,
|
||||
find_provider,
|
||||
find_worker,
|
||||
history_dir,
|
||||
list_revisions,
|
||||
load_registry,
|
||||
registry_to_dict,
|
||||
registry_to_document,
|
||||
rollback_to_revision,
|
||||
save_registry,
|
||||
validate_payload,
|
||||
worker_to_dict,
|
||||
workers_for_provider,
|
||||
)
|
||||
|
||||
_EXPECTED_PROVIDER_IDS = ("claude", "grok", "codex", "agy", "kimi-k")
|
||||
|
||||
|
||||
def _provider(provider_id: str = "claude", **overrides) -> dict:
|
||||
payload = {
|
||||
"id": provider_id,
|
||||
"display_name": "Claude",
|
||||
"vendor": "Anthropic",
|
||||
"executable": "claude",
|
||||
"available": True,
|
||||
"models": ["claude-opus-4-8"],
|
||||
"notes": "",
|
||||
}
|
||||
payload.update(overrides)
|
||||
return payload
|
||||
|
||||
|
||||
def _worker(worker_id: str = "claude-author", **overrides) -> dict:
|
||||
payload = {
|
||||
"id": worker_id,
|
||||
"display_name": "Claude author",
|
||||
"provider": "claude",
|
||||
"model": "claude-opus-4-8",
|
||||
"project": "gitea-tools",
|
||||
"role": "author",
|
||||
"namespace": "gitea-author",
|
||||
"profile": "prgs-author",
|
||||
"workflow": "skills/llm-project-workflow/workflows/work-issue.md",
|
||||
"schedule": {"kind": "cron", "expression": "0 * * * *"},
|
||||
"timeout_seconds": 3600,
|
||||
"enabled": True,
|
||||
"scheduler": {"kind": "launchd", "label": "cc.prgs.claude.author"},
|
||||
"notes": "",
|
||||
}
|
||||
payload.update(overrides)
|
||||
return payload
|
||||
|
||||
|
||||
def _document(providers=None, workers=None, **overrides) -> dict:
|
||||
payload = {
|
||||
"version": SCHEMA_VERSION,
|
||||
"revision": 1,
|
||||
"updated_at": "2026-07-22T00:00:00Z",
|
||||
"providers": providers if providers is not None else [_provider()],
|
||||
"workers": workers if workers is not None else [_worker()],
|
||||
}
|
||||
payload.update(overrides)
|
||||
return payload
|
||||
|
||||
|
||||
class _TempRegistryCase(unittest.TestCase):
|
||||
"""Base case giving each test an isolated registry file."""
|
||||
|
||||
def setUp(self):
|
||||
self._tmp = tempfile.TemporaryDirectory()
|
||||
self.addCleanup(self._tmp.cleanup)
|
||||
self.path = Path(self._tmp.name) / "workers.registry.json"
|
||||
|
||||
def write(self, document: dict) -> Path:
|
||||
self.path.write_text(json.dumps(document, indent=2) + "\n", encoding="utf-8")
|
||||
return self.path
|
||||
|
||||
def parse(self, document: dict) -> WorkerRegistry:
|
||||
return validate_payload(document, source_path=self.path)
|
||||
|
||||
|
||||
class TestPackagedRegistry(unittest.TestCase):
|
||||
"""AC: the declarative registry is the source of truth and ships with the app."""
|
||||
|
||||
def test_default_path_points_at_packaged_data(self):
|
||||
path = default_registry_path()
|
||||
self.assertEqual(path.name, "workers.registry.json")
|
||||
self.assertEqual(path.parent.name, "data")
|
||||
|
||||
def test_packaged_registry_loads_and_validates(self):
|
||||
registry = load_registry()
|
||||
self.assertEqual(registry.version, SCHEMA_VERSION)
|
||||
self.assertGreaterEqual(registry.revision, 1)
|
||||
|
||||
def test_packaged_registry_declares_all_five_providers(self):
|
||||
registry = load_registry()
|
||||
self.assertEqual(
|
||||
tuple(provider.id for provider in registry.providers),
|
||||
_EXPECTED_PROVIDER_IDS,
|
||||
)
|
||||
|
||||
def test_packaged_registry_carries_no_credentials(self):
|
||||
raw = default_registry_path().read_text(encoding="utf-8").lower()
|
||||
for marker in ("token", "password", "secret", "api_key", "credential"):
|
||||
self.assertNotIn(marker, raw)
|
||||
|
||||
|
||||
class TestSeparateEntities(_TempRegistryCase):
|
||||
"""AC: providers and configured workers are separate entities."""
|
||||
|
||||
def test_provider_may_exist_with_no_workers(self):
|
||||
registry = self.parse(
|
||||
_document(providers=[_provider("grok", display_name="Grok")], workers=[])
|
||||
)
|
||||
self.assertEqual(len(registry.providers), 1)
|
||||
self.assertEqual(registry.workers, ())
|
||||
self.assertEqual(workers_for_provider(registry, "grok"), ())
|
||||
|
||||
def test_many_workers_may_share_one_provider(self):
|
||||
registry = self.parse(
|
||||
_document(
|
||||
workers=[
|
||||
_worker("claude-author"),
|
||||
_worker(
|
||||
"claude-reviewer",
|
||||
role="reviewer",
|
||||
namespace="gitea-reviewer",
|
||||
profile="prgs-reviewer",
|
||||
scheduler={"kind": "launchd", "label": "cc.prgs.claude.reviewer"},
|
||||
),
|
||||
]
|
||||
)
|
||||
)
|
||||
self.assertEqual(len(workers_for_provider(registry, "claude")), 2)
|
||||
self.assertEqual(len(registry.providers), 1)
|
||||
|
||||
def test_worker_referencing_unknown_provider_is_refused(self):
|
||||
with self.assertRaises(RegistryValidationError) as ctx:
|
||||
self.parse(_document(workers=[_worker(provider="mystery")]))
|
||||
self.assertIn("unknown provider", str(ctx.exception))
|
||||
|
||||
def test_lookup_helpers(self):
|
||||
registry = self.parse(_document())
|
||||
self.assertIsNotNone(find_worker(registry, "claude-author"))
|
||||
self.assertIsNone(find_worker(registry, "absent"))
|
||||
self.assertIsNotNone(find_provider(registry, "claude"))
|
||||
self.assertIsNone(find_provider(registry, "absent"))
|
||||
|
||||
|
||||
class TestRecordedFields(_TempRegistryCase):
|
||||
"""AC: records provider, model, project, role, namespace/profile, workflow,
|
||||
schedule, timeout, enabled state, and scheduler metadata."""
|
||||
|
||||
def test_every_required_field_is_recorded(self):
|
||||
registry = self.parse(_document())
|
||||
worker = registry.workers[0]
|
||||
self.assertEqual(worker.provider, "claude")
|
||||
self.assertEqual(worker.model, "claude-opus-4-8")
|
||||
self.assertEqual(worker.project, "gitea-tools")
|
||||
self.assertEqual(worker.role, "author")
|
||||
self.assertEqual(worker.namespace, "gitea-author")
|
||||
self.assertEqual(worker.profile, "prgs-author")
|
||||
self.assertEqual(worker.workflow, "skills/llm-project-workflow/workflows/work-issue.md")
|
||||
self.assertEqual(worker.schedule.kind, "cron")
|
||||
self.assertEqual(worker.schedule.expression, "0 * * * *")
|
||||
self.assertEqual(worker.timeout_seconds, 3600)
|
||||
self.assertTrue(worker.enabled)
|
||||
self.assertEqual(worker.scheduler.kind, "launchd")
|
||||
self.assertEqual(worker.scheduler.label, "cc.prgs.claude.author")
|
||||
|
||||
def test_each_required_field_is_individually_required(self):
|
||||
for field in (
|
||||
"provider", "model", "project", "role", "namespace",
|
||||
"profile", "workflow", "schedule", "timeout_seconds",
|
||||
"enabled", "scheduler", "id", "display_name",
|
||||
):
|
||||
with self.subTest(field=field):
|
||||
worker = _worker()
|
||||
worker.pop(field)
|
||||
with self.assertRaises(RegistryValidationError):
|
||||
self.parse(_document(workers=[worker]))
|
||||
|
||||
def test_all_sanctioned_roles_are_accepted(self):
|
||||
for role in ALLOWED_ROLES:
|
||||
with self.subTest(role=role):
|
||||
registry = self.parse(_document(workers=[_worker(role=role)]))
|
||||
self.assertEqual(registry.workers[0].role, role)
|
||||
|
||||
def test_unsanctioned_role_is_refused(self):
|
||||
with self.assertRaises(RegistryValidationError) as ctx:
|
||||
self.parse(_document(workers=[_worker(role="admin")]))
|
||||
self.assertIn("role must be one of", str(ctx.exception))
|
||||
|
||||
def test_worker_dict_round_trips_every_field(self):
|
||||
registry = self.parse(_document())
|
||||
encoded = worker_to_dict(registry.workers[0])
|
||||
self.assertEqual(encoded, _worker())
|
||||
json.dumps(encoded) # must stay JSON-safe for the #799 API
|
||||
|
||||
|
||||
class TestSchemaValidation(_TempRegistryCase):
|
||||
"""AC: supports schema validation — and fails closed."""
|
||||
|
||||
def test_unsupported_version_is_refused(self):
|
||||
with self.assertRaises(RegistryValidationError):
|
||||
self.parse(_document(version=2))
|
||||
|
||||
def test_root_must_be_an_object(self):
|
||||
with self.assertRaises(RegistryValidationError):
|
||||
validate_payload([], source_path=self.path)
|
||||
|
||||
def test_providers_must_be_non_empty(self):
|
||||
with self.assertRaises(RegistryValidationError):
|
||||
self.parse(_document(providers=[]))
|
||||
|
||||
def test_unknown_top_level_field_is_refused(self):
|
||||
with self.assertRaises(RegistryValidationError) as ctx:
|
||||
self.parse(_document(fleet=[]))
|
||||
self.assertIn("unknown fields", str(ctx.exception))
|
||||
|
||||
def test_unknown_worker_field_is_refused_not_ignored(self):
|
||||
# A typo'd field must not be silently dropped: "timeout_second" would
|
||||
# otherwise read as "no timeout declared".
|
||||
worker = _worker()
|
||||
worker["timeout_second"] = 30
|
||||
with self.assertRaises(RegistryValidationError) as ctx:
|
||||
self.parse(_document(workers=[worker]))
|
||||
self.assertIn("timeout_second", str(ctx.exception))
|
||||
|
||||
def test_credentials_are_refused_anywhere_in_the_document(self):
|
||||
for label, mutate in (
|
||||
("provider.api_token", lambda doc: doc["providers"][0].__setitem__("api_token", "x")),
|
||||
("worker.password", lambda doc: doc["workers"][0].__setitem__("password", "x")),
|
||||
("root.secret", lambda doc: doc.__setitem__("secret", "x")),
|
||||
):
|
||||
with self.subTest(field=label):
|
||||
document = _document()
|
||||
mutate(document)
|
||||
with self.assertRaises(ValueError) as ctx:
|
||||
self.parse(document)
|
||||
self.assertIn("credential", str(ctx.exception).lower())
|
||||
|
||||
def test_duplicate_worker_id_is_refused(self):
|
||||
workers = [_worker("dup"), _worker("dup", scheduler={"kind": "manual"})]
|
||||
with self.assertRaises(RegistryValidationError) as ctx:
|
||||
self.parse(_document(workers=workers))
|
||||
self.assertIn("duplicate worker id", str(ctx.exception))
|
||||
|
||||
def test_duplicate_provider_id_is_refused(self):
|
||||
with self.assertRaises(RegistryValidationError) as ctx:
|
||||
self.parse(_document(providers=[_provider("claude"), _provider("claude")], workers=[]))
|
||||
self.assertIn("duplicate provider id", str(ctx.exception))
|
||||
|
||||
def test_duplicate_launchagent_label_is_refused(self):
|
||||
# Two workers sharing a label would silently overwrite each other's agent.
|
||||
workers = [
|
||||
_worker("a", scheduler={"kind": "launchd", "label": "cc.prgs.same"}),
|
||||
_worker("b", scheduler={"kind": "launchd", "label": "cc.prgs.same"}),
|
||||
]
|
||||
with self.assertRaises(RegistryValidationError) as ctx:
|
||||
self.parse(_document(workers=workers))
|
||||
self.assertIn("duplicate scheduler label", str(ctx.exception))
|
||||
|
||||
def test_manual_scheduler_needs_no_label_and_many_may_coexist(self):
|
||||
workers = [
|
||||
_worker("a", scheduler={"kind": "manual"}),
|
||||
_worker("b", scheduler={"kind": "manual"}),
|
||||
]
|
||||
registry = self.parse(_document(workers=workers))
|
||||
self.assertEqual([w.scheduler.label for w in registry.workers], [None, None])
|
||||
|
||||
def test_launchd_scheduler_requires_a_label(self):
|
||||
with self.assertRaises(RegistryValidationError) as ctx:
|
||||
self.parse(_document(workers=[_worker(scheduler={"kind": "launchd"})]))
|
||||
self.assertIn("label is required", str(ctx.exception))
|
||||
|
||||
def test_unknown_scheduler_kind_is_refused(self):
|
||||
with self.assertRaises(RegistryValidationError):
|
||||
self.parse(_document(workers=[_worker(scheduler={"kind": "systemd", "label": "x"})]))
|
||||
|
||||
def test_timeout_must_be_a_positive_bounded_integer(self):
|
||||
for bad in (0, -1, "3600", 1.5, True, 86_401):
|
||||
with self.subTest(timeout=bad):
|
||||
with self.assertRaises(RegistryValidationError):
|
||||
self.parse(_document(workers=[_worker(timeout_seconds=bad)]))
|
||||
|
||||
def test_enabled_must_be_a_real_boolean(self):
|
||||
for bad in ("true", 1, None):
|
||||
with self.subTest(enabled=bad):
|
||||
with self.assertRaises(RegistryValidationError):
|
||||
self.parse(_document(workers=[_worker(enabled=bad)]))
|
||||
|
||||
def test_identifier_shape_is_enforced(self):
|
||||
for bad in ("Claude Author", "-leading", "UPPER", ""):
|
||||
with self.subTest(worker_id=bad):
|
||||
with self.assertRaises(RegistryValidationError):
|
||||
self.parse(_document(workers=[_worker(bad)]))
|
||||
|
||||
|
||||
class TestScheduleValidation(_TempRegistryCase):
|
||||
"""Schedules are declarations; next-run computation belongs to #803."""
|
||||
|
||||
def test_interval_schedule_requires_positive_seconds(self):
|
||||
registry = self.parse(
|
||||
_document(workers=[_worker(schedule={"kind": "interval", "seconds": 900})])
|
||||
)
|
||||
self.assertEqual(registry.workers[0].schedule.seconds, 900)
|
||||
with self.assertRaises(RegistryValidationError):
|
||||
self.parse(_document(workers=[_worker(schedule={"kind": "interval"})]))
|
||||
with self.assertRaises(RegistryValidationError):
|
||||
self.parse(_document(workers=[_worker(schedule={"kind": "interval", "seconds": 0})]))
|
||||
|
||||
def test_cron_schedule_requires_five_fields(self):
|
||||
with self.assertRaises(RegistryValidationError) as ctx:
|
||||
self.parse(_document(workers=[_worker(schedule={"kind": "cron", "expression": "0 *"})]))
|
||||
self.assertIn("five crontab fields", str(ctx.exception))
|
||||
|
||||
def test_manual_schedule_needs_no_timing(self):
|
||||
registry = self.parse(_document(workers=[_worker(schedule={"kind": "manual"})]))
|
||||
schedule = registry.workers[0].schedule
|
||||
self.assertEqual(schedule.kind, "manual")
|
||||
self.assertIsNone(schedule.seconds)
|
||||
self.assertIsNone(schedule.expression)
|
||||
|
||||
def test_fields_from_the_wrong_kind_are_refused(self):
|
||||
with self.assertRaises(RegistryValidationError) as ctx:
|
||||
self.parse(_document(workers=[_worker(schedule={"kind": "manual", "seconds": 60})]))
|
||||
self.assertIn("not valid for kind", str(ctx.exception))
|
||||
|
||||
def test_unknown_schedule_kind_is_refused(self):
|
||||
with self.assertRaises(RegistryValidationError):
|
||||
self.parse(_document(workers=[_worker(schedule={"kind": "hourly"})]))
|
||||
|
||||
|
||||
class TestAtomicPersistence(_TempRegistryCase):
|
||||
"""AC: atomic persistence."""
|
||||
|
||||
def test_save_then_load_round_trips(self):
|
||||
registry = self.parse(_document())
|
||||
save_registry(registry, self.path)
|
||||
reloaded = load_registry(self.path)
|
||||
self.assertEqual(
|
||||
[worker_to_dict(w) for w in reloaded.workers],
|
||||
[worker_to_dict(w) for w in registry.workers],
|
||||
)
|
||||
|
||||
def test_save_leaves_no_temp_files_behind(self):
|
||||
registry = self.parse(_document())
|
||||
save_registry(registry, self.path)
|
||||
save_registry(registry, self.path)
|
||||
leftovers = [p.name for p in self.path.parent.iterdir() if p.name.startswith(".")]
|
||||
self.assertEqual(leftovers, [])
|
||||
|
||||
def test_save_refuses_to_persist_an_invalid_document(self):
|
||||
registry = self.parse(_document())
|
||||
broken = WorkerRegistry(
|
||||
version=registry.version,
|
||||
revision=registry.revision,
|
||||
updated_at=registry.updated_at,
|
||||
providers=registry.providers,
|
||||
# A worker whose provider is not declared in the registry.
|
||||
workers=tuple(
|
||||
type(worker)(**{**worker.__dict__, "provider": "vanished"})
|
||||
for worker in registry.workers
|
||||
),
|
||||
source_path=self.path,
|
||||
)
|
||||
with self.assertRaises(RegistryValidationError):
|
||||
save_registry(broken, self.path)
|
||||
self.assertFalse(self.path.exists(), "invalid save must not create the file")
|
||||
|
||||
def test_document_shape_excludes_local_paths_but_api_shape_includes_it(self):
|
||||
registry = self.parse(_document())
|
||||
self.assertNotIn("source_path", registry_to_document(registry))
|
||||
self.assertEqual(registry_to_dict(registry)["source_path"], str(self.path))
|
||||
|
||||
|
||||
class TestVersioningAndRollback(_TempRegistryCase):
|
||||
"""AC: versioning and rollback."""
|
||||
|
||||
def _seed(self) -> WorkerRegistry:
|
||||
self.write(_document())
|
||||
return load_registry(self.path)
|
||||
|
||||
def test_revision_increments_on_each_save(self):
|
||||
registry = self._seed()
|
||||
self.assertEqual(registry.revision, 1)
|
||||
second = save_registry(registry, self.path)
|
||||
self.assertEqual(second.revision, 2)
|
||||
third = save_registry(second, self.path)
|
||||
self.assertEqual(third.revision, 3)
|
||||
|
||||
def test_updated_at_is_refreshed_and_utc(self):
|
||||
registry = self._seed()
|
||||
saved = save_registry(registry, self.path)
|
||||
self.assertRegex(saved.updated_at, r"^\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}Z$")
|
||||
|
||||
def test_superseded_revisions_are_retained(self):
|
||||
registry = self._seed()
|
||||
second = save_registry(registry, self.path)
|
||||
save_registry(second, self.path)
|
||||
self.assertEqual(list_revisions(self.path), (1, 2))
|
||||
self.assertTrue(history_dir(self.path).is_dir())
|
||||
|
||||
def test_rollback_restores_prior_content_as_a_new_revision(self):
|
||||
self.write(_document(workers=[_worker("original")]))
|
||||
registry = load_registry(self.path)
|
||||
|
||||
changed = WorkerRegistry(
|
||||
version=registry.version,
|
||||
revision=registry.revision,
|
||||
updated_at=registry.updated_at,
|
||||
providers=registry.providers,
|
||||
workers=(), # operator deletes every worker
|
||||
source_path=self.path,
|
||||
)
|
||||
save_registry(changed, self.path)
|
||||
self.assertEqual(load_registry(self.path).workers, ())
|
||||
|
||||
restored = rollback_to_revision(1, self.path)
|
||||
self.assertEqual([w.id for w in restored.workers], ["original"])
|
||||
# Append-only: the rollback publishes a new head rather than rewinding.
|
||||
self.assertGreater(restored.revision, 2)
|
||||
self.assertEqual([w.id for w in load_registry(self.path).workers], ["original"])
|
||||
|
||||
def test_rollback_to_unknown_revision_fails_closed(self):
|
||||
self._seed()
|
||||
with self.assertRaises(RegistryValidationError) as ctx:
|
||||
rollback_to_revision(99, self.path)
|
||||
self.assertIn("not retained", str(ctx.exception))
|
||||
|
||||
def test_revision_must_be_a_positive_integer(self):
|
||||
for bad in (0, -1, "1", None):
|
||||
with self.subTest(revision=bad):
|
||||
with self.assertRaises(RegistryValidationError):
|
||||
self.parse(_document(revision=bad))
|
||||
|
||||
def test_history_is_empty_before_any_save(self):
|
||||
self.write(_document())
|
||||
self.assertEqual(list_revisions(self.path), ())
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -0,0 +1,57 @@
|
||||
{
|
||||
"version": 1,
|
||||
"revision": 1,
|
||||
"updated_at": "2026-07-22T00:00:00Z",
|
||||
"providers": [
|
||||
{
|
||||
"id": "claude",
|
||||
"display_name": "Claude",
|
||||
"vendor": "Anthropic",
|
||||
"executable": "claude",
|
||||
"available": true,
|
||||
"models": [
|
||||
"claude-opus-4-8",
|
||||
"claude-sonnet-5",
|
||||
"claude-haiku-4-5-20251001"
|
||||
],
|
||||
"notes": "Model list is a declaration. Live enumeration and version inspection belong to the provider adapter framework (#800)."
|
||||
},
|
||||
{
|
||||
"id": "grok",
|
||||
"display_name": "Grok",
|
||||
"vendor": "xAI",
|
||||
"executable": "grok",
|
||||
"available": true,
|
||||
"models": [],
|
||||
"notes": "Models enumerated by the provider adapter (#800); not declared here."
|
||||
},
|
||||
{
|
||||
"id": "codex",
|
||||
"display_name": "Codex",
|
||||
"vendor": "OpenAI",
|
||||
"executable": "codex",
|
||||
"available": true,
|
||||
"models": [],
|
||||
"notes": "Models enumerated by the provider adapter (#800); not declared here."
|
||||
},
|
||||
{
|
||||
"id": "agy",
|
||||
"display_name": "AGY",
|
||||
"vendor": "Antigravity",
|
||||
"executable": "agy",
|
||||
"available": true,
|
||||
"models": [],
|
||||
"notes": "MCP allowlist gating applies to this provider; confirm server-side allowlist before configuring a worker."
|
||||
},
|
||||
{
|
||||
"id": "kimi-k",
|
||||
"display_name": "Kimi K",
|
||||
"vendor": "Moonshot AI",
|
||||
"executable": "kimi",
|
||||
"available": true,
|
||||
"models": [],
|
||||
"notes": "Provider id is kimi-k; the executable on PATH is kimi. Models enumerated by the provider adapter (#800)."
|
||||
}
|
||||
],
|
||||
"workers": []
|
||||
}
|
||||
@@ -8,27 +8,7 @@ from dataclasses import dataclass
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
_FORBIDDEN_EXACT_KEYS = frozenset({
|
||||
"token",
|
||||
"password",
|
||||
"secret",
|
||||
"credential",
|
||||
"auth",
|
||||
"api_key",
|
||||
"api-key",
|
||||
})
|
||||
_FORBIDDEN_KEY_PREFIXES = ("auth_", "api_key_", "api-key_")
|
||||
_FORBIDDEN_KEY_SUFFIXES = ("_token", "_secret", "_password", "_credential", "_auth")
|
||||
|
||||
|
||||
def _is_forbidden_key(key: str) -> bool:
|
||||
lowered = key.lower()
|
||||
if lowered in _FORBIDDEN_EXACT_KEYS:
|
||||
return True
|
||||
return (
|
||||
lowered.startswith(_FORBIDDEN_KEY_PREFIXES)
|
||||
or lowered.endswith(_FORBIDDEN_KEY_SUFFIXES)
|
||||
)
|
||||
from webui.registry_safety import reject_credential_keys as _reject_credential_keys
|
||||
|
||||
_REQUIRED_PROJECT_FIELDS = (
|
||||
"id",
|
||||
@@ -79,18 +59,6 @@ def default_registry_path() -> Path:
|
||||
return (Path(__file__).resolve().parent / "data" / "projects.registry.json").resolve()
|
||||
|
||||
|
||||
def _reject_credential_keys(obj: Any, *, path: str = "") -> None:
|
||||
if isinstance(obj, dict):
|
||||
for key, value in obj.items():
|
||||
key_path = f"{path}.{key}" if path else key
|
||||
if _is_forbidden_key(key):
|
||||
raise ValueError(f"registry must not store credentials ({key_path})")
|
||||
_reject_credential_keys(value, path=key_path)
|
||||
elif isinstance(obj, list):
|
||||
for index, item in enumerate(obj):
|
||||
_reject_credential_keys(item, path=f"{path}[{index}]")
|
||||
|
||||
|
||||
def _parse_onboarding(raw: list[dict[str, Any]] | None) -> tuple[OnboardingStep, ...]:
|
||||
if not raw:
|
||||
return ()
|
||||
|
||||
@@ -0,0 +1,52 @@
|
||||
"""Shared credential-rejection guard for web UI registries (#427, #798).
|
||||
|
||||
Registries are operator-editable declarative files that the web UI loads and,
|
||||
for the worker registry, writes back. None of them may ever carry a secret:
|
||||
credentials belong in the keychain and reach worker processes through
|
||||
environment injection, never through a file the browser layer can read.
|
||||
|
||||
The check is structural rather than value-based on purpose. A value scanner has
|
||||
to guess what a secret looks like; a key scanner refuses the *shape* of a
|
||||
credential field, so an operator cannot introduce one by accident and a later
|
||||
loader cannot silently pass one through.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Any
|
||||
|
||||
_FORBIDDEN_EXACT_KEYS = frozenset({
|
||||
"token",
|
||||
"password",
|
||||
"secret",
|
||||
"credential",
|
||||
"auth",
|
||||
"api_key",
|
||||
"api-key",
|
||||
})
|
||||
_FORBIDDEN_KEY_PREFIXES = ("auth_", "api_key_", "api-key_")
|
||||
_FORBIDDEN_KEY_SUFFIXES = ("_token", "_secret", "_password", "_credential", "_auth")
|
||||
|
||||
|
||||
def is_forbidden_key(key: str) -> bool:
|
||||
"""Return True when *key* names a credential field."""
|
||||
lowered = key.lower()
|
||||
if lowered in _FORBIDDEN_EXACT_KEYS:
|
||||
return True
|
||||
return (
|
||||
lowered.startswith(_FORBIDDEN_KEY_PREFIXES)
|
||||
or lowered.endswith(_FORBIDDEN_KEY_SUFFIXES)
|
||||
)
|
||||
|
||||
|
||||
def reject_credential_keys(obj: Any, *, path: str = "", subject: str = "registry") -> None:
|
||||
"""Raise ValueError when *obj* carries a credential-shaped key at any depth."""
|
||||
if isinstance(obj, dict):
|
||||
for key, value in obj.items():
|
||||
key_path = f"{path}.{key}" if path else key
|
||||
if is_forbidden_key(key):
|
||||
raise ValueError(f"{subject} must not store credentials ({key_path})")
|
||||
reject_credential_keys(value, path=key_path, subject=subject)
|
||||
elif isinstance(obj, list):
|
||||
for index, item in enumerate(obj):
|
||||
reject_credential_keys(item, path=f"{path}[{index}]", subject=subject)
|
||||
@@ -0,0 +1,647 @@
|
||||
"""Declarative worker registry and configuration schema (#798, epic #797).
|
||||
|
||||
The registry is the single source of truth for the scheduled multi-LLM worker
|
||||
fleet. It is a versioned JSON document holding two *separate* entity kinds:
|
||||
|
||||
* **Providers** — the LLM runtimes a worker can be built on (Claude, Grok,
|
||||
Codex, AGY, Kimi K). A provider describes the runtime itself: vendor,
|
||||
executable name, models it can serve, and whether it is available on this
|
||||
machine. Providers exist whether or not any worker uses them.
|
||||
* **Workers** — a configured *instance*: one provider, one model, one project,
|
||||
one role, one MCP namespace/profile, one workflow, one schedule. Several
|
||||
workers may share a provider; a worker naming an undeclared provider is
|
||||
refused.
|
||||
|
||||
Keeping them separate is what lets #799 list all five providers even when a
|
||||
provider currently has no configured worker, and it stops provider facts from
|
||||
being copied into (and drifting across) every worker record.
|
||||
|
||||
Scope boundary. This module owns the data model, its validation, and its
|
||||
persistence. It does **not** schedule anything, launch anything, probe provider
|
||||
executables, or serve HTTP. Loading a registry never touches a process; the
|
||||
live fields a dashboard wants (PID, elapsed time, next run) are derived
|
||||
elsewhere (#799, #801, #803, #804) from these declarations.
|
||||
|
||||
Safety invariants:
|
||||
|
||||
* No credential may be stored (:mod:`webui.registry_safety`), so the registry
|
||||
stays safe to render and to hand to a browser layer.
|
||||
* Validation fails closed. Unknown fields are refused rather than ignored, so a
|
||||
typo cannot silently disable a timeout or a role binding.
|
||||
* Writes are atomic and every superseded document is retained as a numbered
|
||||
revision, so a bad edit is recoverable by rollback rather than hand-repair.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
import tempfile
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
from webui.registry_safety import reject_credential_keys
|
||||
|
||||
SCHEMA_VERSION = 1
|
||||
|
||||
#: Roles a worker may hold. These mirror the sanctioned MCP role kinds; a
|
||||
#: worker may not invent one, because the role selects the namespace/profile
|
||||
#: whose capability gates constrain it.
|
||||
ALLOWED_ROLES = ("author", "reviewer", "merger", "reconciler", "cleanup")
|
||||
|
||||
#: Scheduler backends the registry can describe. ``manual`` means the worker is
|
||||
#: only ever started on request and has no recurring trigger.
|
||||
ALLOWED_SCHEDULER_KINDS = ("launchd", "manual")
|
||||
|
||||
#: Schedule kinds. Next-run computation belongs to #803; this module only
|
||||
#: guarantees the declaration is well formed.
|
||||
ALLOWED_SCHEDULE_KINDS = ("interval", "cron", "manual")
|
||||
|
||||
_REQUIRED_PROVIDER_FIELDS = ("id", "display_name", "vendor", "executable", "available")
|
||||
_OPTIONAL_PROVIDER_FIELDS = ("models", "notes")
|
||||
|
||||
_REQUIRED_WORKER_FIELDS = (
|
||||
"id",
|
||||
"display_name",
|
||||
"provider",
|
||||
"model",
|
||||
"project",
|
||||
"role",
|
||||
"namespace",
|
||||
"profile",
|
||||
"workflow",
|
||||
"schedule",
|
||||
"timeout_seconds",
|
||||
"enabled",
|
||||
"scheduler",
|
||||
)
|
||||
_OPTIONAL_WORKER_FIELDS = ("notes",)
|
||||
|
||||
_ID_RE = re.compile(r"^[a-z0-9][a-z0-9._-]*$")
|
||||
|
||||
#: Guards against an operator writing a timeout that would let a worker hold a
|
||||
#: lease effectively forever. 24h is far above any sanctioned cycle.
|
||||
_MAX_TIMEOUT_SECONDS = 86_400
|
||||
|
||||
#: How many superseded revisions to retain beside the live file.
|
||||
_HISTORY_LIMIT = 20
|
||||
|
||||
_TOP_LEVEL_FIELDS = frozenset({"version", "revision", "updated_at", "providers", "workers"})
|
||||
|
||||
|
||||
class RegistryValidationError(ValueError):
|
||||
"""Raised when a registry document violates the schema."""
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class ProviderRecord:
|
||||
id: str
|
||||
display_name: str
|
||||
vendor: str
|
||||
executable: str
|
||||
available: bool
|
||||
models: tuple[str, ...]
|
||||
notes: str
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class ScheduleSpec:
|
||||
kind: str
|
||||
#: Set for ``interval`` schedules.
|
||||
seconds: int | None
|
||||
#: Set for ``cron`` schedules — a five-field crontab expression.
|
||||
expression: str | None
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class SchedulerSpec:
|
||||
kind: str
|
||||
#: LaunchAgent label; required for ``launchd``, absent for ``manual``.
|
||||
label: str | None
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class WorkerRecord:
|
||||
id: str
|
||||
display_name: str
|
||||
provider: str
|
||||
model: str
|
||||
project: str
|
||||
role: str
|
||||
namespace: str
|
||||
profile: str
|
||||
workflow: str
|
||||
schedule: ScheduleSpec
|
||||
timeout_seconds: int
|
||||
enabled: bool
|
||||
scheduler: SchedulerSpec
|
||||
notes: str
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class WorkerRegistry:
|
||||
version: int
|
||||
revision: int
|
||||
updated_at: str
|
||||
providers: tuple[ProviderRecord, ...]
|
||||
workers: tuple[WorkerRecord, ...]
|
||||
source_path: Path
|
||||
|
||||
|
||||
# ── paths ────────────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def default_registry_path() -> Path:
|
||||
"""Location of the packaged worker registry, overridable for tests/deploys."""
|
||||
override = os.environ.get("WEBUI_WORKER_REGISTRY", "").strip()
|
||||
if override:
|
||||
return Path(override).expanduser().resolve()
|
||||
return (Path(__file__).resolve().parent / "data" / "workers.registry.json").resolve()
|
||||
|
||||
|
||||
def history_dir(path: Path | None = None) -> Path:
|
||||
"""Directory holding superseded revisions of *path*."""
|
||||
source = (path or default_registry_path()).resolve()
|
||||
return source.parent / f"{source.name}.history"
|
||||
|
||||
|
||||
# ── field helpers ────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def _require_exact_fields(
|
||||
raw: Any,
|
||||
*,
|
||||
required: tuple[str, ...],
|
||||
optional: tuple[str, ...],
|
||||
subject: str,
|
||||
) -> dict[str, Any]:
|
||||
if not isinstance(raw, dict):
|
||||
raise RegistryValidationError(f"{subject} must be an object")
|
||||
missing = [field for field in required if field not in raw]
|
||||
if missing:
|
||||
raise RegistryValidationError(
|
||||
f"{subject} missing required fields: {', '.join(sorted(missing))}"
|
||||
)
|
||||
unknown = sorted(set(raw) - set(required) - set(optional))
|
||||
if unknown:
|
||||
# Fail closed: silently dropping an unrecognized key is how a typo'd
|
||||
# "timeout_second" ends up meaning "no timeout".
|
||||
raise RegistryValidationError(f"{subject} has unknown fields: {', '.join(unknown)}")
|
||||
return raw
|
||||
|
||||
|
||||
def _require_identifier(value: Any, *, subject: str) -> str:
|
||||
text = str(value).strip()
|
||||
if not _ID_RE.match(text):
|
||||
raise RegistryValidationError(
|
||||
f"{subject} must be lowercase alphanumeric with '.', '_', or '-' (got {value!r})"
|
||||
)
|
||||
return text
|
||||
|
||||
|
||||
def _require_text(value: Any, *, subject: str) -> str:
|
||||
if not isinstance(value, str):
|
||||
raise RegistryValidationError(f"{subject} must be a string (got {value!r})")
|
||||
text = value.strip()
|
||||
if not text:
|
||||
raise RegistryValidationError(f"{subject} must be a non-empty string")
|
||||
return text
|
||||
|
||||
|
||||
def _require_bool(value: Any, *, subject: str) -> bool:
|
||||
if not isinstance(value, bool):
|
||||
raise RegistryValidationError(f"{subject} must be a boolean (got {value!r})")
|
||||
return value
|
||||
|
||||
|
||||
def _require_positive_int(value: Any, *, subject: str, maximum: int | None = None) -> int:
|
||||
if isinstance(value, bool) or not isinstance(value, int):
|
||||
raise RegistryValidationError(f"{subject} must be an integer (got {value!r})")
|
||||
if value <= 0:
|
||||
raise RegistryValidationError(f"{subject} must be greater than zero (got {value})")
|
||||
if maximum is not None and value > maximum:
|
||||
raise RegistryValidationError(f"{subject} must not exceed {maximum} (got {value})")
|
||||
return value
|
||||
|
||||
|
||||
# ── parsing ──────────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def _parse_provider(raw: Any) -> ProviderRecord:
|
||||
data = _require_exact_fields(
|
||||
raw,
|
||||
required=_REQUIRED_PROVIDER_FIELDS,
|
||||
optional=_OPTIONAL_PROVIDER_FIELDS,
|
||||
subject="provider",
|
||||
)
|
||||
provider_id = _require_identifier(data["id"], subject="provider.id")
|
||||
|
||||
models_raw = data.get("models") or []
|
||||
if not isinstance(models_raw, list):
|
||||
raise RegistryValidationError(f"provider[{provider_id}].models must be an array")
|
||||
models = tuple(
|
||||
_require_text(item, subject=f"provider[{provider_id}].models[]") for item in models_raw
|
||||
)
|
||||
|
||||
return ProviderRecord(
|
||||
id=provider_id,
|
||||
display_name=_require_text(
|
||||
data["display_name"], subject=f"provider[{provider_id}].display_name"
|
||||
),
|
||||
vendor=_require_text(data["vendor"], subject=f"provider[{provider_id}].vendor"),
|
||||
executable=_require_text(data["executable"], subject=f"provider[{provider_id}].executable"),
|
||||
available=_require_bool(data["available"], subject=f"provider[{provider_id}].available"),
|
||||
models=models,
|
||||
notes=str(data.get("notes") or "").strip(),
|
||||
)
|
||||
|
||||
|
||||
def _parse_schedule(raw: Any, *, subject: str) -> ScheduleSpec:
|
||||
if not isinstance(raw, dict):
|
||||
raise RegistryValidationError(f"{subject} must be an object")
|
||||
kind = _require_text(raw.get("kind"), subject=f"{subject}.kind")
|
||||
if kind not in ALLOWED_SCHEDULE_KINDS:
|
||||
raise RegistryValidationError(
|
||||
f"{subject}.kind must be one of {', '.join(ALLOWED_SCHEDULE_KINDS)} (got {kind!r})"
|
||||
)
|
||||
|
||||
seconds: int | None = None
|
||||
expression: str | None = None
|
||||
|
||||
if kind == "interval":
|
||||
if "seconds" not in raw:
|
||||
raise RegistryValidationError(f"{subject}.seconds is required for interval schedules")
|
||||
seconds = _require_positive_int(raw["seconds"], subject=f"{subject}.seconds")
|
||||
elif kind == "cron":
|
||||
if "expression" not in raw:
|
||||
raise RegistryValidationError(f"{subject}.expression is required for cron schedules")
|
||||
expression = _require_text(raw["expression"], subject=f"{subject}.expression")
|
||||
if len(expression.split()) != 5:
|
||||
raise RegistryValidationError(
|
||||
f"{subject}.expression must have five crontab fields (got {expression!r})"
|
||||
)
|
||||
|
||||
allowed = {"kind"}
|
||||
if kind == "interval":
|
||||
allowed.add("seconds")
|
||||
elif kind == "cron":
|
||||
allowed.add("expression")
|
||||
unknown = sorted(set(raw) - allowed)
|
||||
if unknown:
|
||||
raise RegistryValidationError(
|
||||
f"{subject} has fields not valid for kind {kind!r}: {', '.join(unknown)}"
|
||||
)
|
||||
|
||||
return ScheduleSpec(kind=kind, seconds=seconds, expression=expression)
|
||||
|
||||
|
||||
def _parse_scheduler(raw: Any, *, subject: str) -> SchedulerSpec:
|
||||
if not isinstance(raw, dict):
|
||||
raise RegistryValidationError(f"{subject} must be an object")
|
||||
kind = _require_text(raw.get("kind"), subject=f"{subject}.kind")
|
||||
if kind not in ALLOWED_SCHEDULER_KINDS:
|
||||
raise RegistryValidationError(
|
||||
f"{subject}.kind must be one of {', '.join(ALLOWED_SCHEDULER_KINDS)} (got {kind!r})"
|
||||
)
|
||||
|
||||
label: str | None = None
|
||||
if kind == "launchd":
|
||||
if "label" not in raw:
|
||||
raise RegistryValidationError(f"{subject}.label is required for launchd schedulers")
|
||||
label = _require_text(raw["label"], subject=f"{subject}.label")
|
||||
|
||||
allowed = {"kind"}
|
||||
if kind == "launchd":
|
||||
allowed.add("label")
|
||||
unknown = sorted(set(raw) - allowed)
|
||||
if unknown:
|
||||
raise RegistryValidationError(
|
||||
f"{subject} has fields not valid for kind {kind!r}: {', '.join(unknown)}"
|
||||
)
|
||||
|
||||
return SchedulerSpec(kind=kind, label=label)
|
||||
|
||||
|
||||
def _parse_worker(raw: Any) -> WorkerRecord:
|
||||
data = _require_exact_fields(
|
||||
raw,
|
||||
required=_REQUIRED_WORKER_FIELDS,
|
||||
optional=_OPTIONAL_WORKER_FIELDS,
|
||||
subject="worker",
|
||||
)
|
||||
worker_id = _require_identifier(data["id"], subject="worker.id")
|
||||
|
||||
role = _require_text(data["role"], subject=f"worker[{worker_id}].role")
|
||||
if role not in ALLOWED_ROLES:
|
||||
raise RegistryValidationError(
|
||||
f"worker[{worker_id}].role must be one of {', '.join(ALLOWED_ROLES)} (got {role!r})"
|
||||
)
|
||||
|
||||
return WorkerRecord(
|
||||
id=worker_id,
|
||||
display_name=_require_text(
|
||||
data["display_name"], subject=f"worker[{worker_id}].display_name"
|
||||
),
|
||||
provider=_require_identifier(data["provider"], subject=f"worker[{worker_id}].provider"),
|
||||
model=_require_text(data["model"], subject=f"worker[{worker_id}].model"),
|
||||
project=_require_text(data["project"], subject=f"worker[{worker_id}].project"),
|
||||
role=role,
|
||||
namespace=_require_text(data["namespace"], subject=f"worker[{worker_id}].namespace"),
|
||||
profile=_require_text(data["profile"], subject=f"worker[{worker_id}].profile"),
|
||||
workflow=_require_text(data["workflow"], subject=f"worker[{worker_id}].workflow"),
|
||||
schedule=_parse_schedule(data["schedule"], subject=f"worker[{worker_id}].schedule"),
|
||||
timeout_seconds=_require_positive_int(
|
||||
data["timeout_seconds"],
|
||||
subject=f"worker[{worker_id}].timeout_seconds",
|
||||
maximum=_MAX_TIMEOUT_SECONDS,
|
||||
),
|
||||
enabled=_require_bool(data["enabled"], subject=f"worker[{worker_id}].enabled"),
|
||||
scheduler=_parse_scheduler(data["scheduler"], subject=f"worker[{worker_id}].scheduler"),
|
||||
notes=str(data.get("notes") or "").strip(),
|
||||
)
|
||||
|
||||
|
||||
def _require_unique(values: list[str], *, subject: str) -> None:
|
||||
seen: set[str] = set()
|
||||
for value in values:
|
||||
if value in seen:
|
||||
raise RegistryValidationError(f"duplicate {subject}: {value}")
|
||||
seen.add(value)
|
||||
|
||||
|
||||
def validate_payload(payload: Any, *, source_path: Path) -> WorkerRegistry:
|
||||
"""Validate a decoded registry document and return the typed registry.
|
||||
|
||||
Raises :class:`RegistryValidationError` on any violation; never partially
|
||||
accepts a document.
|
||||
"""
|
||||
if not isinstance(payload, dict):
|
||||
raise RegistryValidationError("registry root must be an object")
|
||||
|
||||
version = payload.get("version")
|
||||
if version != SCHEMA_VERSION:
|
||||
raise RegistryValidationError(f"unsupported registry version: {version!r}")
|
||||
|
||||
reject_credential_keys(payload, subject="worker registry")
|
||||
|
||||
unknown = sorted(set(payload) - _TOP_LEVEL_FIELDS)
|
||||
if unknown:
|
||||
raise RegistryValidationError(f"registry has unknown fields: {', '.join(unknown)}")
|
||||
|
||||
revision = _require_positive_int(payload.get("revision"), subject="revision")
|
||||
updated_at = _require_text(payload.get("updated_at"), subject="updated_at")
|
||||
|
||||
providers_raw = payload.get("providers")
|
||||
if not isinstance(providers_raw, list) or not providers_raw:
|
||||
raise RegistryValidationError("providers must be a non-empty array")
|
||||
providers = tuple(_parse_provider(item) for item in providers_raw)
|
||||
_require_unique([provider.id for provider in providers], subject="provider id")
|
||||
|
||||
workers_raw = payload.get("workers")
|
||||
if not isinstance(workers_raw, list):
|
||||
raise RegistryValidationError("workers must be an array")
|
||||
workers = tuple(_parse_worker(item) for item in workers_raw)
|
||||
_require_unique([worker.id for worker in workers], subject="worker id")
|
||||
|
||||
# Referential integrity: a worker naming an undeclared provider would look
|
||||
# configured while being unrunnable, which is exactly the ambiguous
|
||||
# ownership the epic requires to fail closed.
|
||||
known_providers = {provider.id for provider in providers}
|
||||
for worker in workers:
|
||||
if worker.provider not in known_providers:
|
||||
raise RegistryValidationError(
|
||||
f"worker[{worker.id}].provider references unknown provider {worker.provider!r}"
|
||||
)
|
||||
|
||||
# A LaunchAgent label identifies a job to launchd; two workers sharing one
|
||||
# would silently overwrite each other's agent.
|
||||
_require_unique(
|
||||
[worker.scheduler.label for worker in workers if worker.scheduler.label],
|
||||
subject="scheduler label",
|
||||
)
|
||||
|
||||
return WorkerRegistry(
|
||||
version=version,
|
||||
revision=revision,
|
||||
updated_at=updated_at,
|
||||
providers=providers,
|
||||
workers=workers,
|
||||
source_path=source_path,
|
||||
)
|
||||
|
||||
|
||||
def load_registry(path: Path | None = None) -> WorkerRegistry:
|
||||
"""Load and validate the worker registry from disk."""
|
||||
source = (path or default_registry_path()).resolve()
|
||||
payload = json.loads(source.read_text(encoding="utf-8"))
|
||||
return validate_payload(payload, source_path=source)
|
||||
|
||||
|
||||
# ── serialization ────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def provider_to_dict(provider: ProviderRecord) -> dict[str, Any]:
|
||||
return {
|
||||
"id": provider.id,
|
||||
"display_name": provider.display_name,
|
||||
"vendor": provider.vendor,
|
||||
"executable": provider.executable,
|
||||
"available": provider.available,
|
||||
"models": list(provider.models),
|
||||
"notes": provider.notes,
|
||||
}
|
||||
|
||||
|
||||
def _schedule_to_dict(schedule: ScheduleSpec) -> dict[str, Any]:
|
||||
payload: dict[str, Any] = {"kind": schedule.kind}
|
||||
if schedule.kind == "interval":
|
||||
payload["seconds"] = schedule.seconds
|
||||
elif schedule.kind == "cron":
|
||||
payload["expression"] = schedule.expression
|
||||
return payload
|
||||
|
||||
|
||||
def _scheduler_to_dict(scheduler: SchedulerSpec) -> dict[str, Any]:
|
||||
payload: dict[str, Any] = {"kind": scheduler.kind}
|
||||
if scheduler.kind == "launchd":
|
||||
payload["label"] = scheduler.label
|
||||
return payload
|
||||
|
||||
|
||||
def worker_to_dict(worker: WorkerRecord) -> dict[str, Any]:
|
||||
return {
|
||||
"id": worker.id,
|
||||
"display_name": worker.display_name,
|
||||
"provider": worker.provider,
|
||||
"model": worker.model,
|
||||
"project": worker.project,
|
||||
"role": worker.role,
|
||||
"namespace": worker.namespace,
|
||||
"profile": worker.profile,
|
||||
"workflow": worker.workflow,
|
||||
"schedule": _schedule_to_dict(worker.schedule),
|
||||
"timeout_seconds": worker.timeout_seconds,
|
||||
"enabled": worker.enabled,
|
||||
"scheduler": _scheduler_to_dict(worker.scheduler),
|
||||
"notes": worker.notes,
|
||||
}
|
||||
|
||||
|
||||
def registry_to_document(registry: WorkerRegistry) -> dict[str, Any]:
|
||||
"""Serialize to the on-disk document shape (no local paths embedded)."""
|
||||
return {
|
||||
"version": registry.version,
|
||||
"revision": registry.revision,
|
||||
"updated_at": registry.updated_at,
|
||||
"providers": [provider_to_dict(provider) for provider in registry.providers],
|
||||
"workers": [worker_to_dict(worker) for worker in registry.workers],
|
||||
}
|
||||
|
||||
|
||||
def registry_to_dict(registry: WorkerRegistry) -> dict[str, Any]:
|
||||
"""Serialize for JSON API responses (adds the resolved source path)."""
|
||||
document = registry_to_document(registry)
|
||||
document["source_path"] = str(registry.source_path)
|
||||
return document
|
||||
|
||||
|
||||
def find_worker(registry: WorkerRegistry, worker_id: str) -> WorkerRecord | None:
|
||||
for worker in registry.workers:
|
||||
if worker.id == worker_id:
|
||||
return worker
|
||||
return None
|
||||
|
||||
|
||||
def find_provider(registry: WorkerRegistry, provider_id: str) -> ProviderRecord | None:
|
||||
for provider in registry.providers:
|
||||
if provider.id == provider_id:
|
||||
return provider
|
||||
return None
|
||||
|
||||
|
||||
def workers_for_provider(registry: WorkerRegistry, provider_id: str) -> tuple[WorkerRecord, ...]:
|
||||
return tuple(worker for worker in registry.workers if worker.provider == provider_id)
|
||||
|
||||
|
||||
# ── persistence ──────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def _utc_now() -> str:
|
||||
return datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
|
||||
|
||||
|
||||
def _atomic_write(path: Path, payload: str) -> None:
|
||||
"""Write *payload* to *path* atomically: temp file in the same dir, fsync, replace."""
|
||||
parent = path.parent
|
||||
parent.mkdir(parents=True, exist_ok=True)
|
||||
fd, temp_path = tempfile.mkstemp(prefix=f".{path.name}-", suffix=".tmp", 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 _revision_path(directory: Path, revision: int) -> Path:
|
||||
return directory / f"rev-{revision:06d}.json"
|
||||
|
||||
|
||||
def _prune_history(path: Path) -> None:
|
||||
directory = history_dir(path)
|
||||
revisions = list_revisions(path)
|
||||
excess = len(revisions) - _HISTORY_LIMIT
|
||||
for revision in revisions[: max(0, excess)]:
|
||||
_revision_path(directory, revision).unlink(missing_ok=True)
|
||||
|
||||
|
||||
def _archive_current(path: Path) -> int | None:
|
||||
"""Copy the live document into the history dir under its own revision number."""
|
||||
if not path.exists():
|
||||
return None
|
||||
try:
|
||||
existing = json.loads(path.read_text(encoding="utf-8"))
|
||||
revision = int(existing.get("revision", 0))
|
||||
except (json.JSONDecodeError, TypeError, ValueError, AttributeError):
|
||||
# An unreadable live file has no trustworthy revision number to file it
|
||||
# under, so it cannot join the history chain.
|
||||
return None
|
||||
if revision <= 0:
|
||||
return None
|
||||
_atomic_write(
|
||||
_revision_path(history_dir(path), revision),
|
||||
json.dumps(existing, indent=2, sort_keys=True) + "\n",
|
||||
)
|
||||
_prune_history(path)
|
||||
return revision
|
||||
|
||||
|
||||
def list_revisions(path: Path | None = None) -> tuple[int, ...]:
|
||||
"""Revision numbers retained in history for *path*, oldest first."""
|
||||
directory = history_dir(path)
|
||||
if not directory.is_dir():
|
||||
return ()
|
||||
revisions: list[int] = []
|
||||
for entry in directory.glob("rev-*.json"):
|
||||
try:
|
||||
revisions.append(int(entry.stem.split("-", 1)[1]))
|
||||
except (IndexError, ValueError):
|
||||
continue
|
||||
return tuple(sorted(revisions))
|
||||
|
||||
|
||||
def save_registry(
|
||||
registry: WorkerRegistry,
|
||||
path: Path | None = None,
|
||||
*,
|
||||
updated_at: str | None = None,
|
||||
) -> WorkerRegistry:
|
||||
"""Validate, archive the superseded revision, then atomically persist a new one.
|
||||
|
||||
The stored revision is always the previous revision plus one, so a reader
|
||||
can tell two documents apart even when their content is otherwise equal.
|
||||
Returns the registry exactly as persisted.
|
||||
"""
|
||||
target = (path or registry.source_path or default_registry_path()).resolve()
|
||||
|
||||
document = registry_to_document(registry)
|
||||
# Re-validate before writing: a registry assembled in memory has not
|
||||
# necessarily been through the loader.
|
||||
validate_payload(document, source_path=target)
|
||||
|
||||
archived = _archive_current(target)
|
||||
document["revision"] = (archived + 1) if archived is not None else registry.revision
|
||||
document["updated_at"] = updated_at or _utc_now()
|
||||
|
||||
persisted = validate_payload(document, source_path=target)
|
||||
_atomic_write(target, json.dumps(document, indent=2, sort_keys=True) + "\n")
|
||||
return persisted
|
||||
|
||||
|
||||
def rollback_to_revision(revision: int, path: Path | None = None) -> WorkerRegistry:
|
||||
"""Restore a retained *revision* as a new head revision.
|
||||
|
||||
History is append-only: rolling back does not delete the revisions in
|
||||
between, it republishes the chosen one under the next revision number, so a
|
||||
rollback is itself reversible.
|
||||
"""
|
||||
target = (path or default_registry_path()).resolve()
|
||||
snapshot_path = _revision_path(history_dir(target), revision)
|
||||
if not snapshot_path.exists():
|
||||
available = ", ".join(str(item) for item in list_revisions(target)) or "(none)"
|
||||
raise RegistryValidationError(
|
||||
f"revision {revision} is not retained for {target.name}; available: {available}"
|
||||
)
|
||||
|
||||
payload = json.loads(snapshot_path.read_text(encoding="utf-8"))
|
||||
restored = validate_payload(payload, source_path=target)
|
||||
return save_registry(restored, target)
|
||||
Reference in New Issue
Block a user