Compare commits

..
8 changed files with 473 additions and 196 deletions
+26 -29
View File
@@ -24,6 +24,8 @@ import subprocess
import uuid import uuid
from datetime import datetime, timedelta, timezone from datetime import datetime, timedelta, timezone
from typing import Any from typing import Any
import hashlib
@@ -2258,7 +2260,12 @@ def _seed_session_context(
expected_username=expected, expected_username=expected,
source=source, source=source,
canonical_repository_root=canonical_root_pin, canonical_repository_root=canonical_root_pin,
cohort_id=_COHORT_ID,
startup_sha=_STARTUP_PARITY.get("startup_head"),
endpoint=host or (profile.get("base_url") or "").strip() or None,
config_fingerprint=_CONFIG_FINGERPRINT,
) )
import issue_work_duplicate_gate # noqa: E402 import issue_work_duplicate_gate # noqa: E402
import issue_workflow_labels # noqa: E402 import issue_workflow_labels # noqa: E402
import terminal_pr_label_cleanup # noqa: E402 # #780 status:pr-open terminal rule import terminal_pr_label_cleanup # noqa: E402 # #780 status:pr-open terminal rule
@@ -2287,6 +2294,11 @@ import stable_control_runtime # noqa: E402
# master has advanced past the running code and fail closed until restart. # master has advanced past the running code and fail closed until restart.
# Read-only operations are never blocked by staleness. # Read-only operations are never blocked by staleness.
_STARTUP_PARITY = master_parity_gate.capture_startup_parity(PROJECT_ROOT) _STARTUP_PARITY = master_parity_gate.capture_startup_parity(PROJECT_ROOT)
_COHORT_ID: str = f"cohort-p{os.getpid()}-{_STARTUP_PARITY.get('startup_head') or 'unknown'}"
_CONFIG_FINGERPRINT: str = hashlib.sha256(
(PROJECT_ROOT + str(_STARTUP_PARITY.get("startup_head"))).encode("utf-8")
).hexdigest()[:16]
# Stable-control runtime facts (#615): which runtime this process serves from. # Stable-control runtime facts (#615): which runtime this process serves from.
# These are the *immutable* facts -- process root, branch, head, checkout-ness -- # These are the *immutable* facts -- process root, branch, head, checkout-ness --
@@ -3325,16 +3337,10 @@ def _effective_remote(remote: str) -> str:
return remote return remote
def _resolve( def _resolve(remote: str, host: str | None, org: str | None, repo: str | None):
remote: str,
host: str | None,
org: str | None,
repo: str | None,
for_mutation: bool = False,
):
"""Resolve remote + overrides to (host, org, repo). """Resolve remote + overrides to (host, org, repo).
#714 / #530 / #707: when the caller omits org and/or repo, prefer the #714 / #530: when the caller omits org and/or repo, prefer the
workspace-aligned git remote over ``REMOTES`` defaults (e.g. bare workspace-aligned git remote over ``REMOTES`` defaults (e.g. bare
``remote=prgs`` must not force ``Timesheet`` when the checkout is ``remote=prgs`` must not force ``Timesheet`` when the checkout is
``Gitea-Tools``). Explicit caller org/repo always win and are validated ``Gitea-Tools``). Explicit caller org/repo always win and are validated
@@ -3420,7 +3426,6 @@ def _resolve(
# Workspace-filled sides are intentional alignment for #530. # Workspace-filled sides are intentional alignment for #530.
org_explicit=org_explicit or filled_org, org_explicit=org_explicit or filled_org,
repo_explicit=repo_explicit or filled_repo, repo_explicit=repo_explicit or filled_repo,
for_mutation=for_mutation,
) )
return (resolved_host, resolved_org, resolved_repo) return (resolved_host, resolved_org, resolved_repo)
@@ -3482,9 +3487,8 @@ def _enforce_remote_repo_guard(
*, *,
org_explicit: bool, org_explicit: bool,
repo_explicit: bool, repo_explicit: bool,
for_mutation: bool = False,
) -> None: ) -> None:
"""Fail closed on a remote/repo mismatch vs. the local git remote (#530/#707). """Fail closed on a remote/repo mismatch vs. the local git remote (#530).
Best-effort: bypassed under pytest unless ``GITEA_FORCE_REMOTE_REPO_CHECK`` is Best-effort: bypassed under pytest unless ``GITEA_FORCE_REMOTE_REPO_CHECK`` is
set, so the unit suite (which calls tools with bare remotes against mocked APIs) set, so the unit suite (which calls tools with bare remotes against mocked APIs)
@@ -3496,12 +3500,6 @@ def _enforce_remote_repo_guard(
): ):
return return
local_remote_url = _local_git_remote_url(remote) local_remote_url = _local_git_remote_url(remote)
primary_org = None
primary_repo = None
ctx = session_ctx.get_session_context()
if ctx:
primary_org = ctx.get("org")
primary_repo = ctx.get("repository")
assessment = remote_repo_guard.assess_remote_repo_match( assessment = remote_repo_guard.assess_remote_repo_match(
remote=remote, remote=remote,
resolved_org=resolved_org, resolved_org=resolved_org,
@@ -3509,9 +3507,6 @@ def _enforce_remote_repo_guard(
local_remote_url=local_remote_url, local_remote_url=local_remote_url,
org_explicit=org_explicit, org_explicit=org_explicit,
repo_explicit=repo_explicit, repo_explicit=repo_explicit,
for_mutation=for_mutation,
primary_org=primary_org,
primary_repo=primary_repo,
) )
if assessment["block"]: if assessment["block"]:
raise RuntimeError(remote_repo_guard.format_remote_repo_guard_error(assessment)) raise RuntimeError(remote_repo_guard.format_remote_repo_guard_error(assessment))
@@ -4080,7 +4075,7 @@ def gitea_lock_issue(
resolved_worktree = issue_lock_worktree.resolve_author_worktree_path( resolved_worktree = issue_lock_worktree.resolve_author_worktree_path(
worktree_path, _canonical_local_git_root() worktree_path, _canonical_local_git_root()
) )
h, o, r = _resolve(remote, host, org, repo, for_mutation=True) h, o, r = _resolve(remote, host, org, repo)
existing_issue_lock = _load_existing_issue_lock( existing_issue_lock = _load_existing_issue_lock(
remote=remote, org=o, repo=r, issue_number=issue_number remote=remote, org=o, repo=r, issue_number=issue_number
) )
@@ -5256,7 +5251,7 @@ def gitea_create_pr(
org=org, org=org,
repo=repo, repo=repo,
) )
h, o, r = _resolve(remote, host, org, repo, for_mutation=True) h, o, r = _resolve(remote, host, org, repo)
# ── Issue Lock Validation (Issue #194 / #196 / #443) ── # ── Issue Lock Validation (Issue #194 / #196 / #443) ──
lock_data = _resolve_issue_lock_for_pr(remote=remote, org=o, repo=r, head=head) lock_data = _resolve_issue_lock_for_pr(remote=remote, org=o, repo=r, head=head)
@@ -9577,7 +9572,7 @@ def gitea_edit_pr(
required_permission="gitea.pr.close", required_permission="gitea.pr.close",
) )
h, o, r = _resolve(remote, host, org, repo, for_mutation=True) h, o, r = _resolve(remote, host, org, repo)
auth = _auth(h) auth = _auth(h)
url = f"{repo_api_url(h, o, r)}/pulls/{pr_number}" url = f"{repo_api_url(h, o, r)}/pulls/{pr_number}"
@@ -9861,7 +9856,7 @@ def gitea_commit_files(
verify_preflight_purity(remote=remote, worktree_path=worktree_path, task="commit_files", org=org, repo=repo) verify_preflight_purity(remote=remote, worktree_path=worktree_path, task="commit_files", org=org, repo=repo)
processed_files, source_proofs = _prepare_commit_payload_files(files) processed_files, source_proofs = _prepare_commit_payload_files(files)
h, o, r = _resolve(remote, host, org, repo, for_mutation=True) h, o, r = _resolve(remote, host, org, repo)
auth = _auth(h) auth = _auth(h)
url = f"{repo_api_url(h, o, r)}/contents" url = f"{repo_api_url(h, o, r)}/contents"
@@ -9997,7 +9992,7 @@ def gitea_publish_unpublished_issue_branch(
repo=repo, repo=repo,
) )
h, o, r = _resolve(remote, host, org, repo, for_mutation=True) h, o, r = _resolve(remote, host, org, repo)
git_remote = (git_remote_name or remote or "").strip() git_remote = (git_remote_name or remote or "").strip()
existing_lock = issue_lock_store.load_issue_lock( existing_lock = issue_lock_store.load_issue_lock(
@@ -10200,7 +10195,7 @@ def gitea_bootstrap_author_issue_worktree(
repo=repo, repo=repo,
) )
h, o, r = _resolve(remote, host, org, repo, for_mutation=True) h, o, r = _resolve(remote, host, org, repo)
canonical_root = _canonical_local_git_root() canonical_root = _canonical_local_git_root()
import author_issue_bootstrap import author_issue_bootstrap
@@ -12148,7 +12143,7 @@ def gitea_reconcile_merged_cleanups(
"blocker_kind": "invalid_pr_number", "blocker_kind": "invalid_pr_number",
} }
h, o, r = _resolve(remote, host, org, repo, for_mutation=True) h, o, r = _resolve(remote, host, org, repo)
auth = _auth(h) auth = _auth(h)
base = repo_api_url(h, o, r) base = repo_api_url(h, o, r)
@@ -12538,7 +12533,7 @@ def gitea_assess_already_landed_reconciliation(
"permission_report": _permission_block_report("gitea.read"), "permission_report": _permission_block_report("gitea.read"),
} }
h, o, r = _resolve(remote, host, org, repo, for_mutation=True) h, o, r = _resolve(remote, host, org, repo)
auth = _auth(h) auth = _auth(h)
pr = api_request("GET", f"{repo_api_url(h, o, r)}/pulls/{pr_number}", auth) pr = api_request("GET", f"{repo_api_url(h, o, r)}/pulls/{pr_number}", auth)
@@ -14214,8 +14209,10 @@ def _current_master_parity() -> dict:
current_head = master_parity_gate.read_git_head(PROJECT_ROOT) current_head = master_parity_gate.read_git_head(PROJECT_ROOT)
live_head = master_parity_gate.read_remote_master_head( live_head = master_parity_gate.read_remote_master_head(
PROJECT_ROOT, remote=_git_default_remote_name(PROJECT_ROOT)) PROJECT_ROOT, remote=_git_default_remote_name(PROJECT_ROOT))
bound_context = session_ctx.get_session_context()
return master_parity_gate.assess_master_parity( return master_parity_gate.assess_master_parity(
_STARTUP_PARITY, current_head, live_remote_head=live_head) _STARTUP_PARITY, current_head, live_remote_head=live_head, bound_cohort=bound_context)
def _current_runtime_mode_report(refresh: bool = False) -> dict: def _current_runtime_mode_report(refresh: bool = False) -> dict:
+70 -9
View File
@@ -183,6 +183,7 @@ def assess_master_parity(
startup: dict | None, startup: dict | None,
current_head: str | None, current_head: str | None,
live_remote_head: str | None = None, live_remote_head: str | None = None,
bound_cohort: dict | None = None,
) -> dict: ) -> dict:
"""Compare the startup baseline against the current on-disk ``HEAD``. """Compare the startup baseline against the current on-disk ``HEAD``.
@@ -192,7 +193,8 @@ def assess_master_parity(
could not be determined, which is not treated as stale). could not be determined, which is not treated as stale).
- ``stale`` -- the on-disk master has definitively advanced past the - ``stale`` -- the on-disk master has definitively advanced past the
running process. running process.
- ``restart_required`` -- ``stale`` or ``live_stale``; the recovery action. - ``restart_required`` -- ``stale``, ``live_stale``, or ``cohort_stale``; the
recovery action.
- ``determinable`` -- whether both local HEADs were known well enough to - ``determinable`` -- whether both local HEADs were known well enough to
compare. compare.
- ``startup_head`` / ``current_head`` / ``reasons``. - ``startup_head`` / ``current_head`` / ``reasons``.
@@ -209,24 +211,69 @@ def assess_master_parity(
- ``live_known`` -- whether the live remote target was resolved. - ``live_known`` -- whether the live remote target was resolved.
- ``live_stale`` -- the live remote master has advanced past the running - ``live_stale`` -- the live remote master has advanced past the running
process (daemon is behind live master) even if local parity is green. process (daemon is behind live master) even if local parity is green.
- ``mutation_safe`` -- the daemon code, local checkout, and live remote - ``bound_cohort`` -- metadata describing the bound MCP cohort.
target all agree; the only state in which a mutation may rely on parity. - ``cohort_parity_match`` -- whether bound cohort startup SHA matches parity.
- ``cohort_stale`` -- bound cohort startup SHA is stale relative to parity.
- ``mutation_safe`` -- the daemon code, local checkout, live remote target,
and bound cohort all agree; the only state in which a mutation may rely
on parity.
""" """
startup_head = (startup or {}).get("startup_head") startup_head = (startup or {}).get("startup_head")
reasons: list[str] = [] reasons: list[str] = []
cohort_info: dict | None = None
cohort_parity_match = True
cohort_stale = False
if bound_cohort:
c_id = str(bound_cohort.get("cohort_id") or "").strip() or None
c_pid = bound_cohort.get("pid")
c_sha = str(
bound_cohort.get("startup_sha")
or bound_cohort.get("git_head")
or ""
).strip() or None
c_endpoint = str(bound_cohort.get("endpoint") or "").strip() or None
c_fingerprint = str(
bound_cohort.get("config_fingerprint") or ""
).strip() or None
cohort_info = {
"cohort_id": c_id,
"pid": c_pid,
"startup_sha": c_sha,
"endpoint": c_endpoint,
"config_fingerprint": c_fingerprint,
}
parity_ref = live_remote_head or current_head or startup_head
if c_sha and parity_ref:
if c_sha.lower() != parity_ref.lower():
cohort_parity_match = False
cohort_stale = True
reasons.append(
f"bound cohort startup SHA '{_short(c_sha)}' does not "
f"match authoritative parity SHA '{_short(parity_ref)}' "
"(stale cohort refused)"
)
if startup_head is None: if startup_head is None:
reasons.append( reasons.append(
"startup commit was not captured; code parity cannot be enforced") "startup commit was not captured; code parity cannot be enforced")
return _result(True, False, False, startup_head, current_head, return _result(True, False, False, startup_head, current_head,
live_remote_head, False, reasons) live_remote_head, False, reasons,
bound_cohort=cohort_info,
cohort_parity_match=cohort_parity_match,
cohort_stale=cohort_stale)
if current_head is None: if current_head is None:
reasons.append( reasons.append(
"current workspace HEAD could not be read; code parity cannot be " "current workspace HEAD could not be read; code parity cannot be "
"enforced") "enforced")
return _result(True, False, False, startup_head, current_head, return _result(True, False, False, startup_head, current_head,
live_remote_head, False, reasons) live_remote_head, False, reasons,
bound_cohort=cohort_info,
cohort_parity_match=cohort_parity_match,
cohort_stale=cohort_stale)
local_in_parity = startup_head == current_head local_in_parity = startup_head == current_head
local_stale = not local_in_parity local_stale = not local_in_parity
@@ -246,18 +293,28 @@ def assess_master_parity(
return _result( return _result(
local_in_parity, local_stale, True, startup_head, current_head, local_in_parity, local_stale, True, startup_head, current_head,
live_remote_head, live_stale, reasons) live_remote_head, live_stale, reasons,
bound_cohort=cohort_info,
cohort_parity_match=cohort_parity_match,
cohort_stale=cohort_stale)
def _result(in_parity, stale, determinable, startup_head, current_head, def _result(in_parity, stale, determinable, startup_head, current_head,
live_remote_head, live_stale, reasons): live_remote_head, live_stale, reasons, bound_cohort=None,
cohort_parity_match=True, cohort_stale=False):
live_known = live_remote_head is not None live_known = live_remote_head is not None
mutation_safe = ( mutation_safe = (
determinable and in_parity and live_known and not live_stale) determinable
and in_parity
and live_known
and not live_stale
and cohort_parity_match
and not cohort_stale
)
return { return {
"in_parity": in_parity, "in_parity": in_parity,
"stale": stale, "stale": stale,
"restart_required": stale or live_stale, "restart_required": stale or live_stale or cohort_stale,
"determinable": determinable, "determinable": determinable,
"startup_head": startup_head, "startup_head": startup_head,
"current_head": current_head, "current_head": current_head,
@@ -267,11 +324,15 @@ def _result(in_parity, stale, determinable, startup_head, current_head,
"live_remote_head": live_remote_head, "live_remote_head": live_remote_head,
"live_known": live_known, "live_known": live_known,
"live_stale": live_stale, "live_stale": live_stale,
"cohort_parity_match": cohort_parity_match,
"cohort_stale": cohort_stale,
"bound_cohort": bound_cohort,
"mutation_safe": mutation_safe, "mutation_safe": mutation_safe,
"reasons": list(reasons), "reasons": list(reasons),
} }
def gate_disabled() -> bool: def gate_disabled() -> bool:
"""Whether the parity gate is disabled by env escape hatch.""" """Whether the parity gate is disabled by env escape hatch."""
return bool((os.environ.get(ENV_DISABLE) or "").strip()) return bool((os.environ.get(ENV_DISABLE) or "").strip())
+45 -3
View File
@@ -104,6 +104,8 @@ def classify_namespace_probe(
profile: str | None = None, profile: str | None = None,
configured: bool = True, configured: bool = True,
probe_source: str | None = None, probe_source: str | None = None,
expected_parity_sha: str | None = None,
bound_cohort: dict[str, Any] | None = None,
) -> dict[str, Any]: ) -> dict[str, Any]:
"""Classify whether a required tool is callable through a live namespace. """Classify whether a required tool is callable through a live namespace.
@@ -135,20 +137,50 @@ def classify_namespace_probe(
else: else:
error_type = "namespace_call_failed" error_type = "namespace_call_failed"
# Extract cohort metadata
cohort_meta = bound_cohort or probe.get("cohort") or probe.get("bound_cohort") or {}
cohort_id = str(
cohort_meta.get("cohort_id") or probe.get("cohort_id") or ""
).strip() or None
startup_sha = str(
cohort_meta.get("startup_sha")
or cohort_meta.get("git_head")
or probe.get("startup_sha")
or probe.get("git_head")
or ""
).strip() or None
endpoint = str(
cohort_meta.get("endpoint") or probe.get("endpoint") or ""
).strip() or None
config_fingerprint = str(
cohort_meta.get("config_fingerprint") or probe.get("config_fingerprint") or ""
).strip() or None
expected_sha = (expected_parity_sha or "").strip().lower() or None
stale_cohort = False
if expected_sha and startup_sha:
if startup_sha.lower() != expected_sha:
stale_cohort = True
error_type = "stale_cohort_refused"
if not configured: if not configured:
error_type = "namespace_not_configured" error_type = "namespace_not_configured"
elif registered is False: elif registered is False:
error_type = "tool_missing" error_type = "tool_missing"
elif not probe_result: elif not probe_result:
error_type = "live_probe_missing" error_type = "live_probe_missing"
elif stale_cohort:
error_type = "stale_cohort_refused"
elif not probe_success and not error_type: elif not probe_success and not error_type:
error_type = "namespace_call_failed" error_type = "namespace_call_failed"
callable_live = bool(configured and probe_result and probe_success) callable_live = bool(configured and probe_result and probe_success and not stale_cohort)
# Probe-path health (spawn or client). IDE-proven only for client path. # Probe-path health (spawn or client). IDE-proven only for client path.
healthy = bool(configured and registered is not False and callable_live) healthy = bool(configured and registered is not False and callable_live)
ide_namespace_proven = bool(healthy and source == PROBE_SOURCE_CLIENT) ide_namespace_proven = bool(healthy and source == PROBE_SOURCE_CLIENT)
process_pid = process.get("pid") if isinstance(process, dict) else None process_pid = process.get("pid") if isinstance(process, dict) else (
cohort_meta.get("pid") if isinstance(cohort_meta, dict) else None
)
profile_name = profile or ( profile_name = profile or (
process.get("profile") if isinstance(process, dict) else None process.get("profile") if isinstance(process, dict) else None
) )
@@ -161,7 +193,13 @@ def classify_namespace_probe(
reasons.append( reasons.append(
f"Required tool '{tool}' is not registered in namespace '{ns}'." f"Required tool '{tool}' is not registered in namespace '{ns}'."
) )
if error_type == "live_probe_missing": if error_type == "stale_cohort_refused":
reasons.append(
f"Bound cohort startup SHA '{startup_sha[:12] if startup_sha else 'unknown'}' "
f"does not match expected parity SHA '{expected_sha[:12] if expected_sha else 'unknown'}' "
"(stale cohort refused)."
)
elif error_type == "live_probe_missing":
reasons.append( reasons.append(
f"No live client invocation proof was supplied for '{ns}.{tool}'." f"No live client invocation proof was supplied for '{ns}.{tool}'."
) )
@@ -248,6 +286,10 @@ def classify_namespace_probe(
"env": env_summary, "env": env_summary,
"config_path": config_path, "config_path": config_path,
"probe_source": source, "probe_source": source,
"cohort_id": cohort_id,
"startup_sha": startup_sha,
"endpoint": endpoint,
"config_fingerprint": config_fingerprint,
}, },
"blocks_merge_workflow": blocks, "blocks_merge_workflow": blocks,
} }
+22 -1
View File
@@ -151,12 +151,16 @@ class RestartCompletionProof:
unresolved_count: int unresolved_count: int
skipped_count: int skipped_count: int
note: str note: str
binding_unchanged: bool = False
prior_reconcile_id: str | None = None
def as_dict(self) -> dict[str, Any]: def as_dict(self) -> dict[str, Any]:
return { return {
"schema_version": self.schema_version, "schema_version": self.schema_version,
"reconcile_version": self.reconcile_version, "reconcile_version": self.reconcile_version,
"reconcile_id": self.reconcile_id, "reconcile_id": self.reconcile_id,
"binding_unchanged": self.binding_unchanged,
"prior_reconcile_id": self.prior_reconcile_id,
"started_at": self.started_at, "started_at": self.started_at,
"finished_at": self.finished_at, "finished_at": self.finished_at,
"boot_head_sha": self.boot_head_sha, "boot_head_sha": self.boot_head_sha,
@@ -173,6 +177,7 @@ class RestartCompletionProof:
"skipped_count": self.skipped_count, "skipped_count": self.skipped_count,
"note": self.note, "note": self.note,
"links": { "links": {
"umbrella": 655, "umbrella": 655,
"vision": 652, "vision": 652,
"roadmap": 653, "roadmap": 653,
@@ -359,7 +364,9 @@ def reconcile_after_restart(
now: datetime | None = None, now: datetime | None = None,
mode: str = MODE_LOG_ONLY, mode: str = MODE_LOG_ONLY,
reconcile_id: str | None = None, reconcile_id: str | None = None,
prior_reconcile_id: str | None = None,
) -> RestartCompletionProof: ) -> RestartCompletionProof:
"""Classify a post-restart inventory into a completion proof (#662). """Classify a post-restart inventory into a completion proof (#662).
Parameters Parameters
@@ -750,12 +757,26 @@ def reconcile_after_restart(
f"Mode={mode_norm}." f"Mode={mode_norm}."
) )
prior_id = str(
prior_reconcile_id
or inventory.get("prior_reconcile_id")
or ""
).strip() or None
target_rec_id = str(reconcile_id or "").strip() or None
binding_unchanged = bool(
prior_id and target_rec_id and target_rec_id == prior_id
)
final_reconcile_id = target_rec_id or f"reconcile-{uuid4().hex[:12]}"
return RestartCompletionProof( return RestartCompletionProof(
schema_version=SCHEMA_VERSION, schema_version=SCHEMA_VERSION,
reconcile_version=RECONCILE_VERSION, reconcile_version=RECONCILE_VERSION,
reconcile_id=(reconcile_id or f"reconcile-{uuid4().hex[:12]}"), reconcile_id=final_reconcile_id,
binding_unchanged=binding_unchanged,
prior_reconcile_id=prior_id,
started_at=_ts(started), started_at=_ts(started),
finished_at=_ts(finished), finished_at=_ts(finished),
boot_head_sha=( boot_head_sha=(
str(inventory.get("boot_head_sha")).strip() str(inventory.get("boot_head_sha")).strip()
if inventory.get("boot_head_sha") if inventory.get("boot_head_sha")
+7 -60
View File
@@ -54,62 +54,20 @@ def assess_remote_repo_match(
local_remote_url: str | None, local_remote_url: str | None,
org_explicit: bool, org_explicit: bool,
repo_explicit: bool, repo_explicit: bool,
for_mutation: bool = False,
primary_org: str | None = None,
primary_repo: str | None = None,
) -> dict: ) -> dict:
"""Fail closed when the resolved org/repo disagrees with the local git remote. """Fail closed when the resolved org/repo disagrees with the local git remote.
The guard enforces two protection levels: The guard is intentionally conservative:
1. Cross-Project Mutation Boundary (#707): * When the caller passed both ``org`` and ``repo`` explicitly, their intent is
When ``for_mutation`` is True (codebase mutation operation: branch, commit, PR, authoritative and the guard never blocks.
merge, branch deletion), the target repository (``resolved_org/resolved_repo``) * When the local git remote URL is unavailable (``None``/empty), corroboration
MUST match the primary authorized project context (``primary_org/primary_repo`` is impossible, so the guard does not block (best-effort only).
or parsed from ``local_remote_url``). Any attempt to mutate a different project * Otherwise, the resolved ``org/repo`` slug must appear in the local remote URL
fails closed, even if explicit org/repo were passed. Metadata operations (such (case-insensitive); if it does not, the guard blocks.
as creating issues or commenting on issues) across projects remain allowed.
2. Workspace Mismatch Guard (#530):
When ``for_mutation`` is False (or targets match), bare remotes must resolve
to an org/repo slug present in the local git remote URL. When explicit org/repo
are supplied for non-mutation operations, the caller's intent is authoritative.
""" """
reasons: list[str] = [] reasons: list[str] = []
eff_primary_org = primary_org
eff_primary_repo = primary_repo
if not eff_primary_org or not eff_primary_repo:
parsed_primary = parse_org_repo_from_remote_url(local_remote_url)
if parsed_primary:
eff_primary_org = eff_primary_org or parsed_primary[0]
eff_primary_repo = eff_primary_repo or parsed_primary[1]
# #707: Enforce cross-project codebase mutation boundary if for_mutation is True
if for_mutation and eff_primary_org and eff_primary_repo:
if (
resolved_org.lower() != eff_primary_org.lower()
or resolved_repo.lower() != eff_primary_repo.lower()
):
reasons.append(
f"Cross-project mutation guard (#707): Attempted codebase mutation targeting "
f"'{resolved_org}/{resolved_repo}' outside of primary authorized project "
f"context '{eff_primary_org}/{eff_primary_repo}'"
)
return {
"proven": False,
"block": True,
"cross_project_mutation_block": True,
"reasons": reasons,
"remote": remote,
"resolved_org": resolved_org,
"resolved_repo": resolved_repo,
"primary_org": eff_primary_org,
"primary_repo": eff_primary_repo,
"local_remote_url": local_remote_url,
"remediation": f"Cross-project codebase work is forbidden. Create an issue in the target repository ('{resolved_org}/{resolved_repo}') instead.",
}
if org_explicit and repo_explicit: if org_explicit and repo_explicit:
return _assessment(True, reasons, remote, resolved_org, resolved_repo, local_remote_url) return _assessment(True, reasons, remote, resolved_org, resolved_repo, local_remote_url)
@@ -130,16 +88,6 @@ def assess_remote_repo_match(
def format_remote_repo_guard_error(assessment: dict) -> str: def format_remote_repo_guard_error(assessment: dict) -> str:
"""Single RuntimeError message for the MCP resolver gate.""" """Single RuntimeError message for the MCP resolver gate."""
if assessment.get("cross_project_mutation_block"):
resolved = f"{assessment.get('resolved_org')}/{assessment.get('resolved_repo')}"
primary = f"{assessment.get('primary_org')}/{assessment.get('primary_repo')}"
return (
f"Cross-project mutation guard (#707): Attempted codebase mutation targeting '{resolved}' "
f"outside of primary authorized project context '{primary}'. "
f"Cross-project codebase work (creating branches, committing files, creating PRs) is forbidden; "
f"create an issue in the target repository ('{resolved}') instead."
)
reasons = "; ".join( reasons = "; ".join(
assessment.get("reasons") or ["remote/repo resolution mismatch"] assessment.get("reasons") or ["remote/repo resolution mismatch"]
) )
@@ -172,4 +120,3 @@ def _assessment(
"local_remote_url": local_remote_url, "local_remote_url": local_remote_url,
"remediation": REMEDIATION, "remediation": REMEDIATION,
} }
+50
View File
@@ -33,6 +33,10 @@ class _SessionContext:
source: str source: str
pid: int pid: int
canonical_repository_root: str | None = None canonical_repository_root: str | None = None
cohort_id: str | None = None
startup_sha: str | None = None
endpoint: str | None = None
config_fingerprint: str | None = None
def as_dict(self) -> dict[str, Any]: def as_dict(self) -> dict[str, Any]:
return { return {
@@ -47,9 +51,14 @@ class _SessionContext:
"source": self.source, "source": self.source,
"pid": self.pid, "pid": self.pid,
"canonical_repository_root": self.canonical_repository_root, "canonical_repository_root": self.canonical_repository_root,
"cohort_id": self.cohort_id,
"startup_sha": self.startup_sha,
"endpoint": self.endpoint,
"config_fingerprint": self.config_fingerprint,
} }
# Process-local only — never a shared file (same rationale as mutation authority). # Process-local only — never a shared file (same rationale as mutation authority).
# The frozen value prevents partial mutation, while the lock makes first-bind and # The frozen value prevents partial mutation, while the lock makes first-bind and
# sanctioned rebind atomic across concurrent MCP calls. # sanctioned rebind atomic across concurrent MCP calls.
@@ -72,6 +81,13 @@ def _reset_session_context_for_testing() -> None:
_SESSION_CONTEXT = None _SESSION_CONTEXT = None
def clear_session_context() -> None:
"""Purge process-session context and cohort bindings on disconnect."""
global _SESSION_CONTEXT
with _SESSION_CONTEXT_LOCK:
_SESSION_CONTEXT = None
def get_session_context() -> dict[str, Any] | None: def get_session_context() -> dict[str, Any] | None:
"""Return a detached snapshot of the bound context, or None if unbound.""" """Return a detached snapshot of the bound context, or None if unbound."""
with _SESSION_CONTEXT_LOCK: with _SESSION_CONTEXT_LOCK:
@@ -146,6 +162,10 @@ def bind_session_context(
expected_username: str | None = None, expected_username: str | None = None,
source: str = "bind", source: str = "bind",
canonical_repository_root: str | None = None, canonical_repository_root: str | None = None,
cohort_id: str | None = None,
startup_sha: str | None = None,
endpoint: str | None = None,
config_fingerprint: str | None = None,
) -> dict[str, Any]: ) -> dict[str, Any]:
"""Atomically bind/re-bind context (the explicit activation path).""" """Atomically bind/re-bind context (the explicit activation path)."""
with _SESSION_CONTEXT_LOCK: with _SESSION_CONTEXT_LOCK:
@@ -160,6 +180,10 @@ def bind_session_context(
expected_username=expected_username, expected_username=expected_username,
source=source, source=source,
canonical_repository_root=canonical_repository_root, canonical_repository_root=canonical_repository_root,
cohort_id=cohort_id,
startup_sha=startup_sha,
endpoint=endpoint,
config_fingerprint=config_fingerprint,
) )
@@ -175,6 +199,10 @@ def _bind_session_context_unlocked(
expected_username: str | None, expected_username: str | None,
source: str, source: str,
canonical_repository_root: str | None = None, canonical_repository_root: str | None = None,
cohort_id: str | None = None,
startup_sha: str | None = None,
endpoint: str | None = None,
config_fingerprint: str | None = None,
) -> dict[str, Any]: ) -> dict[str, Any]:
"""Store a complete immutable context while the caller holds the lock.""" """Store a complete immutable context while the caller holds the lock."""
global _SESSION_CONTEXT global _SESSION_CONTEXT
@@ -190,6 +218,10 @@ def _bind_session_context_unlocked(
source=source, source=source,
pid=os.getpid(), pid=os.getpid(),
canonical_repository_root=(canonical_repository_root or "").strip() or None, canonical_repository_root=(canonical_repository_root or "").strip() or None,
cohort_id=(cohort_id or "").strip() or None,
startup_sha=(startup_sha or "").strip() or None,
endpoint=(endpoint or "").strip() or None,
config_fingerprint=(config_fingerprint or "").strip() or None,
) )
return _SESSION_CONTEXT.as_dict() return _SESSION_CONTEXT.as_dict()
@@ -206,6 +238,10 @@ def seed_session_context_if_unbound(
expected_username: str | None = None, expected_username: str | None = None,
source: str = "seed", source: str = "seed",
canonical_repository_root: str | None = None, canonical_repository_root: str | None = None,
cohort_id: str | None = None,
startup_sha: str | None = None,
endpoint: str | None = None,
config_fingerprint: str | None = None,
) -> dict[str, Any]: ) -> dict[str, Any]:
"""Atomically bind only when this process has no current context. """Atomically bind only when this process has no current context.
@@ -227,10 +263,15 @@ def seed_session_context_if_unbound(
expected_username=expected_username, expected_username=expected_username,
source=source, source=source,
canonical_repository_root=canonical_repository_root, canonical_repository_root=canonical_repository_root,
cohort_id=cohort_id,
startup_sha=startup_sha,
endpoint=endpoint,
config_fingerprint=config_fingerprint,
) )
return _SESSION_CONTEXT.as_dict() return _SESSION_CONTEXT.as_dict()
def assess_session_context( def assess_session_context(
*, *,
profile_name: str | None, profile_name: str | None,
@@ -586,6 +627,10 @@ def mutation_context_audit_fields(
"session_identity": None, "session_identity": None,
"session_repository": None, "session_repository": None,
"session_org": None, "session_org": None,
"session_cohort_id": None,
"session_startup_sha": None,
"session_endpoint": None,
"session_config_fingerprint": None,
} }
return { return {
"session_context_bound": True, "session_context_bound": True,
@@ -598,9 +643,14 @@ def mutation_context_audit_fields(
"session_role_kind": data.get("role_kind"), "session_role_kind": data.get("role_kind"),
"session_context_source": data.get("source"), "session_context_source": data.get("source"),
"session_canonical_repository_root": data.get("canonical_repository_root"), "session_canonical_repository_root": data.get("canonical_repository_root"),
"session_cohort_id": data.get("cohort_id"),
"session_startup_sha": data.get("startup_sha"),
"session_endpoint": data.get("endpoint"),
"session_config_fingerprint": data.get("config_fingerprint"),
} }
def _assessment( def _assessment(
proven: bool, reasons: list[str], ctx: Mapping[str, Any] | None proven: bool, reasons: list[str], ctx: Mapping[str, Any] | None
) -> dict[str, Any]: ) -> dict[str, Any]:
@@ -1,94 +0,0 @@
import os
import sys
import unittest
from unittest import mock
import remote_repo_guard
import gitea_mcp_server as server
from session_context_binding import _reset_session_context_for_testing, bind_session_context
LOCAL_GITEA_TOOLS_URL = "https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools.git"
class TestCrossProjectMutationBoundary(unittest.TestCase):
def setUp(self):
os.environ["PYTEST_CURRENT_TEST"] = "test"
_reset_session_context_for_testing()
def tearDown(self):
_reset_session_context_for_testing()
os.environ.pop("PYTEST_CURRENT_TEST", None)
def test_cross_project_codebase_mutation_blocked(self):
"""Codebase mutations targeting another project must fail closed (#707)."""
assessment = remote_repo_guard.assess_remote_repo_match(
remote="prgs",
resolved_org="Other-Org",
resolved_repo="Other-Repo",
local_remote_url=LOCAL_GITEA_TOOLS_URL,
org_explicit=True,
repo_explicit=True,
for_mutation=True,
)
self.assertTrue(assessment["block"])
self.assertTrue(assessment.get("cross_project_mutation_block"))
self.assertEqual(assessment["primary_org"], "Scaled-Tech-Consulting")
self.assertEqual(assessment["primary_repo"], "Gitea-Tools")
err_msg = remote_repo_guard.format_remote_repo_guard_error(assessment)
self.assertIn("Cross-project mutation guard (#707)", err_msg)
self.assertIn("Attempted codebase mutation targeting 'Other-Org/Other-Repo'", err_msg)
self.assertIn("outside of primary authorized project context 'Scaled-Tech-Consulting/Gitea-Tools'", err_msg)
self.assertIn("create an issue in the target repository ('Other-Org/Other-Repo') instead", err_msg)
def test_same_project_codebase_mutation_allowed(self):
"""Codebase mutations targeting the primary project context are allowed."""
assessment = remote_repo_guard.assess_remote_repo_match(
remote="prgs",
resolved_org="Scaled-Tech-Consulting",
resolved_repo="Gitea-Tools",
local_remote_url=LOCAL_GITEA_TOOLS_URL,
org_explicit=True,
repo_explicit=True,
for_mutation=True,
)
self.assertFalse(assessment["block"])
self.assertFalse(assessment.get("cross_project_mutation_block", False))
def test_cross_project_metadata_operation_allowed(self):
"""Metadata operations (issue creation, comments) across project boundaries are allowed."""
assessment = remote_repo_guard.assess_remote_repo_match(
remote="prgs",
resolved_org="Other-Org",
resolved_repo="Other-Repo",
local_remote_url=LOCAL_GITEA_TOOLS_URL,
org_explicit=True,
repo_explicit=True,
for_mutation=False,
)
self.assertTrue(assessment["proven"])
self.assertFalse(assessment["block"])
@mock.patch.dict(os.environ, {"GITEA_FORCE_REMOTE_REPO_CHECK": "1"})
@mock.patch("gitea_mcp_server._local_git_remote_url", return_value=LOCAL_GITEA_TOOLS_URL)
def test_mcp_server_resolve_cross_project_mutation_fails_closed(self, mock_remote):
"""_resolve with for_mutation=True blocks cross-project targets."""
with self.assertRaises(RuntimeError) as ctx:
server._resolve("prgs", None, "Other-Org", "Other-Repo", for_mutation=True)
err = str(ctx.exception)
self.assertIn("Cross-project mutation guard (#707)", err)
self.assertIn("create an issue in the target repository", err)
@mock.patch.dict(os.environ, {"GITEA_FORCE_REMOTE_REPO_CHECK": "1"})
@mock.patch("gitea_mcp_server._local_git_remote_url", return_value=LOCAL_GITEA_TOOLS_URL)
def test_mcp_server_resolve_cross_project_metadata_succeeds(self, mock_remote):
"""_resolve with for_mutation=False allows explicit cross-project target."""
host, org, repo = server._resolve("prgs", None, "Other-Org", "Other-Repo", for_mutation=False)
self.assertEqual(org, "Other-Org")
self.assertEqual(repo, "Other-Repo")
if __name__ == "__main__":
unittest.main()
@@ -0,0 +1,253 @@
"""Tests for Issue #689: Deterministic MCP namespace attachment.
Verifies cohort identity exposure, stale cohort refusal, parity matching,
reconcile_id freshness, session context cleanup, and regression scenarios.
"""
from __future__ import annotations
import os
import unittest
from unittest.mock import patch
import master_parity_gate
import mcp_namespace_health
import post_restart_reconcile
import session_context_binding as session_ctx
class TestIssue689DeterministicCohortAttachment(unittest.TestCase):
"""Suite covering Issue #689 acceptance criteria."""
def setUp(self) -> None:
session_ctx._reset_session_context_for_testing()
def tearDown(self) -> None:
session_ctx._reset_session_context_for_testing()
def test_ac1_session_context_exposes_cohort_identity(self) -> None:
"""AC1: Session context exposes cohort ID, startup SHA, endpoint, and config fingerprint."""
ctx = session_ctx.bind_session_context(
profile_name="prgs-author",
remote="prgs",
host="gitea.prgs.cc",
identity="jcwalker3",
repository="Gitea-Tools",
org="Scaled-Tech-Consulting",
role_kind="author",
cohort_id="cohort-p1234-abc123456789",
startup_sha="abc123456789def",
endpoint="gitea.prgs.cc",
config_fingerprint="fingerprint12345",
)
self.assertEqual(ctx["cohort_id"], "cohort-p1234-abc123456789")
self.assertEqual(ctx["startup_sha"], "abc123456789def")
self.assertEqual(ctx["endpoint"], "gitea.prgs.cc")
self.assertEqual(ctx["config_fingerprint"], "fingerprint12345")
fetched = session_ctx.get_session_context()
self.assertIsNotNone(fetched)
self.assertEqual(fetched["cohort_id"], "cohort-p1234-abc123456789")
self.assertEqual(fetched["startup_sha"], "abc123456789def")
def test_ac2_stale_cohort_refused_by_probe_classifier(self) -> None:
"""AC2: Probe classifier refuses binding to a stale cohort as stale_cohort_refused."""
probe_res = {
"success": True,
"cohort": {
"cohort_id": "cohort-obsolete-1",
"startup_sha": "22698c1000000000000000000000000000000000",
"endpoint": "gitea.prgs.cc",
"config_fingerprint": "fp-old",
},
}
res = mcp_namespace_health.classify_namespace_probe(
"gitea-author",
probe_result=probe_res,
probe_source="client_namespace",
expected_parity_sha="a4c73766f4b0cc32f7c3808688eceeb6fee74335",
)
self.assertFalse(res["healthy"])
self.assertFalse(res["success"])
self.assertEqual(res["error_type"], "stale_cohort_refused")
self.assertIn("stale cohort refused", " ".join(res["reasons"]))
self.assertEqual(
res["diagnostics"]["startup_sha"],
"22698c1000000000000000000000000000000000",
)
def test_ac3_reconnection_parity_matching_and_fail_closed(self) -> None:
"""AC3: Parity gate fails closed when bound cohort startup SHA mismatches parity."""
startup = {"startup_head": "a4c73766f4b0cc32f7c3808688eceeb6fee74335"}
current = "a4c73766f4b0cc32f7c3808688eceeb6fee74335"
live_remote = "a4c73766f4b0cc32f7c3808688eceeb6fee74335"
# Matching cohort
matching_cohort = {
"cohort_id": "cohort-fresh",
"startup_sha": "a4c73766f4b0cc32f7c3808688eceeb6fee74335",
}
res_matching = master_parity_gate.assess_master_parity(
startup, current, live_remote_head=live_remote, bound_cohort=matching_cohort
)
self.assertTrue(res_matching["cohort_parity_match"])
self.assertFalse(res_matching["cohort_stale"])
self.assertTrue(res_matching["mutation_safe"])
# Mismatched obsolete cohort
obsolete_cohort = {
"cohort_id": "cohort-obsolete-22698c1",
"startup_sha": "22698c1000000000000000000000000000000000",
}
res_stale = master_parity_gate.assess_master_parity(
startup, current, live_remote_head=live_remote, bound_cohort=obsolete_cohort
)
self.assertFalse(res_stale["cohort_parity_match"])
self.assertTrue(res_stale["cohort_stale"])
self.assertTrue(res_stale["restart_required"])
self.assertFalse(res_stale["mutation_safe"])
def test_ac4_reconcile_id_freshness(self) -> None:
"""AC4: Re-attachment distinguishes new reconcile_id from preserved binding."""
inventory = {"inventory_complete": True}
# New attachment generates fresh reconcile_id
proof1 = post_restart_reconcile.reconcile_after_restart(inventory)
self.assertFalse(proof1.binding_unchanged)
self.assertTrue(proof1.reconcile_id.startswith("reconcile-"))
# Preserved binding reports binding_unchanged=True
proof2 = post_restart_reconcile.reconcile_after_restart(
inventory,
reconcile_id=proof1.reconcile_id,
prior_reconcile_id=proof1.reconcile_id,
)
self.assertTrue(proof2.binding_unchanged)
self.assertEqual(proof2.reconcile_id, proof1.reconcile_id)
# Disconnected re-attachment gets new reconcile_id
proof3 = post_restart_reconcile.reconcile_after_restart(
inventory,
prior_reconcile_id=proof1.reconcile_id,
)
self.assertFalse(proof3.binding_unchanged)
self.assertNotEqual(proof3.reconcile_id, proof1.reconcile_id)
def test_ac5_session_disconnect_clears_bindings(self) -> None:
"""AC5: clear_session_context purges session context on disconnect."""
session_ctx.bind_session_context(
profile_name="prgs-author",
remote="prgs",
host="gitea.prgs.cc",
identity="jcwalker3",
cohort_id="cohort-1",
)
self.assertIsNotNone(session_ctx.get_session_context())
session_ctx.clear_session_context()
self.assertIsNone(session_ctx.get_session_context())
def test_ac6_bound_cohort_in_diagnostics(self) -> None:
"""AC6: Bound cohort identity appears in audit diagnostics."""
session_ctx.bind_session_context(
profile_name="prgs-author",
remote="prgs",
host="gitea.prgs.cc",
identity="jcwalker3",
cohort_id="cohort-test-99",
startup_sha="sha99999",
endpoint="gitea.prgs.cc",
config_fingerprint="fp999",
)
audit = session_ctx.mutation_context_audit_fields()
self.assertTrue(audit["session_context_bound"])
self.assertEqual(audit["session_cohort_id"], "cohort-test-99")
self.assertEqual(audit["session_startup_sha"], "sha99999")
self.assertEqual(audit["session_endpoint"], "gitea.prgs.cc")
self.assertEqual(audit["session_config_fingerprint"], "fp999")
def test_ac7_regression_n_reconnects_never_bind_to_obsolete_daemon(self) -> None:
"""AC7: N reconnects against a daemon set containing obsolete daemons never bind obsolete ones."""
live_master = "master-head-latest-12345"
daemons = [
{"id": "d1", "startup_sha": "obsolete-head-11111"},
{"id": "d2", "startup_sha": "obsolete-head-22698c1"},
{"id": "d3", "startup_sha": live_master},
{"id": "d4", "startup_sha": "obsolete-head-33333"},
]
for _ in range(5):
for daemon in daemons:
res = master_parity_gate.assess_master_parity(
{"startup_head": live_master},
live_master,
live_remote_head=live_master,
bound_cohort=daemon,
)
if daemon["startup_sha"] != live_master:
self.assertFalse(res["mutation_safe"])
self.assertTrue(res["cohort_stale"])
else:
self.assertTrue(res["mutation_safe"])
self.assertFalse(res["cohort_stale"])
def test_ac8_regression_incident_shape_reproduction(self) -> None:
"""AC8: Reproduce incident shape — obsolete cohort 22698c1 resident vs newer daemon."""
live_master = "a4c73766f4b0cc32f7c3808688eceeb6fee74335"
obsolete_cohort = {
"cohort_id": "cohort-resident-22698c1",
"startup_sha": "22698c1000000000000000000000000000000000",
}
new_cohort = {
"cohort_id": "cohort-spawned-new",
"startup_sha": live_master,
}
# Obsolete cohort fails parity check
obs_res = mcp_namespace_health.classify_namespace_probe(
"gitea-author",
probe_result={"success": True, "cohort": obsolete_cohort},
probe_source="client_namespace",
expected_parity_sha=live_master,
)
self.assertFalse(obs_res["healthy"])
self.assertEqual(obs_res["error_type"], "stale_cohort_refused")
# Fresh cohort succeeds
new_res = mcp_namespace_health.classify_namespace_probe(
"gitea-author",
probe_result={"success": True, "cohort": new_cohort},
probe_source="client_namespace",
expected_parity_sha=live_master,
)
self.assertTrue(new_res["healthy"])
def test_ac9_regression_bound_cohort_going_stale_detected(self) -> None:
"""AC9: A bound cohort that later goes stale is detected on next attachment check."""
initial_master = "sha-v1-initial"
cohort = {"cohort_id": "c1", "startup_sha": initial_master}
# Initial state: in parity
res1 = master_parity_gate.assess_master_parity(
{"startup_head": initial_master},
initial_master,
live_remote_head=initial_master,
bound_cohort=cohort,
)
self.assertTrue(res1["mutation_safe"])
# Master advances to sha-v2-advanced while cohort remains at sha-v1-initial
advanced_master = "sha-v2-advanced"
res2 = master_parity_gate.assess_master_parity(
{"startup_head": initial_master},
advanced_master,
live_remote_head=advanced_master,
bound_cohort=cohort,
)
self.assertFalse(res2["mutation_safe"])
self.assertTrue(res2["restart_required"])
self.assertTrue(res2["cohort_stale"])
if __name__ == "__main__":
unittest.main()