Compare commits

...
Author SHA1 Message Date
jcwalker3andClaude Opus 4.8 1a38ef95e3 fix(runtime): recognize client identity environment and refresh worker registrations
Two defects left behind by #948 made sanctioned multi-client operation
impossible. They are inseparable: fixing either alone still leaves the
multi-client canary unable to run.

1. Client-identity environment keys were not recognized.

   gitea_mcp_server reads GITEA_MCP_CLIENT, GITEA_MCP_CLIENT_INSTANCE and
   GITEA_MCP_CLIENT_SESSION as the authoritative client-identity inputs for
   worker registration, but none of the three appeared in
   RECOGNIZED_GITEA_ENV_KEYS or matched a recognized prefix. The runtime
   diagnostic scans every peer mcp_server.py process environment and classifies
   any unlisted GITEA_* key as an unsupported override, which is raised as a
   runtime blocker, so gitea_resolve_task_capability returned
   blocker_kind=runtime_reconnect_required with stop_required=true. Because that
   resolver is the mandatory preflight for every author, reviewer and merger
   mutation, setting the very variable #948 requires closed the mutation gate
   for the whole fleet, and reconnecting could not clear it: the variable is
   re-exported from the client's server definition on every launch.

   The three keys are now named individually in the recognized-key set. No
   prefix is added, so an unrecognized GITEA_* override is still refused
   exactly as before.

   A related inconsistency in the same path is also fixed. The diagnostic
   reasons are raised as one RuntimeError, but the preflight re-raise
   recognized only "stale-runtime:", so an "unsupported-env:" reason was
   silently swallowed there while still failing the resolver. Both reason
   families now live in RUNTIME_DIAGNOSTIC_HARD_PREFIXES beside the function
   that produces them, and both propagate identically. This only widens what is
   refused, never what is permitted.

2. WorkerRegistry.heartbeat() had no production caller.

   #948 delivered heartbeat() but only tests called it. The single production
   writer registers once per process behind an attempted-once flag, and
   register() stamps the same timestamp into both started_at and
   last_heartbeat_at. Nothing advanced it afterwards: no lifespan hook, no
   background task, no atexit handler in a process that blocks in mcp.run().
   Since liveness is age against heartbeat_ttl_seconds, that TTL was not a
   liveness window at all but a hard cap on how long any client could stay
   attached; at 900 seconds a healthy, connected, client-managed process became
   session_ownership=unowned with blocker_kind=session_attachment_missing.

   WorkerHeartbeatSupervisor in mcp_worker_identity is the missing caller,
   started from _active_worker_identity() at the moment register() succeeds,
   because that is the only point where identity and fencing_epoch are both
   known. It is a daemon thread rather than an asyncio task or a request-driven
   refresh because renewal must survive an idle session, and because the
   registry performs blocking BEGIN IMMEDIATE sqlite writes that must not run on
   the server's event loop. daemon=True is deliberate: a hard kill takes the
   thread with it, so a dead worker still goes stale on the normal TTL.

   heartbeat_interval_for() returns one third of the TTL, hard-capped at one
   half, so two consecutive beats can be lost without the row expiring and no
   override can produce an interval that outlives the registration it renews.
   heartbeat() gains optional keyword-only expectations (session, generation,
   client name, pid); each supplied one must match the recorded row or the
   renewal is refused with the existing BLOCKER_FENCED literal rather than a new
   blocker_kind, since consumers switch on that value. Omitting them preserves
   the pre-existing behavior exactly. Client names are compared normalized, so
   several namespaces of one application stay one client while separate
   applications stay distinct.

   A terminal refusal stops the supervisor permanently and records why, so a
   fenced session can never beat its way back into ownership. A transient
   failure is counted and beating continues. An atexit hook stops it on orderly
   shutdown. status() is surfaced read-only as worker_heartbeat on
   gitea_get_runtime_context so a stopped heartbeat is diagnosable before the
   TTL turns it into session_attachment_missing; it grants nothing.

   claim_generation() still has no production caller. It bumps fencing_epoch,
   which would fence the supervisor's cached epoch, and the strict refusal is
   left in place deliberately: auto-re-adopting a bumped epoch would defeat
   fencing.

No lock or lease TTL is changed, including the author issue-lock TTL, and no
mutation refused today becomes permitted.

Tests: tests/test_issue_975_client_identity_heartbeat.py adds 40 focused tests
covering all 13 acceptance criteria. Every TTL assertion uses an injected
clock; no test waits for a real TTL. The thread-loop tests use a
millisecond-scale interval with bounded polling.

Focused: 40 passed, 15 subtests passed.
Full suite from inside the branches worktree: 28 failed, 6105 passed, 6 skipped,
1105 subtests — an identical failure set to the 324a0c8a baseline measured in a
sibling branches worktree (28 failed, 6065 passed, 1090 subtests). Zero new
failures; the delta is exactly the added tests.

Closes #975

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-29 20:06:10 -05:00
sysadmin 324a0c8a8d Merge pull request 'fix(runtime): support validated cross-repository canonical roots (#973)' (#974) from fix/issue-973-cross-repo-canonical-roots into master 2026-07-29 17:45:44 -05:00
sysadminandClaude Opus 4.8 34968475c7 fix(runtime): reject unsupported repository-authority modes (#973 B10)
assess_canonical_repository_root declared a public `mode` keyword defaulting
to "validation" and documented exactly two supported values, but never checked
the argument against an allowlist. Both dispatch points were permissive:

* the configured-root path tested `mode == "derivation"` and routed every other
  value into a catch-all `else`, so an unsupported mode silently received
  validation semantics; and
* the single-repository default path tested `mode == "validation"`, so an
  unsupported mode skipped the identity comparison entirely and was strictly
  weaker than validation, not an alias of it.

Measured at the previous head: mode="invalid_mode" with require_binding=True
and matching expected/observed identities returned proven=True, block=False
with no reasons; on the unconfigured path an unsupported mode returned
proven=True where mode="validation" returned proven=False for identical inputs.
Empty string and None behaved like any other unsupported value, and no
assessment ever emitted a mode-specific rejection.

Add an explicit two-value allowlist (SUPPORTED_MODES) and refuse every other
explicitly supplied value — unknown strings, misspellings, case and whitespace
variants, the empty string, None, and non-strings — as the first act of the
function, before any candidate-root existence check, path or symlink
resolution, git top-level discovery, remote-URL or repository-identity
discovery, and before any expected-versus-observed comparison or
validation/derivation behaviour. The refusal reports proven=False,
block=True, a mode-specific reason naming the offending value, and
reason_code=DENY_UNKNOWN_MODE, following the existing
webui.sanctioned_restart.DENY_UNKNOWN_MODE convention. No repository identity
is resolved through a refused mode: resolved_slug and canonical_repo_root are
both None.

Omission continues to select validation, so the documented default is
unchanged. Unsupported modes are refused rather than normalized onto a
supported mode. No mode is exposed through MCP request parameters, environment
variables, repository configuration, session input, or any public reviewer,
merger, issue, PR or lease API; the only production call sites remain an
omitted mode (validation) and the hardcoded "derivation" literal.

resolve_namespace_mutation_context keeps the install checkout as the canonical
root when an assessment resolves none, so a refused mode cannot bind a root
derived through an undefined mode while roots_aligned and the carried
assessment stay fail-closed.

Regressions in tests/test_issue_973_b10_mode_contract.py cover invalid modes
with missing, matching and conflicting identities, empty string, explicit None,
representative non-strings, misspellings and whitespace variants, omission
defaulting to validation, explicit validation and derivation behaviour, the
single-repository default path differential, rejection ordering (both spied and
mock-free), the absence of any request/environment/configuration injection
surface, and fail-closed reviewer, merger, mutation-context and final
mutation-authorization behaviour through the production paths.

Closes #973

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-29 17:52:32 -04:00
sysadmin a6c7d1491e fix(runtime): fix identity derivation disarming in canonical root guard (Closes #973) 2026-07-29 16:02:32 -04:00
sysadmin 9d96cf4cfa fix(runtime): derive expected repository identity independently of configured root (Closes #973) 2026-07-29 14:44:50 -04:00
sysadmin 7978008709 Fix review 646 blockers B1-B7 for cross-repo canonical roots (#973) 2026-07-29 14:01:17 -04:00
sysadmin 6a53308473 fix(runtime): support validated cross-repository canonical roots (#973) 2026-07-29 13:19:29 -04:00
sysadmin 626be8b178 Merge pull request 'fix(reconcile): safely resolve worktree bindings whose paths are missing (Closes #970)' (#972) from fix/issue-970-safely-resolve-missing-worktrees into master 2026-07-29 10:22:20 -05:00
jcwalker3andClaude Opus 4.8 c763161702 fix(reconcile): server-enforced, revalidated missing-worktree cleanup (#970 review 644 B1-B5)
Addresses the five blocking findings of review 644 on PR #972.

B1 — live ownership and status revalidation. resolve_missing_worktree_binding
now re-reads the authoritative lease, session, checkpoint, and issue-lock rows
immediately before mutating and diffs them against the snapshot the audit
recorded (binding path and identity, lease id/status/session/owner pid,
checkpoint path/status, live-session evidence, trustworthy ownership evidence,
issue-lock state). Any drift fails closed without mutation, and the binding is
reclassified from the live values rather than the audit snapshot. A candidate
carrying no audited snapshot is refused rather than trusted.

B2 — server-enforced cleanup authorization. Apply mode no longer accepts a
client-supplied operator_authorized boolean; it is rejected outright at the MCP
tool and in the module (#709 F1 / review 434). Authorization is now the
project's own reconciliation cleanup gate, required at both the task-capability
boundary (new reconciler-only reconcile_missing_worktree_bindings capability,
gitea.branch.delete, role-exclusive) and the production mutation boundary
(an authorized audit_reconciliation_mode cleanup phase, re-checked at the point
of mutation so a forged authorization mapping cannot stand in for the gate).
Dry-run remains available to any gitea.read profile and stays non-mutating.
Existing role, repository, parity, and provenance gates are unchanged.

B3 — expected-path compare-and-swap. retire_session_checkpoint_worktree_path
now requires expected_path and performs a guarded update keyed on the stored
path, refusing without mutation when the stored path was moved, replaced, or
concurrently changed, when the row is unknown, or when a selector matches more
than one checkpoint. retire_lease_worktree_path gains the same treatment plus
optional status/session/owner-pid compare-and-swap, and its UPDATE is keyed on
the audited path. Both report an idempotent already_retired outcome instead of
falsely reporting a retirement.

B4 — live-session and issue-lock evidence. session_active is now derived from
the control-plane sessions table instead of never being set, along two axes:
genuine liveness (recorded active, PID not dead, heartbeat fresh — the rule
reused from restart_coordinator) and weaker but still trustworthy recorded
ownership. A non-terminal lease now protects its binding regardless of whether
the recorded PID is alive, so dead-PID evidence alone can no longer retire a
lease the control plane still holds. The previously unused issue_lock_store is
now read: a live durable issue lock binding the path or branch blocks cleanup,
and locks whose own paths are missing are reported for release through their
own lifecycle rather than retired here.

B5 — adversarial regression coverage. The suite now drives the registered MCP
tools through mcp_server, the real ControlPlaneDB, and the real cleanup gate,
covering lease status/ownership/session/path drift, expected-path mismatch,
concurrent recreation, an unauthorized caller submitting operator_authorized,
wrong profile and missing capability, live-session and trustworthy-owner
evidence, conflicting issue locks, non-mutating dry-run, exact-binding-only
retirement, preservation of unrelated worktrees and git metadata, idempotent
re-execution, and worktrees-dimension resolution.

All original #970 acceptance criteria are preserved, including the distinctions
between deleted paths, moved paths, unavailable hosts or mounts, transient
filesystem failures, live ownership, and concurrent recreation.

Tests: focused #970 suite 47 passed. Adjacent suites (capability role
invariants, audit reconciliation mode, control plane DB, lease lifecycle,
reconciler cleanup integration, delete-branch capability, restart coordinator,
bootstrap lock contract) 252 passed / 93 subtests. Full suite 28 failed /
6000 passed, an exact match of the pre-change baseline's 28 failing test ids
at 3f584352 (28 failed / 5960 passed).

Closes #970

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-29 06:31:52 -05:00
sysadmin 3f584352df fix(reconcile): safely resolve worktree bindings whose paths are missing (Closes #970) 2026-07-29 05:43:52 -04:00
17 changed files with 5838 additions and 68 deletions
+136 -21
View File
@@ -37,6 +37,48 @@ CANONICAL_ROOT_ENV = "GITEA_CANONICAL_REPOSITORY_ROOT"
# Candidate git remote names probed when deriving repository identity.
_IDENTITY_REMOTE_CANDIDATES = ("prgs", "origin", "dadeschools", "mdcps")
# Repository-authority modes (#973 B10). Exactly two values are supported.
# ``mode`` selects how repository authority is established, so an unrecognised
# value must never be normalised onto one of these: aliasing a trusted mode is
# precisely the defect. Omitting the argument keeps the documented safe default,
# ``validation``.
MODE_VALIDATION = "validation"
MODE_DERIVATION = "derivation"
SUPPORTED_MODES: tuple[str, ...] = (MODE_VALIDATION, MODE_DERIVATION)
# Reason code emitted when an unsupported mode is refused. Mirrors the existing
# ``webui.sanctioned_restart.DENY_UNKNOWN_MODE`` convention so callers and tests
# can assert the refusal cause rather than string-matching prose.
DENY_UNKNOWN_MODE = "unknown_mode"
def unsupported_mode_reason(mode: object) -> str | None:
"""Precise rejection reason for *mode*, or None when *mode* is supported.
Only the two documented string values are accepted, compared exactly — no
stripping, no case folding — so misspellings and whitespace variants are
refused rather than coerced. An empty string, ``None``, and any non-string
are all *explicitly supplied* unsupported values and are refused on the same
footing; none of them is normalised to a supported mode. Omitting the
argument entirely never reaches here with an unsupported value because the
parameter default is ``"validation"``.
"""
if isinstance(mode, str) and mode in SUPPORTED_MODES:
return None
supported = ", ".join(repr(m) for m in SUPPORTED_MODES)
if not isinstance(mode, str):
return (
f"unsupported repository-authority mode {mode!r} of type "
f"{type(mode).__name__}: only {supported} are supported; the mode "
"is refused before any repository assessment and no repository "
"identity was resolved through it (fail closed)"
)
return (
f"unsupported repository-authority mode {mode!r}: only {supported} are "
"supported; the mode is refused before any repository assessment and no "
"repository identity was resolved through it (fail closed)"
)
def configured_canonical_root(
profile: Mapping | None,
@@ -132,6 +174,7 @@ def assess_canonical_repository_root(
process_project_root: str,
remote: str | None = None,
require_binding: bool = False,
mode: str = "validation",
) -> dict:
"""Validate the canonical repository root binding, failing closed on forgery.
@@ -140,15 +183,41 @@ def assess_canonical_repository_root(
``configured`` (whether a cross-repo binding was declared),
``resolved_slug`` and ``source``.
Without a configured binding the single-repo default is preserved: the
canonical root is derived from *process_project_root* and never blocks
(unless *require_binding* explicitly demands one).
*mode* accepts exactly the two values in :data:`SUPPORTED_MODES`:
- ``"validation"`` (default): an independently trusted expected repository
slug is known (or required). Candidate root's observed identity must match.
Unprovable or missing expected identity fails closed when *require_binding* is True.
- ``"derivation"``: caller is deriving the canonical repository identity.
No expected slug exists yet by design. Derivation succeeds if the configured
root exists, is a git repository, and carries a resolvable git remote identity.
With a configured binding the path must exist, be a git repository, and —
when *expected_slug* is known — carry a matching repository identity. A
mismatched or (when *require_binding*) unprovable identity is a forged or
conflicting binding and fails closed.
Every other explicitly supplied value — unknown strings, misspellings, the
empty string, ``None``, and non-strings — is refused with ``proven`` False,
``block`` True, and ``reason_code`` :data:`DENY_UNKNOWN_MODE` (#973 B10).
Omitting *mode* entirely keeps the documented ``"validation"`` default.
"""
# #973 B10: refuse an unsupported repository-authority mode as the very first
# act, before any candidate-root existence check, path or symlink resolution,
# git top-level discovery, remote-URL or repository-identity discovery, and
# before any expected-versus-observed comparison or validation/derivation
# behaviour. An unsupported mode previously fell into the catch-all ``else``
# on the configured-root path (silently receiving validation semantics) and
# skipped the identity comparison entirely on the single-repository default
# path (strictly weaker than validation), so matching identities could return
# ``proven`` True. No repository identity may be resolved through a mode the
# contract does not define.
mode_reason = unsupported_mode_reason(mode)
if mode_reason is not None:
return _assessment(
proven=False,
reasons=[mode_reason],
configured=bool((configured_value or "").strip()),
canonical_repo_root=None,
resolved_slug=None,
source=source,
reason_code=DENY_UNKNOWN_MODE,
)
process_root = os.path.realpath(process_project_root)
declared = (configured_value or "").strip()
@@ -168,12 +237,31 @@ def assess_canonical_repository_root(
)
# Single-repo default: canonical root follows the install checkout.
derived = resolve_repo_toplevel(process_root) or process_root
resolved_slug = repository_identity_slug(derived, remote=remote)
reasons: list[str] = []
# #973 B10: the allowlist above guarantees *mode* is one of the two
# supported values here, so an unsupported value can no longer skip this
# identity comparison and end up strictly weaker than validation.
if mode == MODE_VALIDATION and expected_slug:
expected = expected_slug.strip()
if resolved_slug and resolved_slug.lower() != expected.lower():
reasons.append(
f"canonical repository root identity mismatch: '{derived}' resolves "
f"to repository '{resolved_slug}' but expected repository identity "
f"is '{expected}' (forged or conflicting binding, fail closed)"
)
elif not resolved_slug and require_binding:
reasons.append(
f"canonical repository root '{derived}' has no resolvable git "
f"remote identity to confirm authorization for '{expected}' "
"(fail closed)"
)
return _assessment(
proven=True,
reasons=[],
proven=not reasons,
reasons=reasons,
configured=False,
canonical_repo_root=derived,
resolved_slug=None,
resolved_slug=resolved_slug,
source=None,
)
@@ -207,19 +295,36 @@ def assess_canonical_repository_root(
resolved_slug = repository_identity_slug(toplevel, remote=remote)
reasons: list[str] = []
expected = (expected_slug or "").strip() or None
if expected:
if resolved_slug and resolved_slug.lower() != expected.lower():
if mode == MODE_DERIVATION:
if not resolved_slug:
reasons.append(
f"canonical repository root identity mismatch: '{toplevel}' resolves "
f"to repository '{resolved_slug}' but the session is authorized for "
f"'{expected}' (forged or conflicting binding, fail closed)"
f"configured canonical repository root '{toplevel}' has no resolvable "
"git remote identity (fail closed)"
)
elif not resolved_slug and require_binding:
else:
# #973 B10: MODE_VALIDATION only. This arm is no longer a catch-all — the
# allowlist above admits no third value, so an unsupported mode can no
# longer silently receive validation semantics here.
expected = (expected_slug or "").strip() or None
if expected:
if resolved_slug and resolved_slug.lower() != expected.lower():
reasons.append(
f"canonical repository root identity mismatch: '{toplevel}' resolves "
f"to repository '{resolved_slug}' but the session is authorized for "
f"'{expected}' (forged or conflicting binding, fail closed)"
)
elif not resolved_slug and require_binding:
reasons.append(
f"canonical repository root '{toplevel}' has no resolvable git "
f"remote identity to confirm authorization for '{expected}' "
"(fail closed)"
)
elif require_binding:
reasons.append(
f"canonical repository root '{toplevel}' has no resolvable git "
f"remote identity to confirm authorization for '{expected}' "
"(fail closed)"
f"canonical repository root '{toplevel}' has configured value "
f"'{configured_value}' but authoritative expected repository identity "
"is unprovable or missing (fail closed)"
)
return _assessment(
@@ -252,10 +357,19 @@ def _assessment(
proven: bool,
reasons: list[str],
configured: bool,
canonical_repo_root: str,
canonical_repo_root: str | None,
resolved_slug: str | None,
source: str | None,
reason_code: str | None = None,
) -> dict:
"""Build the assessment payload.
``canonical_repo_root`` is None only when the assessment refused to resolve
one at all — today exactly the unsupported-mode refusal (#973 B10), which
must not derive a trusted repository identity through an undefined mode.
``reason_code`` is a machine-checkable refusal cause; None for every
ordinary (non-coded) outcome.
"""
return {
"proven": proven,
"block": not proven,
@@ -264,4 +378,5 @@ def _assessment(
"canonical_repo_root": canonical_repo_root,
"resolved_slug": resolved_slug,
"source": source,
"reason_code": reason_code,
}
+298
View File
@@ -280,6 +280,22 @@ def _ts(dt: datetime | None = None) -> str:
return value.astimezone(timezone.utc).replace(microsecond=0).isoformat().replace("+00:00", "Z")
def _realpath_or_raw(value: str | None) -> str:
"""Normalize a filesystem path for compare-and-swap equality (#970).
Symlinks and ``..`` segments must not make two spellings of the same path
look different, but an unresolvable path must still compare as itself
rather than collapsing to empty — an empty result means "no path given".
"""
text = (value or "").strip()
if not text:
return ""
try:
return os.path.realpath(os.path.abspath(text))
except Exception:
return text
def _parse_ts(value: str | None) -> datetime | None:
if not value:
return None
@@ -1913,6 +1929,168 @@ class ControlPlaneDB:
),
)
def retire_lease_worktree_path(
self,
lease_id: str,
*,
expected_path: str | None = None,
expected_status: str | None = None,
expected_session_id: str | None = None,
expected_owner_pid: int | None = None,
reason: str = "missing_worktree_path_retired",
) -> dict[str, Any]:
"""Retire a missing worktree_path binding from a control-plane lease (#970).
Clears worktree_path on the lease row, updates provenance_json with
durable retirement audit proof, and writes a worktree_binding_retired
event.
The update is a compare-and-swap (#970 review 644 B1/B3): the caller
states the exact path it audited and, when known, the lease status,
owning session, and owner pid it classified against. Every stated value
must still match the stored row, and the ``UPDATE`` itself is keyed on
the stored ``worktree_path``, so a concurrent writer that moved or
replaced the binding between audit and apply loses the race instead of
having its value silently overwritten. A mismatch raises and mutates
nothing.
``expected_path`` is mandatory: a retirement that does not name the path
it intends to clear cannot be safe against concurrent recreation.
"""
now_s = _ts()
expected_norm = _realpath_or_raw(expected_path)
if not expected_norm:
raise ControlPlaneError(
f"cannot retire lease {lease_id} worktree_path: expected_path is "
"required for compare-and-swap retirement (fail closed)"
)
with self._tx() as conn:
cols = self._lease_columns(conn)
row = conn.execute(
"SELECT * FROM leases WHERE lease_id = ?",
(lease_id,),
).fetchone()
if not row:
raise ControlPlaneError(f"unknown lease_id {lease_id}")
record = dict(row)
current_wt = (record.get("worktree_path") or "").strip()
if not current_wt:
# Idempotent: the binding this caller audited is already gone.
return {
"lease_id": lease_id,
"retired": False,
"already_retired": True,
"prior_worktree_path": "",
"expected_worktree_path": expected_path,
"reason": reason,
"compare_and_swap": {
"matched": True,
"outcome": "already_retired",
},
}
if _realpath_or_raw(current_wt) != expected_norm:
raise ControlPlaneError(
f"cannot retire lease {lease_id} worktree_path: expected "
f"'{expected_path}' does not match current '{current_wt}' "
"(fail closed)"
)
for field, expected_value in (
("status", expected_status),
("session_id", expected_session_id),
):
if expected_value is None:
continue
current_value = record.get(field)
if str(current_value or "").strip() != str(expected_value).strip():
raise ControlPlaneError(
f"cannot retire lease {lease_id} worktree_path: lease "
f"{field} changed since audit (expected "
f"'{expected_value}', found '{current_value}'); fail closed"
)
if expected_owner_pid is not None:
current_pid = record.get("owner_pid")
if current_pid is not None and int(current_pid) != int(expected_owner_pid):
raise ControlPlaneError(
f"cannot retire lease {lease_id} worktree_path: lease "
f"owner_pid changed since audit (expected "
f"{expected_owner_pid}, found {current_pid}); fail closed"
)
# Parse and update provenance_json
raw_prov = record.get("provenance_json") or "{}"
try:
prov = json.loads(raw_prov) if isinstance(raw_prov, str) else dict(raw_prov)
except Exception:
prov = {}
if not isinstance(prov, dict):
prov = {}
prior_path = current_wt
prov.update({
"worktree_path_retired": True,
"retired_worktree_path": prior_path,
"retired_at": now_s,
"retirement_reason": reason,
"retired_from_status": record.get("status"),
"retired_from_session_id": record.get("session_id"),
"worktree_path": "",
})
prov_json = json.dumps(prov)
if "worktree_path" in cols:
# CAS: keyed on the exact stored path this caller audited.
cur = conn.execute(
"UPDATE leases SET worktree_path = '', provenance_json = ? "
"WHERE lease_id = ? AND worktree_path = ?",
(prov_json, lease_id, record.get("worktree_path")),
)
if cur.rowcount != 1:
raise ControlPlaneError(
f"cannot retire lease {lease_id} worktree_path: "
"compare-and-swap matched no row (concurrent change); "
"fail closed"
)
else:
conn.execute(
"UPDATE leases SET provenance_json = ? WHERE lease_id = ?",
(prov_json, lease_id),
)
conn.execute(
"""
INSERT INTO events(work_item_id, event_type, message, created_at)
VALUES (?, 'worktree_binding_retired', ?, ?)
""",
(
record["work_item_id"],
f"lease {lease_id} worktree_path '{prior_path}' retired: {reason}",
now_s,
),
)
return {
"lease_id": lease_id,
"retired": True,
"already_retired": False,
"prior_worktree_path": prior_path,
"expected_worktree_path": expected_path,
"retired_at": now_s,
"reason": reason,
"compare_and_swap": {
"matched": True,
"outcome": "retired",
"expected_status": expected_status,
"expected_session_id": expected_session_id,
"expected_owner_pid": expected_owner_pid,
},
}
def abandon_lease(
self,
*,
@@ -3051,3 +3229,123 @@ class ControlPlaneDB:
"live_lease_id": None if live_lease_id is None else str(live_lease_id),
"reconcile_action": "reconcile_required" if stale else "safe_to_resume",
}
def retire_session_checkpoint_worktree_path(
self,
session_id: str,
*,
checkpoint_id: str | None = None,
expected_path: str | None = None,
expected_status: str | None = None,
reason: str = "missing_worktree_path_retired",
) -> dict[str, Any]:
"""Retire a missing worktree_path from session_checkpoints (#970).
Compare-and-swap, mirroring :meth:`retire_lease_worktree_path` (#970
review 644 B3). ``expected_path`` names the exact stored path the caller
audited; the guarded ``UPDATE`` is keyed on that stored value, so a
checkpoint whose path was moved, replaced, or concurrently rewritten
after the audit is refused without mutation rather than blindly cleared.
Exactly one checkpoint row is targeted: by ``checkpoint_id`` when given,
otherwise by ``session_id``, which must identify a single row.
"""
now_s = _ts()
expected_norm = _realpath_or_raw(expected_path)
if not expected_norm:
raise ControlPlaneError(
"cannot retire session checkpoint worktree_path: expected_path "
"is required for compare-and-swap retirement (fail closed)"
)
if not checkpoint_id and not (session_id or "").strip():
raise ControlPlaneError(
"cannot retire session checkpoint worktree_path: checkpoint_id "
"or session_id is required (fail closed)"
)
with self._tx() as conn:
if checkpoint_id:
selector_sql = "SELECT * FROM session_checkpoints WHERE checkpoint_id = ?"
selector_params: tuple[Any, ...] = (checkpoint_id,)
selector_desc = f"checkpoint_id '{checkpoint_id}'"
else:
selector_sql = "SELECT * FROM session_checkpoints WHERE session_id = ?"
selector_params = (session_id,)
selector_desc = f"session_id '{session_id}'"
rows = [dict(r) for r in conn.execute(selector_sql, selector_params).fetchall()]
if not rows:
raise ControlPlaneError(
f"cannot retire session checkpoint worktree_path: no "
f"checkpoint matches {selector_desc} (fail closed)"
)
if len(rows) > 1:
raise ControlPlaneError(
f"cannot retire session checkpoint worktree_path: "
f"{selector_desc} matches {len(rows)} checkpoints; supply an "
"exact checkpoint_id (fail closed)"
)
record = rows[0]
target_checkpoint_id = record.get("checkpoint_id")
current_wt = (record.get("worktree_path") or "").strip()
if not current_wt:
# Idempotent: the binding this caller audited is already gone.
return {
"session_id": session_id,
"checkpoint_id": target_checkpoint_id,
"retired": False,
"already_retired": True,
"prior_worktree_path": "",
"expected_worktree_path": expected_path,
"reason": reason,
"compare_and_swap": {
"matched": True,
"outcome": "already_retired",
},
}
if _realpath_or_raw(current_wt) != expected_norm:
raise ControlPlaneError(
f"cannot retire session checkpoint worktree_path for "
f"{selector_desc}: expected '{expected_path}' does not match "
f"current '{current_wt}' (fail closed)"
)
if expected_status is not None:
current_status = record.get("status")
if str(current_status or "").strip() != str(expected_status).strip():
raise ControlPlaneError(
f"cannot retire session checkpoint worktree_path for "
f"{selector_desc}: status changed since audit (expected "
f"'{expected_status}', found '{current_status}'); fail closed"
)
cur = conn.execute(
"UPDATE session_checkpoints SET worktree_path = '', updated_at = ? "
"WHERE checkpoint_id = ? AND worktree_path = ?",
(now_s, target_checkpoint_id, record.get("worktree_path")),
)
if cur.rowcount != 1:
raise ControlPlaneError(
f"cannot retire session checkpoint worktree_path for "
f"{selector_desc}: compare-and-swap matched no row "
"(concurrent change); fail closed"
)
return {
"session_id": session_id,
"checkpoint_id": target_checkpoint_id,
"retired": True,
"already_retired": False,
"prior_worktree_path": current_wt,
"expected_worktree_path": expected_path,
"retired_at": now_s,
"reason": reason,
"compare_and_swap": {
"matched": True,
"outcome": "retired",
"expected_status": expected_status,
},
}
+5
View File
@@ -65,6 +65,7 @@ that gates each call, not which tools exist.
- `gitea_assess_work_issue_duplicate`
- `gitea_assess_worktree_cleanup_integrity`
- `gitea_audit_config`
- `gitea_audit_missing_worktree_bindings`
- `gitea_audit_runtime_recovery_contamination`
- `gitea_audit_stable_branch_contamination`
- `gitea_audit_worktree_cleanup`
@@ -126,16 +127,20 @@ that gates each call, not which tools exist.
- `gitea_post_heartbeat`
- `gitea_publish_unpublished_issue_branch`
- `gitea_quarantine_contaminated_review`
- `gitea_rebind_dirty_same_claimant_author_session`
- `gitea_reclaim_expired_workflow_lease`
- `gitea_reconcile_after_restart`
- `gitea_reconcile_already_landed_pr`
- `gitea_reconcile_issue_claims`
- `gitea_reconcile_merged_cleanups`
- `gitea_reconcile_missing_worktree_bindings`
- `gitea_reconcile_superseded_by_merged_pr`
- `gitea_record_daemon_process_kill_attempt`
- `gitea_record_irrecoverable_decision_lock_provenance`
- `gitea_record_pre_review_command`
- `gitea_record_shell_spawn_outcome`
- `gitea_record_stable_branch_push_attempt`
- `gitea_recover_dirty_orphaned_issue_worktree`
- `gitea_recover_incomplete_bootstrap_lock`
- `gitea_release_merger_pr_lease`
- `gitea_release_reviewer_pr_lease`
+14
View File
@@ -1180,6 +1180,10 @@ RECOGNIZED_GITEA_ENV_KEYS = frozenset({
"GITEA_SERVER_PROVENANCE",
"GITEA_AUTHOR_WORKTREE",
"GITEA_ACTIVE_WORKTREE",
"GITEA_REVIEWER_WORKTREE",
"GITEA_MERGER_WORKTREE",
"GITEA_CANONICAL_REPOSITORY_ROOT",
"GITEA_MCP_SESSION_STATE_TTL_HOURS",
"GITEA_DISABLE_KEYCHAIN",
"GITEA_CONTROL_PLANE_DB",
"GITEA_DB_PATH",
@@ -1189,6 +1193,16 @@ RECOGNIZED_GITEA_ENV_KEYS = frozenset({
"GITEA_IRRECOVERABLE_HMAC_SECRET",
"GITEA_FORCE_MCP_RUNTIME_CHECK",
"GITEA_FORCE_CLIENT_MANAGED",
# #975: the client-identity inputs the server actually consumes at startup
# (CLIENT_NAME_ENV / CLIENT_INSTANCE_ENV / CLIENT_SESSION_ENV in
# gitea_mcp_server). Production read them while this allowlist omitted them,
# so the peer-env scan classified them as unsupported overrides and the
# capability resolver refused every mutation fleet-wide. Named individually
# on purpose: no prefix is added, so an unrecognised GITEA_* override is
# still refused exactly as it was before.
"GITEA_MCP_CLIENT",
"GITEA_MCP_CLIENT_INSTANCE",
"GITEA_MCP_CLIENT_SESSION",
})
RECOGNIZED_GITEA_ENV_PREFIXES = (
+350 -17
View File
@@ -500,10 +500,66 @@ def _resolve_preflight_workspace_path(worktree_path: str | None = None) -> str:
return workspace
def _resolve_namespace_mutation_context(worktree_path: str | None = None) -> dict:
def _process_root_git_remote_url(remote_name: str) -> str | None:
"""Best-effort local ``git remote get-url`` strictly inside ``PROJECT_ROOT``.
#973 (B8): Must derive expected repository identity from an authority
independent of the candidate configured canonical root.
"""
try:
proc = subprocess.run(
["git", "remote", "get-url", remote_name],
capture_output=True,
text=True,
cwd=PROJECT_ROOT,
)
if proc.returncode != 0:
return None
url = (proc.stdout or "").strip()
return url or None
except Exception:
return None
def _resolve_expected_repository_slug(
remote: str | None = None,
org: str | None = None,
repo: str | None = None,
) -> str | None:
"""Resolve expected repository slug from session context, parameters, or process root.
Must derive expected repository identity ONLY from trusted sources independent
of the candidate configured canonical root (#973 B8).
Bound session context takes precedence over request parameters so request-supplied
coordinates cannot replace a bound canonical root (#741 / #973 B9).
"""
bound = session_ctx.get_session_context() or {}
b_org = bound.get("org")
b_repo = bound.get("repository")
if b_org and b_repo:
return session_ctx.format_repository_slug(b_org, b_repo)
if org and repo:
return session_ctx.format_repository_slug(org, repo)
eff_remote = remote or bound.get("remote") or _effective_remote()
parsed = remote_repo_guard.parse_org_repo_from_remote_url(
_process_root_git_remote_url(eff_remote)
)
if parsed:
return session_ctx.format_repository_slug(parsed[0], parsed[1])
return None
def _resolve_namespace_mutation_context(
worktree_path: str | None = None,
remote: str | None = None,
) -> dict:
"""Canonical namespace workspace + repository root for guards (#460/#510/#706/#618)."""
role = _effective_workspace_role()
configured_root, _source = _configured_canonical_root()
bound = session_ctx.get_session_context() or {}
eff_remote = remote or bound.get("remote") or _effective_remote()
expected_slug = _resolve_expected_repository_slug(eff_remote)
return nwb.resolve_namespace_mutation_context(
role_kind=role,
worktree_path=worktree_path,
@@ -516,9 +572,12 @@ def _resolve_namespace_mutation_context(worktree_path: str | None = None) -> dic
),
profile_name=get_profile().get("profile_name"),
configured_canonical_root=configured_root,
expected_slug=expected_slug,
remote=eff_remote,
)
def _resolve_author_mutation_context(worktree_path: str | None = None) -> dict:
"""Backward-compatible alias for namespace workspace context."""
return _resolve_namespace_mutation_context(worktree_path)
@@ -620,6 +679,9 @@ def _preflight_workspace_details(worktree_path: str | None, dirty_files: list[st
inspected_root = _get_git_root(workspace)
process_root = ctx["process_project_root"]
canonical_root = ctx["canonical_repo_root"]
crr_assessment = ctx.get("canonical_root_assessment") or {}
resolved_slug = crr_assessment.get("resolved_slug")
expected_slug = ctx.get("expected_slug")
active_root = os.path.realpath(inspected_root or workspace)
if active_root == canonical_root:
dirty_scope = "control checkout"
@@ -649,6 +711,10 @@ def _preflight_workspace_details(worktree_path: str | None, dirty_files: list[st
"workspace_healthy": not bool(
ctx.get("bound_worktree_missing") or ctx.get("author_worktree_block")
),
"canonical_root_assessment": crr_assessment,
"expected_repository_slug": expected_slug,
"observed_repository_identity": resolved_slug,
"worktree_registration_result": ctx.get("in_git_worktree_list"),
}
if not ctx["roots_aligned"]:
details["workspace_root_mismatch"] = (
@@ -658,6 +724,7 @@ def _preflight_workspace_details(worktree_path: str | None, dirty_files: list[st
return details
def _format_preflight_workspace_details(details: dict) -> str:
parts = [
f"MCP server process root: {details.get('mcp_server_process_root')}",
@@ -930,10 +997,7 @@ def _enforce_canonical_repository_root(
if not configured_value:
return
bound = session_ctx.get_session_context() or {}
expected_slug = session_ctx.format_repository_slug(
bound.get("org"), bound.get("repository")
)
expected_slug = _resolve_expected_repository_slug(remote)
assessment = crr.assess_canonical_repository_root(
configured_value=configured_value,
source=source,
@@ -1852,8 +1916,16 @@ def _verify_role_mutation_workspace(
if runtime_reasons:
raise RuntimeError("; ".join(runtime_reasons))
except Exception as exc:
if "stale-runtime:" in str(exc):
raise RuntimeError(str(exc))
# #975: ``_check_mcp_runtimes_diagnostics`` raises every reason it
# produces through this one RuntimeError, but only ``stale-runtime:``
# was re-raised here — an ``unsupported-env:`` reason was swallowed
# while still failing the capability resolver, so this preflight and
# the resolver disagreed about the identical diagnostic. Both
# authoritative prefixes now propagate the same way. This can only ever
# widen what is refused, never widen what is permitted.
message = str(exc)
if any(prefix in message for prefix in RUNTIME_DIAGNOSTIC_HARD_PREFIXES):
raise RuntimeError(message)
pass
role = _effective_workspace_role()
@@ -1866,6 +1938,9 @@ def _verify_role_mutation_workspace(
# back to the install checkout and validated Gitea-Tools/branches/ instead
# of the target repository the namespace is actually bound to.
_configured_root, _configured_source = _configured_canonical_root()
bound = session_ctx.get_session_context() or {}
eff_remote = remote or bound.get("remote") or _effective_remote()
expected_slug = _resolve_expected_repository_slug(eff_remote, org=org, repo=repo)
assessment = nwb.assess_namespace_mutation_workspace(
role_kind=role,
worktree_path=worktree_path,
@@ -1880,6 +1955,8 @@ def _verify_role_mutation_workspace(
profile_name=get_profile().get("profile_name"),
current_branch=git_state.get("current_branch"),
configured_canonical_root=_configured_root,
expected_slug=expected_slug,
remote=eff_remote,
)
if assessment["block"]:
raise RuntimeError(
@@ -2184,6 +2261,7 @@ def _canonical_repository_slug(
process_project_root=PROJECT_ROOT,
remote=remote,
require_binding=True,
mode="derivation",
)
slug = assessment.get("resolved_slug")
if assessment.get("block") or not slug:
@@ -3414,7 +3492,7 @@ def cleanup_in_progress_for_pr(
# ── Helpers ───────────────────────────────────────────────────────────────────
def _effective_remote(remote: str) -> str:
def _effective_remote(remote: str = "dadeschools") -> str:
"""If remote is the default ('dadeschools') but the active profile base_url maps to a known remote, use that remote instead."""
try:
profile = get_profile()
@@ -13547,6 +13625,162 @@ def gitea_scan_already_landed_open_prs(
}
@mcp.tool()
def gitea_audit_missing_worktree_bindings(
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
) -> dict:
"""Read-only: audit and classify worktree bindings whose paths are missing on disk (#970).
Correlates missing-path bindings from control-plane leases, session checkpoints,
and issue locks with repository, host, branch, issue/PR, session, and lease state.
Distinguishes deleted worktrees from moved paths, host/mount failures, and live
leases/sessions.
Args:
remote: Known instance 'dadeschools' or 'prgs'.
host: Override the Gitea host.
org: Override the owner/organization.
repo: Override the repository name.
Returns:
dict with audit counts, missing binding classifications, and resolution status.
"""
read_block = _profile_operation_gate("gitea.read")
if read_block:
return {
"success": False,
"performed": False,
"reasons": read_block,
"permission_report": _permission_block_report("gitea.read"),
}
import missing_worktree_reconcile
db, _ = _control_plane_db_or_error()
root = _canonical_local_git_root()
h, o, r = _resolve(remote, host, org, repo)
return missing_worktree_reconcile.audit_missing_worktree_bindings(
db,
project_root=root,
remote=remote,
org=o,
repo=r,
host=h,
)
@mcp.tool()
def gitea_reconcile_missing_worktree_bindings(
dry_run: bool = True,
# Deprecated: retained so callers that still pass it get an explicit deny.
operator_authorized: bool = False,
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
) -> dict:
"""Audit and safely resolve (retire) missing worktree path bindings (#970).
Identifies confirmed stale deleted worktree bindings, re-validates them
against a fresh authoritative read immediately before mutation to prevent
recreation and ownership races, and retires only the exact stale bindings
while preserving unrelated worktrees and Git metadata.
Apply mode (``dry_run=False``) is authorized **server-side** (#970 review
644 B2): it requires the reconciler-only ``gitea.branch.delete``
capability, the resolved ``reconcile_missing_worktree_bindings`` task
capability, and an active cleanup phase minted through
``gitea_authorize_reconciliation_cleanup_phase``. A caller-supplied
``operator_authorized`` is never authorization evidence (#709 F1 /
review 434) and is rejected outright. Dry-run stays available to any
``gitea.read`` profile and never mutates.
Args:
dry_run: If True (default), reports planned mutations without modifying state.
operator_authorized: Rejected. Authorization is a server-side artifact.
remote: Known instance 'dadeschools' or 'prgs'.
host: Override the Gitea host.
org: Override the owner/organization.
repo: Override the repository name.
Returns:
dict with before/after audit state, resolutions, and dimension status.
"""
import missing_worktree_reconcile
# Explicitly reject the self-assertable Boolean before anything else, so it
# can never combine with a legitimate gate to authorize a mutation.
if operator_authorized:
return {
"success": False,
"performed": False,
"reason": "operator_authorized_rejected",
"reasons": [missing_worktree_reconcile.OPERATOR_AUTHORIZED_REJECTION],
"exact_next_action": (
"authorize cleanup via gitea_authorize_reconciliation_cleanup_phase "
"from a reconciler profile, then re-run with dry_run=False"
),
}
read_block = _profile_operation_gate("gitea.read")
if read_block:
return {
"success": False,
"performed": False,
"reasons": read_block,
"permission_report": _permission_block_report("gitea.read"),
}
db, db_errs = _control_plane_db_or_error()
root = _canonical_local_git_root()
h, o, r = _resolve(remote, host, org, repo)
profile = get_profile()
role = (profile.get("role") or profile.get("role_kind") or "").strip().lower()
cleanup_authorization = None
if not dry_run:
cleanup_task = missing_worktree_reconcile.CLEANUP_TASK
required_permission = task_capability_map.required_permission(cleanup_task)
capability_blockers = _profile_operation_gate(required_permission)
role_matches = role == task_capability_map.required_role(cleanup_task)
cleanup_authorization = missing_worktree_reconcile.assess_cleanup_authorization(
role=role,
capability_blockers=capability_blockers,
operator_authorized=False,
task_capability_resolved=(not capability_blockers) and role_matches,
)
if not cleanup_authorization.get("authorized"):
return {
"success": False,
"performed": False,
"reason": "cleanup_authorization_required",
"reasons": cleanup_authorization.get("reasons") or [],
"cleanup_authorization": cleanup_authorization,
"permission_report": _permission_block_report(required_permission),
"exact_next_action": (
"resolve the reconcile_missing_worktree_bindings capability "
"from a reconciler profile and authorize cleanup via "
"gitea_authorize_reconciliation_cleanup_phase"
),
}
return missing_worktree_reconcile.reconcile_missing_worktree_bindings(
db,
project_root=root,
remote=remote,
org=o,
repo=r,
host=h,
dry_run=dry_run,
operator_authorized=False,
cleanup_authorization=cleanup_authorization,
)
@mcp.tool()
def gitea_audit_worktree_cleanup(
remote: str = "dadeschools",
@@ -15121,22 +15355,17 @@ def _current_runtime_mode_report(refresh: bool = False) -> dict:
workspace_root = None
aligned = None
canonical_root = None
resolved_slug = None
try:
ctx = _resolve_namespace_mutation_context(None)
workspace_root = ctx.get("workspace_path")
canonical_root = ctx.get("canonical_repo_root")
# Alignment keeps its established repository-level meaning (#615 F1):
# does this namespace target the repository the process is installed in,
# i.e. canonical_repo_root == process_project_root. It is deliberately
# NOT path equality between the task workspace and the process root --
# the global worktree rule requires task work to live in a branches/
# worktree, so that comparison would classify every correctly bound
# author, reviewer, and merger session as unsafe.
aligned = ctx.get("roots_aligned")
crr_assessment = ctx.get("canonical_root_assessment") or {}
resolved_slug = crr_assessment.get("resolved_slug")
except Exception:
# An unresolvable binding is reported as unknown alignment, never as
# proof of alignment.
aligned = None
resolved_slug = None
try:
profile_name = get_profile()["profile_name"]
except Exception:
@@ -15149,6 +15378,7 @@ def _current_runtime_mode_report(refresh: bool = False) -> dict:
dirty_files=dirty_files,
active_task_workspace=workspace_root,
canonical_repository_root=canonical_root,
repository_slug=resolved_slug,
workspace_roots_aligned=aligned,
profile=profile_name,
declared_mode=stable_control_runtime.declared_runtime_mode(),
@@ -15492,6 +15722,10 @@ _WORKER_REGISTRY = None
_WORKER_IDENTITY: str | None = None
_WORKER_GENERATION: str | None = None
_WORKER_REGISTRATION_ATTEMPTED = False
#: #975: the one heartbeat supervisor for this process's registration. One
#: process registers exactly one worker identity, so there is exactly one
#: supervisor and it is never replaced.
_WORKER_HEARTBEAT_SUPERVISOR = None
#: Env a client launcher may set to name itself and its session. Absent values
#: are reported as unknown; they are never guessed at, because guessing is what
@@ -15578,6 +15812,17 @@ def _active_worker_identity() -> str | None:
if outcome.get("registered"):
_WORKER_IDENTITY = identity
_WORKER_GENERATION = generation
# #975: registration is the only moment identity and fencing
# epoch are both known, so the heartbeat supervisor is started
# here. Without it ``last_heartbeat_at`` never left
# ``started_at`` and every healthy client lost ownership at the
# TTL.
_start_worker_heartbeat(
identity=identity,
fencing_epoch=outcome.get("fencing_epoch"),
hints=hints,
generation=generation,
)
return identity
if not outcome.get("collision"):
return None
@@ -15588,6 +15833,78 @@ def _active_worker_identity() -> str | None:
return None
def _start_worker_heartbeat(
*,
identity: str,
fencing_epoch,
hints: dict,
generation: str | None,
) -> None:
"""Attach a heartbeat supervisor to the registration just created (#975).
Never raises: a supervisor that cannot start leaves the row exactly as
``register()`` wrote it, which is the pre-#975 behaviour, rather than
failing the tool call that happened to trigger lazy registration.
Under pytest the thread is deliberately not started. Tests drive
``WorkerHeartbeatSupervisor`` directly with an injected clock, so the suite
proves the lifecycle without leaving background sqlite writers behind.
"""
global _WORKER_HEARTBEAT_SUPERVISOR
if _WORKER_HEARTBEAT_SUPERVISOR is not None:
return
registry = _worker_registry()
if registry is None or fencing_epoch is None:
return
try:
supervisor = mcp_worker_identity.WorkerHeartbeatSupervisor(
registry,
worker_identity=identity,
fencing_epoch=int(fencing_epoch),
session_id=hints.get("session_id"),
generation_id=generation,
client_name=hints.get("client_name"),
pid=os.getpid(),
ttl_seconds=mcp_worker_identity.DEFAULT_HEARTBEAT_TTL_SECONDS,
interval_seconds=os.environ.get(
mcp_worker_identity.HEARTBEAT_INTERVAL_ENV
),
)
_WORKER_HEARTBEAT_SUPERVISOR = supervisor
if not mcp_daemon_guard.is_pytest_runtime():
supervisor.start()
except Exception:
return
def _worker_heartbeat_status() -> dict:
"""Read-only heartbeat observability for ``gitea_get_runtime_context`` (#975).
Reports the unsupervised case explicitly rather than omitting the key, so an
operator can tell "no supervisor" apart from "supervisor with no beats yet".
"""
supervisor = _WORKER_HEARTBEAT_SUPERVISOR
if supervisor is None:
return {
"supervised": False,
"running": False,
"reasons": [
"no worker heartbeat supervisor is attached to this process; the "
"registration is not being renewed and will go stale at its TTL"
],
}
try:
status = dict(supervisor.status())
except Exception as exc:
return {
"supervised": True,
"running": False,
"reasons": [f"heartbeat status unavailable: {type(exc).__name__}: {exc}"],
}
status.setdefault("reasons", [])
return status
def _active_role_kind_safe() -> str | None:
"""Best-effort role for the registry record; never raises into a tool call.
@@ -19229,6 +19546,11 @@ def gitea_get_runtime_context(
"fencing_epoch": provenance_assessment["fencing_epoch"],
"conflicting_live_sessions": provenance_assessment["conflicting_live_sessions"],
"provenance_assessment": provenance_assessment,
# #975: whether this registration is actually being renewed. Read-only,
# and it grants nothing — ownership still comes from the attachment
# record above. It exists so "my heartbeat stopped" is diagnosable
# before the TTL turns it into session_attachment_missing.
"worker_heartbeat": _worker_heartbeat_status(),
"unconsumed_gitea_env": unconsumed_env,
"preflight_ready": preflight["preflight_ready"],
"preflight_block_reasons": preflight["preflight_block_reasons"],
@@ -21737,6 +22059,17 @@ def gitea_route_task_session(
# self-recovery was removed from the read-only path (was _trigger_mcp_auto_restart).
# Recovery is owned exclusively by the IDE/client reconnect path.
#: Every reason prefix ``_check_mcp_runtimes_diagnostics`` can emit. Callers
#: raise its reasons as one RuntimeError, and #975 found that the preflight
#: re-raise recognised only ``stale-runtime:``, silently dropping
#: ``unsupported-env:`` while the capability resolver still failed on it. This
#: lives beside the producer so a newly added reason family cannot be forgotten
#: by a distant re-raise predicate again.
RUNTIME_DIAGNOSTIC_HARD_PREFIXES: tuple[str, ...] = (
"stale-runtime:",
"unsupported-env:",
)
def _check_mcp_runtimes_diagnostics(task: str, matching_profiles: list[str]) -> list[str]:
"""Read-only: report missing or stale MCP runtimes (no config or process mutation).
+386
View File
@@ -29,6 +29,7 @@ implements:
from __future__ import annotations
import atexit
import hashlib
import os
import re
@@ -123,6 +124,103 @@ DEFAULT_REGISTRY_PATH = os.path.expanduser(
#: worker within a single operator coffee break.
DEFAULT_HEARTBEAT_TTL_SECONDS = 900.0
#: Optional operator override for the beat interval, in seconds. Clamped by
#: :func:`heartbeat_interval_for` so it can never be set at or above the TTL.
HEARTBEAT_INTERVAL_ENV = "GITEA_WORKER_HEARTBEAT_INTERVAL_SECONDS"
#: Never beat faster than this, so a misconfigured interval cannot turn the
#: supervisor into a busy sqlite writer.
MIN_HEARTBEAT_INTERVAL_SECONDS = 1.0
def heartbeat_interval_for(
ttl_seconds: float = DEFAULT_HEARTBEAT_TTL_SECONDS,
override: str | float | None = None,
) -> float:
"""Beat interval for *ttl_seconds*: one third of the TTL, capped at a half.
#975. A third means two consecutive beats can be lost before the row is
allowed to look stale, and the half-TTL cap is a hard ceiling so no
override can produce an interval that expires the registration it is
supposed to renew. That is the whole liveness contract: a healthy worker
stays owned, and a worker that has genuinely stopped beating still expires
on the configured TTL — this function never touches the TTL itself.
"""
try:
ttl = float(ttl_seconds)
except (TypeError, ValueError):
ttl = DEFAULT_HEARTBEAT_TTL_SECONDS
if not ttl > 0:
ttl = DEFAULT_HEARTBEAT_TTL_SECONDS
interval = ttl / 3.0
if override is not None:
candidate = str(override).strip()
if candidate:
try:
parsed = float(candidate)
except (TypeError, ValueError):
parsed = None
if parsed is not None and parsed > 0:
interval = parsed
ceiling = ttl / 2.0
if interval > ceiling:
interval = ceiling
if interval < MIN_HEARTBEAT_INTERVAL_SECONDS:
interval = min(MIN_HEARTBEAT_INTERVAL_SECONDS, ceiling)
return interval
#: Recorded-vs-presented pairs compared by :func:`_heartbeat_expectation_drift`.
_HEARTBEAT_EXPECTATION_FIELDS = (
"session_id",
"generation_id",
"client_name",
"pid",
)
def _heartbeat_expectation_drift(
record: dict[str, Any],
*,
expected_session_id: str | None = None,
expected_generation_id: str | None = None,
expected_client_name: str | None = None,
expected_pid: int | None = None,
) -> list[tuple[str, Any, Any]]:
"""Return ``(field, recorded, presented)`` for every mismatched expectation.
Only supplied expectations are compared, so a caller that presents nothing
gets the pre-#975 behaviour. ``client_name`` is compared through
:func:`normalize_client_name` so one application's namespaces stay one
client while genuinely different clients stay distinct.
"""
presented = {
"session_id": expected_session_id,
"generation_id": expected_generation_id,
"client_name": (
normalize_client_name(expected_client_name)
if expected_client_name is not None
else None
),
"pid": expected_pid,
}
drift: list[tuple[str, Any, Any]] = []
for field in _HEARTBEAT_EXPECTATION_FIELDS:
want = presented[field]
if want is None:
continue
got = record.get(field)
if field == "pid":
match = got is not None and int(got) == int(want)
else:
match = got == want
if not match:
drift.append((field, got, want))
return drift
STATUS_ACTIVE = "active"
STATUS_SUPERSEDED = "superseded"
STATUS_RELEASED = "released"
@@ -786,11 +884,25 @@ class WorkerRegistry:
worker_identity: str,
fencing_epoch: int,
now: datetime | None = None,
expected_session_id: str | None = None,
expected_generation_id: str | None = None,
expected_client_name: str | None = None,
expected_pid: int | None = None,
) -> dict[str, Any]:
"""Renew only the owning registration (#948 AC11).
A stale epoch is refused rather than silently renewed, so a superseded
session that resumes cannot heartbeat its way back into ownership.
#975 adds optional keyword-only *expectations*. Every one supplied must
match the recorded row or the renewal is refused. They exist because
identity plus epoch cannot express "renew the row I registered, and only
that row": a recycled PID, or a second session of the same client, would
otherwise be renewable by the wrong beater. They are fencing tokens,
never assertions — supplying one can only cause a refusal, never grant
anything, and omitting them preserves the pre-#975 behaviour exactly.
Drift reports the existing ``BLOCKER_FENCED`` rather than a new
``blocker_kind``, because consumers switch on that value.
"""
stamp = _ts(now or _utc_now())
with self._tx() as conn:
@@ -807,6 +919,30 @@ class WorkerRegistry:
"reasons": [f"no registration for {worker_identity!r}"],
}
record = self._row_to_record(row)
drift = _heartbeat_expectation_drift(
record,
expected_session_id=expected_session_id,
expected_generation_id=expected_generation_id,
expected_client_name=expected_client_name,
expected_pid=expected_pid,
)
if drift:
return {
"success": False,
"renewed": False,
"mutation_performed": False,
"blocker_kind": BLOCKER_FENCED,
"expectation_drift": drift,
"reasons": [
"presented worker expectations do not match the recorded "
"registration, so this beater does not own the row: "
+ "; ".join(
f"{field} recorded {recorded!r}, presented {presented!r}"
for field, recorded, presented in drift
)
+ " (#975)"
],
}
if record["status"] != STATUS_ACTIVE:
return {
"success": False,
@@ -1430,3 +1566,253 @@ def resolve_bound_remote(
"explicitly to avoid host drift (#948)."
],
}
# --- Production heartbeat lifecycle (#975) --------------------------------
#
# #948 delivered ``WorkerRegistry.heartbeat()`` and nothing ever called it, so
# ``last_heartbeat_at`` stayed pinned to ``started_at`` for every production
# worker and ``heartbeat_ttl_seconds`` stopped being a liveness window at all —
# it became a hard cap on how long any client could stay attached. This
# supervisor is the missing caller.
#
# It is a daemon thread rather than an asyncio task or a per-request refresh
# because the renewal has to survive an *idle* session: a client waiting on a
# lease with no tool call in flight must not lose ownership, which rules out
# request-driven refresh. Registration itself happens lazily inside a tool call
# on the server's event loop, and the registry performs blocking
# ``BEGIN IMMEDIATE`` sqlite writes, which must not run on that loop.
# ``daemon=True`` is deliberate: a hard kill takes the thread down with the
# process, so a dead worker still goes stale on the normal TTL.
#: Refusals that mean this worker no longer owns its row. Beating again could
#: only ever be an attempt to renew ownership it has already lost, so the
#: supervisor stops permanently instead of retrying.
TERMINAL_HEARTBEAT_BLOCKERS = frozenset({BLOCKER_FENCED, BLOCKER_NO_ATTACHMENT})
class WorkerHeartbeatSupervisor:
"""Periodically renew exactly one worker registration.
Contract:
* It never registers. A supervisor exists only for an already-registered
identity, so it cannot create a second registration or a second identity
system.
* Every beat presents the full expectation set, so it can renew only the row
matching this exact client name, worker identity, session, generation and
pid.
* A terminal refusal (fenced or missing) stops it permanently and records
why. A fenced session must never beat its way back into ownership.
* A transient failure (a locked database, say) is counted and the loop
continues, so one contended write does not silently end the heartbeat.
* Nothing here raises into a caller. ``beat_once`` returns its outcome and
the thread body swallows everything, because a heartbeat failure must
degrade to "not renewed" and never crash a tool call.
"""
def __init__(
self,
registry: WorkerRegistry,
*,
worker_identity: str,
fencing_epoch: int,
session_id: str | None = None,
generation_id: str | None = None,
client_name: str | None = None,
pid: int | None = None,
ttl_seconds: float = DEFAULT_HEARTBEAT_TTL_SECONDS,
interval_seconds: float | None = None,
clock=None,
) -> None:
self._registry = registry
self.worker_identity = worker_identity
self.fencing_epoch = int(fencing_epoch)
self.session_id = session_id
self.generation_id = generation_id
self.client_name = (
normalize_client_name(client_name) if client_name is not None else None
)
self.pid = int(pid) if pid is not None else None
self.ttl_seconds = float(ttl_seconds)
self.interval_seconds = (
heartbeat_interval_for(self.ttl_seconds)
if interval_seconds is None
else heartbeat_interval_for(self.ttl_seconds, interval_seconds)
)
#: Injectable so every TTL test uses controlled time and no test waits
#: for a real interval or a real TTL to elapse.
self._clock = clock or _utc_now
self._stop_event = threading.Event()
self._thread: threading.Thread | None = None
self._lock = threading.Lock()
self._atexit_registered = False
self.started = False
self.stopped_reason: str | None = None
self.beats_attempted = 0
self.beats_renewed = 0
self.transient_failures = 0
self.last_beat_at: str | None = None
self.last_result: dict[str, Any] | None = None
# -- one beat --
def beat_once(self, now: datetime | None = None) -> dict[str, Any]:
"""Renew once. Never raises; returns the registry outcome or a failure."""
if self.stopped_reason is not None:
return {
"success": False,
"renewed": False,
"beat_attempted": False,
"reasons": [
f"supervisor already stopped: {self.stopped_reason}"
],
}
stamp = now or self._clock()
self.beats_attempted += 1
try:
result = self._registry.heartbeat(
worker_identity=self.worker_identity,
fencing_epoch=self.fencing_epoch,
now=stamp,
expected_session_id=self.session_id,
expected_generation_id=self.generation_id,
expected_client_name=self.client_name,
expected_pid=self.pid,
)
except Exception as exc:
# Transient by assumption: an unexpected error is not proof this
# worker lost ownership, so it must not silently end the heartbeat.
self.transient_failures += 1
result = {
"success": False,
"renewed": False,
"transient": True,
"blocker_kind": None,
"reasons": [f"heartbeat raised {type(exc).__name__}: {exc}"],
}
self.last_result = result
return result
self.last_result = result
if result.get("renewed"):
self.beats_renewed += 1
self.last_beat_at = result.get("last_heartbeat_at") or _ts(stamp)
return result
blocker = result.get("blocker_kind")
if blocker in TERMINAL_HEARTBEAT_BLOCKERS:
self._stop_internal(
reason=(
f"refused with blocker_kind={blocker!r}: "
+ "; ".join(result.get("reasons") or [])
)
)
else:
self.transient_failures += 1
return result
# -- lifecycle --
def start(self) -> dict[str, Any]:
"""Start the beat thread. Idempotent; safe to call from a tool call."""
with self._lock:
if self.stopped_reason is not None:
return {
"started": False,
"reasons": [f"supervisor stopped: {self.stopped_reason}"],
}
if self._thread is not None and self._thread.is_alive():
return {"started": True, "already_running": True, "reasons": []}
self._stop_event.clear()
thread = threading.Thread(
target=self._run,
name=f"gitea-worker-heartbeat-{self.worker_identity}",
daemon=True,
)
self._thread = thread
self.started = True
if not self._atexit_registered:
# Orderly shutdown stops the heartbeat. A hard kill does not
# run this, which is correct: the row must then go stale.
atexit.register(self._atexit_stop)
self._atexit_registered = True
thread.start()
return {"started": True, "already_running": False, "reasons": []}
def _run(self) -> None:
while not self._stop_event.is_set():
# Wait first: registration already stamped a fresh heartbeat, so an
# immediate beat would be a redundant write on every server launch.
if self._stop_event.wait(self.interval_seconds):
return
try:
self.beat_once()
except Exception:
# beat_once is already total; this is the last-resort guard that
# keeps a supervisor thread from dying silently.
self.transient_failures += 1
if self.stopped_reason is not None:
return
def stop(self, reason: str = "stopped") -> dict[str, Any]:
"""Stop beating. Idempotent, and prompt because the loop waits on an Event."""
self._stop_internal(reason=reason)
thread = self._thread
if thread is not None and thread.is_alive():
thread.join(timeout=max(1.0, min(5.0, self.interval_seconds)))
return {
"stopped": True,
"reason": self.stopped_reason,
"thread_alive": bool(thread is not None and thread.is_alive()),
}
def _stop_internal(self, *, reason: str) -> None:
if self.stopped_reason is None:
self.stopped_reason = reason
self._stop_event.set()
if self._atexit_registered:
try:
atexit.unregister(self._atexit_stop)
except Exception:
pass
self._atexit_registered = False
def _atexit_stop(self) -> None:
try:
self.stop(reason="process exit")
except Exception:
pass
# -- observability --
def status(self) -> dict[str, Any]:
"""Read-only observability payload; safe to embed in a tool result."""
thread = self._thread
return {
"supervised": True,
"worker_identity": self.worker_identity,
"fencing_epoch": self.fencing_epoch,
"session_id": self.session_id,
"generation_id": self.generation_id,
"client_name": self.client_name,
"pid": self.pid,
"heartbeat_ttl_seconds": self.ttl_seconds,
"heartbeat_interval_seconds": self.interval_seconds,
"started": self.started,
"running": bool(
thread is not None
and thread.is_alive()
and self.stopped_reason is None
),
"stopped_reason": self.stopped_reason,
"beats_attempted": self.beats_attempted,
"beats_renewed": self.beats_renewed,
"transient_failures": self.transient_failures,
"last_heartbeat_at": self.last_beat_at,
"last_blocker_kind": (self.last_result or {}).get("blocker_kind"),
}
File diff suppressed because it is too large Load Diff
+153 -27
View File
@@ -8,8 +8,10 @@ poison workspace purity checks in another namespace.
from __future__ import annotations
import os
import subprocess
import author_mutation_worktree as amw
import canonical_repository_root as crr
ACTIVE_WORKTREE_ENV = amw.ACTIVE_WORKTREE_ENV
AUTHOR_WORKTREE_ENV = amw.AUTHOR_WORKTREE_ENV
@@ -152,6 +154,60 @@ def resolve_namespace_workspace(
return os.path.realpath(process_project_root), "MCP server process root (default)"
def verify_git_common_directory_membership(
workspace_path: str,
canonical_repo_root: str,
) -> tuple[bool, str | None]:
"""Verify that workspace_path belongs to canonical_repo_root via git common-dir or branches containment."""
ws = (workspace_path or "").strip()
root = (canonical_repo_root or "").strip()
if not ws or not root:
return False, "empty workspace or canonical root path"
try:
real_ws = os.path.realpath(os.path.abspath(ws))
real_root = os.path.realpath(os.path.abspath(root))
except Exception as exc:
return False, f"invalid workspace or root path: {exc}"
if real_ws == real_root:
return True, None
if not os.path.isdir(real_ws):
return False, f"workspace directory '{real_ws}' does not exist"
try:
res = subprocess.run(
["git", "-C", real_ws, "rev-parse", "--git-common-dir"],
capture_output=True,
text=True,
check=False,
)
if res.returncode == 0:
common_raw = (res.stdout or "").strip()
common_dir = amw._realpath_git_common_dir(real_ws, common_raw)
real_common = os.path.realpath(common_dir)
canonical_git = os.path.realpath(os.path.join(real_root, ".git"))
if real_common in (canonical_git, real_root):
return True, None
return (
False,
f"workspace '{real_ws}' git common directory '{real_common}' does not match "
f"canonical repository root '{real_root}' (.git at '{canonical_git}')"
)
else:
return (
False,
f"workspace '{real_ws}' is not a valid git repository or git rev-parse failed"
)
except Exception as exc:
return (
False,
f"failed to inspect git common directory for workspace '{real_ws}': {exc}"
)
def resolve_namespace_mutation_context(
*,
role_kind: str,
@@ -163,6 +219,8 @@ def resolve_namespace_mutation_context(
worktree: str | None = None,
profile_name: str | None = None,
configured_canonical_root: str | None = None,
expected_slug: str | None = None,
remote: str | None = None,
) -> dict:
"""Shared workspace resolution for runtime_context and mutation guards.
@@ -180,11 +238,41 @@ def resolve_namespace_mutation_context(
env_map = env if env is not None else os.environ
process_root = os.path.realpath(process_project_root)
role = normalize_role_kind(role_kind, profile_name=profile_name)
configured = (configured_canonical_root or "").strip()
if configured:
canonical_root = os.path.realpath(configured)
configured_val = (configured_canonical_root or "").strip()
if configured_val:
crr_assessment = crr.assess_canonical_repository_root(
configured_value=configured_val,
source="configured_canonical_root",
expected_slug=expected_slug,
process_project_root=process_root,
remote=remote,
require_binding=True,
)
# #973 B10: a refused repository-authority mode deliberately resolves no
# canonical root, so downstream guards keep evaluating the install
# checkout rather than an identity derived through an undefined mode.
# ``roots_aligned`` still follows ``proven`` and stays False, and the
# refusal (with its reason_code) rides along in
# ``canonical_root_assessment`` — this is a fail-closed fallback, never a
# normalisation of the mode.
canonical_root = crr_assessment["canonical_repo_root"] or amw.resolve_canonical_repo_root(
process_root, process_root
)
roots_aligned = crr_assessment["proven"]
else:
canonical_root = amw.resolve_canonical_repo_root(process_root, process_root)
crr_assessment = {
"proven": True,
"block": False,
"reasons": [],
"configured": False,
"canonical_repo_root": amw.resolve_canonical_repo_root(process_root, process_root),
"resolved_slug": None,
"source": None,
"reason_code": None,
}
canonical_root = crr_assessment["canonical_repo_root"]
roots_aligned = (canonical_root == process_root)
durable: dict | None = None
if role == "author":
@@ -234,7 +322,10 @@ def resolve_namespace_mutation_context(
"ignored_bindings": demotions + (pollution.get("ignored_bindings") or []),
"process_project_root": process_root,
"canonical_repo_root": canonical_root,
"roots_aligned": canonical_root == process_root,
"roots_aligned": roots_aligned,
"canonical_root_assessment": crr_assessment,
"expected_slug": expected_slug,
"remote": remote,
}
if durable is not None:
result["author_worktree_resolution"] = durable
@@ -246,6 +337,15 @@ def resolve_namespace_mutation_context(
result["author_worktree_reasons"] = list(durable.get("reasons") or [])
result["author_worktree_blocker_kind"] = durable.get("blocker_kind")
result["operator_recovery"] = durable.get("operator_recovery")
else:
path_exists = os.path.exists(workspace)
result["path_exists"] = path_exists
result["in_git_worktree_list"] = (
amw.path_in_git_worktree_list(workspace, canonical_root)
if path_exists
else False
)
result["bound_worktree_missing"] = not path_exists
return result
@@ -378,8 +478,8 @@ def format_namespace_workspace_binding_error(
def assess_namespace_mutation_workspace(
*,
role_kind: str,
worktree_path: str | None,
worktree: str | None,
worktree_path: str | None = None,
worktree: str | None = None,
process_project_root: str,
env: dict[str, str] | os._Environ | None = None,
session_lease_worktree: str | None = None,
@@ -387,6 +487,8 @@ def assess_namespace_mutation_workspace(
profile_name: str | None = None,
current_branch: str | None = None,
configured_canonical_root: str | None = None,
expected_slug: str | None = None,
remote: str | None = None,
) -> dict:
"""Evaluate namespace workspace binding before preflight/mutation."""
ctx = resolve_namespace_mutation_context(
@@ -399,6 +501,8 @@ def assess_namespace_mutation_workspace(
session_lock_worktree=session_lock_worktree,
profile_name=profile_name,
configured_canonical_root=configured_canonical_root,
expected_slug=expected_slug,
remote=remote,
)
mutation_workspace = ctx["workspace_path"]
binding_source = ctx["workspace_binding_source"]
@@ -422,6 +526,27 @@ def assess_namespace_mutation_workspace(
reasons = list(metadata.get("reasons") or [])
operator_recovery = ctx.get("operator_recovery")
crr_reasons = list(ctx.get("canonical_root_assessment", {}).get("reasons") or [])
if crr_reasons:
reasons.extend(crr_reasons)
path_exists = ctx.get("path_exists")
if path_exists is None:
path_exists = os.path.exists(mutation_workspace)
if not path_exists:
if role != "author":
reasons.append(
f"{role} mutation blocked: configured workspace directory '{mutation_workspace}' does not exist (nonexistent worktree)"
)
else:
valid_common, common_err = verify_git_common_directory_membership(
mutation_workspace, ctx["canonical_repo_root"]
)
if not valid_common and common_err:
reasons.append(common_err)
if role == "author":
# #618 durable resolution already validated existence, membership,
# branches/, lock ownership, and traversal safety when present.
@@ -438,26 +563,27 @@ def assess_namespace_mutation_workspace(
)
if branches["block"]:
reasons.extend(branches["reasons"])
elif (
role == "reviewer"
and mutation_workspace == process_root
and not amw.is_path_under_branches(mutation_workspace, ctx["canonical_repo_root"])
):
reasons.append(
f"{role} mutation blocked: workspace is the stable control checkout; "
f"create or reconnect to a session-owned worktree under branches/ "
f"or set {ROLE_WORKTREE_ENVS.get(role, ACTIVE_WORKTREE_ENV)} / "
f"{ACTIVE_WORKTREE_ENV}"
)
elif (
role in {"reviewer", "merger"}
and mutation_workspace != process_root
and not amw.is_path_under_branches(mutation_workspace, ctx["canonical_repo_root"])
):
reasons.append(
f"{role} mutation blocked: workspace '{mutation_workspace}' is not under "
f"'{ctx['canonical_repo_root']}/branches/'"
)
elif role in {"reviewer", "merger"}:
if mutation_workspace == process_root and not amw.is_path_under_branches(mutation_workspace, ctx["canonical_repo_root"]):
reasons.append(
f"{role} mutation blocked: workspace is the stable control checkout; "
f"create or reconnect to a session-owned worktree under branches/ "
f"or set {ROLE_WORKTREE_ENVS.get(role, ACTIVE_WORKTREE_ENV)} / "
f"{ACTIVE_WORKTREE_ENV}"
)
elif mutation_workspace != process_root and not amw.is_path_under_branches(mutation_workspace, ctx["canonical_repo_root"]):
reasons.append(
f"{role} mutation blocked: workspace '{mutation_workspace}' is not under "
f"'{ctx['canonical_repo_root']}/branches/'"
)
if path_exists and amw.is_path_under_branches(mutation_workspace, ctx["canonical_repo_root"]):
in_list = ctx.get("in_git_worktree_list")
if in_list is False:
reasons.append(
f"{role} mutation blocked: workspace '{mutation_workspace}' is under branches/ "
f"but is not registered in git worktree list for '{ctx['canonical_repo_root']}'"
)
block = bool(reasons)
return {
+21
View File
@@ -386,6 +386,25 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"permission": "gitea.branch.delete",
"role": "reconciler",
},
# #970: auditing missing worktree bindings is read-only; retiring one is a
# control-plane cleanup mutation and carries the same reconciler-only
# authority as any other reconciliation cleanup (review 644 B2).
"audit_missing_worktree_bindings": {
"permission": "gitea.read",
"role": "reconciler",
},
"gitea_audit_missing_worktree_bindings": {
"permission": "gitea.read",
"role": "reconciler",
},
"reconcile_missing_worktree_bindings": {
"permission": "gitea.branch.delete",
"role": "reconciler",
},
"gitea_reconcile_missing_worktree_bindings": {
"permission": "gitea.branch.delete",
"role": "reconciler",
},
"work_issue": {
"permission": "gitea.pr.create",
"role": "author",
@@ -653,6 +672,8 @@ ROLE_EXCLUSIVE_TASKS: frozenset[str] = frozenset(
"delete_branch",
"cleanup_merged_pr_branch",
"reconciliation_cleanup",
"reconcile_missing_worktree_bindings",
"gitea_reconcile_missing_worktree_bindings",
"work_issue",
"work-issue",
}
@@ -313,6 +313,7 @@ class TestCanonicalRootGuardBinding(_ServerHarness):
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode="derivation",
)
self.assertFalse(got.get("block"), got.get("reasons"))
self.assertEqual(got["resolved_slug"], TARGET_SLUG)
@@ -241,9 +241,10 @@ class TestNamespaceContextUsesConfiguredRoot(unittest.TestCase):
process_project_root=self.install,
env={},
configured_canonical_root=self.target,
expected_slug="Scaled-Tech-Consulting/mcp-control-plane",
)
self.assertEqual(ctx["canonical_repo_root"], self.target)
self.assertFalse(ctx["roots_aligned"])
self.assertTrue(ctx["roots_aligned"])
def test_target_worktree_is_member_of_target_root(self):
got = nwb.amw.assess_workspace_repo_membership(
@@ -268,6 +269,7 @@ class TestNamespaceContextUsesConfiguredRoot(unittest.TestCase):
env={},
current_branch="feat/issue-1",
configured_canonical_root=self.target,
expected_slug="Scaled-Tech-Consulting/mcp-control-plane",
)
self.assertFalse(assessment["block"], assessment.get("reasons"))
self.assertEqual(assessment["canonical_repo_root"], self.target)
File diff suppressed because it is too large Load Diff
+739
View File
@@ -0,0 +1,739 @@
"""Regression tests for Issue #973 blocker B10: repository-authority mode contract.
Before this repair ``assess_canonical_repository_root`` never checked its
``mode`` argument against an allowlist. Both dispatch points were permissive:
* the configured-root path tested ``mode == "derivation"`` and sent every other
value into a catch-all ``else``, so an unsupported mode silently received
*validation* semantics, and
* the single-repository default path tested ``mode == "validation"``, so an
unsupported mode skipped the identity comparison entirely and was strictly
*weaker* than validation.
The measured consequence was that ``mode="invalid_mode"`` with matching expected
and observed identities returned ``proven: True`` / ``block: False`` with no
reasons, and that an unsupported mode passed on the default path where
``"validation"`` correctly blocked.
These tests exercise the production module directly with real git repositories
no patched stand-in for the function under test and drive the production
enforcement and mutation-context consumers rather than only the intermediate
assessment.
"""
from __future__ import annotations
import ast
import inspect
import os
import subprocess
import tempfile
import unittest
from unittest.mock import patch
import canonical_repository_root as crr
import gitea_config
import gitea_mcp_server as mcp_server
import namespace_workspace_binding as nwb
import stable_control_runtime
INSTALL_SLUG = "Scaled-Tech-Consulting/Gitea-Tools"
TARGET_SLUG = "Scaled-Tech-Consulting/mcp-control-plane"
FOREIGN_SLUG = "Someone-Else/Evil-Repo"
# Explicitly supplied values that must all be refused. Omission is *not* in this
# list: omitting the argument keeps the documented ``"validation"`` default.
UNSUPPORTED_STRING_MODES = (
"invalid_mode",
"",
"validaton", # misspelling
"derivaton", # misspelling
"Validation", # case variant
"DERIVATION", # case variant
" validation", # leading whitespace
"validation ", # trailing whitespace
"derivation\n", # trailing newline
"validation,derivation",
)
UNSUPPORTED_NON_STRING_MODES = (
None,
0,
1,
True,
False,
3.14,
[],
["validation"],
{},
{"mode": "validation"},
("validation",),
object(),
)
def _init_repo(path: str, remote_url: str, *, user: str = "Test User") -> None:
os.makedirs(path, exist_ok=True)
subprocess.run(["git", "init", "-b", "master"], cwd=path, check=True,
capture_output=True)
subprocess.run(["git", "config", "user.email", "[email protected]"], cwd=path,
check=True, capture_output=True)
subprocess.run(["git", "config", "user.name", user], cwd=path, check=True,
capture_output=True)
with open(os.path.join(path, "README.md"), "w") as handle:
handle.write(f"{os.path.basename(path)}\n")
subprocess.run(["git", "add", "README.md"], cwd=path, check=True,
capture_output=True)
subprocess.run(["git", "commit", "-m", "initial"], cwd=path, check=True,
capture_output=True)
subprocess.run(["git", "remote", "add", "prgs", remote_url], cwd=path,
check=True, capture_output=True)
class _CanonicalRootFixture(unittest.TestCase):
"""Real install / target / foreign git repositories, as in the #973 suite."""
def setUp(self):
self._tmp = tempfile.TemporaryDirectory()
self.tmp_dir = os.path.realpath(self._tmp.name)
self.install_root = os.path.join(self.tmp_dir, "Gitea-Tools")
_init_repo(
self.install_root,
"https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools.git",
)
self.target_root = os.path.join(self.tmp_dir, "mcp-control-plane")
_init_repo(
self.target_root,
"https://gitea.prgs.cc/Scaled-Tech-Consulting/mcp-control-plane.git",
)
self.evil_root = os.path.join(self.tmp_dir, "Evil-Repo")
_init_repo(
self.evil_root,
"https://gitea.prgs.cc/Someone-Else/Evil-Repo.git",
user="Evil User",
)
self.target_branches = os.path.join(self.target_root, "branches")
self.target_worktree = os.path.join(self.target_branches, "rev-pr-99")
subprocess.run(
["git", "worktree", "add", "-b", "rev-pr-99", self.target_worktree],
cwd=self.target_root, check=True, capture_output=True,
)
def tearDown(self):
self._tmp.cleanup()
def assertRefusedForMode(self, assessment: dict, mode) -> None:
"""Assert a fail-closed refusal attributable to *mode* and nothing else."""
self.assertFalse(assessment["proven"], assessment)
self.assertTrue(assessment["block"], assessment)
self.assertEqual(assessment["reason_code"], crr.DENY_UNKNOWN_MODE, assessment)
self.assertEqual(len(assessment["reasons"]), 1, assessment)
reason = assessment["reasons"][0]
self.assertIn("unsupported repository-authority mode", reason)
self.assertIn(repr(mode), reason)
# No trusted repository identity may be derived through an invalid mode.
self.assertIsNone(assessment["resolved_slug"], assessment)
self.assertIsNone(assessment["canonical_repo_root"], assessment)
class TestB10UnsupportedModeIsRejected(_CanonicalRootFixture):
"""Direct assessment tests for invalid-mode parsing."""
def test_invalid_string_with_missing_expected_identity(self):
"""G1: blocks for the mode, not incidentally for a missing identity."""
assessment = crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="env",
expected_slug=None,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode="invalid_mode",
)
self.assertRefusedForMode(assessment, "invalid_mode")
self.assertFalse(
any("unprovable or missing" in r for r in assessment["reasons"]),
"must block because the mode is unsupported, not because the expected "
"identity happened to be missing",
)
def test_invalid_string_with_matching_identities(self):
"""G2: the contract violation — matching identities used to return proven."""
assessment = crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="env",
expected_slug=TARGET_SLUG,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode="invalid_mode",
)
self.assertRefusedForMode(assessment, "invalid_mode")
def test_invalid_string_with_conflicting_identities(self):
"""G3: refused for the mode, not for the incidental identity mismatch."""
assessment = crr.assess_canonical_repository_root(
configured_value=self.evil_root,
source="env",
expected_slug=INSTALL_SLUG,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode="invalid_mode",
)
self.assertRefusedForMode(assessment, "invalid_mode")
self.assertFalse(
any("identity mismatch" in r for r in assessment["reasons"]),
"the identity comparison must not have run at all",
)
def test_empty_string_mode(self):
"""G4a: an explicitly supplied empty string is an unsupported value."""
assessment = crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="env",
expected_slug=TARGET_SLUG,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode="",
)
self.assertRefusedForMode(assessment, "")
def test_explicit_none_mode(self):
"""G4b: explicit None is refused; it is not treated as omission."""
assessment = crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="env",
expected_slug=TARGET_SLUG,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode=None,
)
self.assertRefusedForMode(assessment, None)
self.assertIn("of type NoneType", assessment["reasons"][0])
def test_representative_non_string_modes(self):
for mode in UNSUPPORTED_NON_STRING_MODES:
with self.subTest(mode=repr(mode)):
assessment = crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="env",
expected_slug=TARGET_SLUG,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode=mode,
)
self.assertRefusedForMode(assessment, mode)
self.assertIn(
f"of type {type(mode).__name__}", assessment["reasons"][0]
)
def test_unknown_strings_misspellings_and_whitespace_variants(self):
for mode in UNSUPPORTED_STRING_MODES:
with self.subTest(mode=repr(mode)):
assessment = crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="env",
expected_slug=TARGET_SLUG,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode=mode,
)
self.assertRefusedForMode(assessment, mode)
def test_supported_modes_are_exactly_two(self):
self.assertEqual(
crr.SUPPORTED_MODES, ("validation", "derivation")
)
self.assertIsNone(crr.unsupported_mode_reason("validation"))
self.assertIsNone(crr.unsupported_mode_reason("derivation"))
self.assertIsNotNone(crr.unsupported_mode_reason("invalid_mode"))
class TestB10SupportedModesUnchanged(_CanonicalRootFixture):
"""The repair must not disturb the two documented modes."""
def test_omitted_mode_defaults_to_validation(self):
"""G4c/G4d: omission still selects validation, proven by both outcomes."""
matching = crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="env",
expected_slug=TARGET_SLUG,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
)
self.assertTrue(matching["proven"], matching)
self.assertFalse(matching["block"], matching)
self.assertIsNone(matching["reason_code"], matching)
self.assertEqual(matching["resolved_slug"], TARGET_SLUG)
conflicting = crr.assess_canonical_repository_root(
configured_value=self.evil_root,
source="env",
expected_slug=INSTALL_SLUG,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
)
self.assertFalse(conflicting["proven"], conflicting)
self.assertTrue(conflicting["block"], conflicting)
self.assertIsNone(conflicting["reason_code"], conflicting)
self.assertTrue(
any("identity mismatch" in r for r in conflicting["reasons"]),
"omission must behave exactly like explicit validation",
)
def test_explicit_validation_retains_strict_behavior(self):
"""G5a/G5b/G5c."""
ok = crr.assess_canonical_repository_root(
configured_value=self.target_root, source="env",
expected_slug=TARGET_SLUG, process_project_root=self.install_root,
remote="prgs", require_binding=True, mode="validation",
)
self.assertTrue(ok["proven"], ok)
mismatch = crr.assess_canonical_repository_root(
configured_value=self.evil_root, source="env",
expected_slug=INSTALL_SLUG, process_project_root=self.install_root,
remote="prgs", require_binding=True, mode="validation",
)
self.assertTrue(mismatch["block"], mismatch)
self.assertTrue(any("identity mismatch" in r for r in mismatch["reasons"]))
unprovable = crr.assess_canonical_repository_root(
configured_value=self.target_root, source="env",
expected_slug=None, process_project_root=self.install_root,
remote="prgs", require_binding=True, mode="validation",
)
self.assertTrue(unprovable["block"], unprovable)
self.assertTrue(
any("unprovable or missing" in r for r in unprovable["reasons"])
)
def test_explicit_derivation_retains_trusted_derivation(self):
"""G5d: derivation still resolves identity with no expected slug."""
derived = crr.assess_canonical_repository_root(
configured_value=self.target_root, source="env",
expected_slug=None, process_project_root=self.install_root,
remote="prgs", require_binding=True, mode="derivation",
)
self.assertTrue(derived["proven"], derived)
self.assertFalse(derived["block"], derived)
self.assertIsNone(derived["reason_code"], derived)
self.assertEqual(derived["resolved_slug"], TARGET_SLUG)
def test_derivation_without_resolvable_remote_still_fails_closed(self):
no_remote = os.path.join(self.tmp_dir, "no-remote-target")
os.makedirs(no_remote)
subprocess.run(["git", "init", "-b", "master"], cwd=no_remote, check=True,
capture_output=True)
assessment = crr.assess_canonical_repository_root(
configured_value=no_remote, source="env", expected_slug=None,
process_project_root=self.install_root, remote="prgs",
require_binding=True, mode="derivation",
)
self.assertTrue(assessment["block"], assessment)
self.assertTrue(
any("no resolvable" in r for r in assessment["reasons"]), assessment
)
class TestB10SingleRepositoryDefaultPath(_CanonicalRootFixture):
"""G6: on the unconfigured path an invalid mode used to be weaker than validation."""
def _assess(self, mode_kwargs: dict) -> dict:
return crr.assess_canonical_repository_root(
configured_value=None,
source=None,
expected_slug=FOREIGN_SLUG,
process_project_root=self.install_root,
remote="prgs",
require_binding=False,
**mode_kwargs,
)
def test_validation_blocks_a_foreign_expected_identity(self):
"""G6b: the reference behaviour the invalid mode must not undercut."""
got = self._assess({"mode": "validation"})
self.assertFalse(got["proven"], got)
self.assertTrue(got["block"], got)
self.assertTrue(any("identity mismatch" in r for r in got["reasons"]))
def test_invalid_mode_no_longer_passes_where_validation_blocks(self):
"""G6a: identical inputs, only the mode differs — must not fail open."""
invalid = self._assess({"mode": "invalid_mode"})
self.assertRefusedForMode(invalid, "invalid_mode")
validation = self._assess({"mode": "validation"})
self.assertEqual(
invalid["proven"], validation["proven"],
"an unsupported mode must never be more permissive than validation",
)
self.assertTrue(invalid["block"] and validation["block"])
def test_empty_string_mode_on_default_path(self):
"""G6c."""
self.assertRefusedForMode(self._assess({"mode": ""}), "")
def test_omitted_mode_on_default_path_still_validates(self):
got = self._assess({})
self.assertFalse(got["proven"], got)
self.assertTrue(any("identity mismatch" in r for r in got["reasons"]), got)
class TestB10RejectionOrdering(_CanonicalRootFixture):
"""Refusal must precede every form of candidate-root or Git inspection."""
def test_no_git_or_identity_discovery_runs_for_an_unsupported_mode(self):
with patch.object(crr, "resolve_repo_toplevel") as toplevel, \
patch.object(crr, "repository_identity_slug") as identity:
assessment = crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="env",
expected_slug=TARGET_SLUG,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode="invalid_mode",
)
self.assertRefusedForMode(assessment, "invalid_mode")
toplevel.assert_not_called()
identity.assert_not_called()
def test_the_same_spies_do_fire_for_a_supported_mode(self):
"""Control: proves the previous test's assertions are not vacuous."""
with patch.object(crr, "resolve_repo_toplevel",
wraps=crr.resolve_repo_toplevel) as toplevel, \
patch.object(crr, "repository_identity_slug",
wraps=crr.repository_identity_slug) as identity:
crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="env",
expected_slug=TARGET_SLUG,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode="validation",
)
toplevel.assert_called()
identity.assert_called()
def test_nonexistent_candidate_root_still_reports_the_mode_refusal(self):
"""No mocks: the existence check cannot have run before the refusal."""
nonexistent = os.path.join(self.tmp_dir, "no-such-repository")
assessment = crr.assess_canonical_repository_root(
configured_value=nonexistent,
source="env",
expected_slug=TARGET_SLUG,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode="invalid_mode",
)
self.assertRefusedForMode(assessment, "invalid_mode")
self.assertFalse(
any("does not exist" in r for r in assessment["reasons"]), assessment
)
def test_symlinked_candidate_root_is_not_resolved_for_an_unsupported_mode(self):
link = os.path.join(self.tmp_dir, "target-alias")
os.symlink(self.target_root, link)
assessment = crr.assess_canonical_repository_root(
configured_value=link,
source="env",
expected_slug=TARGET_SLUG,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode="invalid_mode",
)
self.assertRefusedForMode(assessment, "invalid_mode")
# The refusal payload must not leak a resolved path for the candidate.
self.assertIsNone(assessment["canonical_repo_root"], assessment)
class TestB10NoModeInjectionSurface(unittest.TestCase):
"""``mode`` must not be reachable from requests, environment, or config."""
PRODUCTION_MODULES = ("gitea_mcp_server.py", "namespace_workspace_binding.py")
def _repo_root(self) -> str:
return os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
def test_production_call_sites_only_pass_allowlisted_literal_modes(self):
"""Valid hardcoded modes reach the correct path; nothing else is passed."""
seen: list[tuple[str, str | None]] = []
for name in self.PRODUCTION_MODULES:
path = os.path.join(self._repo_root(), name)
with open(path) as handle:
tree = ast.parse(handle.read())
for node in ast.walk(tree):
if not isinstance(node, ast.Call):
continue
func = node.func
called = (
func.attr if isinstance(func, ast.Attribute)
else getattr(func, "id", None)
)
if called != "assess_canonical_repository_root":
continue
supplied = [k for k in node.keywords if k.arg == "mode"]
if not supplied:
seen.append((name, None))
continue
value = supplied[0].value
self.assertIsInstance(
value, ast.Constant,
f"{name}: mode must be a literal, never a variable or expression",
)
self.assertIn(
value.value, crr.SUPPORTED_MODES,
f"{name}: unsupported mode literal {value.value!r}",
)
seen.append((name, value.value))
self.assertTrue(seen, "expected production call sites to be found")
# Both documented modes are exercised by production, and omission is used.
self.assertIn(None, [mode for _, mode in seen])
self.assertIn("derivation", [mode for _, mode in seen])
def test_no_public_entry_point_exposes_a_mode_parameter(self):
for func in (
nwb.resolve_namespace_mutation_context,
nwb.assess_namespace_mutation_workspace,
mcp_server._resolve_namespace_mutation_context,
mcp_server._enforce_canonical_repository_root,
mcp_server._canonical_repository_slug,
mcp_server._trusted_session_repository,
mcp_server._resolve_expected_repository_slug,
):
with self.subTest(func=func.__name__):
self.assertNotIn("mode", inspect.signature(func).parameters)
def test_no_environment_key_selects_a_repository_authority_mode(self):
for key in gitea_config.RECOGNIZED_GITEA_ENV_KEYS:
self.assertNotIn(
"CANONICAL_REPOSITORY_MODE", key.upper(),
f"{key} would expose a repository-authority mode selector",
)
# An invented mode-ish variable is simply not consumed by crr.
env = {
"GITEA_CANONICAL_REPOSITORY_ROOT": "/some/path",
"GITEA_CANONICAL_REPOSITORY_MODE": "invalid_mode",
}
value, source = crr.configured_canonical_root(None, env)
self.assertEqual(value, "/some/path")
self.assertNotIn("mode", (source or "").lower())
self.assertIn(
"GITEA_CANONICAL_REPOSITORY_MODE",
gitea_config.get_unconsumed_gitea_env_overrides(env),
"an unknown GITEA_* key must still be rejected as unrecognised",
)
def test_repository_configuration_carries_no_mode_field(self):
profile = {
"canonical_repository_root": "/some/path",
"mode": "invalid_mode",
}
value, source = crr.configured_canonical_root(profile, {})
self.assertEqual(value, "/some/path")
self.assertEqual(source, "profile canonical_repository_root")
class TestB10ProductionEnforcementPaths(_CanonicalRootFixture):
"""Production enforcement, mutation-context, reviewer and merger consumers."""
def _force_invalid_mode(self):
"""Simulate a future call site threading an unsupported mode.
The real ``assess_canonical_repository_root`` still executes only the
caller-side argument is substituted so the refusal under test is
produced by production code, not by a stand-in. No caller-controlled
``mode`` parameter is added to any production signature to achieve this.
"""
real = crr.assess_canonical_repository_root
def _wrapper(**kwargs):
kwargs["mode"] = "invalid_mode"
return real(**kwargs)
return patch.object(crr, "assess_canonical_repository_root", _wrapper)
def test_mutation_context_fails_closed_under_a_refused_mode(self):
with self._force_invalid_mode():
ctx = nwb.resolve_namespace_mutation_context(
role_kind="reviewer",
worktree_path=self.target_worktree,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.target_root,
expected_slug=TARGET_SLUG,
remote="prgs",
)
assessment = ctx["canonical_root_assessment"]
self.assertFalse(ctx["roots_aligned"], ctx)
self.assertTrue(assessment["block"], assessment)
self.assertEqual(assessment["reason_code"], crr.DENY_UNKNOWN_MODE)
self.assertIsNone(assessment["resolved_slug"])
# The refused mode must not yield the candidate root as canonical.
self.assertNotEqual(ctx["canonical_repo_root"], self.target_root)
self.assertEqual(ctx["canonical_repo_root"], self.install_root)
def test_mutation_context_unchanged_for_the_supported_default(self):
ctx = nwb.resolve_namespace_mutation_context(
role_kind="reviewer",
worktree_path=self.target_worktree,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.target_root,
expected_slug=TARGET_SLUG,
remote="prgs",
)
self.assertTrue(ctx["roots_aligned"], ctx)
self.assertEqual(ctx["canonical_repo_root"], self.target_root)
self.assertIsNone(ctx["canonical_root_assessment"]["reason_code"])
def test_reviewer_and_merger_authorization_fails_closed_under_a_refused_mode(self):
for role in ("reviewer", "merger"):
with self.subTest(role=role), self._force_invalid_mode():
assessment = nwb.assess_namespace_mutation_workspace(
role_kind=role,
worktree_path=self.target_worktree,
worktree=None,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.target_root,
expected_slug=TARGET_SLUG,
remote="prgs",
)
self.assertTrue(assessment["block"], assessment)
self.assertTrue(
any("unsupported repository-authority mode" in r
for r in assessment["reasons"]),
assessment,
)
def test_reviewer_and_merger_authorization_unchanged_for_supported_modes(self):
for role in ("reviewer", "merger"):
with self.subTest(role=role):
assessment = nwb.assess_namespace_mutation_workspace(
role_kind=role,
worktree_path=self.target_worktree,
worktree=None,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.target_root,
expected_slug=TARGET_SLUG,
remote="prgs",
)
self.assertFalse(assessment["block"], assessment)
def test_final_mutation_gate_blocks_when_the_mode_refusal_unaligns_roots(self):
with self._force_invalid_mode():
ctx = nwb.resolve_namespace_mutation_context(
role_kind="reviewer",
worktree_path=self.target_worktree,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.target_root,
expected_slug=TARGET_SLUG,
remote="prgs",
)
report = stable_control_runtime.build_runtime_report(
process_root=self.install_root,
checkout_branch="master",
runtime_head="abcdef123456",
active_task_workspace=ctx["workspace_path"],
canonical_repository_root=ctx["canonical_repo_root"],
workspace_roots_aligned=ctx["roots_aligned"],
)
gate = stable_control_runtime.assess_runtime_mutation_gate(report)
self.assertTrue(gate["block"], gate)
def test_enforce_canonical_repository_root_uses_the_validation_default(self):
"""B8 boundary intact: a foreign configured root raises through production."""
with patch.object(mcp_server, "PROJECT_ROOT", self.install_root), \
patch.object(mcp_server, "_configured_canonical_root",
return_value=(self.evil_root, "env")), \
patch.object(mcp_server.session_ctx, "get_session_context",
return_value=None):
with self.assertRaises(RuntimeError) as raised:
mcp_server._enforce_canonical_repository_root(remote="prgs")
self.assertIn("identity mismatch", str(raised.exception))
def test_enforce_canonical_repository_root_refuses_a_threaded_invalid_mode(self):
bound = {"org": "Scaled-Tech-Consulting",
"repository": "mcp-control-plane", "remote": "prgs"}
with patch.object(mcp_server, "PROJECT_ROOT", self.install_root), \
patch.object(mcp_server, "_configured_canonical_root",
return_value=(self.target_root, "env")), \
patch.object(mcp_server.session_ctx, "get_session_context",
return_value=bound), \
self._force_invalid_mode():
with self.assertRaises(RuntimeError) as raised:
mcp_server._enforce_canonical_repository_root(remote="prgs")
self.assertIn("unsupported repository-authority mode", str(raised.exception))
def test_enforce_canonical_repository_root_passes_for_a_valid_binding(self):
bound = {"org": "Scaled-Tech-Consulting",
"repository": "mcp-control-plane", "remote": "prgs"}
with patch.object(mcp_server, "PROJECT_ROOT", self.install_root), \
patch.object(mcp_server, "_configured_canonical_root",
return_value=(self.target_root, "env")), \
patch.object(mcp_server.session_ctx, "get_session_context",
return_value=bound), \
patch.object(mcp_server.session_ctx, "assess_session_context",
return_value={"block": False, "reasons": []}), \
patch.object(mcp_server, "get_profile",
return_value={"profile_name": "prgs-reviewer"}):
mcp_server._enforce_canonical_repository_root(remote="prgs")
def test_canonical_repository_slug_derivation_unbroken(self):
"""B9 boundary intact: legitimate cross-repository derivation still works."""
profile = {
"profile_name": "prgs-author",
"allowed_repositories": [TARGET_SLUG],
"canonical_repository_root": self.target_root,
}
with patch.object(mcp_server, "PROJECT_ROOT", self.install_root), \
patch.object(mcp_server.session_ctx, "get_session_context",
return_value=None):
slug, reasons = mcp_server._canonical_repository_slug(profile, "prgs")
self.assertEqual(slug, TARGET_SLUG, reasons)
self.assertEqual(reasons, [])
def test_canonical_repository_slug_fails_closed_under_a_refused_mode(self):
profile = {
"profile_name": "prgs-author",
"allowed_repositories": [TARGET_SLUG],
"canonical_repository_root": self.target_root,
}
with patch.object(mcp_server, "PROJECT_ROOT", self.install_root), \
patch.object(mcp_server.session_ctx, "get_session_context",
return_value=None), \
self._force_invalid_mode():
slug, reasons = mcp_server._canonical_repository_slug(profile, "prgs")
self.assertIsNone(slug)
self.assertTrue(
any("unsupported repository-authority mode" in r for r in reasons),
reasons,
)
result = mcp_server._trusted_session_repository(
profile, "prgs", for_mutation=True
)
self.assertIsNone(result["org"])
self.assertIsNone(result["repository"])
self.assertTrue(result["reasons"])
if __name__ == "__main__":
unittest.main()
@@ -0,0 +1,546 @@
"""Regression tests for Issue #973: validated cross-repository canonical roots."""
from __future__ import annotations
import os
import shutil
import tempfile
import unittest
from unittest.mock import patch, MagicMock
from pathlib import Path
import subprocess
import gitea_config
import namespace_workspace_binding as nwb
import canonical_repository_root as crr
import stable_control_runtime
import gitea_mcp_server as mcp_server
class TestIssue973RecognizedEnvKeys(unittest.TestCase):
"""Test recognized environment variable keys under #973."""
def test_recognized_gitea_env_keys(self):
for key in (
"GITEA_CANONICAL_REPOSITORY_ROOT",
"GITEA_REVIEWER_WORKTREE",
"GITEA_MERGER_WORKTREE",
"GITEA_MCP_SESSION_STATE_TTL_HOURS",
):
self.assertIn(key, gitea_config.RECOGNIZED_GITEA_ENV_KEYS)
def test_get_unconsumed_gitea_env_overrides_ignores_recognized(self):
env = {
"GITEA_CANONICAL_REPOSITORY_ROOT": "/some/path",
"GITEA_REVIEWER_WORKTREE": "/some/reviewer/path",
"GITEA_MERGER_WORKTREE": "/some/merger/path",
"GITEA_MCP_SESSION_STATE_TTL_HOURS": "24",
"GITEA_UNRECOGNIZED_FOO_VAR": "bar",
}
unconsumed = gitea_config.get_unconsumed_gitea_env_overrides(env)
self.assertNotIn("GITEA_CANONICAL_REPOSITORY_ROOT", unconsumed)
self.assertNotIn("GITEA_REVIEWER_WORKTREE", unconsumed)
self.assertNotIn("GITEA_MERGER_WORKTREE", unconsumed)
self.assertNotIn("GITEA_MCP_SESSION_STATE_TTL_HOURS", unconsumed)
self.assertIn("GITEA_UNRECOGNIZED_FOO_VAR", unconsumed)
class TestIssue973CrossRepoCanonicalRoots(unittest.TestCase):
"""Test workspace binding and canonical root validation for cross-repo namespaces."""
def setUp(self):
self._tmp = tempfile.TemporaryDirectory()
self.tmp_dir = os.path.realpath(self._tmp.name)
# Create simulated installation root
self.install_root = os.path.join(self.tmp_dir, "Gitea-Tools")
os.makedirs(self.install_root)
subprocess.run(["git", "init", "-b", "master"], cwd=self.install_root, check=True)
subprocess.run(["git", "config", "user.email", "[email protected]"], cwd=self.install_root, check=True)
subprocess.run(["git", "config", "user.name", "Test User"], cwd=self.install_root, check=True)
with open(os.path.join(self.install_root, "README.md"), "w") as f:
f.write("install\n")
subprocess.run(["git", "add", "README.md"], cwd=self.install_root, check=True)
subprocess.run(["git", "commit", "-m", "initial"], cwd=self.install_root, check=True)
# Create simulated target repository root
self.target_root = os.path.join(self.tmp_dir, "mcp-control-plane")
os.makedirs(self.target_root)
subprocess.run(["git", "init", "-b", "master"], cwd=self.target_root, check=True)
subprocess.run(["git", "config", "user.email", "[email protected]"], cwd=self.target_root, check=True)
subprocess.run(["git", "config", "user.name", "Test User"], cwd=self.target_root, check=True)
with open(os.path.join(self.target_root, "README.md"), "w") as f:
f.write("target\n")
subprocess.run(["git", "add", "README.md"], cwd=self.target_root, check=True)
subprocess.run(["git", "commit", "-m", "initial"], cwd=self.target_root, check=True)
# Add remotes to simulate real git repositories with identities
subprocess.run(["git", "remote", "add", "prgs", "https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools.git"], cwd=self.install_root, check=True)
subprocess.run(["git", "remote", "add", "prgs", "https://gitea.prgs.cc/Scaled-Tech-Consulting/mcp-control-plane.git"], cwd=self.target_root, check=True)
# Create simulated foreign repository root
self.evil_root = os.path.join(self.tmp_dir, "Evil-Repo")
os.makedirs(self.evil_root)
subprocess.run(["git", "init", "-b", "master"], cwd=self.evil_root, check=True)
subprocess.run(["git", "config", "user.email", "[email protected]"], cwd=self.evil_root, check=True)
subprocess.run(["git", "config", "user.name", "Evil User"], cwd=self.evil_root, check=True)
with open(os.path.join(self.evil_root, "README.md"), "w") as f:
f.write("evil\n")
subprocess.run(["git", "add", "README.md"], cwd=self.evil_root, check=True)
subprocess.run(["git", "commit", "-m", "initial"], cwd=self.evil_root, check=True)
subprocess.run(["git", "remote", "add", "prgs", "https://gitea.prgs.cc/Someone-Else/Evil-Repo.git"], cwd=self.evil_root, check=True)
# Create branches/ directory and a valid registered worktree in target repository
self.target_branches = os.path.join(self.target_root, "branches")
self.target_worktree = os.path.join(self.target_branches, "rev-pr-99")
subprocess.run(["git", "worktree", "add", "-b", "rev-pr-99", self.target_worktree], cwd=self.target_root, check=True)
def tearDown(self):
self._tmp.cleanup()
def test_valid_same_repository_configuration(self):
ctx = nwb.resolve_namespace_mutation_context(
role_kind="reviewer",
worktree_path=None,
process_project_root=self.install_root,
env={},
configured_canonical_root=None,
)
self.assertEqual(ctx["canonical_repo_root"], self.install_root)
self.assertTrue(ctx["roots_aligned"])
self.assertTrue(ctx["canonical_root_assessment"]["proven"])
def test_valid_cross_repo_canonical_root(self):
ctx = nwb.resolve_namespace_mutation_context(
role_kind="reviewer",
worktree_path=self.target_worktree,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.target_root,
expected_slug="Scaled-Tech-Consulting/mcp-control-plane",
)
self.assertEqual(ctx["canonical_repo_root"], self.target_root)
self.assertTrue(ctx["roots_aligned"])
self.assertTrue(ctx["canonical_root_assessment"]["proven"])
def test_expected_repository_identity_match(self):
assessment = crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="test",
expected_slug="Scaled-Tech-Consulting/mcp-control-plane",
process_project_root=self.install_root,
remote="prgs",
)
self.assertTrue(assessment["proven"])
self.assertFalse(assessment["block"])
def test_foreign_repository_identity_mismatch(self):
assessment = crr.assess_canonical_repository_root(
configured_value=self.evil_root,
source="test",
expected_slug="Scaled-Tech-Consulting/Gitea-Tools",
process_project_root=self.install_root,
remote="prgs",
)
self.assertFalse(assessment["proven"])
self.assertTrue(assessment["block"])
self.assertTrue(any("identity mismatch" in r for r in assessment["reasons"]))
def test_native_repository_binding_mismatch(self):
ctx = nwb.resolve_namespace_mutation_context(
role_kind="reviewer",
worktree_path=self.target_worktree,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.evil_root,
expected_slug="Scaled-Tech-Consulting/Gitea-Tools",
remote="prgs",
)
self.assertFalse(ctx["roots_aligned"])
self.assertFalse(ctx["canonical_root_assessment"]["proven"])
self.assertTrue(any("identity mismatch" in r for r in ctx["canonical_root_assessment"]["reasons"]))
def test_unpatched_foreign_configured_root_derives_expected_from_process_root_and_blocks(self):
"""B8: Production path test where foreign configured root cannot self-authorize."""
with patch.object(mcp_server, "PROJECT_ROOT", self.install_root), \
patch.object(mcp_server, "_configured_canonical_root", return_value=(self.evil_root, "env")):
expected_slug = mcp_server._resolve_expected_repository_slug("prgs")
self.assertEqual(expected_slug, "Scaled-Tech-Consulting/Gitea-Tools")
assessment = crr.assess_canonical_repository_root(
configured_value=self.evil_root,
source="env",
expected_slug=expected_slug,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
)
self.assertFalse(assessment["proven"])
self.assertTrue(assessment["block"])
self.assertEqual(assessment["resolved_slug"], "Someone-Else/Evil-Repo")
self.assertTrue(any("identity mismatch" in r for r in assessment["reasons"]))
def test_unpatched_valid_cross_repo_matching_session_context(self):
"""B8: Valid cross-repo namespace matches when session context is bound to target repo."""
bound_ctx = {"org": "Scaled-Tech-Consulting", "repository": "mcp-control-plane", "remote": "prgs"}
with patch.object(mcp_server, "PROJECT_ROOT", self.install_root), \
patch.object(mcp_server.session_ctx, "get_session_context", return_value=bound_ctx):
expected_slug = mcp_server._resolve_expected_repository_slug("prgs")
self.assertEqual(expected_slug, "Scaled-Tech-Consulting/mcp-control-plane")
assessment = crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="env",
expected_slug=expected_slug,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
)
self.assertTrue(assessment["proven"])
self.assertFalse(assessment["block"])
self.assertEqual(assessment["resolved_slug"], "Scaled-Tech-Consulting/mcp-control-plane")
def test_unprovable_expected_identity_fails_closed(self):
"""B8: If expected repository identity is unprovable for a configured root, fail closed."""
no_remote_root = os.path.join(self.tmp_dir, "no-remote-process-root")
os.makedirs(no_remote_root)
subprocess.run(["git", "init", "-b", "master"], cwd=no_remote_root, check=True)
with patch.object(mcp_server, "PROJECT_ROOT", no_remote_root), \
patch.object(mcp_server.session_ctx, "get_session_context", return_value=None):
expected_slug = mcp_server._resolve_expected_repository_slug("prgs")
self.assertIsNone(expected_slug)
assessment = crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="env",
expected_slug=expected_slug,
process_project_root=no_remote_root,
remote="prgs",
require_binding=True,
)
self.assertFalse(assessment["proven"])
self.assertTrue(assessment["block"])
self.assertTrue(any("unprovable or missing" in r for r in assessment["reasons"]))
def test_missing_canonical_root(self):
ctx = nwb.resolve_namespace_mutation_context(
role_kind="author",
worktree_path=None,
process_project_root=self.install_root,
env={},
configured_canonical_root="",
)
self.assertEqual(ctx["canonical_repo_root"], self.install_root)
self.assertTrue(ctx["roots_aligned"])
def test_nonexistent_configured_canonical_root(self):
nonexistent = os.path.join(self.tmp_dir, "nonexistent-repo")
ctx = nwb.resolve_namespace_mutation_context(
role_kind="reviewer",
worktree_path=self.target_worktree,
process_project_root=self.install_root,
env={},
configured_canonical_root=nonexistent,
)
self.assertFalse(ctx["roots_aligned"])
self.assertFalse(ctx["canonical_root_assessment"]["proven"])
self.assertTrue(any("does not exist" in r for r in ctx["canonical_root_assessment"]["reasons"]))
def test_non_git_configured_canonical_root(self):
non_git = os.path.join(self.tmp_dir, "non-git-dir")
os.makedirs(non_git)
ctx = nwb.resolve_namespace_mutation_context(
role_kind="reviewer",
worktree_path=self.target_worktree,
process_project_root=self.install_root,
env={},
configured_canonical_root=non_git,
)
self.assertFalse(ctx["roots_aligned"])
self.assertFalse(ctx["canonical_root_assessment"]["proven"])
self.assertTrue(any("not a git repository" in r for r in ctx["canonical_root_assessment"]["reasons"]))
def test_git_common_directory_membership_matching(self):
valid, err = nwb.verify_git_common_directory_membership(
self.target_worktree, self.target_root
)
self.assertTrue(valid, err)
self.assertIsNone(err)
def test_foreign_git_common_directory(self):
# Foreign worktree created under install_root
install_branches = os.path.join(self.install_root, "branches")
foreign_wt = os.path.join(install_branches, "foreign-wt")
subprocess.run(["git", "worktree", "add", "-b", "foreign-wt", foreign_wt], cwd=self.install_root, check=True)
valid, err = nwb.verify_git_common_directory_membership(
foreign_wt, self.target_root
)
self.assertFalse(valid)
self.assertIn("does not match", err)
def test_normalized_path_aliases(self):
alias_path = self.target_worktree + "/../rev-pr-99/./"
valid, err = nwb.verify_git_common_directory_membership(
alias_path, self.target_root
)
self.assertTrue(valid, err)
def test_safe_symlink_identity(self):
link_path = os.path.join(self.target_branches, "symlink-rev-99")
try:
os.symlink(self.target_worktree, link_path)
valid, err = nwb.verify_git_common_directory_membership(
link_path, self.target_root
)
self.assertTrue(valid, err)
finally:
if os.path.exists(link_path):
os.unlink(link_path)
def test_symlink_escape_or_foreign_alias(self):
outside_dir = os.path.join(self.tmp_dir, "outside-target")
os.makedirs(outside_dir)
link_escape = os.path.join(self.target_branches, "escape-link")
try:
os.symlink(outside_dir, link_escape)
assessment = nwb.assess_namespace_mutation_workspace(
role_kind="reviewer",
worktree_path=link_escape,
worktree=None,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.target_root,
expected_slug="Scaled-Tech-Consulting/mcp-control-plane",
)
self.assertTrue(assessment["block"])
finally:
if os.path.exists(link_escape):
os.unlink(link_escape)
def test_reviewer_worktree_registered_and_valid(self):
assessment = nwb.assess_namespace_mutation_workspace(
role_kind="reviewer",
worktree_path=self.target_worktree,
worktree=None,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.target_root,
expected_slug="Scaled-Tech-Consulting/mcp-control-plane",
)
self.assertFalse(assessment["block"])
def test_reviewer_worktree_unregistered_blocks(self):
unreg_wt = os.path.join(self.target_branches, "unregistered-reviewer")
os.makedirs(unreg_wt)
assessment = nwb.assess_namespace_mutation_workspace(
role_kind="reviewer",
worktree_path=unreg_wt,
worktree=None,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.target_root,
expected_slug="Scaled-Tech-Consulting/mcp-control-plane",
)
self.assertTrue(assessment["block"])
self.assertTrue(any("is not registered in git worktree list" in r for r in assessment["reasons"]))
def test_merger_worktree_registered_and_valid(self):
assessment = nwb.assess_namespace_mutation_workspace(
role_kind="merger",
worktree_path=self.target_worktree,
worktree=None,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.target_root,
expected_slug="Scaled-Tech-Consulting/mcp-control-plane",
)
self.assertFalse(assessment["block"])
def test_merger_worktree_unregistered_blocks(self):
unreg_wt = os.path.join(self.target_branches, "unregistered-merger")
os.makedirs(unreg_wt)
assessment = nwb.assess_namespace_mutation_workspace(
role_kind="merger",
worktree_path=unreg_wt,
worktree=None,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.target_root,
expected_slug="Scaled-Tech-Consulting/mcp-control-plane",
)
self.assertTrue(assessment["block"])
self.assertTrue(any("is not registered in git worktree list" in r for r in assessment["reasons"]))
def test_reviewer_or_merger_worktree_outside_branches_blocks(self):
outside_wt = os.path.join(self.target_root, "outside_branches_wt")
os.makedirs(outside_wt)
for r_kind in ("reviewer", "merger"):
assessment = nwb.assess_namespace_mutation_workspace(
role_kind=r_kind,
worktree_path=outside_wt,
worktree=None,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.target_root,
expected_slug="Scaled-Tech-Consulting/mcp-control-plane",
)
self.assertTrue(assessment["block"])
self.assertTrue(any("is not under" in r for r in assessment["reasons"]))
def test_nonexistent_reviewer_or_merger_worktree_blocks(self):
nonexistent_wt = os.path.join(self.target_branches, "nonexistent-wt")
for r_kind in ("reviewer", "merger"):
assessment = nwb.assess_namespace_mutation_workspace(
role_kind=r_kind,
worktree_path=nonexistent_wt,
worktree=None,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.target_root,
expected_slug="Scaled-Tech-Consulting/mcp-control-plane",
)
self.assertTrue(assessment["block"])
self.assertTrue(any("does not exist" in r for r in assessment["reasons"]))
def test_safe_and_unsafe_mutation_alignment_outcomes(self):
# Safe alignment (same repo)
report_safe = stable_control_runtime.build_runtime_report(
process_root=self.install_root,
checkout_branch="master",
runtime_head="abcdef123456",
active_task_workspace=self.install_root,
canonical_repository_root=self.install_root,
workspace_roots_aligned=True,
)
gate_safe = stable_control_runtime.assess_runtime_mutation_gate(report_safe)
self.assertFalse(gate_safe["block"])
# Unsafe alignment
report_unsafe = stable_control_runtime.build_runtime_report(
process_root=self.install_root,
checkout_branch="master",
runtime_head="abcdef123456",
active_task_workspace=self.target_worktree,
canonical_repository_root=self.target_root,
workspace_roots_aligned=False,
)
gate_unsafe = stable_control_runtime.assess_runtime_mutation_gate(report_unsafe)
self.assertTrue(gate_unsafe["block"])
def test_reviewer_lease_lifecycle_production_path(self):
"""B5: Automated regression for reviewer lease acquire and release through production path."""
mock_whoami = {
"authenticated": True,
"username": "sysadmin",
"remote": "prgs",
"profile": {
"profile_name": "prgs-reviewer",
"role": "reviewer",
"role_kind": "reviewer",
"allowed_operations": ["gitea.read", "gitea.pr.comment", "gitea.pr.approve", "gitea.pr.request_changes"],
"forbidden_operations": [],
},
}
def mock_api_request(method, url, auth=None, json_data=None):
if method == "GET":
return {"number": 99, "head": {"sha": "abc1234"}, "state": "open", "merged": False, "merged_at": None}
elif method == "POST":
return {"id": 9999, "body": (json_data or {}).get("body", "")}
return {}
with patch.object(mcp_server, "gitea_whoami", return_value=mock_whoami), \
patch.object(mcp_server, "get_profile", return_value=mock_whoami["profile"]), \
patch.object(mcp_server, "_effective_workspace_role", return_value="reviewer"), \
patch.object(mcp_server, "_configured_canonical_root", return_value=(self.target_root, "env")), \
patch.object(mcp_server, "_reviewer_session_worktree", return_value=self.target_worktree), \
patch.object(mcp_server, "_auth", return_value="token mock-token"), \
patch.object(mcp_server, "_fetch_pr_comments", return_value=[]), \
patch.object(mcp_server, "api_request", side_effect=mock_api_request), \
patch("reviewer_pr_lease.assess_acquire_lease", return_value={"acquire_allowed": True, "reasons": [], "lease_body": "<!-- LEASE -->"}), \
patch("reviewer_pr_lease.find_active_reviewer_lease", return_value={"session_id": "sid-123", "reviewer": "sysadmin"}), \
patch("reviewer_pr_lease.get_session_lease", return_value={"session_id": "sid-123", "reviewer": "sysadmin"}), \
patch("reviewer_pr_lease.clear_session_lease") as mock_clear:
acq_res = mcp_server.gitea_acquire_reviewer_pr_lease(
pr_number=99,
remote="prgs",
worktree=self.target_worktree,
org="Scaled-Tech-Consulting",
repo="mcp-control-plane",
)
self.assertTrue(acq_res.get("success"), acq_res)
rel_res = mcp_server.gitea_release_reviewer_pr_lease(
pr_number=99,
worktree=self.target_worktree,
remote="prgs",
org="Scaled-Tech-Consulting",
repo="mcp-control-plane",
)
self.assertTrue(rel_res.get("success"), rel_res)
mock_clear.assert_called_once()
def test_unbound_session_derives_and_retains_configured_target_repository(self):
"""B9: Unbound session with a legitimate configured target root derives and retains target repository."""
prof = {
"profile_name": "prgs-author",
"allowed_repositories": ["Scaled-Tech-Consulting/mcp-control-plane"],
"canonical_repository_root": self.target_root,
}
with patch.object(mcp_server, "PROJECT_ROOT", self.install_root), \
patch.object(mcp_server.session_ctx, "get_session_context", return_value=None):
slug, reasons = mcp_server._canonical_repository_slug(prof, "prgs")
self.assertEqual(slug, "Scaled-Tech-Consulting/mcp-control-plane", f"reasons: {reasons}")
self.assertEqual(reasons, [])
res = mcp_server._trusted_session_repository(prof, "prgs")
self.assertEqual(res["org"], "Scaled-Tech-Consulting")
self.assertEqual(res["repository"], "mcp-control-plane")
def test_process_root_a_configured_target_b_retains_b(self):
"""B9: Process root A (Gitea-Tools) and configured target B (mcp-control-plane) keep B as expected identity when bound."""
bound_ctx = {"org": "Scaled-Tech-Consulting", "repository": "mcp-control-plane", "remote": "prgs"}
with patch.object(mcp_server, "PROJECT_ROOT", self.install_root), \
patch.object(mcp_server.session_ctx, "get_session_context", return_value=bound_ctx):
expected = mcp_server._resolve_expected_repository_slug("prgs")
self.assertEqual(expected, "Scaled-Tech-Consulting/mcp-control-plane")
def test_request_supplied_repository_c_cannot_replace_derived_or_bound_repository_b(self):
"""B9: Request-supplied repository C (Timesheet) cannot replace bound repository B (mcp-control-plane)."""
bound_ctx = {"org": "Scaled-Tech-Consulting", "repository": "mcp-control-plane", "remote": "prgs"}
with patch.object(mcp_server, "PROJECT_ROOT", self.install_root), \
patch.object(mcp_server.session_ctx, "get_session_context", return_value=bound_ctx):
expected = mcp_server._resolve_expected_repository_slug("prgs", org="Scaled-Tech-Consulting", repo="Timesheet")
self.assertEqual(expected, "Scaled-Tech-Consulting/mcp-control-plane")
def test_known_expected_repo_a_plus_candidate_root_b_fails_closed(self):
"""B9: Known expected repository A plus candidate root B fails closed in validation mode."""
assessment = crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="env",
expected_slug="Scaled-Tech-Consulting/Gitea-Tools",
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode="validation",
)
self.assertFalse(assessment["proven"])
self.assertTrue(assessment["block"])
self.assertTrue(any("identity mismatch" in r for r in assessment["reasons"]))
def test_validation_mode_with_unprovable_expected_identity_fails_closed(self):
"""B9: Validation mode with expected_slug=None and require_binding=True fails closed."""
assessment = crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="env",
expected_slug=None,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode="validation",
)
self.assertFalse(assessment["proven"])
self.assertTrue(assessment["block"])
self.assertTrue(any("unprovable or missing" in r for r in assessment["reasons"]))
if __name__ == "__main__":
unittest.main()
@@ -0,0 +1,759 @@
"""Client-identity env recognition and the production heartbeat lifecycle (#975).
Two defects left behind by #948, covered here against the issue's acceptance
criteria:
1. ``GITEA_MCP_CLIENT`` / ``GITEA_MCP_CLIENT_INSTANCE`` /
``GITEA_MCP_CLIENT_SESSION`` were consumed by production while absent from
the recognized-key allowlist, so the peer-env scan classified them as
unsupported overrides and the capability resolver refused every mutation
fleet-wide.
2. ``WorkerRegistry.heartbeat()`` had no production caller, so
``last_heartbeat_at`` never left ``started_at`` and ``heartbeat_ttl_seconds``
became a hard cap on attachment rather than a liveness window.
Every TTL assertion here uses an injected clock (AC13). No test waits for a
real TTL. The one test that exercises the real thread loop uses a
millisecond-scale interval and polls with a bounded timeout, never a TTL.
All client, session, and generation identifiers are synthetic.
"""
from __future__ import annotations
import os
import tempfile
import threading
import time
import unittest
from datetime import datetime, timedelta, timezone
import gitea_config
import mcp_worker_identity as mwi
NOW = datetime(2026, 7, 30, 0, 40, 9, tzinfo=timezone.utc)
TTL = 900.0
#: The three keys production reads as client identity.
CLIENT_IDENTITY_ENV_KEYS = (
"GITEA_MCP_CLIENT",
"GITEA_MCP_CLIENT_INSTANCE",
"GITEA_MCP_CLIENT_SESSION",
)
def _registry() -> mwi.WorkerRegistry:
"""A registry on a throwaway path; never the operator's real one."""
handle, path = tempfile.mkstemp(suffix=".sqlite3")
os.close(handle)
os.unlink(path)
return mwi.WorkerRegistry(path)
def _attach(
registry: mwi.WorkerRegistry,
*,
client: str = "claude_code",
session: str = "sess-0001",
generation: str = "gen-0001",
profile: str = "prgs-author",
role: str = "author",
pid: int = 4242,
now: datetime = NOW,
ttl: float = TTL,
) -> dict:
"""Register one synthetic worker and return the outcome plus its identity."""
identity = mwi.generate_worker_identity(client, session, now=now)
outcome = registry.register(
worker_identity=identity,
client_name=client,
client_instance_id=f"inst-{session}",
session_id=session,
generation_id=generation,
role=role,
profile=profile,
pid=pid,
heartbeat_ttl_seconds=ttl,
now=now,
)
outcome["identity"] = identity
return outcome
def _supervisor(
registry: mwi.WorkerRegistry,
outcome: dict,
*,
client: str = "claude_code",
session: str = "sess-0001",
generation: str = "gen-0001",
pid: int = 4242,
ttl: float = TTL,
interval_seconds: float | None = None,
clock=None,
) -> mwi.WorkerHeartbeatSupervisor:
"""A supervisor bound to *outcome*'s registration. Never started here."""
return mwi.WorkerHeartbeatSupervisor(
registry,
worker_identity=outcome["identity"],
fencing_epoch=outcome["fencing_epoch"],
session_id=session,
generation_id=generation,
client_name=client,
pid=pid,
ttl_seconds=ttl,
interval_seconds=interval_seconds,
clock=clock,
)
class _ManualClock:
"""Controlled time. Nothing in this suite sleeps to reach a TTL."""
def __init__(self, start: datetime = NOW) -> None:
self.now = start
def __call__(self) -> datetime:
return self.now
def advance(self, seconds: float) -> datetime:
self.now = self.now + timedelta(seconds=seconds)
return self.now
# --- AC1-AC3: environment recognition ------------------------------------
class ClientIdentityEnvRecognitionTests(unittest.TestCase):
"""AC1-AC3: the keys production consumes are the keys the gate accepts."""
def test_each_client_identity_key_is_recognized(self):
# AC1: individually asserted, so a partial fix cannot pass.
for key in CLIENT_IDENTITY_ENV_KEYS:
with self.subTest(key=key):
self.assertIn(
key,
gitea_config.RECOGNIZED_GITEA_ENV_KEYS,
f"{key} is read by the server but not recognized by the gate",
)
def test_recognized_client_identity_keys_are_not_unconsumed(self):
# AC3 at the classification layer the resolver reads: with only these
# keys set there is no unsupported override to raise a blocker from.
env = {
"GITEA_MCP_CLIENT": "claude_code",
"GITEA_MCP_CLIENT_INSTANCE": "inst-0001",
"GITEA_MCP_CLIENT_SESSION": "sess-0001",
}
self.assertEqual(gitea_config.get_unconsumed_gitea_env_overrides(env), {})
def test_unknown_gitea_key_is_still_unsupported(self):
# AC2: the gate is not broadened. An unrecognised override is refused
# exactly as before, and the recognised siblings do not shield it.
env = {
"GITEA_MCP_CLIENT": "claude_code",
"GITEA_TOTALLY_INVENTED_OVERRIDE": "1",
}
unconsumed = gitea_config.get_unconsumed_gitea_env_overrides(env)
self.assertIn("GITEA_TOTALLY_INVENTED_OVERRIDE", unconsumed)
self.assertNotIn("GITEA_MCP_CLIENT", unconsumed)
def test_no_new_prefix_admits_arbitrary_client_keys(self):
# A prefix would have been the lazy fix; it would admit anything
# beginning GITEA_MCP_CLIENT, including typos and future unknowns.
env = {"GITEA_MCP_CLIENT_SOMETHING_ELSE": "x"}
self.assertIn(
"GITEA_MCP_CLIENT_SOMETHING_ELSE",
gitea_config.get_unconsumed_gitea_env_overrides(env),
)
def test_server_consumes_exactly_the_recognized_key_names(self):
# Guards the actual defect class: the constants production reads and the
# allowlist drifting apart again.
import gitea_mcp_server as mcp_server
for name in (
mcp_server.CLIENT_NAME_ENV,
mcp_server.CLIENT_INSTANCE_ENV,
mcp_server.CLIENT_SESSION_ENV,
):
with self.subTest(env_key=name):
self.assertIn(name, gitea_config.RECOGNIZED_GITEA_ENV_KEYS)
class RuntimeDiagnosticReasonPropagationTests(unittest.TestCase):
"""`unsupported-env:` must not be swallowed where `stale-runtime:` is not."""
def test_both_reason_families_are_hard_prefixes(self):
import gitea_mcp_server as mcp_server
self.assertIn("stale-runtime:", mcp_server.RUNTIME_DIAGNOSTIC_HARD_PREFIXES)
self.assertIn("unsupported-env:", mcp_server.RUNTIME_DIAGNOSTIC_HARD_PREFIXES)
def test_unsupported_env_message_is_recognized_as_hard(self):
import gitea_mcp_server as mcp_server
message = (
"unsupported-env: Unsupported GITEA_* environment variable override(s) "
"detected: GITEA_INVENTED=1. Unknown env overrides are unsupported."
)
self.assertTrue(
any(p in message for p in mcp_server.RUNTIME_DIAGNOSTIC_HARD_PREFIXES)
)
# --- Interval policy ------------------------------------------------------
class HeartbeatIntervalPolicyTests(unittest.TestCase):
"""The interval must always leave room for lost beats inside the TTL."""
def test_default_interval_is_one_third_of_ttl(self):
self.assertAlmostEqual(mwi.heartbeat_interval_for(900.0), 300.0)
def test_interval_is_capped_at_half_the_ttl(self):
# An override may not exceed the ceiling, or it would expire the very
# registration it exists to renew.
self.assertAlmostEqual(mwi.heartbeat_interval_for(900.0, 5000), 450.0)
def test_override_below_the_ceiling_is_honoured(self):
self.assertAlmostEqual(mwi.heartbeat_interval_for(900.0, 120), 120.0)
def test_unparsable_or_nonpositive_override_falls_back_to_policy(self):
for bad in ("", "abc", "0", "-5", None):
with self.subTest(override=bad):
self.assertAlmostEqual(mwi.heartbeat_interval_for(900.0, bad), 300.0)
def test_interval_is_always_strictly_below_the_ttl(self):
for ttl in (1.0, 30.0, 900.0, 14400.0):
with self.subTest(ttl=ttl):
self.assertLess(mwi.heartbeat_interval_for(ttl), ttl)
# --- AC4-AC6: a healthy worker stays owned --------------------------------
class HealthyWorkerRemainsOwnedTests(unittest.TestCase):
"""AC4-AC6: the TTL becomes a liveness window instead of an attachment cap."""
def test_unsupervised_worker_goes_unowned_at_the_ttl(self):
# The regression itself, asserted so the fix cannot be mistaken for a
# TTL change: without a beater the row still expires exactly as before.
registry = _registry()
outcome = _attach(registry)
record = registry.get(outcome["identity"])
self.assertEqual(record["started_at"], record["last_heartbeat_at"])
after_ttl = NOW + timedelta(seconds=TTL + 1)
liveness = registry.is_live(record, now=after_ttl, pid_alive=True)
self.assertFalse(liveness["live"])
assessment = mwi.assess_provenance(
registry=registry,
worker_identity=outcome["identity"],
generation_id="gen-0001",
env={"GITEA_CLIENT_MANAGED": "1"},
native_transport_bound=True,
now=after_ttl,
pid_alive_probe=lambda _pid: True,
)
self.assertEqual(assessment["session_ownership"], "unowned")
self.assertEqual(assessment["blocker_kind"], mwi.BLOCKER_NO_ATTACHMENT)
def test_supervised_worker_stays_owned_past_two_full_ttls(self):
# AC4: beyond two complete TTL periods of simulated time.
registry = _registry()
outcome = _attach(registry)
clock = _ManualClock()
supervisor = _supervisor(registry, outcome, clock=clock)
interval = supervisor.interval_seconds
self.assertLess(interval, TTL)
elapsed = 0.0
target = TTL * 2 + interval
while elapsed < target:
clock.advance(interval)
elapsed += interval
result = supervisor.beat_once()
self.assertTrue(result["renewed"], result.get("reasons"))
assessment = mwi.assess_provenance(
registry=registry,
worker_identity=outcome["identity"],
generation_id="gen-0001",
env={"GITEA_CLIENT_MANAGED": "1"},
native_transport_bound=True,
now=clock.now,
pid_alive_probe=lambda _pid: True,
)
self.assertEqual(
assessment["session_ownership"],
"owned",
f"ownership lost after {elapsed:.0f}s of simulated time",
)
self.assertGreater(elapsed, TTL * 2)
def test_last_heartbeat_advances_and_diverges_from_started_at(self):
# AC5.
registry = _registry()
outcome = _attach(registry)
clock = _ManualClock()
supervisor = _supervisor(registry, outcome, clock=clock)
before = registry.get(outcome["identity"])
clock.advance(300)
supervisor.beat_once()
after = registry.get(outcome["identity"])
self.assertEqual(before["started_at"], after["started_at"])
self.assertNotEqual(after["started_at"], after["last_heartbeat_at"])
self.assertGreater(after["last_heartbeat_at"], before["last_heartbeat_at"])
def test_repeated_heartbeats_create_no_duplicate_rows(self):
# AC6: heartbeating is not a second registration path.
registry = _registry()
outcome = _attach(registry)
clock = _ManualClock()
supervisor = _supervisor(registry, outcome, clock=clock)
for _ in range(5):
clock.advance(300)
self.assertTrue(supervisor.beat_once()["renewed"])
workers = registry.list_workers()
self.assertEqual(len(workers), 1)
self.assertEqual(workers[0]["worker_identity"], outcome["identity"])
self.assertEqual(supervisor.status()["beats_renewed"], 5)
# --- AC7, AC8, AC12: only the owner may renew -----------------------------
class HeartbeatOwnershipFencingTests(unittest.TestCase):
"""AC7/AC8/AC12: one worker can never refresh another worker's row."""
def test_two_client_names_stay_distinct_rows_and_identities(self):
# AC7.
registry = _registry()
first = _attach(registry, client="claude_code", session="sess-a")
second = _attach(registry, client="gemini", session="sess-b")
self.assertNotEqual(first["identity"], second["identity"])
rows = {w["worker_identity"]: w for w in registry.list_workers()}
self.assertEqual(len(rows), 2)
self.assertEqual(rows[first["identity"]]["client_name"], "claude_code")
self.assertEqual(rows[second["identity"]]["client_name"], "gemini")
def test_wrong_session_cannot_refresh_the_row(self):
# AC8, session dimension.
registry = _registry()
outcome = _attach(registry, session="sess-real")
result = registry.heartbeat(
worker_identity=outcome["identity"],
fencing_epoch=outcome["fencing_epoch"],
expected_session_id="sess-impostor",
)
self.assertFalse(result["renewed"])
self.assertEqual(result["blocker_kind"], mwi.BLOCKER_FENCED)
self.assertIn(
"session_id", [field for field, _r, _p in result["expectation_drift"]]
)
def test_wrong_generation_cannot_refresh_the_row(self):
# AC8, generation dimension.
registry = _registry()
outcome = _attach(registry, generation="gen-real")
result = registry.heartbeat(
worker_identity=outcome["identity"],
fencing_epoch=outcome["fencing_epoch"],
expected_generation_id="gen-other",
)
self.assertFalse(result["renewed"])
self.assertEqual(result["blocker_kind"], mwi.BLOCKER_FENCED)
def test_wrong_client_name_cannot_refresh_the_row(self):
registry = _registry()
outcome = _attach(registry, client="claude_code")
result = registry.heartbeat(
worker_identity=outcome["identity"],
fencing_epoch=outcome["fencing_epoch"],
expected_client_name="gemini",
)
self.assertFalse(result["renewed"])
self.assertEqual(result["blocker_kind"], mwi.BLOCKER_FENCED)
def test_wrong_pid_cannot_refresh_the_row(self):
registry = _registry()
outcome = _attach(registry, pid=4242)
result = registry.heartbeat(
worker_identity=outcome["identity"],
fencing_epoch=outcome["fencing_epoch"],
expected_pid=9999,
)
self.assertFalse(result["renewed"])
self.assertEqual(result["blocker_kind"], mwi.BLOCKER_FENCED)
def test_namespaces_of_one_application_share_a_client_name(self):
# Several namespaces of one application must stay one client, while
# separate applications stay distinct. Aliases normalise, and the
# expectation check compares normalised values, so an author and a
# reviewer namespace of the same product are not treated as impostors.
registry = _registry()
outcome = _attach(registry, client="claude")
self.assertEqual(
registry.get(outcome["identity"])["client_name"], "claude_code"
)
result = registry.heartbeat(
worker_identity=outcome["identity"],
fencing_epoch=outcome["fencing_epoch"],
expected_client_name="claude_desktop",
)
self.assertTrue(result["renewed"], result.get("reasons"))
def test_matching_expectations_renew_normally(self):
registry = _registry()
outcome = _attach(registry)
result = registry.heartbeat(
worker_identity=outcome["identity"],
fencing_epoch=outcome["fencing_epoch"],
expected_session_id="sess-0001",
expected_generation_id="gen-0001",
expected_client_name="claude_code",
expected_pid=4242,
)
self.assertTrue(result["renewed"], result.get("reasons"))
def test_omitted_expectations_preserve_pre_fix_behaviour(self):
# AC12: the #948 contract is unchanged for callers presenting nothing.
registry = _registry()
outcome = _attach(registry)
self.assertTrue(
registry.heartbeat(
worker_identity=outcome["identity"],
fencing_epoch=outcome["fencing_epoch"],
)["renewed"]
)
def test_stale_fencing_epoch_is_still_refused(self):
# AC12: the pre-existing epoch fence is untouched.
registry = _registry()
outcome = _attach(registry)
result = registry.heartbeat(
worker_identity=outcome["identity"],
fencing_epoch=int(outcome["fencing_epoch"]) + 1,
)
self.assertFalse(result["renewed"])
self.assertEqual(result["blocker_kind"], mwi.BLOCKER_FENCED)
def test_fenced_supervisor_stops_permanently(self):
# A fenced session must never beat its way back into ownership.
registry = _registry()
outcome = _attach(registry)
supervisor = mwi.WorkerHeartbeatSupervisor(
registry,
worker_identity=outcome["identity"],
fencing_epoch=int(outcome["fencing_epoch"]) + 1,
session_id="sess-0001",
generation_id="gen-0001",
client_name="claude_code",
pid=4242,
ttl_seconds=TTL,
clock=_ManualClock(),
)
first = supervisor.beat_once()
self.assertFalse(first["renewed"])
self.assertIsNotNone(supervisor.stopped_reason)
# A second attempt does not even reach the registry.
second = supervisor.beat_once()
self.assertFalse(second["renewed"])
self.assertFalse(second["beat_attempted"])
self.assertEqual(supervisor.status()["beats_attempted"], 1)
def test_missing_registration_stops_the_supervisor(self):
registry = _registry()
supervisor = mwi.WorkerHeartbeatSupervisor(
registry,
worker_identity=mwi.generate_worker_identity("grok", "sess-gone", now=NOW),
fencing_epoch=1,
ttl_seconds=TTL,
clock=_ManualClock(),
)
result = supervisor.beat_once()
self.assertEqual(result["blocker_kind"], mwi.BLOCKER_NO_ATTACHMENT)
self.assertIsNotNone(supervisor.stopped_reason)
# --- AC9, AC10, AC11: stopping and expiry ---------------------------------
class ShutdownAndExpiryTests(unittest.TestCase):
"""AC9-AC11: orderly shutdown stops beating; dead workers still expire."""
def test_stop_ends_the_heartbeat(self):
# AC9.
registry = _registry()
outcome = _attach(registry)
clock = _ManualClock()
supervisor = _supervisor(registry, outcome, clock=clock)
clock.advance(300)
self.assertTrue(supervisor.beat_once()["renewed"])
supervisor.stop(reason="orderly shutdown")
self.assertFalse(supervisor.status()["running"])
self.assertEqual(supervisor.stopped_reason, "orderly shutdown")
clock.advance(300)
after_stop = supervisor.beat_once()
self.assertFalse(after_stop["renewed"])
self.assertFalse(after_stop["beat_attempted"])
def test_stop_is_idempotent_and_keeps_the_first_reason(self):
registry = _registry()
outcome = _attach(registry)
supervisor = _supervisor(registry, outcome, clock=_ManualClock())
supervisor.stop(reason="orderly shutdown")
supervisor.stop(reason="second call")
self.assertEqual(supervisor.stopped_reason, "orderly shutdown")
def test_stopped_worker_becomes_unowned_once_the_ttl_elapses(self):
# AC10: stopping beating restores normal expiration; the fix does not
# make a worker immortal.
registry = _registry()
outcome = _attach(registry)
clock = _ManualClock()
supervisor = _supervisor(registry, outcome, clock=clock)
clock.advance(300)
supervisor.beat_once()
supervisor.stop(reason="orderly shutdown")
last_beat = clock.now
clock.advance(TTL + 1)
record = registry.get(outcome["identity"])
liveness = registry.is_live(record, now=clock.now, pid_alive=True)
self.assertFalse(liveness["live"])
self.assertGreater(liveness["heartbeat_age_seconds"], TTL)
assessment = mwi.assess_provenance(
registry=registry,
worker_identity=outcome["identity"],
generation_id="gen-0001",
env={"GITEA_CLIENT_MANAGED": "1"},
native_transport_bound=True,
now=clock.now,
pid_alive_probe=lambda _pid: True,
)
self.assertEqual(assessment["session_ownership"], "unowned")
self.assertLess(last_beat, clock.now)
def test_dead_worker_expires_even_with_a_live_pid(self):
# A recycled or long-lived pid must not keep a silent worker owned.
registry = _registry()
outcome = _attach(registry)
record = registry.get(outcome["identity"])
liveness = registry.is_live(
record, now=NOW + timedelta(seconds=TTL + 1), pid_alive=True
)
self.assertFalse(liveness["live"])
self.assertFalse(liveness["heartbeat_fresh"])
def test_one_worker_expiring_does_not_stale_healthy_workers(self):
# AC11.
registry = _registry()
healthy = _attach(registry, client="claude_code", session="sess-live")
abandoned = _attach(registry, client="gemini", session="sess-dead")
clock = _ManualClock()
supervisor = _supervisor(
registry,
healthy,
client="claude_code",
session="sess-live",
clock=clock,
)
elapsed = 0.0
while elapsed < TTL + 300:
clock.advance(supervisor.interval_seconds)
elapsed += supervisor.interval_seconds
self.assertTrue(supervisor.beat_once()["renewed"])
healthy_live = registry.is_live(
registry.get(healthy["identity"]), now=clock.now, pid_alive=True
)
abandoned_live = registry.is_live(
registry.get(abandoned["identity"]), now=clock.now, pid_alive=True
)
self.assertTrue(healthy_live["live"])
self.assertFalse(abandoned_live["live"])
# --- Failure handling and observability -----------------------------------
class HeartbeatFailureHandlingTests(unittest.TestCase):
"""A heartbeat failure never grants ownership and never crashes a caller."""
def test_transient_exception_is_counted_and_beating_continues(self):
registry = _registry()
outcome = _attach(registry)
clock = _ManualClock()
supervisor = _supervisor(registry, outcome, clock=clock)
calls = {"n": 0}
real = registry.heartbeat
def flaky(**kwargs):
calls["n"] += 1
if calls["n"] == 1:
raise RuntimeError("database is locked")
return real(**kwargs)
registry.heartbeat = flaky # type: ignore[method-assign]
try:
clock.advance(300)
first = supervisor.beat_once()
self.assertFalse(first["renewed"])
self.assertTrue(first["transient"])
self.assertIsNone(supervisor.stopped_reason)
clock.advance(300)
second = supervisor.beat_once()
self.assertTrue(second["renewed"], second.get("reasons"))
finally:
registry.heartbeat = real # type: ignore[method-assign]
status = supervisor.status()
self.assertEqual(status["transient_failures"], 1)
self.assertEqual(status["beats_renewed"], 1)
self.assertTrue(status["supervised"])
def test_status_reports_an_unstarted_supervisor_honestly(self):
registry = _registry()
outcome = _attach(registry)
status = _supervisor(registry, outcome, clock=_ManualClock()).status()
self.assertTrue(status["supervised"])
self.assertFalse(status["started"])
self.assertFalse(status["running"])
self.assertEqual(status["beats_renewed"], 0)
self.assertIsNone(status["stopped_reason"])
def test_runtime_context_reports_unsupervised_when_none_attached(self):
# The observability surface must distinguish "no supervisor" from
# "supervisor with no beats yet", or the regression is invisible again.
import gitea_mcp_server as mcp_server
saved = mcp_server._WORKER_HEARTBEAT_SUPERVISOR
mcp_server._WORKER_HEARTBEAT_SUPERVISOR = None
try:
status = mcp_server._worker_heartbeat_status()
finally:
mcp_server._WORKER_HEARTBEAT_SUPERVISOR = saved
self.assertFalse(status["supervised"])
self.assertFalse(status["running"])
self.assertTrue(status["reasons"])
def test_runtime_context_surfaces_an_attached_supervisor(self):
import gitea_mcp_server as mcp_server
registry = _registry()
outcome = _attach(registry)
supervisor = _supervisor(registry, outcome, clock=_ManualClock())
saved = mcp_server._WORKER_HEARTBEAT_SUPERVISOR
mcp_server._WORKER_HEARTBEAT_SUPERVISOR = supervisor
try:
status = mcp_server._worker_heartbeat_status()
finally:
mcp_server._WORKER_HEARTBEAT_SUPERVISOR = saved
self.assertTrue(status["supervised"])
self.assertEqual(status["worker_identity"], outcome["identity"])
self.assertEqual(status["heartbeat_ttl_seconds"], TTL)
self.assertLess(status["heartbeat_interval_seconds"], TTL)
def test_status_failure_degrades_instead_of_raising(self):
import gitea_mcp_server as mcp_server
class Exploding:
def status(self):
raise RuntimeError("boom")
saved = mcp_server._WORKER_HEARTBEAT_SUPERVISOR
mcp_server._WORKER_HEARTBEAT_SUPERVISOR = Exploding()
try:
status = mcp_server._worker_heartbeat_status()
finally:
mcp_server._WORKER_HEARTBEAT_SUPERVISOR = saved
self.assertTrue(status["supervised"])
self.assertFalse(status["running"])
self.assertTrue(status["reasons"])
class SupervisorThreadLoopTests(unittest.TestCase):
"""The real thread, on a millisecond interval. Never waits for a TTL."""
def test_thread_beats_repeatedly_then_stops_on_request(self):
registry = _registry()
outcome = _attach(registry)
beats = threading.Event()
seen = {"n": 0}
real = registry.heartbeat
def counting(**kwargs):
result = real(**kwargs)
seen["n"] += 1
if seen["n"] >= 3:
beats.set()
return result
registry.heartbeat = counting # type: ignore[method-assign]
supervisor = _supervisor(registry, outcome, ttl=1.0, interval_seconds=0.01)
try:
started = supervisor.start()
self.assertTrue(started["started"])
self.assertTrue(
beats.wait(timeout=10.0), "heartbeat thread produced no beats"
)
self.assertTrue(supervisor.status()["running"])
finally:
supervisor.stop(reason="test teardown")
registry.heartbeat = real # type: ignore[method-assign]
self.assertGreaterEqual(supervisor.status()["beats_renewed"], 3)
self.assertFalse(supervisor.status()["running"])
# Stopping is prompt: the loop waits on an Event, not a sleep.
deadline = time.monotonic() + 5.0
while supervisor._thread is not None and supervisor._thread.is_alive():
if time.monotonic() > deadline:
self.fail("heartbeat thread did not exit after stop()")
time.sleep(0.01)
def test_start_is_idempotent(self):
registry = _registry()
outcome = _attach(registry)
supervisor = _supervisor(registry, outcome, ttl=1.0, interval_seconds=0.05)
try:
self.assertFalse(supervisor.start()["already_running"])
self.assertTrue(supervisor.start()["already_running"])
finally:
supervisor.stop(reason="test teardown")
def test_start_refuses_after_a_terminal_stop(self):
registry = _registry()
outcome = _attach(registry)
supervisor = _supervisor(registry, outcome, clock=_ManualClock())
supervisor.stop(reason="orderly shutdown")
self.assertFalse(supervisor.start()["started"])
if __name__ == "__main__":
unittest.main()
+7 -2
View File
@@ -245,10 +245,15 @@ class TestNamespaceWorkspaceIntegration(unittest.TestCase):
def test_pr487_style_merge_binds_clean_merger_workspace(
self, _exists, _isdir, mock_run
):
mock_run.return_value = MagicMock(returncode=0, stdout=f"{CONTROL_ROOT}/.git\n")
def mock_git(cmd, *args, **kwargs):
if "rev-parse" in cmd:
return MagicMock(returncode=0, stdout=f"{CONTROL_ROOT}/.git\n")
return MagicMock(returncode=0, stdout=f"worktree {CONTROL_ROOT}\nworktree {MERGER_CLEAN}\n")
mock_run.side_effect = mock_git
os.environ[nwb.AUTHOR_WORKTREE_ENV] = AUTHOR_DIRTY
os.environ[nwb.MERGER_WORKTREE_ENV] = MERGER_CLEAN
srv._preflight_resolved_role = "reviewer"
with mock.patch.object(srv, "PROJECT_ROOT", MCP_PROCESS_ROOT):
with mock.patch("gitea_mcp_server.get_profile", return_value=self._merger_profile()):
resolved = srv._verify_role_mutation_workspace("prgs")
self.assertEqual(resolved, os.path.realpath(MCP_PROCESS_ROOT))
self.assertEqual(resolved, os.path.realpath(MERGER_CLEAN))
@@ -153,6 +153,12 @@ EXPECTED_ROLE_EXCLUSIVE_TASKS = frozenset(
"delete_branch",
"cleanup_merged_pr_branch",
"reconciliation_cleanup",
# #970 review 644 B2: retiring a missing worktree binding is a
# control-plane cleanup mutation, so it carries the same reconciler-only
# authority as every other reconciliation cleanup. Permission alone must
# not authorize it.
"reconcile_missing_worktree_bindings",
"gitea_reconcile_missing_worktree_bindings",
"work_issue",
"work-issue",
}