Compare commits

..
Author SHA1 Message Date
sysadminandClaude Opus 5 92615f474b feat(control-plane): add authoritative native MCP fleet inventory (#949)
Every pre-existing runtime surface is per-process. gitea_get_runtime_context
and gitea_assess_master_parity describe only the server answering the call,
and gitea_assess_mcp_namespace_health accepts process, probe_result and
registered_tools from the caller, so it cannot constrain the caller. The
control-plane sessions table records allocator task sessions, not server
processes. Five self-reports of one revision therefore never proved that five
processes exist, that no sixth exists, or that all five share one client
cohort.

Add gitea_assess_fleet_inventory, a strictly read-only capability that takes
no evidence parameters. It combines two sources that must agree:

- a control-plane runtime registry row each server writes about itself,
  from the official entrypoint, immediately after the native transport bind;
- a process observation the answering server performs, never the caller.

A member is running only when both agree and the observed process predates
its registration, so configuration alone never counts and a recycled PID
cannot impersonate an exited server. Classification is a pure function of the
snapshot, so gitea-controller and gitea-reconciler return the same verdict.

Missing, duplicate, unexpected, stale and unregistered members are reported
in separate fields rather than collapsed into one error, because each needs a
different operator action. single_cohort is derived only from recorded cohort
identity and stays null when unknown, so matching revisions never establish a
cohort. Any unreadable registry, unavailable listing, unregistered process,
unknown cohort or unknown revision sets inventory_complete and
mutation_gate_satisfied false with a specific blocked_reason.

The capability performs no restart, reconnect, drain, lease mutation, issue
mutation or process termination, never terminates a duplicate, and never
manufactures restart evidence. The only signal sent is signal 0.

Also resolves the fleet namespace from the expected roster:
role_namespace_gate.infer_mcp_namespace recognises only author and reviewer
and echoes the profile name otherwise, which would have mislabelled the
controller, merger and reconciler members. Widening that shared helper is
controller role-metadata work owned by #950 and is deliberately not done here.

Schema v5 to v6 is additive and idempotent: one new table, no existing table,
tool signature or result field changed. Two control-plane tests that pinned
the schema version to the literal 5 now assert against SCHEMA_VERSION.

#950 (controller role metadata), #951 (restart receipts) and #952 (stale-lease
consistency) remain separate and untouched.

Tests: 118 new across tests/test_mcp_fleet_inventory.py,
tests/test_control_plane_db_server_runtimes.py and
tests/test_fleet_inventory_tool.py, covering every acceptance criterion
including the multi-LLM duplicate-server regression that motivated the issue.

Closes #949

Co-Authored-By: Claude Opus 5 (1M context) <[email protected]>
Claude-Session: https://claude.ai/code/session_01NBaQKcRRKeKbZuhk72MhEf
2026-07-27 23:25:45 -04:00
sysadmin 82d71b7702 Merge pull request 'fix(author bootstrap): restore missing runtime identity and session helpers (Closes #943)' (#944) from fix/issue-943-runtime-context-helpers into master 2026-07-27 20:03:08 -05:00
sysadmin 47bfae07d2 fix(author bootstrap): resolve allocator session ownership, authority alignment, and test coverage (#943) 2026-07-27 19:52:27 -04:00
sysadmin 35ed8a2fcb Merge pull request 'fix(gate): preserve exact-owner renewal evidence across duplicate rechecks (Closes #945)' (#946) from fix/issue-945-owning-pr-renewal-evidence into master 2026-07-27 18:27:35 -05:00
sysadmin ed9414ebda Merge pull request 'fix(test): add pytest.ini to restrict testpaths and ignore worktree branches (#927)' (#928) from fix/issue-927-pytest-config into master 2026-07-27 15:38:45 -05:00
sysadmin 5fea326988 Merge pull request 'feat(mcp): expose sanctioned Codex MCP reconnect request (Closes #678)' (#918) from fix/issue-678-codex-mcp-reconnect into master 2026-07-27 15:32:10 -05:00
sysadmin dc5d2c8caa Merge pull request 'feat(webui): Sentry/GlitchTip observability console & correlation view (Closes #649)' (#917) from feat/issue-649-sentry-console-correlation into master 2026-07-27 15:23:36 -05:00
sysadmin c30b381eb2 Merge pull request 'fix(mcp): add active IDE vs global config-drift diagnostic (Closes #672)' (#920) from fix/issue-672-mcp-config-drift into master 2026-07-27 05:40:04 -05:00
jcwalker3andClaude Opus 5 f49e781102 fix(author bootstrap): restore missing runtime identity and session helpers (Closes #943)
gitea_bootstrap_author_issue_worktree referenced four globals that commit
a942afe (#850) introduced without ever defining:

    _active_username        1 reference, 0 definitions
    _active_profile_name    1 reference, 0 definitions
    _current_session_id     1 reference, 0 definitions
    _author_mutation_block  1 reference, 0 definitions

Evaluating the call arguments therefore raised

    NameError: name '_active_username' is not defined

before author_issue_bootstrap.bootstrap_author_issue_worktree was entered, so
the capability was unusable for every caller including dry_run=true. The fourth
name, _author_mutation_block, sits on the reviewer-stop refusal path and was
found by the generalised regression test rather than by the original report.

The defect was unreachable until PR #942 (#941) wired the bootstrap scope into
workflow_scope_guard: before that, verify_preflight_purity refused first with
missing_issue_worktree, masking it.

Changes:

* _active_username reads the immutable #714 session context that gitea_whoami
  seeds — the identity pin every other mutation gate already consults. An
  unbound context returns None so callers fail closed rather than acting as an
  unverified actor; a profile's expected_username is never substituted.
* _active_profile_name prefers the live get_profile() and falls back to the
  bound session context only when the profile cannot be read.
* _current_session_id mints the same "<profile>-<pid>-<hex>" shape as the three
  pre-existing lease call sites, bound once per process so repeated calls
  describe one session instead of a fresh owner per call, which would make
  lease-ownership comparisons unsatisfiable. None is never memoised.
* _author_mutation_block returns the uniform refusal shape the other author
  mutations already return for the same check_author_mutation_after_reviewer_stop
  block.

No guard, signature, or permission changes. Wrong-role, wrong-profile,
wrong-identity, stale-runtime, expected-base and workflow-scope enforcement all
still gate the call; the helpers only supply values the service then validates
fail-closed (missing_active_identity / missing_active_profile /
missing_owner_session / stale_concurrency_pin).

Regression: tests/test_issue_943_runtime_context_helpers.py (27 tests). The
generalised test resolves every global the wrapper's body references against
module globals and builtins, so the next missing reference fails too rather
than only the three named here — that test is what found
_author_mutation_block. Coverage also coversdry-run reaching and completing the
service with helper-produced bindings, dry-run leaving no branch, worktree,
assignment or lease, apply reaching its intended transition, each fail-closed
mismatch, expected-base mismatch, and the #941/PR #942 scope wiring.

Pre-fix 20 failed / 13 passed against unmodified aab54d48; post-fix 27 passed.
Targeted bootstrap and guard suites: 245 passed, 59 subtests.
Full suite: 28 failed, 5552 passed, 6 skipped, 1006 subtests — the 28 are the
standing baseline, every one of which also fails on the unmodified base.

Closes #943

Co-Authored-By: Claude Opus 5 (1M context) <[email protected]>
Claude-Session: https://claude.ai/code/session_013ygVZQLbbhJTChuVuLaJWb
2026-07-26 07:56:14 -05:00
sysadmin 17dd05ec9d fix(test): add pytest.ini to restrict testpaths and ignore worktree branches (#927) 2026-07-25 23:07:15 -04:00
sysadmin fad44669d9 fix(mcp): add active IDE vs global config-drift diagnostic (Closes #672) 2026-07-25 19:07:10 -04:00
jcwalker3andGrok 4.5 29ad93d145 feat(mcp): expose sanctioned Codex MCP reconnect request (Closes #678)
Add gitea_request_mcp_reconnect as a report-only callable surface so Codex
and other agent hosts can request host/IDE reconnect with typed blockers
and exact UI steps. Never kills processes or edits config.

Co-Authored-By: Grok 4.5 <[email protected]>
2026-07-25 17:52:08 -05:00
sysadmin f2dbf30e81 feat(webui): Sentry/GlitchTip observability console & correlation view (Closes #649) 2026-07-25 18:48:36 -04:00
32 changed files with 5905 additions and 27 deletions
+52
View File
@@ -196,6 +196,58 @@ def _verify_assignment_and_lease_ids(
# Some lease rows may not yet have an assignment join; still require
# the lease itself to exist and bind to the claimed session/issue.
pass
# #943 review 622 B2: a lease that is no longer live confers no ownership.
# Existence alone previously satisfied this gate, so a released or expired
# lease could still authorize a bootstrap for a claim its session had given
# up. Checked before the session comparison so the reason names the real
# problem rather than reporting a mismatch.
from datetime import datetime, timezone
lease_status = str(lease.get("status") or "").strip().lower()
if lease_status and lease_status != "active":
return {
"success": False,
"reason_code": "lease_not_live",
"message": (
f"lease_id '{lid}' is '{lease_status}', not active; a lease that "
"is not live confers no ownership (fail closed)."
),
"exact_next_action": (
"Re-allocate the work item and pass the live assignment/lease pair."
),
}
expires_raw = str(lease.get("expires_at") or "").strip()
if expires_raw:
try:
expires_at = datetime.fromisoformat(expires_raw.replace("Z", "+00:00"))
except ValueError:
return {
"success": False,
"reason_code": "lease_not_live",
"message": (
f"lease_id '{lid}' records an unparseable expiry "
f"'{expires_raw}' (fail closed)."
),
"exact_next_action": (
"Re-allocate the work item and pass the live "
"assignment/lease pair."
),
}
if expires_at.tzinfo is None:
expires_at = expires_at.replace(tzinfo=timezone.utc)
if expires_at <= datetime.now(timezone.utc):
return {
"success": False,
"reason_code": "lease_not_live",
"message": (
f"lease_id '{lid}' expired at {expires_raw}; an expired lease "
"confers no ownership (fail closed)."
),
"exact_next_action": (
"Reclaim or re-allocate the lease, then retry with the live pair."
),
}
lease_session = str(lease.get("session_id") or "").strip()
if lease_session and lease_session != owner_session:
return {
+190 -1
View File
@@ -32,7 +32,7 @@ from typing import Any, Iterator, Sequence
import dependency_graph
import gitea_audit
SCHEMA_VERSION = 5
SCHEMA_VERSION = 6
# Assignable work kinds only — raw monitoring incidents are never work items.
WORK_KINDS = frozenset({"issue", "pr"})
@@ -65,6 +65,33 @@ CREATE TABLE IF NOT EXISTS sessions (
status TEXT NOT NULL DEFAULT 'active'
);
-- Fleet-level MCP server runtime registry (#949). Distinct from ``sessions``:
-- a session is one allocator *task*, while a row here is one *server process*
-- that bound native MCP transport. Each server writes its own row and no other,
-- so a caller can never supply this evidence. Creating the table is itself the
-- v5→v6 migration: additive, idempotent, and it touches no existing table.
CREATE TABLE IF NOT EXISTS mcp_server_runtimes (
runtime_id TEXT PRIMARY KEY,
namespace TEXT NOT NULL,
profile TEXT,
role TEXT,
remote TEXT,
org TEXT,
repo TEXT,
repository_root TEXT,
pid INTEGER NOT NULL,
cohort_id TEXT,
cohort_source TEXT,
client_provenance TEXT,
boot_id TEXT,
startup_head TEXT,
daemon_start_head TEXT,
transport TEXT,
registered_at TEXT NOT NULL,
last_heartbeat_at TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'running'
);
CREATE TABLE IF NOT EXISTS work_items (
work_item_id INTEGER PRIMARY KEY AUTOINCREMENT,
remote TEXT NOT NULL,
@@ -910,6 +937,142 @@ class ControlPlaneDB:
rows = conn.execute(sql, params).fetchall()
return [dict(r) for r in rows]
# ── MCP server runtime registry (#949) ────────────────────────────────
_RUNTIME_COLUMNS: tuple[str, ...] = (
"runtime_id",
"namespace",
"profile",
"role",
"remote",
"org",
"repo",
"repository_root",
"pid",
"cohort_id",
"cohort_source",
"client_provenance",
"boot_id",
"startup_head",
"daemon_start_head",
"transport",
"registered_at",
"last_heartbeat_at",
"status",
)
def register_mcp_server_runtime(
self,
record: dict[str, Any],
*,
retention_seconds: int | None = None,
now: str | None = None,
) -> dict[str, Any]:
"""Record that *this* server process is running (#949).
Only the described process ever calls this, at native transport bind.
Two housekeeping deletions happen here — on the startup write path, never
on the read path — so the registry cannot grow without bound:
* rows that recorded the **same PID**, which the calling process now
owns and which therefore cannot still be a live server;
* rows older than *retention_seconds*.
Neither can hide a live duplicate: a concurrently running server holds a
different PID and writes a fresher row.
"""
stamp = now or _ts()
row = {key: record.get(key) for key in self._RUNTIME_COLUMNS}
row["runtime_id"] = str(record.get("runtime_id") or "").strip()
if not row["runtime_id"]:
raise ValueError("runtime_id is required (fail closed)")
row["namespace"] = str(record.get("namespace") or "").strip()
if not row["namespace"]:
raise ValueError("namespace is required (fail closed)")
if record.get("pid") is None:
raise ValueError("pid is required (fail closed)")
row["pid"] = int(record["pid"])
row["registered_at"] = record.get("registered_at") or stamp
row["last_heartbeat_at"] = record.get("last_heartbeat_at") or stamp
row["status"] = record.get("status") or "running"
columns = ", ".join(self._RUNTIME_COLUMNS)
placeholders = ", ".join("?" for _ in self._RUNTIME_COLUMNS)
values = [row[key] for key in self._RUNTIME_COLUMNS]
with self._tx() as conn:
conn.execute(
"DELETE FROM mcp_server_runtimes WHERE pid = ? AND runtime_id != ?",
(row["pid"], row["runtime_id"]),
)
if retention_seconds and retention_seconds > 0:
cutoff = _ts(
_utc_now() - timedelta(seconds=int(retention_seconds))
)
conn.execute(
"DELETE FROM mcp_server_runtimes WHERE registered_at < ?",
(cutoff,),
)
conn.execute(
f"INSERT OR REPLACE INTO mcp_server_runtimes({columns}) "
f"VALUES ({placeholders})",
values,
)
stored = conn.execute(
"SELECT * FROM mcp_server_runtimes WHERE runtime_id = ?",
(row["runtime_id"],),
).fetchone()
return dict(stored)
def heartbeat_mcp_server_runtime(self, runtime_id: str) -> None:
"""Refresh a runtime row's liveness timestamp.
Never called by the inventory read path: liveness is established from
process evidence so that reading the fleet mutates nothing (#949 AC11).
"""
with self._tx() as conn:
conn.execute(
"UPDATE mcp_server_runtimes SET last_heartbeat_at = ? "
"WHERE runtime_id = ?",
(_ts(), runtime_id),
)
def mark_mcp_server_runtime_stopped(self, runtime_id: str) -> None:
"""Mark a runtime row stopped (graceful shutdown bookkeeping)."""
with self._tx() as conn:
conn.execute(
"UPDATE mcp_server_runtimes SET status = 'stopped', "
"last_heartbeat_at = ? WHERE runtime_id = ?",
(_ts(), runtime_id),
)
def list_mcp_server_runtimes(
self,
*,
statuses: Sequence[str] | None = None,
limit: int = 500,
) -> list[dict[str, Any]]:
"""List runtime registry rows, newest registration first (read-only).
Ordering is fully deterministic: ties on ``registered_at`` break on
``runtime_id`` so two callers reading one snapshot see one order.
"""
clauses: list[str] = []
params: list[Any] = []
if statuses:
placeholders = ", ".join("?" for _ in statuses)
clauses.append(f"status IN ({placeholders})")
params.extend(statuses)
where = ("WHERE " + " AND ".join(clauses)) if clauses else ""
params.append(max(1, int(limit)))
with self._tx(immediate=False) as conn:
rows = conn.execute(
f"SELECT * FROM mcp_server_runtimes {where} "
"ORDER BY registered_at DESC, runtime_id ASC LIMIT ?",
params,
).fetchall()
return [dict(r) for r in rows]
# ── work items ────────────────────────────────────────────────────────
def upsert_work_item(
@@ -1568,6 +1731,32 @@ class ControlPlaneDB:
).fetchone()
return dict(row) if row else None
def list_incident_links(
self,
*,
provider: str | None = None,
gitea_org: str | None = None,
gitea_repo: str | None = None,
limit: int = 100,
) -> list[dict[str, Any]]:
"""List stored incident_links rows, optionally filtered by provider/repo (#612 / #649)."""
query = "SELECT * FROM incident_links WHERE 1=1"
params: list[Any] = []
if provider:
query += " AND provider = ?"
params.append(provider.strip().lower())
if gitea_org:
query += " AND gitea_org = ?"
params.append(_norm_scope(gitea_org))
if gitea_repo:
query += " AND gitea_repo = ?"
params.append(_norm_scope(gitea_repo))
query += " ORDER BY link_id DESC LIMIT ?"
params.append(max(1, limit))
with self._tx(immediate=False) as conn:
rows = conn.execute(query, params).fetchall()
return [dict(r) for r in rows]
# ── lease lifecycle (#601) ────────────────────────────────────────────
+70
View File
@@ -0,0 +1,70 @@
# MCP Config Drift Diagnostic & Sanctioned Repair Runbook (#672)
This document describes the diagnostic framework for detecting configuration drift between the active IDE MCP configuration (`~/.gemini/antigravity-ide/mcp_config.json`) and the offline/global canonical configuration (`~/.gemini/config/mcp_config.json`), and establishes the **sanctioned repair runbook**.
## Background & Problem Statement
Offline tools like `test_mcp_conn.py` test the global configuration (`~/.gemini/config/mcp_config.json`) via `subprocess.Popen`. However, the active IDE/client namespace uses `~/.gemini/antigravity-ide/mcp_config.json`. When required Gitea role servers (`gitea-author`, `gitea-reviewer`, `gitea-merger`, `gitea-reconciler`, `gitea-controller`, `gitea-tools`) are missing or carry mismatched profile environments in the active IDE config:
1. Offline tests pass (`test_mcp_conn.py` green).
2. The IDE client returns `EOF` / `transport closed` when attempting role-scoped mutations.
3. Operators misdiagnose missing server definitions as stale runtimes, leading to forbidden `pkill` attempts (#630) or `mtime` hacks (#655).
## Diagnostic Tool: `mcp_config_drift.py`
Run the diagnostic tool directly to compare configurations:
```bash
python3 mcp_config_drift.py --json
```
Or specify custom config locations:
```bash
python3 mcp_config_drift.py \
--active-config ~/.gemini/antigravity-ide/mcp_config.json \
--global-config ~/.gemini/config/mcp_config.json
```
### Key Diagnostic Outputs
- `in_sync`: Boolean indicating if all required Gitea role servers exist in the active IDE config with matching profile declarations.
- `missing_role_servers`: List of role servers present in global config but missing from active IDE config.
- `profile_mismatches`: List of profile environment mismatches per server.
- `reasons`: Explicit, human-readable list of drift causes.
All returned payloads automatically redact secret tokens, DSNs, Authorization headers, and private keys.
---
## Sanctioned Repair Path (Step-by-Step)
When `mcp_config_drift.py` reports drift (`in_sync: false`), execute the following **sanctioned repair steps**:
1. **Backup Active IDE Config:**
```bash
cp ~/.gemini/antigravity-ide/mcp_config.json ~/.gemini/antigravity-ide/mcp_config.json.bak
```
2. **Patch Active IDE Config:**
Copy the missing Gitea role server JSON blocks (`gitea-author`, `gitea-reviewer`, etc.) from `~/.gemini/config/mcp_config.json` into `~/.gemini/antigravity-ide/mcp_config.json`.
3. **Reconnect via IDE/Client:**
Use the IDE / client UI reconnection control (or restart the IDE client app).
4. **Verify Active Namespace Health:**
Invoke `gitea_whoami` (and optional `gitea_resolve_task_capability`) through the active IDE client on each required role namespace.
---
## FORBIDDEN Repair Actions (#630 / #655)
The following actions are **strictly forbidden** for config drift repair:
- ❌ **`pkill` or manual daemon process kill commands:** Process kills cause contamination and break active session leases.
- ❌ **`mtime` touch edits:** Artificial mtime modifications mask stale runtimes without updating configuration.
- ❌ **Source code edits:** Mutating python tool logic to bypass missing server entries.
- ❌ **Session-state edits:** Direct database or lock-file state mutation.
---
## Final Report Guidelines
A workflow final report **must not** rely on offline `test_mcp_conn.py` output alone. Final reports must include active-config evidence from live `gitea_whoami` calls on the active IDE namespaces.
+134
View File
@@ -0,0 +1,134 @@
# Authoritative MCP fleet inventory (#949)
`gitea_assess_fleet_inventory` is the read-only native capability that answers a
question no other surface could: **is exactly one server running for each
configured PRGS profile, and do they all belong to one client cohort?**
## Why the existing surfaces were not enough
| Surface | What it proves | Why it cannot prove the fleet |
| --- | --- | --- |
| `gitea_get_runtime_context` | Profile, identity, workspace binding of **the process answering the call** | Says nothing about the other four namespaces |
| `gitea_assess_master_parity` | Startup vs current vs live revision of **that same process** | Five self-reports of one revision do not establish five processes, nor the absence of a sixth |
| `gitea_assess_mcp_namespace_health` | Whether one named namespace can invoke one tool | Accepts `process`, `probe_result` and `registered_tools` **from the caller** — a capability whose inputs come from the party it constrains is not evidence |
| control-plane `sessions` table | Allocator **task** sessions | A session is a unit of work, not a server process; nothing recorded that a server exists |
Before this capability, satisfying a strict five-process/single-cohort gate
required shell process inspection, cached JSON, or source reading — none of
which are sanctioned workflow evidence.
## Evidence model
Two independent sources must agree before a member counts as running.
**1. Control-plane runtime registry** — table `mcp_server_runtimes`.
Each server writes exactly one row *about itself*, from the official entrypoint,
immediately after the native MCP transport bind. Only a transport-bound process
reaches that line, and no MCP caller can reach it at all. The row is
authoritative for identity: namespace, profile, role, repository binding,
cohort, startup revision, transport, and PID.
**2. Server-side process observation** — a process listing performed by the
server answering the inventory call, never by the caller. It is authoritative
for existence and liveness, and it is the only source that can reveal a running
server the registry does not know about.
A member is `live` only when a registry row has a matching running process whose
start time precedes the registration — so a recycled PID cannot impersonate a
server that has since exited.
### Deliberate non-inferences
* **Configuration is not existence.** A configured profile with no live
corroborated row is `missing`, never `running` (AC8).
* **Matching revisions are not a cohort.** `single_cohort` is derived only from
recorded cohort identity. Five members at one revision with an unknown cohort
yield `single_cohort: null` and a closed gate, never `true` (AC7).
* **Parent-client status is not member health.** Each member is classified from
its own evidence.
Cohort identity comes from `GITEA_MCP_CLIENT_COHORT_ID` when the client sets it,
otherwise from the parent process that launched the server. A server whose
parent has gone away (reparented to init) reports an unknown cohort rather than
guessing.
## Result shape
Top-level verdict fields:
| Field | Meaning |
| --- | --- |
| `inventory_complete` | All evidence was obtainable. False whenever anything below is unknown. |
| `incomplete_reasons` | Every distinct reason completeness failed. |
| `configured_members` | The five expected members with their instance counts and health. |
| `running_members` | Live, corroborated instances of expected profiles. |
| `missing_members` | Expected profiles with no live instance. |
| `duplicate_members` | Expected profiles with more than one live instance, with every PID. |
| `unexpected_members` | Live members outside the expected roster. Never folded into duplicates. |
| `stale_members` | Registry rows whose process is dead, unobserved, PID-recycled, or unknown. |
| `unregistered_processes` | Running MCP server processes with no registry row. |
| `repository_binding_mismatches` | Live members bound to another repository, or with an incomplete binding. |
| `role_mismatches` | Live members whose declared role or namespace contradicts the configured profile. |
| `single_cohort` / `mixed_cohort` | `true`/`false`, or `null` when cohort evidence is unknown. |
| `mixed_revision` / `startup_revisions` | Revision spread across live members. |
| `exactly_one_per_profile` | No missing and no duplicate expected members. |
| `mutation_gate_satisfied` | The full invariant held. |
| `blocked_reason` / `blocked_reasons` | Why the gate is closed; the first is the headline. |
| `mutations_performed` | Always `[]`. |
Ordering is deterministic — every list sorts by namespace, profile, PID, then
runtime id — so two callers reading one snapshot see identical structure.
## How controller and reconciler gates consume it
Classification is a **pure function of the snapshot**. The answering namespace is
reported as `answering_namespace` metadata and never affects the verdict, so
`gitea-controller` and `gitea-reconciler` return the same result for the same
fleet. Neither is privileged over the other.
Consume it like this:
1. Call `gitea_assess_fleet_inventory` from `gitea-controller` **or**
`gitea-reconciler` (any namespace holding `gitea.read` may call it).
2. If `mutation_gate_satisfied` is `true`, the exact-one-instance-per-profile,
single-cohort, single-revision invariant is proven; proceed.
3. Otherwise **stop and report `blocked_reason` verbatim**. Distinguish the
cases — they need different operator actions:
* `missing_members` — the client did not launch that namespace; reconnect it.
* `duplicate_members` — a second client or a manual launch is running that
profile; the **operator** quits the extra client. This capability never
terminates a process.
* `unexpected_members` — an unconfigured PRGS server is live; investigate
before trusting any gate.
* `unregistered_processes` — a server is running code that predates this
capability, or failed to register; the inventory is incomplete by
construction and must not be reported as healthy.
* `mixed_cohort` / `mixed_revision` — the fleet is not one coherent unit.
* `single_cohort: null` — cohort evidence is missing; this is *unknown*, not
*healthy*.
Never treat `inventory_complete: false` as a soft warning. It means the
inventory cannot describe the whole fleet, which is exactly the state the
mutation gate exists to refuse.
## Read-only guarantees
The capability performs no restart, reconnect, drain, lease mutation, issue
mutation, repository write, or process termination, and it never manufactures
historical restart evidence. The only signal it sends is `signal 0`, a
liveness/permission check that delivers nothing to the target process. Registry
writes happen exclusively on the startup path of the process being described,
never on this read path.
## Scope boundaries
This capability is observation only. Related concerns live elsewhere and are
deliberately not absorbed here:
* **#950** — controller role metadata and capability-routing consistency.
* **#951** — durable synchronization and restart receipts.
* **#952** — stale-lease inspection, dashboard, and executor consistency.
* **#900** — cohort lifecycle supervision (draining and reaping superseded
cohorts) is a *mutating* transition that would consume this evidence.
* **#948** — the client/session/generation ownership model that this evidence
feeds.
+21 -6
View File
@@ -47,18 +47,33 @@ Do the steps in order. Stop as soon as a live **client-namespace** call succeeds
- Only the Gitea namespace fails → single-namespace transport close. Continue.
- Every server fails → restart the whole MCP client, not just one namespace.
2. **Reconnect the namespace through the client, not the shell.** Use the IDE /
client MCP-reconnect action for that server entry (in Claude Code:
`/mcp` → reconnect the affected `gitea-*` server). Reconnecting forces the
client to spawn a fresh subprocess and re-open the pipe. This clears the
closed-client state that a bare `kill`/respawn from a terminal does **not**.
2. **Request the sanctioned reconnect surface (#678), then reconnect through
the client — not the shell.** From a still-reachable Gitea MCP namespace
(or after host auto-reconnect), call:
```text
gitea_request_mcp_reconnect(
namespace="gitea-author", # or gitea-reviewer / gitea-merger / …
reason="transport_eof",
client="codex", # or claude_code / generic
)
```
The tool is **report-only**: it never restarts a process. It returns
namespace, profile, pid/session, startup SHA, current master SHA, boundary
status, and a **typed blocker** with exact operator UI steps for Codex
(Reload Developer Tools / per-server reconnect) or Claude Code (`/mcp`).
Then perform the host reconnect those steps describe so the client spawns a
fresh subprocess and re-opens the pipe. That clears the closed-client state
that a bare `kill`/respawn from a terminal does **not**.
3. **Do not "fix" it by importing the server or poking the process.** Reaching
for `python -c 'import gitea_mcp_server ...'`, raw JSON-RPC from a shell,
killing PIDs to force a respawn, or touching MCP config mtimes does **not**
restore the *client's* view of the namespace and violates the daemon-import
guard (#558, `docs/mcp-daemon-import-guard.md`). The only sanctioned repair
is a **client reconnect / relaunch**.
is a **client reconnect / relaunch** (or the typed operator path returned by
`gitea_request_mcp_reconnect`).
4. **Verify through the same path the workflow will use.** After reconnect, call
the specific tool the blocked workflow needs — not just any tool — through
+3 -2
View File
@@ -45,10 +45,11 @@ and *fails closed*.
| `legacy_auto_restart_helper` | removed | A helper (`_trigger_mcp_auto_restart`) that actively restarted the server from the read-only resolver path. | Removed in #685; kept absent by `assert_auto_restart_helper_absent()`. | #685, #657 |
| `config_touch_reload` | removed | Touching (utime) the MCP client config to make the host reload the server. | Removed from the resolver in #685: stale detection is report-only, never mutating config, spawning threads, or calling `os._exit`. | #685, #657 |
| `master_advance_auto_restart` | guarded_fail_closed | On-disk master advancing past the running code. | `master_parity_gate` captures startup parity and blocks mutations while stale, emitting restart guidance; the process never self-restarts. | #420, #591, #657 |
| `stale_runtime_resolver_reconnect` | guarded_fail_closed | The capability resolver detecting a stale serving process. | Report-only (#685): returns `restart_required`/`stop_required` and an exact reconnect action; no restart, thread, config touch, or `os._exit`. | #685, #657 |
| `stale_runtime_resolver_reconnect` | guarded_fail_closed | The capability resolver detecting a stale serving process. | Report-only (#685): returns `restart_required`/`stop_required` and an exact reconnect action; no restart, thread, config touch, or `os._exit`. | #685, #657, #678 |
| `codex_client_reconnect_request` | guarded_fail_closed | `gitea_request_mcp_reconnect` report-only tool for Codex/LLM sessions. | Report-only (#678): returns namespace/profile/pid/startup SHA/master SHA/boundary status and a typed operator blocker with exact client UI steps; never restarts or kills. | #678, #630, #685, #657 |
| `manual_daemon_kill` | forbidden | Shell kills of the daemon: `pkill -f mcp_server.py`, `killall`, broad `pkill -f python` sweeps, or `kill <pid>` of a daemon pid. | Forbidden (#630): `runtime_recovery_guard` classifies these as contamination and `gitea_record_daemon_process_kill_attempt` writes a durable marker that fails later mutations closed. Operator maintenance authorization is read only from the environment. | #630, #657 |
| `conflict_marker_infra_stop` | guarded_fail_closed | The daemon entrypoint scans for unresolved merge-conflict markers at startup and stops (`sys.exit(1)`). | Fail-closed startup stop, not a restart: the process exits and waits for the operator to resolve conflicts and relaunch; never loops. | #657 |
| `ide_client_reconnect` | host_residual | A manual `/mcp reconnect` (or equivalent host action) that recreates the MCP client connection. | Outside this process's control; the sanctioned recovery the gates point operators toward. No in-process code initiates it. | #584, #656, #657 |
| `ide_client_reconnect` | host_residual | A manual `/mcp reconnect` (or equivalent host action) that recreates the MCP client connection. Agents obtain exact UI steps via `gitea_request_mcp_reconnect` (#678). | Outside this process's control; the sanctioned recovery the gates point operators toward. No in-process code initiates it. | #584, #656, #657, #678 |
| `profile_switch_runtime` | sanctioned_narrow_recovery | Switching the active execution profile at runtime (dynamic-profile mode). | In-process and restart-free: `runtime_switching_supported` is true, so a switch rebinds capability without recreating the process. | #656, #657 |
## Guards enforced in CI
+2
View File
@@ -54,6 +54,7 @@ that gates each call, not which tools exist.
- `gitea_assess_already_landed_reconciliation`
- `gitea_assess_conflict_fix_classification`
- `gitea_assess_conflict_fix_push`
- `gitea_assess_fleet_inventory`
- `gitea_assess_gitea_operation_path`
- `gitea_assess_master_parity`
- `gitea_assess_mcp_namespace_health`
@@ -137,6 +138,7 @@ that gates each call, not which tools exist.
- `gitea_release_merger_pr_lease`
- `gitea_release_reviewer_pr_lease`
- `gitea_release_workflow_lease`
- `gitea_request_mcp_reconnect`
- `gitea_request_mcp_restart`
- `gitea_resolve_task_capability`
- `gitea_resume_review_draft`
@@ -0,0 +1,35 @@
# Web Console: Sentry/GlitchTip Observability & Incident Bridge Console (#649)
This document describes the Phase 4 observability console surface integrated into the MCP Control Plane Web Console (`webui/`), backed by the #612 incident bridge and the #613 control-plane DB substrate.
## Architectural Authority Model (ADR Alignment)
Per the Web Console Architecture ADR (`docs/architecture/webui-control-plane-console-architecture-adr.md`):
| Layer | Responsibility | Authority |
|---|---|---|
| **Gitea** | Durable work record | Issues, PRs, comments, reviews, labels, merges |
| **Control-plane DB** | Live coordination & linkage | `incident_links` table, session leases, allocations |
| **Sentry / GlitchTip** | Observability input | Unresolved incidents, error events, stack traces |
| **Incident Bridge (#612)** | Reconciliation engine | Reconciles provider observations into Gitea issues |
| **Web Console (`webui/`)** | Read-only projection & gated actions | Projects connection health & correlation links; gates writes |
> **Key Rule:** Raw monitoring incidents are **never** assignable control-plane `work_items`. They remain observation input only.
## Redaction Boundary Invariants
1. **No secrets in returns or rendering:** Auth tokens (`SENTRY_AUTH_TOKEN`, `GLITCHTIP_AUTH_TOKEN`), DSNs, `Authorization` headers, and sensitive local file paths are passed through `webui.console_redaction` before leaving the server.
2. **Safe projection:** Connection objects report `credentials_present: true/false` rather than exposing raw keys or headers.
## Console Endpoints
- **HTML Surface:** `GET /observability` — Renders provider connection cards, error correlation tables, and gated reconcile controls.
- **Versioned API:** `GET /api/v1/observability` — Returns structured JSON snapshot with `schema_version`, `providers`, `links`, and `metrics`.
- **Legacy Compatibility Alias:** `GET /api/observability` — Read-only compatibility alias for Phase 4.
## Gated Actions
- `observability_reconcile_incident` (`gitea_observability_reconcile_incident`): Triggers or previews dry-run issue reconciliation for a provider incident.
- `observability_link_issue` (`gitea_observability_link_issue`): Links a provider incident to an existing Gitea tracking issue.
Both actions require `operator` role and gate through `task_capability_map`. Execution fails closed in read-only MVP mode.
+2
View File
@@ -98,6 +98,8 @@ already define, and a regression test asserts each mapping matches.
| `system.rebind_session_worktree` | operator | gated_write | `gitea.read` | Yes | No | No | 2 |
| `system.reconcile_cleanups` | controller | privileged | `gitea.pr.close` | Yes | No | No | 2 |
| `initiate_workflow` | operator | gated_write | `gitea.read` | Yes | No | No | 2 |
| `observability_reconcile_incident` | operator | gated_write | `gitea.read` | Yes | No | No | 4 |
| `observability_link_issue` | operator | gated_write | `gitea.read` | Yes | No | No | 4 |
**Dual control** means the acting principal may not be the sole authority: a
second distinct principal must confirm. **Break-glass** means the action is
+679 -11
View File
@@ -196,6 +196,7 @@ RECONCILER_WORKTREE_ENV = "GITEA_RECONCILER_WORKTREE"
import namespace_workspace_binding as nwb # noqa: E402
import canonical_repository_root as crr # noqa: E402 # #706 cross-repo canonical root
import mcp_namespace_health # noqa: E402
import mcp_fleet_inventory # noqa: E402 # #949 authoritative fleet inventory
import stale_binding_recovery # noqa: E402
# Worktree env bindings inherited from the parent environment at daemon boot
@@ -2123,6 +2124,7 @@ import root_checkout_guard # noqa: E402
import workflow_scope_guard # noqa: E402 # #683 production scope / force-on guards
import stable_branch_push_guard # noqa: E402
import runtime_recovery_guard # noqa: E402 # #630 manual daemon-kill contamination
import mcp_client_reconnect # noqa: E402 # #678 sanctioned Codex reconnect request
import remote_repo_guard # noqa: E402
import anti_stomp_preflight # noqa: E402
import issue_claim_heartbeat # noqa: E402
@@ -3664,6 +3666,293 @@ def _authenticated_username(host: str):
return user
# #943 review 622 F3/F4: one coherent authority snapshot per mutation claim.
#
# The reviewed implementation read the identity from the pinned #714 session
# context and the profile from the live ``get_profile()``, so a sanctioned
# rebind could produce a claimant pair whose two halves came from different
# snapshots — and that pair is written durably into the issue lock. The
# canonical pairing is ``record_mutation_authority``: profile from
# ``get_profile()``, identity from ``_authenticated_username(host)``. This
# resolver reproduces that pairing, and uses the pinned session context only for
# drift detection, never as a value source.
_AUTHORITY_PROFILE_UNRESOLVED = "authority_profile_unresolved"
_AUTHORITY_IDENTITY_UNRESOLVED = "authority_identity_unresolved"
_AUTHORITY_IDENTITY_DRIFT = "authority_identity_drift"
_AUTHORITY_PROFILE_DRIFT = "authority_profile_drift"
def _active_mutation_authority(host: str | None) -> dict:
"""Resolve one coherent (identity, profile) authority pair, or a refusal.
Both halves come from the same live snapshot. The immutable #714 session
context is consulted only to detect drift: when it disagrees with the live
snapshot the result is a fail-closed refusal, never a silently blended pair.
Returns a dict with ``ok`` True plus ``identity``/``profile_name``, or ``ok``
False plus ``reason_code``, ``reasons`` and ``expected``/``actual`` when a
drift or resolution failure is detected. Never raises for an unresolved
profile: the caller converts the refusal into a structured author block so
reason code, retryability and transport survival are preserved (F4).
"""
try:
profile = get_profile() or {}
except (RuntimeError, ValueError, TypeError, KeyError, OSError) as exc:
# Narrow, and never a silent fallback to a previously cached name: an
# unresolvable profile is a fail-closed condition, matching
# record_mutation_authority's "active profile unresolved (fail closed)".
return {
"ok": False,
"reason_code": _AUTHORITY_PROFILE_UNRESOLVED,
"reasons": [
"active profile could not be resolved: "
f"{_redact(str(exc))} (fail closed)"
],
}
profile_name = (profile.get("profile_name") or "").strip()
if not profile_name:
return {
"ok": False,
"reason_code": _AUTHORITY_PROFILE_UNRESOLVED,
"reasons": [
"active profile carries no profile_name (fail closed)"
],
}
identity = ""
if host:
identity = (_authenticated_username(host) or "").strip()
if not identity:
# A profile's expected_username is configuration, not proof of an
# authenticated actor, so it is never substituted here.
return {
"ok": False,
"reason_code": _AUTHORITY_IDENTITY_UNRESOLVED,
"reasons": [
"authenticated identity could not be resolved for the active "
"host (fail closed); call gitea_whoami and retry"
],
}
ctx = session_ctx.get_session_context() or {}
pinned_identity = (ctx.get("identity") or "").strip()
pinned_profile = (ctx.get("profile_name") or "").strip()
if pinned_identity and pinned_identity != identity:
return {
"ok": False,
"reason_code": _AUTHORITY_IDENTITY_DRIFT,
"reasons": [
"live authenticated identity disagrees with the bound session "
"context identity (fail closed)"
],
"expected": pinned_identity,
"actual": identity,
}
if pinned_profile and pinned_profile != profile_name:
return {
"ok": False,
"reason_code": _AUTHORITY_PROFILE_DRIFT,
"reasons": [
"live active profile disagrees with the bound session context "
"profile (fail closed)"
],
"expected": pinned_profile,
"actual": profile_name,
}
return {
"ok": True,
"identity": identity,
"profile_name": profile_name,
"session_context_bound": bool(pinned_identity or pinned_profile),
}
def _active_username(host: str | None = None) -> str | None:
"""Authenticated identity for the active mutation authority, else None.
Refuses unbound, blank and whitespace-only identities, and never substitutes
a profile's ``expected_username`` for an authenticated one.
"""
authority = _active_mutation_authority(host)
return authority.get("identity") if authority.get("ok") else None
def _active_profile_name(host: str | None = None) -> str | None:
"""Active profile name for the mutation authority, else None.
Shares the single snapshot with :func:`_active_username`, so the two can
never describe different authority states.
"""
authority = _active_mutation_authority(host)
return authority.get("profile_name") if authority.get("ok") else None
# #943 review 622 B1: the workflow session that owns the work — never a
# process-lifetime or PID-derived value.
#
# The reviewed implementation minted "<profile>-<pid>-<hex>" once per process.
# The MCP daemon outlives every task it serves, so that identifier conflates
# sequential author tasks and can never equal the control-plane session that
# owns an allocator-created lease; ``_verify_assignment_and_lease_ids`` therefore
# refused with lease_session_mismatch on the canonical allocated path.
# ``issue_lock_store.mint_task_session_id`` states the rule directly: an
# ownership key "deliberately contains no process identifier", because the PID
# "cannot identify *which* task holds a claim".
_SESSION_UNVERIFIED = "workflow_session_unverified"
_SESSION_REQUIRED = "workflow_session_required_for_allocated_work"
_SESSION_LOCK_OWNER_MISMATCH = "issue_lock_owner_mismatch"
def _resolve_owner_workflow_session(
*,
issue_number: int,
assignment_id: str | None,
lease_id: str | None,
session_id: str | None,
identity: str,
profile_name: str,
remote: str,
org: str | None,
repo: str | None,
) -> dict:
"""Resolve the authoritative workflow session that owns this work.
Precedence, each step fail-closed:
1. An explicit ``session_id`` is verified against the control-plane
``sessions`` table: it must exist, be active, and match this role and
profile. An unverifiable identifier is refused, never trusted.
2. Otherwise an existing issue lock for this issue supplies its per-task
``task_session_id``, but only when the lock's recorded claimant matches
the resolved authority pair.
3. Otherwise, when allocator identifiers are supplied, the request is
refused: ownership of an allocated assignment cannot be established
without its owning session, and the presence of the identifiers is not
itself evidence of ownership.
4. Otherwise no allocator identifiers and no existing lock a fresh
per-task key is minted through the canonical
``issue_lock_store.mint_task_session_id``. It carries no process
identifier and is never memoised, so sequential tasks on one daemon never
share an owner.
Whether the resolved session actually owns a supplied lease stays the
decision of ``author_issue_bootstrap._verify_assignment_and_lease_ids``;
this function never pre-empts, duplicates or bypasses that gate.
"""
declared = (session_id or "").strip()
if declared:
db, errs = _control_plane_db_or_error()
if db is None:
return {
"ok": False,
"reason_code": _SESSION_UNVERIFIED,
"reasons": errs,
}
try:
rows = db.list_sessions(statuses=("active",))
except Exception as exc: # noqa: BLE001 — surface structured
return {
"ok": False,
"reason_code": _SESSION_UNVERIFIED,
"reasons": [
"could not read control-plane sessions to verify "
f"session_id: {_redact(str(exc))} (fail closed)"
],
}
match = next(
(r for r in rows if str(r.get("session_id") or "") == declared), None
)
if match is None:
return {
"ok": False,
"reason_code": _SESSION_UNVERIFIED,
"reasons": [
f"session_id '{declared}' is not an active control-plane "
"session (fail closed)"
],
}
row_role = (match.get("role") or "").strip().lower()
row_profile = (match.get("profile") or "").strip()
if row_role and row_role != "author":
return {
"ok": False,
"reason_code": _SESSION_UNVERIFIED,
"reasons": [
f"session_id '{declared}' is recorded for role "
f"'{row_role}', not author (fail closed)"
],
"expected": "author",
"actual": row_role,
}
if row_profile and row_profile != profile_name:
return {
"ok": False,
"reason_code": _SESSION_UNVERIFIED,
"reasons": [
f"session_id '{declared}' is recorded for profile "
f"'{row_profile}', not '{profile_name}' (fail closed)"
],
"expected": row_profile,
"actual": profile_name,
}
return {"ok": True, "session_id": declared, "session_source": "declared"}
lock = None
try:
lock = issue_lock_store.load_issue_lock(
remote=remote,
org=org or "",
repo=repo or "",
issue_number=int(issue_number),
)
except Exception: # noqa: BLE001 — absent/unreadable lock is not fatal here
lock = None
lock_session = issue_lock_store.lease_task_session_id(lock) if lock else ""
if lock_session:
lease = (lock or {}).get("work_lease") or {}
claimant = lease.get("claimant") or {}
lock_identity = (claimant.get("username") or "").strip()
lock_profile = (claimant.get("profile") or "").strip()
if (lock_identity and lock_identity != identity) or (
lock_profile and lock_profile != profile_name
):
return {
"ok": False,
"reason_code": _SESSION_LOCK_OWNER_MISMATCH,
"reasons": [
f"issue #{int(issue_number)} is locked by "
f"'{lock_identity or 'unknown'}' ({lock_profile or 'unknown'}), "
f"not '{identity}' ({profile_name}) (fail closed)"
],
"expected": f"{lock_identity}/{lock_profile}",
"actual": f"{identity}/{profile_name}",
}
return {
"ok": True,
"session_id": lock_session,
"session_source": "issue_lock",
}
if (assignment_id or "").strip() or (lease_id or "").strip():
return {
"ok": False,
"reason_code": _SESSION_REQUIRED,
"reasons": [
"assignment_id/lease_id were supplied but no owning workflow "
"session could be established (fail closed); pass the "
"session_id that holds the lease, or bind the issue lock first"
],
}
return {
"ok": True,
"session_id": issue_lock_store.mint_task_session_id(
issue_lock_store.AUTHOR_ISSUE_WORK_LEASE
),
"session_source": "minted_task_key",
}
def _authenticated_actor(host: str) -> dict:
"""Resolve the authenticated actor's stable identity (#709 F7 review 438).
@@ -9987,6 +10276,24 @@ def gitea_commit_files(
}
def _author_mutation_block(reasons: list[str], **extra) -> dict:
"""Uniform fail-closed shape for an author mutation refused after a reviewer stop.
#943: referenced by ``gitea_bootstrap_author_issue_worktree`` and never
defined, so the reviewer-stop refusal path raised ``NameError`` instead of
returning its refusal. Mirrors the inline shape the other author mutations
return for the same ``check_author_mutation_after_reviewer_stop`` block.
"""
payload = {
"success": False,
"performed": False,
"outcome": "REFUSED",
"reasons": reasons,
}
payload.update(extra)
return payload
def _publication_block(reasons: list[str], **extra) -> dict:
"""Uniform fail-closed shape for publication refusals (#812 AC20)."""
payload = {
@@ -10243,6 +10550,7 @@ def gitea_bootstrap_author_issue_worktree(
branch_name: str | None = None,
worktree_path: str | None = None,
idempotency_key: str | None = None,
session_id: str | None = None,
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
@@ -10264,6 +10572,14 @@ def gitea_bootstrap_author_issue_worktree(
branch_name: Optional custom branch name (must match issue-<N> pattern).
worktree_path: Optional custom worktree path under branches/.
idempotency_key: Optional key for idempotent replay/resume.
session_id: Optional workflow session that owns the assignment/lease.
Required when assignment_id/lease_id are supplied, because ownership
of an allocated lease is compared by session identifier and cannot
be derived from the daemon process (#943 review 622 B1). Verified
against the control-plane sessions table; an unverifiable value is
refused rather than trusted. Omit for an unallocated bootstrap: the
owning session is then taken from an existing issue lock, or a fresh
per-task key is minted.
remote: Known instance 'dadeschools' or 'prgs'.
host: Override Gitea host.
org: Override Org.
@@ -10302,6 +10618,44 @@ def gitea_bootstrap_author_issue_worktree(
h, o, r = _resolve(remote, host, org, repo)
canonical_root = _canonical_local_git_root()
# One authority snapshot supplies both halves of the claimant pair (F3), and
# an unresolvable profile becomes a structured refusal rather than a silent
# fallback (F4).
authority = _active_mutation_authority(h)
if not authority.get("ok"):
return _author_mutation_block(
authority.get("reasons") or ["mutation authority unresolved"],
reason_code=authority.get("reason_code"),
retryable=False,
transport_survives=True,
expected=authority.get("expected"),
actual=authority.get("actual"),
issue_number=int(issue_number),
)
# The owning workflow session — never process- or PID-derived (B1).
session = _resolve_owner_workflow_session(
issue_number=issue_number,
assignment_id=assignment_id,
lease_id=lease_id,
session_id=session_id,
identity=authority["identity"],
profile_name=authority["profile_name"],
remote=remote,
org=o,
repo=r,
)
if not session.get("ok"):
return _author_mutation_block(
session.get("reasons") or ["owning workflow session unresolved"],
reason_code=session.get("reason_code"),
retryable=False,
transport_survives=True,
expected=session.get("expected"),
actual=session.get("actual"),
issue_number=int(issue_number),
)
import author_issue_bootstrap
return author_issue_bootstrap.bootstrap_author_issue_worktree(
@@ -10317,9 +10671,9 @@ def gitea_bootstrap_author_issue_worktree(
host=h,
org=o,
repo=r,
active_identity=_active_username(),
active_profile=_active_profile_name(),
owner_session=_current_session_id(),
active_identity=authority["identity"],
active_profile=authority["profile_name"],
owner_session=session["session_id"],
dry_run=dry_run,
)
@@ -14140,11 +14494,16 @@ def _classify_operation_gate_reasons(reasons: list[str]) -> dict:
def _stale_runtime_reconnect_action() -> str:
"""Sanctioned recovery for a stale daemon — reconnect only (#685/#897)."""
"""Sanctioned recovery for a stale daemon — reconnect only (#685/#897/#678)."""
return (
"Reconnect the IDE/client MCP session so the server reloads at the "
"current master head. Do not call gitea_activate_profile or switch "
"MCP role sessions — profile switching does not clear a stale daemon."
"blocker_kind=runtime_reconnect_required: call "
"gitea_request_mcp_reconnect(namespace=<active gitea-* namespace>, "
"reason='stale-runtime', client='codex') for a typed operator "
"reconnect blocker with exact UI steps, then reconnect the IDE/client "
"MCP session so the server reloads at the current master head. Do not "
"call gitea_activate_profile, pkill, touch configs, or switch MCP role "
"sessions — profile switching does not clear a stale daemon. After "
"reconnect restart from gitea_whoami → gitea_resolve_task_capability."
)
@@ -18453,6 +18812,116 @@ def gitea_assess_master_parity(
return out
def _fleet_namespace_for_profile(profile_name: str | None) -> str | None:
"""Fleet namespace for *profile_name* (#949).
``role_namespace_gate.infer_mcp_namespace`` recognises only author and
reviewer and returns the profile name for everything else, so it cannot
name the controller, merger, or reconciler namespaces. The fleet roster
carries that mapping; this defers to it and keeps the existing inference as
the fallback for profiles outside the roster. Widening the shared helper is
controller role-metadata work owned by #950.
"""
return mcp_fleet_inventory.namespace_for_profile(
profile_name,
default=role_namespace_gate.infer_mcp_namespace(profile_name),
)
@mcp.tool()
def gitea_assess_fleet_inventory(
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
) -> dict:
"""Read-only: authoritative inventory of the running PRGS MCP fleet (#949).
Every other runtime surface is per-process. ``gitea_get_runtime_context``
and ``gitea_assess_master_parity`` describe only the server answering the
call; ``gitea_assess_mcp_namespace_health`` takes ``process`` and
``probe_result`` from the caller and so cannot constrain it. This tool takes
**no evidence parameters at all**: it combines the control-plane runtime
registry, which each server writes about itself at native transport bind,
with a process observation performed by this server. A member counts as
running only when both agree, so configuration alone never counts as a
running server and a caller cannot supply the answer.
Classification is a pure function of the snapshot, so ``gitea-controller``
and ``gitea-reconciler`` return the same verdict for the same fleet;
``answering_namespace`` is metadata and never changes it.
Fails closed. Any unreadable registry, unavailable process listing,
unregistered server process, unknown cohort, or unknown startup revision
sets ``inventory_complete`` false and ``mutation_gate_satisfied`` false with
a specific ``blocked_reason``. Matching Git revisions never establish a
single cohort ``single_cohort`` is derived only from recorded cohort
identity and stays ``null`` when that is unknown.
Strictly read-only: no restart, reconnect, lease mutation, issue mutation,
repository write, or process termination, and duplicates are reported and
never terminated. The only signal sent is ``signal 0`` liveness probing,
which delivers nothing to the target process. ``mutations_performed`` is
always an empty list.
Args:
remote: Known instance 'dadeschools' or 'prgs'. Declares which
repository binding fleet members are expected to carry.
host: Override the Gitea host.
org: Override the expected owner/organization binding.
repo: Override the expected repository binding.
Returns:
dict with 'inventory_complete', 'configured_members', 'running_members',
'missing_members', 'duplicate_members', 'unexpected_members',
'stale_members', 'unregistered_processes', 'single_cohort',
'mixed_cohort', 'mixed_revision', 'exactly_one_per_profile',
'mutation_gate_satisfied', 'blocked_reason', and per-member evidence.
"""
read_block = _profile_operation_gate("gitea.read")
if read_block:
return {
"success": False,
"read_only": True,
"inventory_complete": False,
"mutation_gate_satisfied": False,
"blocked_reason": "the active profile may not read control-plane state",
"reasons": read_block,
"permission_report": _permission_block_report("gitea.read"),
"mutations_performed": [],
}
ctx = session_ctx.get_session_context() or {}
expected_binding = {
"remote": remote or ctx.get("remote"),
"org": org or ctx.get("org"),
"repo": repo or ctx.get("repository"),
}
db, db_errors = _control_plane_db_or_error()
runtime_rows: list[dict] = []
registry_error: str | None = None
if db is None:
registry_error = "; ".join(db_errors) or "control-plane DB unavailable"
else:
try:
runtime_rows = db.list_mcp_server_runtimes(statuses=("running",))
except Exception as exc: # noqa: BLE001 - unreadable registry is unknown evidence
registry_error = f"runtime registry could not be read: {_redact(str(exc))}"
result = mcp_fleet_inventory.classify_fleet_inventory(
runtime_rows=runtime_rows,
process_scan=mcp_fleet_inventory.scan_mcp_server_processes(),
expected_binding=expected_binding,
registry_available=registry_error is None,
registry_error=registry_error,
answering_namespace=_fleet_namespace_for_profile(_active_profile_name(host)),
)
result["expected_repository_binding"] = expected_binding
result["summary"] = mcp_fleet_inventory.summarize(result)
return result
# #781: documented in the canonical review workflow as a tool reviewers call
# before workflow load, but the registration decorator had been lost, so the
# documented inventory named something no namespace could reach.
@@ -21230,10 +21699,15 @@ def gitea_resolve_task_capability(
# serving process/profile inventory is stale — even if permission is OK.
if runtime_stale_blocker:
next_safe_action = (
"blocker_kind=runtime_reconnect_required: reconnect/restart the "
"IDE-managed Gitea MCP server for this profile so it reloads current "
"master. Do not edit mcp_config.json by hand; the resolver does not "
"touch config, spawn recovery threads, or terminate the process."
"blocker_kind=runtime_reconnect_required: call "
"gitea_request_mcp_reconnect(namespace=<active gitea-* namespace>, "
"reason='stale-runtime', client='codex') for a typed operator "
"blocker with exact UI steps, then reconnect/reload the IDE-managed "
"Gitea MCP server for this profile so it reloads current master. "
"Do not edit mcp_config.json by hand, pkill, or touch configs; the "
"resolver does not touch config, spawn recovery threads, or "
"terminate the process. After reconnect restart from gitea_whoami → "
"gitea_resolve_task_capability."
)
# Task/role alignment guards (#167): the requested task, not the
@@ -22791,6 +23265,129 @@ def gitea_workflow_dashboard(
return payload
@mcp.tool()
def gitea_request_mcp_reconnect(
namespace: str | None = None,
reason: str | None = None,
client: str = "codex",
remote: str = "dadeschools",
host: str | None = None,
session_id: str | None = None,
) -> dict:
"""Request a sanctioned host/IDE MCP reconnect for a named namespace (#678).
Codex and other agent hosts can detect stale or closed Gitea MCP runtimes
(``stop_required`` / ``restart_required`` from capability resolution, transport
EOF, missing namespace attachment). The **host owns the transport** this
process cannot reopen the client's stdio pipe. This tool is the callable
surface agents use to:
1. Report reconnect status fields (namespace, profile, pid/session,
startup SHA, current master SHA, boundary status).
2. Return a **typed blocker** with exact operator UI steps for Codex (or
another client) when reconnect is required.
This tool **never** restarts, kills, reloads, or reconfigures an MCP
process. Forbidden recovery paths (pkill, touch/mtime hacks, config/.env
edits, session-state edits, raw API) are never recommended.
After the operator reconnects, workflows must restart from preflight:
``gitea_whoami`` ``gitea_resolve_task_capability`` task.
Args:
namespace: MCP namespace to reconnect (e.g. ``gitea-author``). Defaults
to the active profile's inferred namespace.
reason: Why reconnect is requested: ``stale-runtime``, ``transport_eof``,
``missing_namespace``, ``not_required``, or free-form (normalized).
client: Operator UI surface ``codex`` (default), ``claude_code``, or
``generic``.
remote: Known instance ``dadeschools`` or ``prgs`` (parity context).
host: Optional host override for parity context.
session_id: Optional session id to echo in the report.
Returns:
dict with reconnect report fields, ``reconnect_performed=False``,
``typed_blocker`` when reconnect is required, and
``forbidden_recovery_paths``.
"""
# Read-only: gitea.read is sufficient. Never a mutation.
read_block = _profile_operation_gate("gitea.read")
if read_block:
return {
"success": False,
"read_only": True,
"reconnect_performed": False,
"mutation_performed": False,
"reasons": read_block,
"permission_report": _permission_block_report("gitea.read"),
"forbidden_recovery_paths": list(
mcp_client_reconnect.FORBIDDEN_RECOVERY_PATHS
),
}
profile = get_profile()
profile_name = (profile.get("profile_name") or "").strip() or None
inferred_ns = role_namespace_gate.infer_mcp_namespace(profile_name)
ns = (namespace or "").strip() or inferred_ns or "gitea-tools"
parity = _current_master_parity()
startup_sha = (
parity.get("daemon_start_head")
or parity.get("startup_head")
or _process_boot_head_sha
)
current_sha = parity.get("local_head") or parity.get("current_head")
if not current_sha:
try:
current_sha = master_parity_gate.read_git_head(PROJECT_ROOT)
except Exception: # noqa: BLE001
current_sha = None
boundary = mcp_client_reconnect.classify_boundary_status(
startup_sha=startup_sha if isinstance(startup_sha, str) else None,
current_master_sha=current_sha if isinstance(current_sha, str) else None,
live_stale=bool(parity.get("live_stale")) if parity.get("live_known") else None,
in_parity=parity.get("in_parity") if parity.get("determinable") else None,
)
# Infer reason from parity when caller left it unspecified.
effective_reason = reason
if not (effective_reason or "").strip():
if parity.get("restart_required") or parity.get("live_stale"):
effective_reason = mcp_client_reconnect.REASON_STALE_RUNTIME
elif boundary == mcp_client_reconnect.BOUNDARY_CLEAN:
effective_reason = mcp_client_reconnect.REASON_NOT_REQUIRED
else:
effective_reason = mcp_client_reconnect.REASON_UNSPECIFIED
payload = mcp_client_reconnect.build_reconnect_request(
namespace=ns,
profile=profile_name,
pid=os.getpid(),
session_id=session_id
or f"{(profile_name or 'session')}-{os.getpid()}",
startup_sha=startup_sha if isinstance(startup_sha, str) else None,
current_master_sha=current_sha if isinstance(current_sha, str) else None,
boundary_status=boundary,
reason=effective_reason,
client=client,
live_stale=bool(parity.get("live_stale")) if parity.get("live_known") else None,
in_parity=parity.get("in_parity") if parity.get("determinable") else None,
restart_required=bool(parity.get("restart_required")),
stop_required=bool(parity.get("restart_required")),
extra={
"remote": remote if remote in REMOTES else remote,
"host": host,
"session_context_audit": session_ctx.mutation_context_audit_fields(),
"parity_summary": master_parity_gate.format_parity(parity),
"live_stale": parity.get("live_stale"),
"live_known": parity.get("live_known"),
"in_parity": parity.get("in_parity"),
},
)
return payload
@mcp.tool()
def gitea_request_mcp_restart(
remote: str = "dadeschools",
@@ -23858,6 +24455,72 @@ def gitea_quarantine_contaminated_review(
}
def _register_fleet_runtime(transport: str = "stdio") -> dict | None:
"""Record this server process in the control-plane runtime registry (#949).
Called once from the official entrypoint immediately after the native
transport bind, so the row can only ever describe the process writing it.
This is the evidence ``gitea_assess_fleet_inventory`` reads; without it the
fleet is unprovable, but a failure here must never prevent the server from
serving, so every error is swallowed after being logged to stderr.
"""
try:
profile_name = gitea_config.selected_profile_name()
try:
profile = get_profile()
except Exception: # noqa: BLE001 - identity resolution must not block boot
profile = {}
allowed = (profile or {}).get("allowed_operations") or []
forbidden = (profile or {}).get("forbidden_operations") or []
role = _role_kind(allowed, forbidden) if allowed else None
ctx = session_ctx.get_session_context() or {}
env = os.environ
provenance = (
"client_managed"
if (
(env.get("GITEA_CLIENT_MANAGED") or "").strip().lower()
in {"1", "true", "yes", "client_managed"}
or (env.get("GITEA_MCP_CLIENT_MANAGED") or "").strip().lower()
in {"1", "true", "yes", "client_managed"}
or (env.get("GITEA_SERVER_PROVENANCE") or "").strip()
== "client_managed"
)
else "manual_launch"
)
record = mcp_fleet_inventory.build_process_runtime_record(
namespace=_fleet_namespace_for_profile(profile_name),
profile=profile_name,
role=role,
remote=ctx.get("remote"),
org=ctx.get("org"),
repo=ctx.get("repository"),
repository_root=PROJECT_ROOT,
startup_head=_process_boot_head_sha,
daemon_start_head=_process_boot_head_sha,
transport=transport,
client_provenance=provenance,
env=env,
)
db, errors = _control_plane_db_or_error()
if db is None:
sys.stderr.write(
f"--- fleet runtime registration skipped: {'; '.join(errors)} ---\n"
)
return None
db.register_mcp_server_runtime(
record,
retention_seconds=mcp_fleet_inventory.RUNTIME_RETENTION_SECONDS,
)
sys.stderr.write(
f"--- fleet runtime registered: {record['runtime_id']} "
f"(cohort {record['cohort_id']}) ---\n"
)
return record
except Exception as exc: # noqa: BLE001 - never block startup
sys.stderr.write(f"--- fleet runtime registration failed: {exc} ---\n")
return None
# ── Entry point ───────────────────────────────────────────────────────────────
if __name__ == "__main__":
@@ -23867,6 +24530,11 @@ if __name__ == "__main__":
# native transport; offline imports / standalone scripts fail closed.
mcp_daemon_guard.mark_sanctioned_daemon()
mcp_daemon_guard.bind_native_mcp_transport(transport="stdio")
# #949: record this process in the control-plane runtime registry. Only a
# transport-bound server reaches this point, so the row is authoritative
# evidence that this namespace is actually running — evidence no caller of
# gitea_assess_fleet_inventory can supply.
_register_fleet_runtime(transport="stdio")
# Lock this session's launch profile into the environment so child CLI
# processes (e.g. review_pr.py) can detect and refuse profile
# side-channel overrides (#199).
+328
View File
@@ -0,0 +1,328 @@
"""Sanctioned MCP client reconnect request surface for Codex/LLM sessions (#678).
Codex and other agent hosts can detect stale or closed Gitea MCP runtimes, but
the host owns the transport. This module never restarts, kills, or reloads a
daemon. It builds:
1. A **callable reconnect request** result agents can invoke via
``gitea_request_mcp_reconnect`` (report-only, side-effect free).
2. A **typed blocker** with exact operator UI steps when recovery must be
performed by the host/operator.
Forbidden recovery paths (must never be recommended):
* ``pkill`` / ``kill`` / ``killall`` of MCP daemons
* ``touch`` / mtime config reload hacks
* ``.env`` or MCP config edits as recovery
* session-state file edits
* raw Gitea API / direct server-import fallbacks
After the operator reconnects, workflows restart from identity / runtime /
capability preflight (``gitea_whoami`` ``gitea_resolve_task_capability``
task).
"""
from __future__ import annotations
from typing import Any, Mapping
# --- Reason vocabulary -------------------------------------------------------
REASON_STALE_RUNTIME = "stale-runtime"
REASON_TRANSPORT_EOF = "transport_eof"
REASON_MISSING_NAMESPACE = "missing_namespace"
REASON_NOT_REQUIRED = "not_required"
REASON_UNSPECIFIED = "unspecified"
VALID_REASONS = frozenset(
{
REASON_STALE_RUNTIME,
REASON_TRANSPORT_EOF,
REASON_MISSING_NAMESPACE,
REASON_NOT_REQUIRED,
REASON_UNSPECIFIED,
}
)
# Boundary statuses reported to callers (match review_workflow_boundary style).
BOUNDARY_CLEAN = "clean"
BOUNDARY_MISMATCH = "mismatch"
BOUNDARY_STALE = "stale"
BOUNDARY_UNKNOWN = "unknown"
# Typed blocker kinds
BLOCKER_OPERATOR_RECONNECT = "operator_mcp_reconnect_required"
BLOCKER_NONE = "none"
FORBIDDEN_RECOVERY_PATHS: tuple[str, ...] = (
"pkill / kill / killall of mcp_server.py, gitea_mcp_server, or broad python sweeps",
"touch / mtime-based MCP config reload hacks",
".env edits as recovery",
"MCP config file edits as recovery",
"session-state file edits as recovery",
"raw Gitea API or direct MCP server-import fallbacks",
)
# Client-specific operator UI steps. Keep Codex first (issue title surface).
OPERATOR_UI_STEPS: dict[str, tuple[str, ...]] = {
"codex": (
"In Codex, open the MCP / Developer tools panel for this workspace.",
"Locate the named Gitea MCP server entry (namespace) that needs reconnect "
"(e.g. gitea-author, gitea-reviewer, gitea-merger, gitea-tools, "
"gitea-controller, gitea-reconciler).",
"Click 'Reload Developer Tools' or the server reconnect/reload control "
"for that entry so the client spawns a fresh MCP subprocess.",
"If per-server reconnect is unavailable, fully restart the Codex client "
"(quit and relaunch) so all MCP namespaces reattach.",
"After reconnect, rerun the blocked workflow from preflight: "
"gitea_whoami → gitea_resolve_task_capability → the original task. "
"Do not resume mid-mutation.",
),
"claude_code": (
"Run `/mcp` (or open the MCP servers UI) in Claude Code.",
"Reconnect the affected gitea-* server entry so the client reopens stdio.",
"If reconnect fails, relaunch the Claude Code session entirely.",
"After reconnect, restart the workflow from gitea_whoami → "
"gitea_resolve_task_capability → task.",
),
"generic": (
"Use the host/IDE MCP reconnect or reload control for the named namespace.",
"If no per-namespace control exists, restart the MCP client/editor.",
"After reconnect, restart the workflow from identity/capability preflight.",
),
}
DEFAULT_CLIENT = "codex"
def normalize_reason(reason: str | None) -> str:
"""Map free-form reason strings onto the closed vocabulary."""
raw = (reason or "").strip().lower()
if not raw:
return REASON_UNSPECIFIED
if raw in VALID_REASONS:
return raw
text = raw.replace(" ", "_").replace("-", "_")
aliases = {
"stale_runtime": REASON_STALE_RUNTIME,
"staleruntime": REASON_STALE_RUNTIME,
"runtime_stale": REASON_STALE_RUNTIME,
"stale": REASON_STALE_RUNTIME,
"transport_eof": REASON_TRANSPORT_EOF,
"transport_closed": REASON_TRANSPORT_EOF,
"eof": REASON_TRANSPORT_EOF,
"client_is_closing": REASON_TRANSPORT_EOF,
"missing_namespace": REASON_MISSING_NAMESPACE,
"namespace_missing": REASON_MISSING_NAMESPACE,
"not_required": REASON_NOT_REQUIRED,
"healthy": REASON_NOT_REQUIRED,
"ok": REASON_NOT_REQUIRED,
"unspecified": REASON_UNSPECIFIED,
}
if text in aliases:
return aliases[text]
hyphenated = text.replace("_", "-")
if hyphenated in VALID_REASONS:
return hyphenated
return REASON_UNSPECIFIED
def normalize_client(client: str | None) -> str:
"""Return a known client key for operator UI steps."""
text = (client or "").strip().lower().replace(" ", "_").replace("-", "_")
if text in ("codex", "openai_codex", "openai"):
return "codex"
if text in ("claude", "claude_code", "claude_desktop", "anthropic"):
return "claude_code"
if text in OPERATOR_UI_STEPS:
return text
return DEFAULT_CLIENT
def classify_boundary_status(
*,
startup_sha: str | None,
current_master_sha: str | None,
live_stale: bool | None = None,
in_parity: bool | None = None,
) -> str:
"""Derive boundary_status from parity evidence."""
if live_stale is True or in_parity is False:
return BOUNDARY_STALE
start = (startup_sha or "").strip().lower()
current = (current_master_sha or "").strip().lower()
if start and current and start != current:
return BOUNDARY_MISMATCH
if start and current and start == current:
return BOUNDARY_CLEAN
if in_parity is True:
return BOUNDARY_CLEAN
return BOUNDARY_UNKNOWN
def operator_ui_steps(client: str | None, *, namespace: str | None = None) -> list[str]:
"""Exact operator UI steps for the named client."""
key = normalize_client(client)
steps = list(OPERATOR_UI_STEPS.get(key) or OPERATOR_UI_STEPS[DEFAULT_CLIENT])
ns = (namespace or "").strip()
if ns:
steps = [
s.replace("named Gitea MCP server entry (namespace)", f"namespace '{ns}'")
.replace("affected gitea-* server entry", f"server entry '{ns}'")
.replace("named namespace", f"namespace '{ns}'")
for s in steps
]
return steps
def build_reconnect_request(
*,
namespace: str,
profile: str | None = None,
pid: int | str | None = None,
session_id: str | None = None,
startup_sha: str | None = None,
current_master_sha: str | None = None,
boundary_status: str | None = None,
reason: str | None = None,
client: str | None = DEFAULT_CLIENT,
live_stale: bool | None = None,
in_parity: bool | None = None,
restart_required: bool | None = None,
stop_required: bool | None = None,
extra: Mapping[str, Any] | None = None,
) -> dict[str, Any]:
"""Build the structured reconnect-request / typed-blocker payload (#678).
Never mutates process, config, or session state. Always side-effect free.
"""
ns = (namespace or "").strip() or "unknown"
normalized_reason = normalize_reason(reason)
boundary = (boundary_status or "").strip() or classify_boundary_status(
startup_sha=startup_sha,
current_master_sha=current_master_sha,
live_stale=live_stale,
in_parity=in_parity,
)
reconnect_needed = True
if normalized_reason == REASON_NOT_REQUIRED and boundary == BOUNDARY_CLEAN:
reconnect_needed = False
if restart_required is False and stop_required is False and boundary == BOUNDARY_CLEAN:
# Explicit healthy probe
if normalized_reason in (REASON_NOT_REQUIRED, REASON_UNSPECIFIED):
reconnect_needed = False
normalized_reason = REASON_NOT_REQUIRED
if restart_required is True or stop_required is True:
reconnect_needed = True
if normalized_reason in (REASON_NOT_REQUIRED, REASON_UNSPECIFIED):
normalized_reason = REASON_STALE_RUNTIME
client_key = normalize_client(client)
steps = operator_ui_steps(client_key, namespace=ns)
result: dict[str, Any] = {
"success": True,
"read_only": True,
"reconnect_performed": False,
"mutation_performed": False,
"reconnect_needed": reconnect_needed,
"namespace": ns,
"profile": (profile or "").strip() or None,
"pid": pid,
"session_id": (session_id or "").strip() or None,
"startup_sha": (startup_sha or "").strip() or None,
"current_master_sha": (current_master_sha or "").strip() or None,
"boundary_status": boundary,
"reason": normalized_reason,
"client": client_key,
"forbidden_recovery_paths": list(FORBIDDEN_RECOVERY_PATHS),
"post_reconnect_preflight": [
"gitea_whoami",
"gitea_resolve_task_capability",
"original_task",
],
"exact_safe_next_action": None,
"blocker_kind": BLOCKER_NONE,
"operator_ui_steps": steps,
"typed_blocker": None,
}
if reconnect_needed:
result["blocker_kind"] = BLOCKER_OPERATOR_RECONNECT
result["stop_required"] = True
result["restart_required"] = True
result["exact_safe_next_action"] = (
f"blocker_kind={BLOCKER_OPERATOR_RECONNECT}: operator must reconnect "
f"MCP namespace '{ns}' via the host UI (client={client_key}). "
"Do not pkill, touch configs, edit session state, or use raw API. "
"After reconnect, restart from gitea_whoami → "
"gitea_resolve_task_capability → task."
)
result["typed_blocker"] = {
"blocker_kind": BLOCKER_OPERATOR_RECONNECT,
"namespaces": [ns],
"why_reconnect_required": normalized_reason,
"operator_ui_steps": steps,
"client": client_key,
"forbidden_recovery_paths": list(FORBIDDEN_RECOVERY_PATHS),
"instruction_after_reconnect": (
"Rerun the blocked workflow from preflight "
"(gitea_whoami → gitea_resolve_task_capability → task). "
"Do not continue mid-mutation from pre-reconnect state."
),
}
else:
result["stop_required"] = False
result["restart_required"] = False
result["exact_safe_next_action"] = (
f"Reconnect not required for namespace '{ns}' "
f"(boundary_status={boundary}). Proceed with the original task."
)
if extra:
for key, value in extra.items():
if key not in result:
result[key] = value
return result
def reasons_never_suggest_forbidden(text: str) -> bool:
"""Return True when *text* does not recommend a forbidden recovery path.
Mentions that *ban* a path (e.g. ``Do not pkill`` / ``never edit session
state``) are allowed. Positive recommendations such as ``use pkill`` or
``run killall`` fail.
"""
import re
lowered = (text or "").lower()
# Strip common ban prefixes so "do not pkill" does not trip positive checks.
scrubbed = re.sub(
r"\b(?:do not|don't|never|must not|forbid(?:den)?|ban(?:ned)?)\b"
r"[^.!;\n]{0,80}",
" ",
lowered,
)
# Positive imperative / advisory forms that would tell an agent to do harm.
positive_suggestions = (
"use pkill",
"run pkill",
"try pkill",
"pkill -f",
"use killall",
"run killall",
"killall mcp",
"use kill ",
"run kill ",
"touch the mcp",
"touch mcp config",
"utime(",
"edit the mcp config to recover",
"edit .env to recover",
"import gitea_mcp_server",
"python -c 'import gitea_mcp",
)
return not any(frag in scrubbed for frag in positive_suggestions)
+237
View File
@@ -0,0 +1,237 @@
"""Antigravity IDE vs Global MCP Config Drift Diagnostic (#672).
Diagnoses config drift between the active IDE MCP configuration
(e.g. ``~/.gemini/antigravity-ide/mcp_config.json``) and the offline/global
canonical configuration (e.g. ``~/.gemini/config/mcp_config.json``).
Hard rules (#672 / #630 / #655):
* Distinguish offline/global success from active IDE namespace availability.
* Never print tokens, DSNs, Authorization headers, or secret-bearing env vars.
* Sanctioned repair path is: backup active config -> patch active config from canonical
-> reconnect through IDE/client -> verify with live ``gitea_whoami``.
* FORBIDDEN: ``pkill``, mtime edits, source edits, or session-state edits for repair.
"""
from __future__ import annotations
import argparse
import json
import os
import sys
from datetime import datetime, timezone
from pathlib import Path
from typing import Any
from webui import console_redaction
DEFAULT_ACTIVE_IDE_CONFIG = "~/.gemini/antigravity-ide/mcp_config.json"
DEFAULT_GLOBAL_CONFIG = "~/.gemini/config/mcp_config.json"
REQUIRED_GITEA_ROLE_SERVERS = (
"gitea-author",
"gitea-reviewer",
"gitea-merger",
"gitea-reconciler",
"gitea-controller",
"gitea-tools",
)
SANCTIONED_REPAIR_RUNBOOK: tuple[str, ...] = (
"1. Backup active IDE config: cp ~/.gemini/antigravity-ide/mcp_config.json ~/.gemini/antigravity-ide/mcp_config.json.bak",
"2. Patch active IDE config: copy required missing Gitea role server entries from global config (~/.gemini/config/mcp_config.json) into active IDE config.",
"3. Reconnect via IDE/client UI or client restart (do NOT use host process kill).",
"4. Verify active namespace health using live gitea_whoami and gitea_resolve_task_capability on each role namespace.",
"FORBIDDEN REPAIR PATHS: pkill / host process kill, mtime touch edits, source code edits, or session-state edits.",
)
def resolve_config_path(path_str: str) -> Path:
"""Expand user and resolve absolute path."""
return Path(os.path.expanduser(path_str)).resolve()
def load_mcp_config(config_path: str | Path) -> tuple[dict[str, Any] | None, str | None]:
"""Load and parse JSON MCP configuration from file.
Returns (config_dict, error_message).
"""
resolved = resolve_config_path(str(config_path))
if not resolved.exists():
return None, f"file_not_found: {resolved}"
try:
with open(resolved, "r", encoding="utf-8") as f:
data = json.load(f)
if not isinstance(data, dict):
return None, f"invalid_schema: root is not a JSON object in {resolved}"
return data, None
except Exception as exc:
return None, f"unreadable_json: {exc} in {resolved}"
def extract_mcp_servers(config: dict[str, Any] | None) -> dict[str, dict[str, Any]]:
"""Extract the mcpServers or mcp_servers mapping safely."""
if not config:
return {}
servers = config.get("mcpServers") or config.get("mcp_servers") or {}
if isinstance(servers, dict):
return {str(k): v for k, v in servers.items() if isinstance(v, dict)}
return {}
def _safe_redact_server_config(srv_cfg: dict[str, Any]) -> dict[str, Any]:
"""Redact secrets from environment variables and command line args."""
safe = {}
if "command" in srv_cfg:
safe["command"] = str(srv_cfg["command"])
if "args" in srv_cfg and isinstance(srv_cfg["args"], list):
safe["args"] = [console_redaction.redact_text(str(a)) for a in srv_cfg["args"]]
if "env" in srv_cfg and isinstance(srv_cfg["env"], dict):
safe_env = {}
for k, v in srv_cfg["env"].items():
if any(secret_kw in k.lower() for secret_kw in ("token", "secret", "pass", "key", "auth")):
safe_env[k] = "[REDACTED]"
else:
safe_env[k] = console_redaction.redact_text(str(v))
safe["env"] = safe_env
return safe
def analyze_config_drift(
active_config_path: str = DEFAULT_ACTIVE_IDE_CONFIG,
global_config_path: str = DEFAULT_GLOBAL_CONFIG,
) -> dict[str, Any]:
"""Analyze MCP configuration drift between active IDE config and global config.
Returns structured diagnostic output.
"""
active_resolved = resolve_config_path(active_config_path)
global_resolved = resolve_config_path(global_config_path)
active_cfg, active_err = load_mcp_config(active_resolved)
global_cfg, global_err = load_mcp_config(global_resolved)
active_servers = extract_mcp_servers(active_cfg)
global_servers = extract_mcp_servers(global_cfg)
missing_role_servers: list[str] = []
present_role_servers: list[str] = []
profile_mismatches: list[dict[str, Any]] = []
reasons: list[str] = []
if active_err:
reasons.append(f"Active IDE config error: {active_err}")
if global_err:
reasons.append(f"Global canonical config error: {global_err}")
# Check Gitea role servers
for srv_name in REQUIRED_GITEA_ROLE_SERVERS:
in_active = srv_name in active_servers
in_global = srv_name in global_servers
if in_active:
present_role_servers.append(srv_name)
elif in_global:
missing_role_servers.append(srv_name)
reasons.append(
f"Missing Gitea role server '{srv_name}' in active IDE config ({active_resolved})"
)
if in_active and in_global:
# Compare profiles & environments
act_env = active_servers[srv_name].get("env", {}) if isinstance(active_servers[srv_name], dict) else {}
glo_env = global_servers[srv_name].get("env", {}) if isinstance(global_servers[srv_name], dict) else {}
act_prof = act_env.get("GITEA_MCP_PROFILE") or act_env.get("GITEA_PROFILE_NAME")
glo_prof = glo_env.get("GITEA_MCP_PROFILE") or glo_env.get("GITEA_PROFILE_NAME")
if act_prof != glo_prof:
mismatch_item = {
"server": srv_name,
"active_profile": act_prof,
"global_profile": glo_prof,
}
profile_mismatches.append(mismatch_item)
reasons.append(
f"Profile mismatch for '{srv_name}': active='{act_prof}' != global='{glo_prof}'"
)
in_sync = bool(
not active_err
and not global_err
and not missing_role_servers
and not profile_mismatches
)
report = {
"timestamp": datetime.now(timezone.utc).isoformat(),
"in_sync": in_sync,
"active_config_path": str(active_resolved),
"active_config_exists": active_cfg is not None,
"global_config_path": str(global_resolved),
"global_config_exists": global_cfg is not None,
"required_role_servers": list(REQUIRED_GITEA_ROLE_SERVERS),
"present_role_servers": present_role_servers,
"missing_role_servers": missing_role_servers,
"profile_mismatches": profile_mismatches,
"reasons": reasons,
"sanctioned_repair_runbook": list(SANCTIONED_REPAIR_RUNBOOK),
"forbidden_repair_methods": [
"pkill / host process kill",
"mtime touch edits",
"source code edits",
"session-state edits",
],
}
return console_redaction.redact_payload(report)
def main() -> None:
parser = argparse.ArgumentParser(
description="Diagnose Gitea MCP role server config drift between active IDE and global config."
)
parser.add_argument(
"--active-config",
default=DEFAULT_ACTIVE_IDE_CONFIG,
help="Path to active IDE MCP config JSON",
)
parser.add_argument(
"--global-config",
default=DEFAULT_GLOBAL_CONFIG,
help="Path to global/canonical MCP config JSON",
)
parser.add_argument(
"--json", action="store_true", help="Print raw JSON report"
)
args = parser.parse_args()
report = analyze_config_drift(args.active_config, args.global_config)
if args.json:
print(json.dumps(report, indent=2))
else:
print("=== MCP Config Drift Diagnostic Report ===")
print(f"Timestamp: {report['timestamp']}")
print(f"In Sync: {report['in_sync']}")
print(f"Active IDE Config: {report['active_config_path']} (exists={report['active_config_exists']})")
print(f"Global Config: {report['global_config_path']} (exists={report['global_config_exists']})")
print(f"Present Role Servers: {', '.join(report['present_role_servers']) if report['present_role_servers'] else 'None'}")
print(f"Missing Role Servers: {', '.join(report['missing_role_servers']) if report['missing_role_servers'] else 'None'}")
if report['profile_mismatches']:
print("Profile Mismatches:")
for m in report['profile_mismatches']:
print(f" - {m['server']}: active={m['active_profile']} vs global={m['global_profile']}")
if report['reasons']:
print("Drift Reasons:")
for r in report['reasons']:
print(f" - {r}")
print("\nSanctioned Repair Runbook:")
for step in report['sanctioned_repair_runbook']:
print(f" {step}")
sys.exit(0 if report["in_sync"] else 1)
if __name__ == "__main__":
main()
+785
View File
@@ -0,0 +1,785 @@
"""Authoritative, read-only PRGS MCP fleet inventory (#949).
Every pre-existing runtime surface is *per process*. ``gitea_get_runtime_context``
and ``gitea_assess_master_parity`` describe only the server answering the call.
``gitea_assess_mcp_namespace_health`` accepts ``process``, ``probe_result`` and
``registered_tools`` **from the caller**, so it cannot constrain the caller. The
control-plane ``sessions`` table records allocator *task* sessions, not server
processes. Five independent self-reports of the same revision therefore never
proved that exactly five processes exist, that no sixth exists, or that all five
belong to one client cohort.
Evidence model
--------------
Two independent sources must agree before a fleet member counts as running:
``control-plane runtime registry``
A row each server writes **about itself** at native transport bind
(:func:`build_process_runtime_record`). No caller can supply it. It is
authoritative for identity: namespace, profile, role, repository binding,
cohort, and the revision the process started at.
``server-side process observation``
A process listing performed by the server answering the inventory call
(:func:`scan_mcp_server_processes`), never by the caller. It is
authoritative for existence and liveness, and it is the only source that
can show a process the registry does not know about.
A member is ``live`` only when a registry row has a matching, still-running
process whose start time precedes the registration (so a recycled PID cannot
impersonate a dead server). Anything the two sources cannot jointly establish
is reported as unknown and fails the mutation gate closed configuration alone
never counts as a running member, and matching Git revisions never establish a
single cohort.
This module performs no restart, reconnect, lease mutation, issue mutation, or
process termination. The only signal it ever sends is ``signal 0`` liveness
probing, which delivers nothing to the target process.
"""
from __future__ import annotations
import os
import secrets
import subprocess
from datetime import datetime, timezone
from typing import Any, Iterable, Mapping, Sequence
# ── expected fleet ────────────────────────────────────────────────────────────
# The configured PRGS fleet. Each entry is one expected member; the roster is
# the definition of "expected" for missing/unexpected classification.
EXPECTED_PRGS_FLEET: tuple[dict[str, str], ...] = (
{"namespace": "gitea-author", "profile": "prgs-author", "role": "author"},
{
"namespace": "gitea-controller",
"profile": "prgs-controller",
"role": "controller",
},
{"namespace": "gitea-reviewer", "profile": "prgs-reviewer", "role": "reviewer"},
{"namespace": "gitea-merger", "profile": "prgs-merger", "role": "merger"},
{
"namespace": "gitea-reconciler",
"profile": "prgs-reconciler",
"role": "reconciler",
},
)
# Cohort identity supplied by a client that manages the whole fleet.
COHORT_ID_ENV = "GITEA_MCP_CLIENT_COHORT_ID"
# Registry rows older than this are pruned at *registration* time (a startup
# write), never on the read path. Dead rows inside the window are still reported
# as stale evidence rather than silently dropped.
RUNTIME_RETENTION_SECONDS = 7 * 24 * 3600
# Liveness classifications for a registry row.
LIVENESS_LIVE = "live"
LIVENESS_DEAD = "dead"
LIVENESS_PID_RECYCLED = "pid_recycled"
LIVENESS_UNOBSERVED = "unobserved"
LIVENESS_UNKNOWN = "unknown"
# Per-member health classifications.
HEALTH_RUNNING = "running"
HEALTH_MISSING = "missing"
HEALTH_DUPLICATE = "duplicate"
HEALTH_UNEXPECTED = "unexpected"
HEALTH_STALE = "stale"
HEALTH_UNKNOWN = "unknown"
EVIDENCE_AUTHORITY = "control_plane_runtime_registry+server_process_observation"
_MCP_PROCESS_MARKER = "mcp_server.py"
_LSTART_FORMAT = "%a %b %d %H:%M:%S %Y"
_ISO_FORMAT = "%Y-%m-%dT%H:%M:%SZ"
# ── time helpers ──────────────────────────────────────────────────────────────
def _utcnow() -> datetime:
return datetime.now(timezone.utc)
def iso_now() -> str:
return _utcnow().strftime(_ISO_FORMAT)
def _parse_iso(value: Any) -> datetime | None:
text = (str(value) if value is not None else "").strip()
if not text:
return None
if text.endswith("Z"):
text = text[:-1] + "+00:00"
try:
parsed = datetime.fromisoformat(text)
except ValueError:
return None
if parsed.tzinfo is None:
parsed = parsed.replace(tzinfo=timezone.utc)
return parsed.astimezone(timezone.utc)
def _int_or_none(value: Any) -> int | None:
try:
return int(value)
except (TypeError, ValueError):
return None
def _clean(value: Any) -> str | None:
text = (str(value) if value is not None else "").strip()
return text or None
# ── process-level evidence (server side only) ─────────────────────────────────
def probe_pid_alive(pid: int | None) -> bool | None:
"""Return whether *pid* exists. ``None`` when it cannot be determined.
Uses ``signal 0``, which performs a permission/existence check and delivers
nothing to the target. This module never sends a terminating signal.
"""
resolved = _int_or_none(pid)
if resolved is None or resolved <= 0:
return None
try:
os.kill(resolved, 0)
except ProcessLookupError:
return False
except PermissionError:
# The process exists but belongs to another user.
return True
except OSError:
return None
return True
def scan_mcp_server_processes(*, runner=subprocess.run) -> dict[str, Any]:
"""Observe running Gitea MCP server processes from this server process.
This is deliberately performed by the answering server, never by the caller:
a caller-supplied process list is exactly the input a fleet gate must not
trust. When the listing cannot be obtained the result reports
``available=False`` so the inventory fails closed instead of assuming that
no unregistered process exists.
"""
try:
proc = runner(
["ps", "-o", "pid,lstart,command", "-ax"],
capture_output=True,
text=True,
check=True,
)
except Exception as exc: # noqa: BLE001 - any failure means unknown evidence
return {
"available": False,
"processes": [],
"reason": f"process listing unavailable: {exc}",
}
processes: list[dict[str, Any]] = []
for raw_line in (proc.stdout or "").splitlines()[1:]:
line = raw_line.strip()
if not line or _MCP_PROCESS_MARKER not in line:
continue
parts = line.split(None, 6)
if len(parts) < 7:
continue
pid = _int_or_none(parts[0])
if pid is None:
continue
try:
naive = datetime.strptime(" ".join(parts[1:6]), _LSTART_FORMAT)
started_at = naive.astimezone(timezone.utc)
except (ValueError, OSError):
started_at = None
processes.append(
{
"pid": pid,
"started_at": started_at.strftime(_ISO_FORMAT) if started_at else None,
"command": parts[6],
}
)
processes.sort(key=lambda item: item["pid"])
return {"available": True, "processes": processes, "reason": None}
# ── cohort / registration record ──────────────────────────────────────────────
def namespace_for_profile(profile: str | None, *, default: str | None = None) -> str | None:
"""Map a configured profile to its fleet namespace.
Resolved from the expected roster rather than from name-shape heuristics.
``role_namespace_gate.infer_mcp_namespace`` only recognises author and
reviewer, so it returns the *profile* name for controller, merger, and
reconciler which would label three of the five members with a namespace
that does not exist. Widening that helper is controller role-metadata work
and belongs to #950; the roster already carries the mapping this inventory
needs, so it is read from there.
A profile outside the roster falls back to *default* (typically the
caller's existing inference), so an unexpected member is still described
rather than dropped.
"""
cleaned = _clean(profile)
for entry in EXPECTED_PRGS_FLEET:
if entry["profile"] == cleaned:
return entry["namespace"]
return default if default is not None else cleaned
def derive_cohort_identity(env: Mapping[str, str] | None = None) -> dict[str, Any]:
"""Derive the client/cohort identity of this server process.
A cohort is the set of servers a single client launched together. The
parent process is the durable expression of that: an IDE/CLI client spawns
every namespace as its own child. A process whose parent has gone away
(reparented to init) cannot prove which cohort it belongs to, and says so
rather than guessing.
Revisions are deliberately not consulted here. Two servers built from the
same commit are not thereby one cohort (#949 AC7).
"""
source_env = os.environ if env is None else env
explicit = _clean(source_env.get(COHORT_ID_ENV))
if explicit:
return {"cohort_id": explicit, "cohort_source": "explicit_env"}
try:
ppid = os.getppid()
except OSError:
ppid = 0
if ppid and ppid > 1:
return {"cohort_id": f"ppid:{ppid}", "cohort_source": "parent_process"}
return {
"cohort_id": None,
"cohort_source": "unknown",
"cohort_reason": (
"parent process is unavailable or reparented to init; this server "
"cannot prove which client cohort launched it"
),
}
def build_process_runtime_record(
*,
namespace: str,
profile: str | None,
role: str | None,
remote: str | None = None,
org: str | None = None,
repo: str | None = None,
repository_root: str | None = None,
pid: int | None = None,
startup_head: str | None = None,
daemon_start_head: str | None = None,
transport: str | None = None,
client_provenance: str | None = None,
env: Mapping[str, str] | None = None,
boot_id: str | None = None,
registered_at: str | None = None,
) -> dict[str, Any]:
"""Build the row a server writes about itself at native transport bind.
Every field describes the *calling* process. Nothing here is caller-supplied
in the MCP sense: the only code that reaches this function is the official
entrypoint of the process being described.
"""
cohort = derive_cohort_identity(env)
resolved_pid = _int_or_none(pid)
if resolved_pid is None:
resolved_pid = os.getpid()
token = boot_id or secrets.token_hex(8)
return {
"runtime_id": f"{namespace}:{resolved_pid}:{token}",
"namespace": namespace,
"profile": _clean(profile),
"role": _clean(role),
"remote": _clean(remote),
"org": _clean(org),
"repo": _clean(repo),
"repository_root": _clean(repository_root),
"pid": resolved_pid,
"cohort_id": cohort["cohort_id"],
"cohort_source": cohort["cohort_source"],
"client_provenance": _clean(client_provenance) or "unknown",
"boot_id": token,
"startup_head": _clean(startup_head),
"daemon_start_head": _clean(daemon_start_head),
"transport": _clean(transport),
"registered_at": registered_at or iso_now(),
"status": "running",
}
# ── classification ────────────────────────────────────────────────────────────
def _expected_index(
expected_fleet: Sequence[Mapping[str, str]],
) -> dict[str, dict[str, str]]:
index: dict[str, dict[str, str]] = {}
for entry in expected_fleet:
profile = _clean(entry.get("profile"))
if profile:
index[profile] = dict(entry)
return index
def _sort_key(member: Mapping[str, Any]) -> tuple:
return (
str(member.get("namespace") or ""),
str(member.get("profile") or ""),
_int_or_none(member.get("pid")) or 0,
str(member.get("runtime_id") or ""),
)
def _normalize_row(
row: Mapping[str, Any],
*,
observed_by_pid: Mapping[int, Mapping[str, Any]],
process_scan_available: bool,
now: datetime,
) -> dict[str, Any]:
pid = _int_or_none(row.get("pid"))
member: dict[str, Any] = {
"runtime_id": _clean(row.get("runtime_id")),
"namespace": _clean(row.get("namespace")),
"profile": _clean(row.get("profile")),
"role": _clean(row.get("role")),
"remote": _clean(row.get("remote")),
"org": _clean(row.get("org")),
"repo": _clean(row.get("repo")),
"repository_root": _clean(row.get("repository_root")),
"pid": pid,
"cohort_id": _clean(row.get("cohort_id")),
"cohort_source": _clean(row.get("cohort_source")) or "unknown",
"client_provenance": _clean(row.get("client_provenance")) or "unknown",
"boot_id": _clean(row.get("boot_id")),
"startup_head": _clean(row.get("startup_head")),
"daemon_start_head": _clean(row.get("daemon_start_head")),
"transport": _clean(row.get("transport")),
"registered_at": _clean(row.get("registered_at")),
"last_heartbeat_at": _clean(row.get("last_heartbeat_at")),
"recorded_status": _clean(row.get("status")) or "unknown",
}
pid_alive = probe_pid_alive(pid)
member["pid_alive"] = pid_alive
observed = observed_by_pid.get(pid) if pid is not None else None
member["process_observed"] = bool(observed) if process_scan_available else None
if not process_scan_available:
# Existence cannot be corroborated; never upgrade to live on the
# registry's word alone.
member["liveness"] = LIVENESS_UNKNOWN
member["liveness_reason"] = (
"process observation unavailable; registry rows cannot be corroborated"
)
elif pid_alive is False:
member["liveness"] = LIVENESS_DEAD
member["liveness_reason"] = "recorded PID is not running"
elif pid_alive is None:
member["liveness"] = LIVENESS_UNKNOWN
member["liveness_reason"] = "PID liveness could not be determined"
elif observed is None:
member["liveness"] = LIVENESS_UNOBSERVED
member["liveness_reason"] = (
"recorded PID is not a running Gitea MCP server process"
)
else:
started_at = _parse_iso(observed.get("started_at"))
registered_at = _parse_iso(member["registered_at"])
if started_at and registered_at and started_at > registered_at:
member["liveness"] = LIVENESS_PID_RECYCLED
member["liveness_reason"] = (
"the process now holding this PID started after the registry row "
"was written; the registered server is gone"
)
else:
member["liveness"] = LIVENESS_LIVE
member["liveness_reason"] = None
heartbeat = _parse_iso(member["last_heartbeat_at"])
member["heartbeat_age_seconds"] = (
int((now - heartbeat).total_seconds()) if heartbeat else None
)
return member
def _binding_matches(
member: Mapping[str, Any], expected_binding: Mapping[str, Any] | None
) -> bool | None:
if not expected_binding:
return None
for field in ("remote", "org", "repo"):
expected = _clean(expected_binding.get(field))
if expected is None:
continue
actual = _clean(member.get(field))
if actual is None:
return None
if actual != expected:
return False
return True
def classify_fleet_inventory(
*,
runtime_rows: Iterable[Mapping[str, Any]],
process_scan: Mapping[str, Any] | None = None,
expected_fleet: Sequence[Mapping[str, str]] = EXPECTED_PRGS_FLEET,
expected_binding: Mapping[str, Any] | None = None,
registry_available: bool = True,
registry_error: str | None = None,
now: datetime | None = None,
answering_namespace: str | None = None,
) -> dict[str, Any]:
"""Classify a fleet snapshot. Pure: identical input yields identical output.
The verdict never depends on which namespace asked, so controller and
reconciler agree by construction; ``answering_namespace`` is reported as
metadata only.
"""
moment = now or _utcnow()
scan = dict(process_scan or {"available": False, "processes": [], "reason": None})
scan_available = bool(scan.get("available"))
observed_processes = list(scan.get("processes") or [])
observed_by_pid: dict[int, Mapping[str, Any]] = {}
for proc in observed_processes:
observed_pid = _int_or_none(proc.get("pid"))
if observed_pid is not None:
observed_by_pid[observed_pid] = proc
expected_index = _expected_index(expected_fleet)
members = [
_normalize_row(
row,
observed_by_pid=observed_by_pid,
process_scan_available=scan_available,
now=moment,
)
for row in (runtime_rows or [])
]
live = [m for m in members if m["liveness"] == LIVENESS_LIVE]
not_live = [m for m in members if m["liveness"] != LIVENESS_LIVE]
live_by_profile: dict[str, list[dict[str, Any]]] = {}
for member in live:
live_by_profile.setdefault(member["profile"] or "", []).append(member)
running_members: list[dict[str, Any]] = []
missing_members: list[dict[str, Any]] = []
duplicate_members: list[dict[str, Any]] = []
unexpected_members: list[dict[str, Any]] = []
binding_mismatches: list[dict[str, Any]] = []
role_mismatches: list[dict[str, Any]] = []
configured_members: list[dict[str, Any]] = []
for entry in expected_fleet:
profile = _clean(entry.get("profile")) or ""
instances = sorted(live_by_profile.get(profile, []), key=_sort_key)
configured_members.append(
{
"namespace": _clean(entry.get("namespace")),
"profile": profile,
"role": _clean(entry.get("role")),
"instance_count": len(instances),
"health": (
HEALTH_MISSING
if not instances
else HEALTH_RUNNING
if len(instances) == 1
else HEALTH_DUPLICATE
),
}
)
if not instances:
missing_members.append(
{
"namespace": _clean(entry.get("namespace")),
"profile": profile,
"role": _clean(entry.get("role")),
"health": HEALTH_MISSING,
"reason": (
"no live registry row corroborated by a running server "
"process"
),
}
)
continue
for instance in instances:
instance["health"] = (
HEALTH_RUNNING if len(instances) == 1 else HEALTH_DUPLICATE
)
instance["expected"] = True
running_members.append(instance)
if len(instances) > 1:
duplicate_members.append(
{
"namespace": _clean(entry.get("namespace")),
"profile": profile,
"role": _clean(entry.get("role")),
"health": HEALTH_DUPLICATE,
"instance_count": len(instances),
"pids": sorted(
i["pid"] for i in instances if i["pid"] is not None
),
"instances": instances,
"reason": (
"more than one live server is registered for this profile"
),
}
)
for member in sorted(live, key=_sort_key):
profile = member["profile"] or ""
expected_entry = expected_index.get(profile)
if expected_entry is not None:
expected_role = _clean(expected_entry.get("role"))
actual_role = member["role"]
if expected_role and actual_role and actual_role != expected_role:
role_mismatches.append(
{
"namespace": member["namespace"],
"profile": profile,
"pid": member["pid"],
"expected_role": expected_role,
"declared_role": actual_role,
"reason": "declared role does not match the configured profile",
}
)
expected_namespace = _clean(expected_entry.get("namespace"))
if (
expected_namespace
and member["namespace"]
and member["namespace"] != expected_namespace
):
role_mismatches.append(
{
"namespace": member["namespace"],
"profile": profile,
"pid": member["pid"],
"expected_namespace": expected_namespace,
"declared_role": member["role"],
"reason": (
"profile is served from a namespace it is not "
"configured for"
),
}
)
else:
member["health"] = HEALTH_UNEXPECTED
member["expected"] = False
unexpected_members.append(member)
match = _binding_matches(member, expected_binding)
member["repository_binding_matches"] = match
if match is False:
binding_mismatches.append(
{
"namespace": member["namespace"],
"profile": profile,
"pid": member["pid"],
"remote": member["remote"],
"org": member["org"],
"repo": member["repo"],
"expected": dict(expected_binding or {}),
"reason": "member is bound to a different repository",
}
)
elif match is None and expected_binding:
binding_mismatches.append(
{
"namespace": member["namespace"],
"profile": profile,
"pid": member["pid"],
"remote": member["remote"],
"org": member["org"],
"repo": member["repo"],
"expected": dict(expected_binding or {}),
"reason": "member did not record a complete repository binding",
}
)
stale_members: list[dict[str, Any]] = []
for member in sorted(not_live, key=_sort_key):
member["health"] = (
HEALTH_UNKNOWN if member["liveness"] == LIVENESS_UNKNOWN else HEALTH_STALE
)
stale_members.append(member)
registered_pids = {m["pid"] for m in live if m["pid"] is not None}
unregistered_processes: list[dict[str, Any]] = []
if scan_available:
for proc in observed_processes:
observed_pid = _int_or_none(proc.get("pid"))
if observed_pid is None or observed_pid in registered_pids:
continue
unregistered_processes.append(
{"pid": observed_pid, "started_at": proc.get("started_at")}
)
unregistered_processes.sort(key=lambda item: item["pid"])
# Cohort. Derived only from recorded cohort identity — never from revisions.
cohort_ids = {m["cohort_id"] for m in live}
cohort_unknown = any(cohort_id is None for cohort_id in cohort_ids)
known_cohorts = sorted(c for c in cohort_ids if c is not None)
if not live or cohort_unknown:
single_cohort: bool | None = None
mixed_cohort: bool | None = None
else:
single_cohort = len(known_cohorts) == 1
mixed_cohort = len(known_cohorts) > 1
# Revision spread. Reported independently of cohort; never used to infer it.
heads = {m["startup_head"] for m in live}
head_unknown = any(head is None for head in heads)
known_heads = sorted(h for h in heads if h is not None)
mixed_revision = (len(known_heads) > 1) if known_heads else None
exactly_one_per_profile = not missing_members and not duplicate_members
no_unexpected_members = not unexpected_members
incomplete_reasons: list[str] = []
if not registry_available:
incomplete_reasons.append(
registry_error or "the control-plane runtime registry could not be read"
)
if not scan_available:
incomplete_reasons.append(
str(scan.get("reason") or "server-side process observation unavailable")
)
if unregistered_processes:
pids = ", ".join(str(p["pid"]) for p in unregistered_processes)
incomplete_reasons.append(
f"running Gitea MCP server process(es) with no runtime registry row "
f"(PIDs: {pids}); the fleet contains members this inventory cannot "
f"describe"
)
if live and cohort_unknown:
incomplete_reasons.append(
"one or more live members did not record a client cohort identity; "
"matching revisions do not establish a single cohort"
)
if live and head_unknown:
incomplete_reasons.append(
"one or more live members did not record a startup revision"
)
if any(m["liveness"] == LIVENESS_UNKNOWN for m in members):
incomplete_reasons.append(
"liveness of one or more registry rows could not be determined"
)
inventory_complete = not incomplete_reasons
blocked_reasons: list[str] = list(incomplete_reasons)
if missing_members:
names = ", ".join(sorted(m["profile"] for m in missing_members))
blocked_reasons.append(f"expected fleet member(s) not running: {names}")
if duplicate_members:
names = ", ".join(sorted(d["profile"] for d in duplicate_members))
blocked_reasons.append(
f"duplicate server(s) registered for profile(s): {names}"
)
if unexpected_members:
names = ", ".join(
sorted(
str(m["profile"] or m["namespace"] or "?") for m in unexpected_members
)
)
blocked_reasons.append(f"unexpected fleet member(s) running: {names}")
if mixed_cohort:
blocked_reasons.append(
"live members span more than one client cohort: "
+ ", ".join(known_cohorts)
)
if mixed_revision:
blocked_reasons.append(
"live members started at different revisions: " + ", ".join(known_heads)
)
if binding_mismatches:
blocked_reasons.append(
"one or more live members are not bound to the expected repository"
)
if role_mismatches:
blocked_reasons.append(
"one or more live members declare a role or namespace that does not "
"match the configured profile"
)
mutation_gate_satisfied = bool(
inventory_complete
and exactly_one_per_profile
and no_unexpected_members
and single_cohort is True
and mixed_revision is False
and not binding_mismatches
and not role_mismatches
)
if not mutation_gate_satisfied and not blocked_reasons:
blocked_reasons.append(
"the fleet snapshot did not establish the exact-one-instance-per-"
"profile, single-cohort invariant"
)
return {
"success": True,
"read_only": True,
"evidence_authority": EVIDENCE_AUTHORITY,
"answering_namespace": _clean(answering_namespace),
"generated_at": moment.strftime(_ISO_FORMAT),
"inventory_complete": inventory_complete,
"incomplete_reasons": incomplete_reasons,
"configured_members": configured_members,
"running_members": sorted(running_members, key=_sort_key),
"missing_members": sorted(missing_members, key=_sort_key),
"duplicate_members": sorted(duplicate_members, key=_sort_key),
"unexpected_members": sorted(unexpected_members, key=_sort_key),
"stale_members": stale_members,
"unregistered_processes": unregistered_processes,
"repository_binding_mismatches": sorted(binding_mismatches, key=_sort_key),
"role_mismatches": sorted(role_mismatches, key=_sort_key),
"cohort_ids": known_cohorts,
"single_cohort": single_cohort,
"mixed_cohort": mixed_cohort,
"startup_revisions": known_heads,
"mixed_revision": mixed_revision,
"exactly_one_per_profile": exactly_one_per_profile,
"no_unexpected_members": no_unexpected_members,
"expected_member_count": len(expected_fleet),
"running_member_count": len(running_members),
"mutation_gate_satisfied": mutation_gate_satisfied,
"blocked_reason": blocked_reasons[0] if blocked_reasons else None,
"blocked_reasons": blocked_reasons,
"process_observation": {
"available": scan_available,
"observed_process_count": len(observed_processes),
"reason": scan.get("reason"),
},
"registry": {
"available": registry_available,
"row_count": len(members),
"error": registry_error,
},
"mutations_performed": [],
}
def summarize(result: Mapping[str, Any]) -> str:
"""One-line human summary of a classification result."""
if result.get("mutation_gate_satisfied"):
return (
f"fleet healthy: {result.get('running_member_count')} of "
f"{result.get('expected_member_count')} members running in a single "
f"cohort at one revision"
)
return f"fleet not provable: {result.get('blocked_reason')}"
+35 -5
View File
@@ -225,8 +225,33 @@ _RESTART_PATHS: tuple[RestartPath, ...] = (
"exact_safe_next_action pointing at IDE/client reconnect; performs "
"no restart, thread spawn, config touch, or os._exit."
),
locations=("gitea_mcp_server.py (gitea_resolve_task_capability)",),
references=("#685", "#657"),
locations=(
"gitea_mcp_server.py (gitea_resolve_task_capability)",
"gitea_mcp_server.py (gitea_request_mcp_reconnect)",
"mcp_client_reconnect.py",
),
references=("#685", "#657", "#678"),
),
RestartPath(
path_id="codex_client_reconnect_request",
title="Sanctioned Codex/LLM reconnect request tool",
mechanism=(
"gitea_request_mcp_reconnect: agents invoke a report-only tool that "
"returns namespace/profile/pid/startup SHA/master SHA/boundary "
"status plus a typed operator blocker with exact client UI steps."
),
classification=CLASS_GUARDED_FAIL_CLOSED,
guard=(
"Report-only (#678): never restarts, kills, reloads, or edits "
"config; recovery is always host/operator reconnect. Forbidden "
"paths (pkill, touch, .env/config/session-state hacks) are listed "
"and never recommended."
),
locations=(
"mcp_client_reconnect.py",
"gitea_mcp_server.py (gitea_request_mcp_reconnect)",
),
references=("#678", "#630", "#685", "#657"),
),
RestartPath(
path_id="manual_daemon_kill",
@@ -271,7 +296,8 @@ _RESTART_PATHS: tuple[RestartPath, ...] = (
title="Host/IDE MCP reconnect",
mechanism=(
"A manual `/mcp reconnect` (or equivalent host action) that the "
"IDE performs to recreate the MCP client connection."
"IDE performs to recreate the MCP client connection. Agents obtain "
"exact UI steps via gitea_request_mcp_reconnect (#678)."
),
classification=CLASS_HOST_RESIDUAL,
guard=(
@@ -279,8 +305,12 @@ _RESTART_PATHS: tuple[RestartPath, ...] = (
"gates point operators toward; documented as residual host "
"behavior. No in-process code initiates it."
),
locations=("host/IDE",),
references=("#584", "#656", "#657"),
locations=(
"host/IDE",
"mcp_client_reconnect.py",
"gitea_mcp_server.py (gitea_request_mcp_reconnect)",
),
references=("#584", "#656", "#657", "#678"),
residual_host=True,
),
RestartPath(
+1
View File
@@ -40,6 +40,7 @@ NON_TOOL_IDENTIFIERS: frozenset[str] = frozenset(
"gitea_auth",
"gitea_config",
"gitea_mcp_server",
"mcp_fleet_inventory",
"mcp_server",
"offline_mcp_helper",
"offline_mcp_runner",
+3
View File
@@ -0,0 +1,3 @@
[pytest]
testpaths = tests
norecursedirs = branches .git venv __pycache__ graphify-out
+25
View File
@@ -252,6 +252,31 @@ Helpers: `scripts/worktree-start`, `scripts/worktree-review`,
- Never place raw tokens in LLM/MCP config.
- Use `gitea_whoami` and `gitea_resolve_task_capability` before mutating.
## Fleet inventory
`gitea_whoami`, `gitea_get_runtime_context` and `gitea_assess_master_parity` each
describe only the server answering the call. Five namespaces independently
reporting the same revision never proved that five processes exist, that no sixth
exists, or that all five belong to one client cohort.
`gitea_assess_fleet_inventory` is the read-only capability that does prove it. It
takes no evidence parameters: it combines the control-plane runtime registry,
which each server writes about itself at native transport bind, with a process
observation the answering server performs. Classification is a pure function of
that snapshot, so `gitea-controller` and `gitea-reconciler` return the same
verdict for the same fleet.
Consume `mutation_gate_satisfied`. When it is false, report `blocked_reason`
verbatim and stop — `missing_members`, `duplicate_members`, `unexpected_members`,
`unregistered_processes`, `mixed_cohort` and `mixed_revision` are reported
separately because each needs a different operator action. Treat
`inventory_complete: false` and `single_cohort: null` as *unknown*, never as
healthy. The capability never terminates a duplicate process, restarts,
reconnects, or touches a lease.
Details, field meanings, and the gate-consumption sequence:
[`docs/mcp-fleet-inventory.md`](../../docs/mcp-fleet-inventory.md) (#949).
## Tool inventory
[`docs/mcp-tool-inventory.md`](../../docs/mcp-tool-inventory.md) is the canonical
+13
View File
@@ -159,6 +159,19 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"permission": "gitea.read",
"role": "reconciler",
},
# #949: fleet inventory is strictly read-only evidence. It is deliberately
# not role-exclusive — the controller and reconciler gates are its named
# consumers, but any namespace holding gitea.read must be able to prove the
# exact-five-process/single-cohort invariant before it acts, and the verdict
# is a pure function of the snapshot so every namespace agrees.
"assess_fleet_inventory": {
"permission": "gitea.read",
"role": "reconciler",
},
"gitea_assess_fleet_inventory": {
"permission": "gitea.read",
"role": "reconciler",
},
# PR synchronization lifecycle: assess is read-only (any role with gitea.read);
# update-by-merge is author-only and mutates the PR head via Gitea API.
"assess_pr_sync_status": {
+6 -2
View File
@@ -14,6 +14,7 @@ from control_plane_db import (
ControlPlaneError,
InvalidWorkKindError,
LeaseRequiredError,
SCHEMA_VERSION,
WORK_KINDS,
_ts,
_utc_now,
@@ -37,7 +38,10 @@ class ControlPlaneDBTest(unittest.TestCase):
rows = dict(conn.execute("SELECT key, value FROM schema_meta").fetchall())
finally:
conn.close()
self.assertEqual(rows["schema_version"], "5")
# Pinned to the constant, not a literal: every additive migration bumps
# SCHEMA_VERSION, and the invariant under test is that the meta row
# records the version the code actually wrote (#949 added v6).
self.assertEqual(rows["schema_version"], str(SCHEMA_VERSION))
self.assertIn("DB coordinates", rows["architecture"])
self.assertIn("bridge", rows["architecture"].lower())
@@ -868,7 +872,7 @@ class SessionCheckpointTest(unittest.TestCase):
conn.close()
self.assertIn("session_checkpoints", names)
record = self._write()
self.assertEqual(record["checkpoint_schema_version"], 5)
self.assertEqual(record["checkpoint_schema_version"], SCHEMA_VERSION)
# AC2 — checkpoints written for multi-role session fixtures.
def test_multi_role_fixtures_each_get_a_row(self) -> None:
@@ -0,0 +1,282 @@
"""Control-plane MCP server runtime registry (#949).
The registry is the half of the fleet evidence a caller cannot supply: each
server writes exactly one row about itself at native transport bind. These tests
pin the storage contract additive migration, deterministic ordering, and
PID-reuse/retention pruning that cannot hide a live duplicate.
"""
import os
import tempfile
import unittest
from datetime import datetime, timedelta, timezone
from unittest import mock
import control_plane_db
import mcp_fleet_inventory as mfi
def _stamp(delta_seconds=0):
moment = datetime.now(timezone.utc) + timedelta(seconds=delta_seconds)
return moment.replace(microsecond=0).strftime("%Y-%m-%dT%H:%M:%SZ")
class _DBCase(unittest.TestCase):
def setUp(self):
self._tmp = tempfile.TemporaryDirectory()
self.addCleanup(self._tmp.cleanup)
self.db_path = os.path.join(self._tmp.name, "control_plane.sqlite3")
self.db = control_plane_db.ControlPlaneDB(db_path=self.db_path)
def record(self, namespace, profile, role, pid, **overrides):
base = mfi.build_process_runtime_record(
namespace=namespace,
profile=profile,
role=role,
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
repository_root="/checkout/Gitea-Tools",
pid=pid,
startup_head="82d71b77028a7abd4f8ab4a4e4d89658a187f73d",
transport="stdio",
client_provenance="client_managed",
env={mfi.COHORT_ID_ENV: "ppid:40990"},
boot_id=f"boot{pid}",
)
base.update(overrides)
return base
class TestSchema(_DBCase):
def test_schema_version_is_bumped(self):
self.assertGreaterEqual(control_plane_db.SCHEMA_VERSION, 6)
def test_runtime_table_exists_on_a_fresh_database(self):
self.assertEqual(self.db.list_mcp_server_runtimes(), [])
def test_migration_is_additive_and_idempotent(self):
"""Re-opening an existing DB must not disturb the other tables."""
self.db.upsert_session(session_id="s-a", role="author", profile="prgs-author")
self.db.register_mcp_server_runtime(
self.record("gitea-author", "prgs-author", "author", 41000)
)
reopened = control_plane_db.ControlPlaneDB(db_path=self.db_path)
self.assertEqual(len(reopened.list_sessions()), 1)
self.assertEqual(len(reopened.list_mcp_server_runtimes()), 1)
class TestRegistration(_DBCase):
def test_registration_round_trips_every_field(self):
record = self.record("gitea-controller", "prgs-controller", "controller", 41001)
stored = self.db.register_mcp_server_runtime(record)
for key in (
"runtime_id",
"namespace",
"profile",
"role",
"remote",
"org",
"repo",
"repository_root",
"pid",
"cohort_id",
"cohort_source",
"client_provenance",
"boot_id",
"startup_head",
"transport",
"status",
):
self.assertEqual(stored[key], record[key], key)
def test_registration_requires_a_runtime_id(self):
record = self.record("gitea-author", "prgs-author", "author", 41002)
record["runtime_id"] = ""
with self.assertRaises(ValueError):
self.db.register_mcp_server_runtime(record)
def test_registration_requires_a_namespace(self):
record = self.record("gitea-author", "prgs-author", "author", 41003)
record["namespace"] = " "
with self.assertRaises(ValueError):
self.db.register_mcp_server_runtime(record)
def test_registration_requires_a_pid(self):
record = self.record("gitea-author", "prgs-author", "author", 41004)
record["pid"] = None
with self.assertRaises(ValueError):
self.db.register_mcp_server_runtime(record)
def test_five_members_register_independently(self):
for index, entry in enumerate(mfi.EXPECTED_PRGS_FLEET):
self.db.register_mcp_server_runtime(
self.record(
entry["namespace"], entry["profile"], entry["role"], 41100 + index
)
)
rows = self.db.list_mcp_server_runtimes()
self.assertEqual(len(rows), 5)
self.assertEqual(
sorted(r["profile"] for r in rows),
[
"prgs-author",
"prgs-controller",
"prgs-merger",
"prgs-reconciler",
"prgs-reviewer",
],
)
class TestPruning(_DBCase):
def test_reusing_a_pid_replaces_the_stale_row(self):
first = self.record("gitea-author", "prgs-author", "author", 41200)
self.db.register_mcp_server_runtime(first)
second = self.record("gitea-reviewer", "prgs-reviewer", "reviewer", 41200)
self.db.register_mcp_server_runtime(second)
rows = self.db.list_mcp_server_runtimes()
self.assertEqual(len(rows), 1)
self.assertEqual(rows[0]["runtime_id"], second["runtime_id"])
def test_pruning_never_removes_a_live_duplicate_on_another_pid(self):
"""The defect this must not have: hiding a second running server."""
first = self.record("gitea-author", "prgs-author", "author", 41300)
self.db.register_mcp_server_runtime(first)
second = self.record("gitea-author", "prgs-author", "author", 41301)
self.db.register_mcp_server_runtime(
second, retention_seconds=mfi.RUNTIME_RETENTION_SECONDS
)
rows = self.db.list_mcp_server_runtimes()
self.assertEqual(len(rows), 2)
self.assertEqual(sorted(r["pid"] for r in rows), [41300, 41301])
def test_retention_drops_rows_older_than_the_window(self):
old = self.record(
"gitea-merger",
"prgs-merger",
"merger",
41400,
registered_at=_stamp(-(mfi.RUNTIME_RETENTION_SECONDS + 3600)),
)
self.db.register_mcp_server_runtime(old)
fresh = self.record("gitea-author", "prgs-author", "author", 41401)
self.db.register_mcp_server_runtime(
fresh, retention_seconds=mfi.RUNTIME_RETENTION_SECONDS
)
rows = self.db.list_mcp_server_runtimes()
self.assertEqual([r["pid"] for r in rows], [41401])
def test_retention_is_skipped_when_not_requested(self):
old = self.record(
"gitea-merger",
"prgs-merger",
"merger",
41500,
registered_at=_stamp(-(mfi.RUNTIME_RETENTION_SECONDS + 3600)),
)
self.db.register_mcp_server_runtime(old)
self.db.register_mcp_server_runtime(
self.record("gitea-author", "prgs-author", "author", 41501)
)
self.assertEqual(len(self.db.list_mcp_server_runtimes()), 2)
class TestListing(_DBCase):
def test_listing_is_deterministically_ordered(self):
for index in range(5):
self.db.register_mcp_server_runtime(
self.record(
"gitea-author",
"prgs-author",
"author",
41600 + index,
registered_at="2026-07-28T02:00:00Z",
)
)
first = [r["runtime_id"] for r in self.db.list_mcp_server_runtimes()]
second = [r["runtime_id"] for r in self.db.list_mcp_server_runtimes()]
self.assertEqual(first, second)
self.assertEqual(first, sorted(first))
def test_status_filter_selects_running_rows(self):
running = self.record("gitea-author", "prgs-author", "author", 41700)
self.db.register_mcp_server_runtime(running)
stopped = self.record("gitea-merger", "prgs-merger", "merger", 41701)
self.db.register_mcp_server_runtime(stopped)
self.db.mark_mcp_server_runtime_stopped(stopped["runtime_id"])
rows = self.db.list_mcp_server_runtimes(statuses=("running",))
self.assertEqual([r["pid"] for r in rows], [41700])
def test_listing_without_a_filter_returns_stopped_rows_too(self):
stopped = self.record("gitea-merger", "prgs-merger", "merger", 41800)
self.db.register_mcp_server_runtime(stopped)
self.db.mark_mcp_server_runtime_stopped(stopped["runtime_id"])
rows = self.db.list_mcp_server_runtimes()
self.assertEqual(len(rows), 1)
self.assertEqual(rows[0]["status"], "stopped")
class TestHeartbeat(_DBCase):
def test_heartbeat_advances_only_the_timestamp(self):
record = self.record(
"gitea-author",
"prgs-author",
"author",
41900,
last_heartbeat_at="2026-07-28T02:00:00Z",
)
self.db.register_mcp_server_runtime(record)
self.db.heartbeat_mcp_server_runtime(record["runtime_id"])
row = self.db.list_mcp_server_runtimes()[0]
self.assertNotEqual(row["last_heartbeat_at"], "2026-07-28T02:00:00Z")
self.assertEqual(row["status"], "running")
self.assertEqual(row["pid"], 41900)
def test_heartbeat_for_an_unknown_runtime_is_a_no_op(self):
self.db.heartbeat_mcp_server_runtime("does-not-exist")
self.assertEqual(self.db.list_mcp_server_runtimes(), [])
class TestEndToEndSnapshot(_DBCase):
"""Registry rows feed the classifier without reshaping."""
def test_registered_fleet_classifies_as_healthy(self):
pids = []
for index, entry in enumerate(mfi.EXPECTED_PRGS_FLEET):
pid = 42000 + index
pids.append(pid)
self.db.register_mcp_server_runtime(
self.record(
entry["namespace"],
entry["profile"],
entry["role"],
pid,
registered_at="2026-07-28T02:00:00Z",
)
)
rows = self.db.list_mcp_server_runtimes(statuses=("running",))
scan = {
"available": True,
"processes": [
{"pid": pid, "started_at": "2026-07-28T01:00:00Z"} for pid in pids
],
"reason": None,
}
with mock.patch.object(mfi, "probe_pid_alive", return_value=True):
result = mfi.classify_fleet_inventory(
runtime_rows=rows,
process_scan=scan,
expected_binding={
"remote": "prgs",
"org": "Scaled-Tech-Consulting",
"repo": "Gitea-Tools",
},
)
self.assertTrue(result["inventory_complete"])
self.assertTrue(result["mutation_gate_satisfied"])
self.assertEqual(result["running_member_count"], 5)
if __name__ == "__main__":
unittest.main()
+337
View File
@@ -0,0 +1,337 @@
"""Native exposure of the fleet inventory capability (#949).
Covers the parts the pure classifier cannot: that the tool is registered and
documented, that it is gated on ``gitea.read`` rather than a role, that it
accepts no caller-supplied evidence, that it fails closed on an unreadable
registry, and that the workflow documentation explains gate consumption.
"""
import inspect
import os
import pathlib
import unittest
from unittest import mock
import gitea_mcp_server
import mcp_fleet_inventory as mfi
import task_capability_map
REPO_ROOT = pathlib.Path(__file__).resolve().parent.parent
TOOL_NAME = "gitea_assess_fleet_inventory"
class TestCapabilityMapExposure(unittest.TestCase):
"""AC14: the invariant is provable through sanctioned native calls."""
def test_task_resolves_to_a_read_permission(self):
self.assertEqual(
task_capability_map.required_permission("assess_fleet_inventory"),
"gitea.read",
)
self.assertEqual(
task_capability_map.required_permission(TOOL_NAME), "gitea.read"
)
def test_task_and_tool_keys_agree(self):
self.assertEqual(
task_capability_map.TASK_CAPABILITY_MAP["assess_fleet_inventory"],
task_capability_map.TASK_CAPABILITY_MAP[TOOL_NAME],
)
def test_task_is_not_role_exclusive(self):
"""Controller and reconciler both need it; so does any read-only gate."""
self.assertNotIn(
"assess_fleet_inventory", task_capability_map.ROLE_EXCLUSIVE_TASKS
)
self.assertNotIn(TOOL_NAME, task_capability_map.ROLE_EXCLUSIVE_TASKS)
def test_task_is_not_registered_as_an_issue_mutation(self):
self.assertNotIn(TOOL_NAME, task_capability_map.ISSUE_MUTATION_TOOL_TASKS)
def test_unknown_task_still_fails_closed(self):
with self.assertRaises(KeyError):
task_capability_map.required_permission("assess_fleet_inventory_typo")
class TestToolRegistration(unittest.TestCase):
def test_tool_is_registered_on_the_server(self):
self.assertTrue(hasattr(gitea_mcp_server, TOOL_NAME))
def test_tool_accepts_no_caller_supplied_evidence(self):
"""The defect #949 names: evidence parameters the caller controls."""
signature = inspect.signature(gitea_mcp_server.gitea_assess_fleet_inventory)
self.assertEqual(sorted(signature.parameters), ["host", "org", "remote", "repo"])
for forbidden in ("process", "probe_result", "registered_tools", "processes"):
self.assertNotIn(forbidden, signature.parameters)
def test_tool_is_documented_in_the_canonical_inventory(self):
doc = (REPO_ROOT / "docs" / "mcp-tool-inventory.md").read_text()
self.assertIn(f"`{TOOL_NAME}`", doc)
def test_docstring_states_the_read_only_guarantee(self):
doc = (gitea_mcp_server.gitea_assess_fleet_inventory.__doc__ or "").lower()
self.assertIn("read-only", doc)
self.assertIn("fails closed", doc)
class TestNamespaceResolution(unittest.TestCase):
"""Every configured profile must resolve to its real fleet namespace.
``role_namespace_gate.infer_mcp_namespace`` recognises only author and
reviewer and echoes the profile name for the rest, which would label the
controller, merger, and reconciler members with namespaces that do not
exist and then flag all three as role mismatches.
"""
def test_every_expected_profile_maps_to_its_namespace(self):
for entry in mfi.EXPECTED_PRGS_FLEET:
self.assertEqual(
mfi.namespace_for_profile(entry["profile"]),
entry["namespace"],
entry["profile"],
)
def test_server_helper_agrees_with_the_roster(self):
for entry in mfi.EXPECTED_PRGS_FLEET:
self.assertEqual(
gitea_mcp_server._fleet_namespace_for_profile(entry["profile"]),
entry["namespace"],
entry["profile"],
)
def test_unknown_profile_falls_back_to_the_supplied_default(self):
self.assertEqual(
mfi.namespace_for_profile("prgs-shadow", default="gitea-shadow"),
"gitea-shadow",
)
def test_unknown_profile_without_a_default_returns_the_profile(self):
self.assertEqual(mfi.namespace_for_profile("prgs-shadow"), "prgs-shadow")
def test_registered_row_uses_the_roster_namespace(self):
fake_db = mock.Mock()
with mock.patch.object(
gitea_mcp_server, "_control_plane_db_or_error", return_value=(fake_db, [])
), mock.patch.object(
gitea_mcp_server, "get_profile", return_value={"allowed_operations": []}
), mock.patch.object(
gitea_mcp_server.gitea_config,
"selected_profile_name",
return_value="prgs-reconciler",
):
record = gitea_mcp_server._register_fleet_runtime(transport="stdio")
self.assertEqual(record["namespace"], "gitea-reconciler")
class TestToolBehavior(unittest.TestCase):
"""The tool wires the registry and the process scan into the classifier."""
@staticmethod
def _healthy_rows():
return [
{
"runtime_id": f"{entry['namespace']}:{43000 + index}:boot",
"namespace": entry["namespace"],
"profile": entry["profile"],
"role": entry["role"],
"remote": "prgs",
"org": "Scaled-Tech-Consulting",
"repo": "Gitea-Tools",
"repository_root": "/checkout/Gitea-Tools",
"pid": 43000 + index,
"cohort_id": "ppid:40990",
"cohort_source": "parent_process",
"client_provenance": "client_managed",
"boot_id": "boot",
"startup_head": "82d71b77028a7abd4f8ab4a4e4d89658a187f73d",
"daemon_start_head": "82d71b77028a7abd4f8ab4a4e4d89658a187f73d",
"transport": "stdio",
"registered_at": "2026-07-28T02:00:00Z",
"last_heartbeat_at": "2026-07-28T02:00:00Z",
"status": "running",
}
for index, entry in enumerate(mfi.EXPECTED_PRGS_FLEET)
]
def _call(self, rows, *, scan=None, db=None, db_errors=None):
fake_db = db
if fake_db is None and db_errors is None:
fake_db = mock.Mock()
fake_db.list_mcp_server_runtimes.return_value = rows
scan = scan or {
"available": True,
"processes": [
{"pid": r["pid"], "started_at": "2026-07-28T01:00:00Z"} for r in rows
],
"reason": None,
}
with mock.patch.object(
gitea_mcp_server, "_profile_operation_gate", return_value=[]
), mock.patch.object(
gitea_mcp_server,
"_control_plane_db_or_error",
return_value=(fake_db, db_errors or []),
), mock.patch.object(
gitea_mcp_server, "_active_profile_name", return_value="prgs-controller"
), mock.patch.object(
mfi, "scan_mcp_server_processes", return_value=scan
), mock.patch.object(
mfi, "probe_pid_alive", return_value=True
):
return gitea_mcp_server.gitea_assess_fleet_inventory(
remote="prgs", org="Scaled-Tech-Consulting", repo="Gitea-Tools"
)
def test_healthy_fleet_satisfies_the_gate_through_the_tool(self):
result = self._call(self._healthy_rows())
self.assertTrue(result["inventory_complete"])
self.assertTrue(result["mutation_gate_satisfied"])
self.assertEqual(result["running_member_count"], 5)
self.assertEqual(result["mutations_performed"], [])
def test_tool_reports_the_expected_repository_binding(self):
result = self._call(self._healthy_rows())
self.assertEqual(
result["expected_repository_binding"],
{"remote": "prgs", "org": "Scaled-Tech-Consulting", "repo": "Gitea-Tools"},
)
def test_tool_reads_only_running_rows(self):
rows = self._healthy_rows()
fake_db = mock.Mock()
fake_db.list_mcp_server_runtimes.return_value = rows
self._call(rows, db=fake_db)
fake_db.list_mcp_server_runtimes.assert_called_once_with(statuses=("running",))
def test_tool_never_writes_to_the_registry(self):
rows = self._healthy_rows()
fake_db = mock.Mock()
fake_db.list_mcp_server_runtimes.return_value = rows
self._call(rows, db=fake_db)
fake_db.register_mcp_server_runtime.assert_not_called()
fake_db.heartbeat_mcp_server_runtime.assert_not_called()
fake_db.mark_mcp_server_runtime_stopped.assert_not_called()
def test_unavailable_control_plane_fails_closed(self):
result = self._call([], db_errors=["control-plane DB substrate unavailable"])
self.assertFalse(result["inventory_complete"])
self.assertFalse(result["mutation_gate_satisfied"])
self.assertFalse(result["registry"]["available"])
self.assertIn("unavailable", result["blocked_reason"])
def test_registry_read_failure_fails_closed(self):
rows = self._healthy_rows()
fake_db = mock.Mock()
fake_db.list_mcp_server_runtimes.side_effect = RuntimeError("db locked")
result = self._call(rows, db=fake_db)
self.assertFalse(result["inventory_complete"])
self.assertFalse(result["mutation_gate_satisfied"])
self.assertIn("could not be read", result["registry"]["error"])
def test_missing_read_permission_blocks_without_touching_the_registry(self):
fake_db = mock.Mock()
with mock.patch.object(
gitea_mcp_server,
"_profile_operation_gate",
return_value=["profile may not read"],
), mock.patch.object(
gitea_mcp_server, "_control_plane_db_or_error", return_value=(fake_db, [])
):
result = gitea_mcp_server.gitea_assess_fleet_inventory(remote="prgs")
self.assertFalse(result["success"])
self.assertFalse(result["mutation_gate_satisfied"])
self.assertIn("permission_report", result)
self.assertEqual(result["mutations_performed"], [])
fake_db.list_mcp_server_runtimes.assert_not_called()
def test_answering_namespace_is_reported(self):
result = self._call(self._healthy_rows())
self.assertEqual(result["answering_namespace"], "gitea-controller")
def test_summary_is_present(self):
result = self._call(self._healthy_rows())
self.assertIn("5 of 5", result["summary"])
class TestStartupRegistration(unittest.TestCase):
"""The registry row is written by the process it describes."""
def test_registration_helper_exists_on_the_entrypoint_module(self):
self.assertTrue(hasattr(gitea_mcp_server, "_register_fleet_runtime"))
def test_registration_writes_one_row_for_this_process(self):
fake_db = mock.Mock()
with mock.patch.object(
gitea_mcp_server, "_control_plane_db_or_error", return_value=(fake_db, [])
), mock.patch.object(
gitea_mcp_server, "get_profile", return_value={"allowed_operations": []}
):
record = gitea_mcp_server._register_fleet_runtime(transport="stdio")
self.assertIsNotNone(record)
fake_db.register_mcp_server_runtime.assert_called_once()
written = fake_db.register_mcp_server_runtime.call_args.args[0]
self.assertEqual(written["pid"], os.getpid())
self.assertEqual(written["transport"], "stdio")
def test_registration_failure_never_blocks_startup(self):
with mock.patch.object(
gitea_mcp_server,
"_control_plane_db_or_error",
side_effect=RuntimeError("boom"),
):
self.assertIsNone(gitea_mcp_server._register_fleet_runtime())
def test_registration_is_skipped_when_the_control_plane_is_unavailable(self):
with mock.patch.object(
gitea_mcp_server,
"_control_plane_db_or_error",
return_value=(None, ["unavailable"]),
):
self.assertIsNone(gitea_mcp_server._register_fleet_runtime())
def test_entrypoint_registers_after_binding_native_transport(self):
"""Order matters: only a transport-bound process may claim a row."""
source = (REPO_ROOT / "gitea_mcp_server.py").read_text()
bind_at = source.index('bind_native_mcp_transport(transport="stdio")')
register_at = source.index('_register_fleet_runtime(transport="stdio")')
run_at = source.index('mcp.run(transport="stdio")')
self.assertLess(bind_at, register_at)
self.assertLess(register_at, run_at)
class TestWorkflowDocumentation(unittest.TestCase):
"""AC13: documentation explains how the gates consume the result."""
def setUp(self):
self.doc = (REPO_ROOT / "docs" / "mcp-fleet-inventory.md").read_text()
self.skill = (
REPO_ROOT / "skills" / "llm-project-workflow" / "SKILL.md"
).read_text()
def test_dedicated_document_exists(self):
self.assertIn("# Authoritative MCP fleet inventory", self.doc)
def test_document_names_both_consuming_namespaces(self):
self.assertIn("gitea-controller", self.doc)
self.assertIn("gitea-reconciler", self.doc)
def test_document_explains_gate_consumption(self):
self.assertIn("mutation_gate_satisfied", self.doc)
self.assertIn("blocked_reason", self.doc)
def test_document_states_the_non_inferences(self):
self.assertIn("Configuration is not existence", self.doc)
self.assertIn("Matching revisions are not a cohort", self.doc)
def test_document_preserves_the_neighbouring_issue_boundaries(self):
for issue in ("#950", "#951", "#952"):
self.assertIn(issue, self.doc)
def test_canonical_workflow_skill_links_the_document(self):
self.assertIn("docs/mcp-fleet-inventory.md", self.skill)
self.assertIn(TOOL_NAME, self.skill)
if __name__ == "__main__":
unittest.main()
@@ -0,0 +1,323 @@
"""Tests for sanctioned Codex MCP reconnect request surface (#678)."""
from __future__ import annotations
import os
import unittest
from unittest import mock
import mcp_client_reconnect as mcr
class NormalizeReasonTests(unittest.TestCase):
def test_stale_runtime_aliases(self):
self.assertEqual(mcr.normalize_reason("stale-runtime"), mcr.REASON_STALE_RUNTIME)
self.assertEqual(mcr.normalize_reason("stale_runtime"), mcr.REASON_STALE_RUNTIME)
self.assertEqual(mcr.normalize_reason("STALE"), mcr.REASON_STALE_RUNTIME)
def test_transport_eof_aliases(self):
self.assertEqual(mcr.normalize_reason("transport_eof"), mcr.REASON_TRANSPORT_EOF)
self.assertEqual(mcr.normalize_reason("EOF"), mcr.REASON_TRANSPORT_EOF)
self.assertEqual(
mcr.normalize_reason("client_is_closing"), mcr.REASON_TRANSPORT_EOF
)
def test_missing_namespace(self):
self.assertEqual(
mcr.normalize_reason("missing_namespace"), mcr.REASON_MISSING_NAMESPACE
)
def test_empty_is_unspecified(self):
self.assertEqual(mcr.normalize_reason(None), mcr.REASON_UNSPECIFIED)
self.assertEqual(mcr.normalize_reason(""), mcr.REASON_UNSPECIFIED)
class BoundaryClassificationTests(unittest.TestCase):
def test_clean_when_shas_match(self):
self.assertEqual(
mcr.classify_boundary_status(
startup_sha="abc", current_master_sha="abc"
),
mcr.BOUNDARY_CLEAN,
)
def test_mismatch_when_shas_differ(self):
self.assertEqual(
mcr.classify_boundary_status(
startup_sha="aaa", current_master_sha="bbb"
),
mcr.BOUNDARY_MISMATCH,
)
def test_stale_when_live_stale(self):
self.assertEqual(
mcr.classify_boundary_status(
startup_sha="aaa",
current_master_sha="aaa",
live_stale=True,
),
mcr.BOUNDARY_STALE,
)
class BuildReconnectRequestTests(unittest.TestCase):
def test_stale_runtime_returns_typed_blocker_with_codex_steps(self):
result = mcr.build_reconnect_request(
namespace="gitea-author",
profile="prgs-author",
pid=1234,
session_id="sess-1",
startup_sha="aaa111",
current_master_sha="bbb222",
reason="stale-runtime",
client="codex",
restart_required=True,
stop_required=True,
)
self.assertTrue(result["success"])
self.assertTrue(result["read_only"])
self.assertFalse(result["reconnect_performed"])
self.assertFalse(result["mutation_performed"])
self.assertTrue(result["reconnect_needed"])
self.assertEqual(result["namespace"], "gitea-author")
self.assertEqual(result["profile"], "prgs-author")
self.assertEqual(result["pid"], 1234)
self.assertEqual(result["session_id"], "sess-1")
self.assertEqual(result["startup_sha"], "aaa111")
self.assertEqual(result["current_master_sha"], "bbb222")
self.assertEqual(result["boundary_status"], mcr.BOUNDARY_MISMATCH)
self.assertEqual(result["blocker_kind"], mcr.BLOCKER_OPERATOR_RECONNECT)
self.assertIsNotNone(result["typed_blocker"])
blocker = result["typed_blocker"]
self.assertEqual(blocker["namespaces"], ["gitea-author"])
self.assertEqual(blocker["why_reconnect_required"], mcr.REASON_STALE_RUNTIME)
self.assertTrue(any("Codex" in s or "Reload" in s for s in blocker["operator_ui_steps"]))
self.assertIn("pkill", " ".join(result["forbidden_recovery_paths"]).lower())
self.assertTrue(
mcr.reasons_never_suggest_forbidden(result["exact_safe_next_action"] or "")
)
# Must not recommend forbidden recovery.
for step in blocker["operator_ui_steps"]:
self.assertTrue(mcr.reasons_never_suggest_forbidden(step), step)
def test_transport_eof_typed_blocker(self):
result = mcr.build_reconnect_request(
namespace="gitea-reviewer",
reason="transport_eof",
client="claude_code",
)
self.assertTrue(result["reconnect_needed"])
self.assertEqual(result["reason"], mcr.REASON_TRANSPORT_EOF)
self.assertEqual(result["client"], "claude_code")
steps = " ".join(result["operator_ui_steps"]).lower()
self.assertIn("/mcp", steps)
def test_missing_namespace_typed_blocker(self):
result = mcr.build_reconnect_request(
namespace="gitea-merger",
reason="missing_namespace",
client="codex",
)
self.assertTrue(result["reconnect_needed"])
self.assertEqual(result["reason"], mcr.REASON_MISSING_NAMESPACE)
self.assertEqual(
result["typed_blocker"]["blocker_kind"], mcr.BLOCKER_OPERATOR_RECONNECT
)
def test_healthy_not_required(self):
result = mcr.build_reconnect_request(
namespace="gitea-tools",
startup_sha="deadbeef",
current_master_sha="deadbeef",
reason="not_required",
client="codex",
in_parity=True,
restart_required=False,
stop_required=False,
)
self.assertFalse(result["reconnect_needed"])
self.assertEqual(result["blocker_kind"], mcr.BLOCKER_NONE)
self.assertIsNone(result["typed_blocker"])
self.assertFalse(result["stop_required"])
self.assertFalse(result["restart_required"])
self.assertIn("not required", (result["exact_safe_next_action"] or "").lower())
def test_successful_reconnect_report_fields_present(self):
"""AC2: reconnect result reports required fields (even when needed)."""
result = mcr.build_reconnect_request(
namespace="gitea-controller",
profile="prgs-controller",
pid=99,
session_id="sid",
startup_sha="s" * 40,
current_master_sha="c" * 40,
reason="stale-runtime",
)
for key in (
"namespace",
"profile",
"pid",
"session_id",
"startup_sha",
"current_master_sha",
"boundary_status",
):
self.assertIn(key, result)
self.assertIsNotNone(result[key], key)
class ToolSurfaceTests(unittest.TestCase):
"""Exercise gitea_request_mcp_reconnect with a stubbed server context."""
def test_tool_is_registered_and_side_effect_free(self):
import gitea_mcp_server as srv
self.assertTrue(hasattr(srv, "gitea_request_mcp_reconnect"))
with mock.patch.object(srv, "_profile_operation_gate", return_value=None):
with mock.patch.object(
srv,
"get_profile",
return_value={
"profile_name": "prgs-author",
"role_kind": "author",
"role": "author",
},
):
with mock.patch.object(
srv,
"_current_master_parity",
return_value={
"startup_head": "a" * 40,
"current_head": "a" * 40,
"daemon_start_head": "a" * 40,
"local_head": "a" * 40,
"in_parity": True,
"stale": False,
"restart_required": False,
"determinable": True,
"live_stale": False,
"live_known": True,
"reasons": [],
},
):
with mock.patch.object(
srv.master_parity_gate,
"format_parity",
return_value="in parity",
):
with mock.patch.object(
srv.role_namespace_gate,
"infer_mcp_namespace",
return_value="gitea-author",
):
with mock.patch.object(
srv.session_ctx,
"mutation_context_audit_fields",
return_value={"session_profile": "prgs-author"},
):
result = srv.gitea_request_mcp_reconnect(
namespace="gitea-author",
reason="not_required",
client="codex",
remote="prgs",
)
self.assertTrue(result.get("success"))
self.assertFalse(result.get("reconnect_performed"))
self.assertFalse(result.get("mutation_performed"))
self.assertEqual(result.get("namespace"), "gitea-author")
self.assertEqual(result.get("profile"), "prgs-author")
self.assertEqual(result.get("pid"), os.getpid())
self.assertIn("startup_sha", result)
self.assertIn("current_master_sha", result)
self.assertIn("boundary_status", result)
self.assertTrue(
mcr.reasons_never_suggest_forbidden(
result.get("exact_safe_next_action") or ""
)
)
def test_tool_stale_returns_typed_blocker(self):
import gitea_mcp_server as srv
with mock.patch.object(srv, "_profile_operation_gate", return_value=None):
with mock.patch.object(
srv,
"get_profile",
return_value={
"profile_name": "prgs-reconciler",
"role_kind": "reconciler",
"role": "reconciler",
},
):
with mock.patch.object(
srv,
"_current_master_parity",
return_value={
"startup_head": "a" * 40,
"current_head": "b" * 40,
"daemon_start_head": "a" * 40,
"local_head": "b" * 40,
"in_parity": False,
"stale": True,
"restart_required": True,
"determinable": True,
"live_stale": True,
"live_known": True,
"reasons": ["stale"],
},
):
with mock.patch.object(
srv.master_parity_gate,
"format_parity",
return_value="stale",
):
with mock.patch.object(
srv.role_namespace_gate,
"infer_mcp_namespace",
return_value="gitea-reconciler",
):
with mock.patch.object(
srv.session_ctx,
"mutation_context_audit_fields",
return_value={},
):
result = srv.gitea_request_mcp_reconnect(
reason="stale-runtime",
client="codex",
)
self.assertTrue(result["reconnect_needed"])
self.assertEqual(
result["blocker_kind"], mcr.BLOCKER_OPERATOR_RECONNECT
)
self.assertIsNotNone(result["typed_blocker"])
self.assertIn("gitea-reconciler", result["typed_blocker"]["namespaces"])
self.assertTrue(result["stop_required"])
self.assertTrue(result["restart_required"])
self.assertTrue(
mcr.reasons_never_suggest_forbidden(
result.get("exact_safe_next_action") or ""
)
)
class InventoryRegistrationTests(unittest.TestCase):
def test_reconnect_path_in_restart_inventory(self):
import mcp_restart_paths as mrp
ids = {p.path_id for p in mrp.iter_restart_paths()}
self.assertIn("codex_client_reconnect_request", ids)
self.assertIn("ide_client_reconnect", ids)
def test_tool_name_in_documented_inventory(self):
import mcp_tool_inventory as inv
doc_path = os.path.join(
os.path.dirname(os.path.dirname(__file__)), inv.INVENTORY_DOC_PATH
)
with open(doc_path, encoding="utf-8") as handle:
documented = inv.parse_documented_inventory(handle.read())
self.assertIn("gitea_request_mcp_reconnect", documented)
if __name__ == "__main__":
unittest.main()
@@ -0,0 +1,747 @@
"""Regression: author bootstrap runtime authority and session ownership (#943).
Two rounds of defects live here.
**Round 1 (#943 as filed).** ``gitea_bootstrap_author_issue_worktree`` passed
four values down to the bootstrap service that were never defined:
``_active_username``, ``_active_profile_name``, ``_current_session_id`` and
``_author_mutation_block``. Every call dry-run included raised
``NameError`` while evaluating the arguments, before the service was entered.
**Round 2 (review 622 on PR #944).** The first fix defined all four but made
``_current_session_id`` mint ``<profile>-<pid>-<hex>`` once per process. The MCP
daemon outlives every task it serves, so that value conflates sequential author
tasks and can never equal the control-plane session that owns an
allocator-created lease: ``_verify_assignment_and_lease_ids`` refused the whole
allocated path with ``lease_session_mismatch``. The reviewed round also read the
identity from the pinned session context while reading the profile from the live
profile, so a rebind could produce a mixed claimant pair, and it swallowed every
``get_profile()`` exception.
These tests therefore drive real state, not mocks of internals: a temporary
control-plane SQLite database and a temporary issue-lock directory, both
redirected through the same environment variables production uses
(``GITEA_CONTROL_PLANE_DB``, ``GITEA_ISSUE_LOCK_DIR``). The ownership gate that
runs is the real one.
``test_every_global_referenced_by_the_wrapper_resolves`` remains: it is what
found ``_author_mutation_block``, and it generalises to the next missing
reference. It supplements the runtime coverage below rather than standing in for
it.
"""
from __future__ import annotations
import ast
import builtins
import os
import re
import subprocess
import tempfile
import unittest
from unittest import mock
import author_issue_bootstrap as aib
import control_plane_db
import create_issue_bootstrap as cib
import gitea_mcp_server as gms
import issue_lock_store
import workflow_scope_guard
BOOTSTRAP_TASK = "bootstrap_author_issue_worktree"
WRAPPER_NAME = "gitea_bootstrap_author_issue_worktree"
RUNTIME_HELPERS = (
"_active_mutation_authority",
"_active_username",
"_active_profile_name",
"_resolve_owner_workflow_session",
"_author_mutation_block",
)
ORG = "Scaled-Tech-Consulting"
REPO = "Gitea-Tools"
IDENTITY = "jcwalker3"
PROFILE = "prgs-author"
# A per-task ownership key must carry no process identifier (#790).
TASK_KEY_RE = re.compile(r"^author_issue_work-[0-9a-f]{16}$")
def _make_control_repo(tmp: str) -> tuple[str, str]:
"""Create a clean control checkout on master and return (path, head)."""
repo = os.path.join(tmp, "repo")
os.makedirs(os.path.join(repo, "branches"))
subprocess.check_call(
["git", "init", "-b", "master", repo],
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
)
subprocess.check_call(
[
"git", "-C", repo,
"-c", "user.email=t@t", "-c", "user.name=t",
"commit", "--allow-empty", "-m", "init",
],
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
)
head = subprocess.check_output(
["git", "-C", repo, "rev-parse", "HEAD"], text=True
).strip()
return repo, head
def _wrapper_ast() -> ast.FunctionDef:
"""Return the AST of the bootstrap wrapper as it exists on disk."""
path = os.path.join(
os.path.dirname(os.path.dirname(os.path.abspath(__file__))),
"gitea_mcp_server.py",
)
with open(path, encoding="utf-8") as fh:
tree = ast.parse(fh.read())
for node in ast.walk(tree):
if isinstance(node, ast.FunctionDef) and node.name == WRAPPER_NAME:
return node
raise AssertionError(f"{WRAPPER_NAME} not found in gitea_mcp_server.py")
class _IsolatedControlPlane(unittest.TestCase):
"""Temp control-plane DB and temp issue-lock dir, via production env vars."""
def setUp(self):
self._tmp = tempfile.TemporaryDirectory()
self.addCleanup(self._tmp.cleanup)
self.tmp = self._tmp.name
self.db_path = os.path.join(self.tmp, "control-plane.sqlite3")
self.lock_dir = os.path.join(self.tmp, "issue-locks")
self.journals = os.path.join(self.tmp, "journals")
os.makedirs(self.lock_dir)
os.makedirs(self.journals)
env = mock.patch.dict(
os.environ,
{
control_plane_db.DB_PATH_ENV: self.db_path,
issue_lock_store.LOCK_DIR_ENV: self.lock_dir,
},
)
env.start()
self.addCleanup(env.stop)
self.db = control_plane_db.ControlPlaneDB(self.db_path)
def _allocate(self, session_id: str, *, issue: int = 943):
"""Create a real assignment + lease owned by *session_id*."""
self.db.upsert_session(
session_id=session_id, role="author", profile=PROFILE, pid=os.getpid()
)
res = self.db.assign_and_lease(
session_id=session_id, role="author", remote="prgs",
org=ORG, repo=REPO, kind="issue", number=issue,
)
self.assertEqual(res.outcome, "assigned", res)
return res.assignment_id, res.lease_id
def _authority(self):
"""A resolved authority pair, as the wrapper would compute it."""
return {"ok": True, "identity": IDENTITY, "profile_name": PROFILE}
def _resolve_session(self, **over):
kwargs = dict(
issue_number=943,
assignment_id=None,
lease_id=None,
session_id=None,
identity=IDENTITY,
profile_name=PROFILE,
remote="prgs",
org=ORG,
repo=REPO,
)
kwargs.update(over)
return gms._resolve_owner_workflow_session(**kwargs)
class OwnershipGateTests(_IsolatedControlPlane):
"""B2: the allocator-driven ownership path, against a real control plane."""
def setUp(self):
super().setUp()
self.repo, self.head = _make_control_repo(self.tmp)
def _bootstrap(self, **over):
kwargs = dict(
issue_number=943,
canonical_repo_root=self.repo,
expected_base_sha=self.head,
branch_name="fix/issue-943-runtime-context-helpers",
remote="prgs",
org=ORG,
repo=REPO,
active_identity=IDENTITY,
active_profile=PROFILE,
lock_dir=self.journals,
idempotency_key="test-943",
dry_run=True,
)
kwargs.update(over)
return aib.bootstrap_author_issue_worktree(**kwargs)
def test_true_owning_session_passes_the_ownership_gate(self):
"""The canonical owner reaches and completes the service."""
session = "prgs-author-task-a"
assignment_id, lease_id = self._allocate(session)
res = self._bootstrap(
assignment_id=assignment_id, lease_id=lease_id, owner_session=session
)
self.assertTrue(res.get("success"), res)
self.assertTrue(res.get("dry_run"))
self.assertEqual(res.get("base_sha"), self.head)
def test_different_session_is_refused(self):
session = "prgs-author-task-a"
assignment_id, lease_id = self._allocate(session)
res = self._bootstrap(
assignment_id=assignment_id,
lease_id=lease_id,
owner_session="prgs-author-task-b",
)
self.assertFalse(res.get("success"))
self.assertEqual(res.get("reason_code"), "lease_session_mismatch")
def test_process_derived_session_would_be_refused(self):
"""The reviewed round-1 value shape can never own an allocated lease."""
session = "prgs-author-task-a"
assignment_id, lease_id = self._allocate(session)
round_one_value = f"{PROFILE}-{os.getpid()}-deadbeef"
self.assertNotEqual(round_one_value, session)
res = self._bootstrap(
assignment_id=assignment_id,
lease_id=lease_id,
owner_session=round_one_value,
)
self.assertFalse(res.get("success"))
self.assertEqual(res.get("reason_code"), "lease_session_mismatch")
def test_unknown_lease_fails_closed(self):
session = "prgs-author-task-a"
assignment_id, _ = self._allocate(session)
res = self._bootstrap(
assignment_id=assignment_id,
lease_id="lease-does-not-exist",
owner_session=session,
)
self.assertFalse(res.get("success"))
self.assertEqual(res.get("reason_code"), "unknown_lease_id")
def test_released_lease_fails_closed(self):
session = "prgs-author-task-a"
assignment_id, lease_id = self._allocate(session)
self.db.release_lease(lease_id, session_id=session)
res = self._bootstrap(
assignment_id=assignment_id, lease_id=lease_id, owner_session=session
)
self.assertFalse(res.get("success"))
self.assertEqual(res.get("reason_code"), "lease_not_live")
def test_force_expired_lease_fails_closed(self):
session = "prgs-author-task-a"
assignment_id, lease_id = self._allocate(session)
self.db.force_expire_lease(lease_id, reason="test")
res = self._bootstrap(
assignment_id=assignment_id, lease_id=lease_id, owner_session=session
)
self.assertFalse(res.get("success"))
self.assertEqual(res.get("reason_code"), "lease_not_live")
def test_replacement_lease_does_not_inherit_prior_ownership(self):
"""A second task's lease is not ownable by the first task's session."""
first = "prgs-author-task-a"
assignment_a, lease_a = self._allocate(first)
self.db.release_lease(lease_a, session_id=first)
second = "prgs-author-task-b"
assignment_b, lease_b = self._allocate(second)
self.assertNotEqual(lease_a, lease_b)
res = self._bootstrap(
assignment_id=assignment_b, lease_id=lease_b, owner_session=first
)
self.assertFalse(res.get("success"))
self.assertEqual(res.get("reason_code"), "lease_session_mismatch")
def test_assignment_lease_identifier_mismatch_fails_closed(self):
session = "prgs-author-task-a"
_, lease_id = self._allocate(session)
res = self._bootstrap(
assignment_id="asn-not-the-recorded-one",
lease_id=lease_id,
owner_session=session,
)
self.assertFalse(res.get("success"))
self.assertEqual(res.get("reason_code"), "assignment_lease_mismatch")
def test_lease_id_without_assignment_id_fails_closed(self):
session = "prgs-author-task-a"
_, lease_id = self._allocate(session)
res = self._bootstrap(lease_id=lease_id, owner_session=session)
self.assertFalse(res.get("success"))
self.assertEqual(res.get("reason_code"), "incomplete_assignment_lease_ids")
def test_dry_run_with_valid_allocator_bindings_leaves_no_durable_state(self):
session = "prgs-author-task-a"
assignment_id, lease_id = self._allocate(session)
res = self._bootstrap(
assignment_id=assignment_id, lease_id=lease_id, owner_session=session
)
self.assertTrue(res.get("success"), res)
branches = subprocess.check_output(
["git", "-C", self.repo, "branch", "--list"], text=True
)
self.assertNotIn("issue-943", branches)
worktrees = subprocess.check_output(
["git", "-C", self.repo, "worktree", "list"], text=True
)
self.assertNotIn("issue-943", worktrees)
self.assertFalse(
os.path.exists(
os.path.join(self.repo, "branches",
"fix-issue-943-runtime-context-helpers")
)
)
journal = res.get("phase_journal") or {}
self.assertFalse(journal.get("completed"))
self.assertFalse(any((journal.get("artifacts_created") or {}).values()))
# The dry run must not have created an issue lock in the isolated dir.
self.assertEqual(os.listdir(self.lock_dir), [])
def test_apply_reaches_the_intended_transition_with_valid_bindings(self):
session = "prgs-author-task-a"
assignment_id, lease_id = self._allocate(session)
res = self._bootstrap(
assignment_id=assignment_id,
lease_id=lease_id,
owner_session=session,
dry_run=False,
)
self.assertTrue(res.get("success"), res)
self.assertNotEqual(res.get("dry_run"), True)
branches = subprocess.check_output(
["git", "-C", self.repo, "branch", "--list"], text=True
)
self.assertIn("issue-943", branches)
self.assertTrue(os.path.isdir(res.get("worktree_path") or ""))
class WorkflowSessionResolutionTests(_IsolatedControlPlane):
"""B1: the wrapper resolves the owning session, never a process identifier."""
def test_declared_session_is_verified_against_the_control_plane(self):
session = "prgs-author-task-a"
assignment_id, lease_id = self._allocate(session)
res = self._resolve_session(
session_id=session, assignment_id=assignment_id, lease_id=lease_id
)
self.assertTrue(res.get("ok"), res)
self.assertEqual(res.get("session_id"), session)
self.assertEqual(res.get("session_source"), "declared")
def test_unknown_declared_session_is_refused_not_trusted(self):
res = self._resolve_session(session_id="prgs-author-not-a-session")
self.assertFalse(res.get("ok"))
self.assertEqual(res.get("reason_code"), "workflow_session_unverified")
def test_declared_session_for_another_role_is_refused(self):
self.db.upsert_session(
session_id="prgs-reviewer-x", role="reviewer", profile="prgs-reviewer",
pid=os.getpid(),
)
res = self._resolve_session(session_id="prgs-reviewer-x")
self.assertFalse(res.get("ok"))
self.assertEqual(res.get("reason_code"), "workflow_session_unverified")
def test_declared_session_for_another_profile_is_refused(self):
self.db.upsert_session(
session_id="other-profile-session", role="author",
profile="prgs-controller", pid=os.getpid(),
)
res = self._resolve_session(session_id="other-profile-session")
self.assertFalse(res.get("ok"))
self.assertEqual(res.get("reason_code"), "workflow_session_unverified")
def test_allocated_work_without_a_session_is_refused(self):
"""Supplying a lease is not itself evidence of ownership."""
session = "prgs-author-task-a"
assignment_id, lease_id = self._allocate(session)
res = self._resolve_session(assignment_id=assignment_id, lease_id=lease_id)
self.assertFalse(res.get("ok"))
self.assertEqual(
res.get("reason_code"), "workflow_session_required_for_allocated_work"
)
def test_existing_issue_lock_supplies_its_per_task_session(self):
lock_session = issue_lock_store.mint_task_session_id(
issue_lock_store.AUTHOR_ISSUE_WORK_LEASE
)
path = issue_lock_store.lock_file_path(
remote="prgs", org=ORG, repo=REPO, issue_number=943,
lock_dir=self.lock_dir,
)
issue_lock_store.write_lock_file(
path,
{
"issue_number": 943,
"branch_name": "fix/issue-943-runtime-context-helpers",
"work_lease": {
"task_session_id": lock_session,
"claimant": {"username": IDENTITY, "profile": PROFILE},
},
},
) if hasattr(issue_lock_store, "write_lock_file") else _write_json(
path,
{
"issue_number": 943,
"branch_name": "fix/issue-943-runtime-context-helpers",
"work_lease": {
"task_session_id": lock_session,
"claimant": {"username": IDENTITY, "profile": PROFILE},
},
},
)
res = self._resolve_session()
self.assertTrue(res.get("ok"), res)
self.assertEqual(res.get("session_id"), lock_session)
self.assertEqual(res.get("session_source"), "issue_lock")
def test_issue_lock_owned_by_another_identity_is_refused(self):
path = issue_lock_store.lock_file_path(
remote="prgs", org=ORG, repo=REPO, issue_number=943,
lock_dir=self.lock_dir,
)
_write_json(
path,
{
"issue_number": 943,
"work_lease": {
"task_session_id": "author_issue_work-" + "0" * 16,
"claimant": {"username": "someone-else", "profile": PROFILE},
},
},
)
res = self._resolve_session()
self.assertFalse(res.get("ok"))
self.assertEqual(res.get("reason_code"), "issue_lock_owner_mismatch")
def test_unallocated_bootstrap_mints_a_per_task_key(self):
res = self._resolve_session()
self.assertTrue(res.get("ok"), res)
self.assertEqual(res.get("session_source"), "minted_task_key")
self.assertRegex(res["session_id"], TASK_KEY_RE)
def test_minted_key_contains_no_process_identifier(self):
res = self._resolve_session()
self.assertNotIn(str(os.getpid()), res["session_id"])
self.assertNotIn(PROFILE, res["session_id"])
def test_sequential_tasks_on_one_daemon_do_not_share_ownership(self):
"""The round-1 defect: one identifier per process for every task."""
first = self._resolve_session()["session_id"]
second = self._resolve_session()["session_id"]
third = self._resolve_session()["session_id"]
self.assertNotEqual(first, second)
self.assertNotEqual(second, third)
self.assertEqual(len({first, second, third}), 3)
def test_no_process_lifetime_cache_remains(self):
self.assertFalse(hasattr(gms, "_ACTIVE_SESSION_ID"))
self.assertFalse(hasattr(gms, "_current_session_id"))
class MutationAuthorityTests(unittest.TestCase):
"""F3/F4: one coherent authority pair, drift detected, no silent fallback."""
def _ctx(self, **over):
base = {"identity": IDENTITY, "profile_name": PROFILE}
base.update(over)
return base
def test_matching_live_and_pinned_authority_resolves(self):
with mock.patch.object(gms, "get_profile",
return_value={"profile_name": PROFILE}), \
mock.patch.object(gms, "_authenticated_username",
return_value=IDENTITY), \
mock.patch.object(gms.session_ctx, "get_session_context",
return_value=self._ctx()):
res = gms._active_mutation_authority("gitea.prgs.cc")
self.assertTrue(res.get("ok"), res)
self.assertEqual(res["identity"], IDENTITY)
self.assertEqual(res["profile_name"], PROFILE)
def test_identity_drift_fails_closed(self):
with mock.patch.object(gms, "get_profile",
return_value={"profile_name": PROFILE}), \
mock.patch.object(gms, "_authenticated_username",
return_value="someone-else"), \
mock.patch.object(gms.session_ctx, "get_session_context",
return_value=self._ctx()):
res = gms._active_mutation_authority("gitea.prgs.cc")
self.assertFalse(res.get("ok"))
self.assertEqual(res.get("reason_code"), "authority_identity_drift")
self.assertEqual(res.get("expected"), IDENTITY)
self.assertEqual(res.get("actual"), "someone-else")
def test_profile_drift_fails_closed(self):
with mock.patch.object(gms, "get_profile",
return_value={"profile_name": "prgs-controller"}), \
mock.patch.object(gms, "_authenticated_username",
return_value=IDENTITY), \
mock.patch.object(gms.session_ctx, "get_session_context",
return_value=self._ctx()):
res = gms._active_mutation_authority("gitea.prgs.cc")
self.assertFalse(res.get("ok"))
self.assertEqual(res.get("reason_code"), "authority_profile_drift")
def test_identity_and_profile_never_come_from_different_snapshots(self):
"""Round 2's mixed pair: pinned identity plus live profile."""
with mock.patch.object(gms, "get_profile",
return_value={"profile_name": "prgs-controller"}), \
mock.patch.object(gms, "_authenticated_username",
return_value="new-identity"), \
mock.patch.object(gms.session_ctx, "get_session_context",
return_value=self._ctx()):
res = gms._active_mutation_authority("gitea.prgs.cc")
self.assertFalse(res.get("ok"))
self.assertIsNone(gms._active_username("gitea.prgs.cc"))
self.assertIsNone(gms._active_profile_name("gitea.prgs.cc"))
def test_unresolvable_profile_is_a_structured_refusal_not_a_fallback(self):
"""F4: no bare-except fallback to a previously pinned profile name."""
with mock.patch.object(gms, "get_profile",
side_effect=RuntimeError("profile disabled")), \
mock.patch.object(gms, "_authenticated_username",
return_value=IDENTITY), \
mock.patch.object(gms.session_ctx, "get_session_context",
return_value=self._ctx()):
res = gms._active_mutation_authority("gitea.prgs.cc")
self.assertFalse(res.get("ok"))
self.assertEqual(res.get("reason_code"), "authority_profile_unresolved")
self.assertNotEqual(res.get("profile_name"), PROFILE)
def test_malformed_profile_without_name_fails_closed(self):
with mock.patch.object(gms, "get_profile", return_value={}), \
mock.patch.object(gms, "_authenticated_username",
return_value=IDENTITY), \
mock.patch.object(gms.session_ctx, "get_session_context",
return_value=None):
res = gms._active_mutation_authority("gitea.prgs.cc")
self.assertFalse(res.get("ok"))
self.assertEqual(res.get("reason_code"), "authority_profile_unresolved")
def test_unresolved_identity_fails_closed(self):
for value in (None, "", " "):
with self.subTest(identity=value):
with mock.patch.object(gms, "get_profile",
return_value={"profile_name": PROFILE}), \
mock.patch.object(gms, "_authenticated_username",
return_value=value), \
mock.patch.object(gms.session_ctx, "get_session_context",
return_value=None):
res = gms._active_mutation_authority("gitea.prgs.cc")
self.assertFalse(res.get("ok"))
self.assertEqual(
res.get("reason_code"), "authority_identity_unresolved"
)
def test_missing_host_cannot_yield_an_identity(self):
with mock.patch.object(gms, "get_profile",
return_value={"profile_name": PROFILE}), \
mock.patch.object(gms.session_ctx, "get_session_context",
return_value=None):
res = gms._active_mutation_authority(None)
self.assertFalse(res.get("ok"))
self.assertEqual(res.get("reason_code"), "authority_identity_unresolved")
def test_expected_username_is_never_substituted_for_authentication(self):
with mock.patch.object(
gms, "get_profile",
return_value={"profile_name": PROFILE, "username": IDENTITY},
), mock.patch.object(gms, "_authenticated_username", return_value=None), \
mock.patch.object(gms.session_ctx, "get_session_context",
return_value={"expected_username": IDENTITY}):
self.assertIsNone(gms._active_username("gitea.prgs.cc"))
def test_accessors_share_one_snapshot(self):
with mock.patch.object(gms, "get_profile",
return_value={"profile_name": PROFILE}), \
mock.patch.object(gms, "_authenticated_username",
return_value=IDENTITY), \
mock.patch.object(gms.session_ctx, "get_session_context",
return_value=self._ctx()):
self.assertEqual(gms._active_username("gitea.prgs.cc"), IDENTITY)
self.assertEqual(gms._active_profile_name("gitea.prgs.cc"), PROFILE)
class AuthorMutationBlockTests(unittest.TestCase):
"""Preserved: the structured refusal shape review 622 confirmed correct."""
def test_matches_the_sibling_refusal_shape(self):
res = gms._author_mutation_block(["stopped"])
self.assertIs(res["success"], False)
self.assertIs(res["performed"], False)
self.assertEqual(res["outcome"], "REFUSED")
self.assertEqual(res["reasons"], ["stopped"])
def test_carries_reason_code_and_transport_fields(self):
res = gms._author_mutation_block(
["nope"], reason_code="authority_identity_drift",
retryable=False, transport_survives=True,
expected="a", actual="b", issue_number=943,
)
self.assertEqual(res["reason_code"], "authority_identity_drift")
self.assertIs(res["retryable"], False)
self.assertIs(res["transport_survives"], True)
self.assertEqual((res["expected"], res["actual"]), ("a", "b"))
self.assertEqual(res["issue_number"], 943)
self.assertIs(res["success"], False)
class RuntimeHelperResolutionTests(unittest.TestCase):
"""Every runtime helper the wrapper references is defined and callable.
Supplements the runtime coverage above; it does not replace it.
"""
def test_named_helpers_are_defined_and_callable(self):
for name in RUNTIME_HELPERS:
with self.subTest(helper=name):
self.assertTrue(hasattr(gms, name), f"{name} is not defined")
self.assertTrue(callable(getattr(gms, name)))
def test_every_global_referenced_by_the_wrapper_resolves(self):
"""The generalised form of the round-1 defect: an unresolvable global."""
fn = _wrapper_ast()
bound: set[str] = {a.arg for a in fn.args.args}
bound |= {a.arg for a in fn.args.kwonlyargs}
if fn.args.vararg:
bound.add(fn.args.vararg.arg)
if fn.args.kwarg:
bound.add(fn.args.kwarg.arg)
for node in ast.walk(fn):
if isinstance(node, ast.Name) and isinstance(
node.ctx, (ast.Store, ast.Del)
):
bound.add(node.id)
elif isinstance(node, (ast.Import, ast.ImportFrom)):
for alias in node.names:
bound.add((alias.asname or alias.name).split(".")[0])
elif isinstance(node, ast.ExceptHandler) and node.name:
bound.add(node.name)
unresolved = sorted(
node.id
for node in ast.walk(fn)
if isinstance(node, ast.Name)
and isinstance(node.ctx, ast.Load)
and node.id not in bound
and not hasattr(gms, node.id)
and not hasattr(builtins, node.id)
)
self.assertEqual(
unresolved, [],
f"{WRAPPER_NAME} references undefined globals: {unresolved}",
)
def test_wrapper_wires_the_authority_and_session_resolvers(self):
fn = _wrapper_ast()
called = {
node.func.id
for node in ast.walk(fn)
if isinstance(node, ast.Call) and isinstance(node.func, ast.Name)
}
self.assertIn("_active_mutation_authority", called)
self.assertIn("_resolve_owner_workflow_session", called)
self.assertIn("_author_mutation_block", called)
def test_wrapper_accepts_an_optional_session_id(self):
"""ABI addition stays backward compatible: optional, defaulting to None."""
fn = _wrapper_ast()
names = [a.arg for a in fn.args.args]
self.assertIn("session_id", names)
offset = len(names) - len(fn.args.defaults)
default = fn.args.defaults[names.index("session_id") - offset]
self.assertIsInstance(default, ast.Constant)
self.assertIsNone(default.value)
class Issue941ScopeGuardNotRegressedTests(unittest.TestCase):
"""Preserved: PR #942's bootstrap-scope wiring still holds."""
def setUp(self):
self._tmp = tempfile.TemporaryDirectory()
self.addCleanup(self._tmp.cleanup)
self.repo, self.head = _make_control_repo(self._tmp.name)
def _assessment(self, task: str = BOOTSTRAP_TASK) -> dict:
return aib.assess_author_issue_bootstrap(
workspace_path=self.repo,
canonical_repo_root=self.repo,
current_branch="master",
head_sha=self.head,
porcelain_status="",
remote_master_sha=self.head,
task=task,
)
def test_bootstrap_task_still_permitted_from_clean_control_checkout(self):
res = workflow_scope_guard.assess_root_source_mutation(
workspace_path=self.repo,
canonical_repo_root=self.repo,
role_kind="author",
mutation_task=BOOTSTRAP_TASK,
porcelain_status="",
bootstrap_assessment=self._assessment(),
)
self.assertFalse(res.get("block"), res)
self.assertNotEqual(
res.get("blocker_kind"), workflow_scope_guard.BLOCKER_MISSING_WORKTREE
)
def test_bootstrap_task_still_blocked_without_evidence(self):
res = workflow_scope_guard.assess_root_source_mutation(
workspace_path=self.repo,
canonical_repo_root=self.repo,
role_kind="author",
mutation_task=BOOTSTRAP_TASK,
porcelain_status="",
)
self.assertTrue(res.get("block"))
self.assertEqual(
res.get("blocker_kind"), workflow_scope_guard.BLOCKER_MISSING_WORKTREE
)
def test_ordinary_author_mutation_still_blocked_from_control_checkout(self):
res = workflow_scope_guard.assess_root_source_mutation(
workspace_path=self.repo,
canonical_repo_root=self.repo,
role_kind="author",
mutation_task="commit_files",
porcelain_status="",
bootstrap_assessment=self._assessment(),
)
self.assertTrue(res.get("block"))
self.assertEqual(
res.get("blocker_kind"), workflow_scope_guard.BLOCKER_MISSING_WORKTREE
)
def test_create_issue_bootstrap_unchanged(self):
self.assertTrue(cib.is_create_issue_task("create_issue"))
self.assertFalse(cib.is_create_issue_task(BOOTSTRAP_TASK))
def _write_json(path: str, payload: dict) -> None:
"""Write an issue-lock file directly, for lock-precedence tests."""
import json
os.makedirs(os.path.dirname(path), exist_ok=True)
with open(path, "w", encoding="utf-8") as fh:
json.dump(payload, fh)
if __name__ == "__main__":
unittest.main()
+149
View File
@@ -0,0 +1,149 @@
"""Unit tests for mcp_config_drift.py (#672)."""
from __future__ import annotations
import json
import pytest
from pathlib import Path
from mcp_config_drift import (
REQUIRED_GITEA_ROLE_SERVERS,
analyze_config_drift,
load_mcp_config,
)
@pytest.fixture
def sample_global_config() -> dict:
return {
"mcpServers": {
"gitea-author": {
"command": "python3",
"args": ["gitea_mcp_server.py"],
"env": {"GITEA_MCP_PROFILE": "prgs-author", "SENTRY_AUTH_TOKEN": "secret-token-999"},
},
"gitea-reviewer": {
"command": "python3",
"args": ["gitea_mcp_server.py"],
"env": {"GITEA_MCP_PROFILE": "prgs-reviewer"},
},
"gitea-merger": {
"command": "python3",
"args": ["gitea_mcp_server.py"],
"env": {"GITEA_MCP_PROFILE": "prgs-merger"},
},
"gitea-reconciler": {
"command": "python3",
"args": ["gitea_mcp_server.py"],
"env": {"GITEA_MCP_PROFILE": "prgs-reconciler"},
},
"gitea-controller": {
"command": "python3",
"args": ["gitea_mcp_server.py"],
"env": {"GITEA_MCP_PROFILE": "prgs-controller"},
},
"gitea-tools": {
"command": "python3",
"args": ["gitea_mcp_server.py"],
"env": {"GITEA_MCP_PROFILE": "prgs-author"},
},
}
}
def write_json(path: Path, data: dict) -> str:
path.write_text(json.dumps(data, indent=2), encoding="utf-8")
return str(path)
def test_drift_detection_in_sync(tmp_path, sample_global_config):
glob_file = tmp_path / "global_mcp.json"
act_file = tmp_path / "active_mcp.json"
write_json(glob_file, sample_global_config)
write_json(act_file, sample_global_config)
report = analyze_config_drift(active_config_path=str(act_file), global_config_path=str(glob_file))
assert report["in_sync"] is True
assert report["missing_role_servers"] == []
assert report["profile_mismatches"] == []
assert set(report["present_role_servers"]) == set(REQUIRED_GITEA_ROLE_SERVERS)
def test_drift_detection_missing_author(tmp_path, sample_global_config):
glob_file = tmp_path / "global_mcp.json"
act_file = tmp_path / "active_mcp.json"
active_config = json.loads(json.dumps(sample_global_config))
del active_config["mcpServers"]["gitea-author"]
write_json(glob_file, sample_global_config)
write_json(act_file, active_config)
report = analyze_config_drift(active_config_path=str(act_file), global_config_path=str(glob_file))
assert report["in_sync"] is False
assert "gitea-author" in report["missing_role_servers"]
assert "gitea-author" not in report["present_role_servers"]
def test_drift_detection_missing_reviewer(tmp_path, sample_global_config):
glob_file = tmp_path / "global_mcp.json"
act_file = tmp_path / "active_mcp.json"
active_config = json.loads(json.dumps(sample_global_config))
del active_config["mcpServers"]["gitea-reviewer"]
write_json(glob_file, sample_global_config)
write_json(act_file, active_config)
report = analyze_config_drift(active_config_path=str(act_file), global_config_path=str(glob_file))
assert report["in_sync"] is False
assert "gitea-reviewer" in report["missing_role_servers"]
def test_drift_detection_profile_mismatch(tmp_path, sample_global_config):
glob_file = tmp_path / "global_mcp.json"
act_file = tmp_path / "active_mcp.json"
active_config = json.loads(json.dumps(sample_global_config))
active_config["mcpServers"]["gitea-author"]["env"]["GITEA_MCP_PROFILE"] = "dadeschools-author"
write_json(glob_file, sample_global_config)
write_json(act_file, active_config)
report = analyze_config_drift(active_config_path=str(act_file), global_config_path=str(glob_file))
assert report["in_sync"] is False
assert len(report["profile_mismatches"]) == 1
mismatch = report["profile_mismatches"][0]
assert mismatch["server"] == "gitea-author"
assert mismatch["active_profile"] == "dadeschools-author"
assert mismatch["global_profile"] == "prgs-author"
def test_secret_redaction_in_drift_report(tmp_path, sample_global_config):
glob_file = tmp_path / "global_mcp.json"
act_file = tmp_path / "active_mcp.json"
write_json(glob_file, sample_global_config)
write_json(act_file, sample_global_config)
report = analyze_config_drift(active_config_path=str(act_file), global_config_path=str(glob_file))
serialized = str(report)
assert "secret-token-999" not in serialized
def test_sanctioned_runbook_forbids_pkill():
report = analyze_config_drift(active_config_path="/nonexistent/path/active.json", global_config_path="/nonexistent/path/global.json")
runbook_text = " ".join(report["sanctioned_repair_runbook"]).lower()
forbidden_text = " ".join(report["forbidden_repair_methods"]).lower()
assert "pkill" in forbidden_text
assert "mtime" in forbidden_text
assert "source" in forbidden_text
assert "session-state" in forbidden_text
+788
View File
@@ -0,0 +1,788 @@
"""Authoritative fleet inventory classification (#949).
One test group per acceptance criterion. The classifier is pure, so every
scenario is expressed as a snapshot: registry rows plus a process observation.
PID liveness is the one impure input, so it is patched per test rather than
depending on whatever happens to be running on the machine.
"""
import os
import unittest
from datetime import datetime, timedelta, timezone
from unittest import mock
import mcp_fleet_inventory as mfi
NOW = datetime(2026, 7, 28, 3, 0, 0, tzinfo=timezone.utc)
HEAD_A = "82d71b77028a7abd4f8ab4a4e4d89658a187f73d"
HEAD_B = "35ed8a2fcb11134a37c862ca6eaca26e3028902a"
COHORT_A = "ppid:40990"
COHORT_B = "ppid:51022"
BINDING = {"remote": "prgs", "org": "Scaled-Tech-Consulting", "repo": "Gitea-Tools"}
_BASE_PID = 41000
def row(
namespace,
profile,
role,
pid,
*,
cohort_id=COHORT_A,
startup_head=HEAD_A,
registered_at="2026-07-28T02:00:00Z",
**overrides,
):
"""One control-plane runtime registry row."""
record = {
"runtime_id": f"{namespace}:{pid}:boot{pid}",
"namespace": namespace,
"profile": profile,
"role": role,
"remote": BINDING["remote"],
"org": BINDING["org"],
"repo": BINDING["repo"],
"repository_root": "/checkout/Gitea-Tools",
"pid": pid,
"cohort_id": cohort_id,
"cohort_source": "parent_process",
"client_provenance": "client_managed",
"boot_id": f"boot{pid}",
"startup_head": startup_head,
"daemon_start_head": startup_head,
"transport": "stdio",
"registered_at": registered_at,
"last_heartbeat_at": registered_at,
"status": "running",
}
record.update(overrides)
return record
def healthy_rows():
"""One live row per expected member, all one cohort, all one revision."""
return [
row(entry["namespace"], entry["profile"], entry["role"], _BASE_PID + index)
for index, entry in enumerate(mfi.EXPECTED_PRGS_FLEET)
]
def scan_for(rows, *, available=True, extra_pids=(), started_at="2026-07-28T01:00:00Z"):
"""A process observation that corroborates *rows* (plus any extra PIDs)."""
processes = [
{"pid": r["pid"], "started_at": started_at, "command": "python mcp_server.py"}
for r in rows
]
processes.extend(
{"pid": pid, "started_at": started_at, "command": "python mcp_server.py"}
for pid in extra_pids
)
return {"available": available, "processes": processes, "reason": None}
def classify(rows, scan=None, **kwargs):
kwargs.setdefault("expected_binding", BINDING)
kwargs.setdefault("now", NOW)
return mfi.classify_fleet_inventory(
runtime_rows=rows,
process_scan=scan if scan is not None else scan_for(rows),
**kwargs,
)
class _AliveMixin:
"""Treat every PID as alive unless a test declares it dead or unknown."""
def setUp(self):
super().setUp()
self.dead_pids = set()
self.unknown_pids = set()
def probe(pid):
if pid in self.unknown_pids:
return None
return pid not in self.dead_pids
patcher = mock.patch.object(mfi, "probe_pid_alive", side_effect=probe)
patcher.start()
self.addCleanup(patcher.stop)
class TestHealthyFleet(_AliveMixin, unittest.TestCase):
"""AC1: a healthy five-server fleet reports all five members exactly once."""
def test_five_members_each_reported_once(self):
result = classify(healthy_rows())
self.assertEqual(len(result["configured_members"]), 5)
self.assertEqual(len(result["running_members"]), 5)
self.assertEqual(result["running_member_count"], 5)
self.assertEqual(result["missing_members"], [])
self.assertEqual(result["duplicate_members"], [])
self.assertEqual(result["unexpected_members"], [])
self.assertTrue(result["exactly_one_per_profile"])
for member in result["configured_members"]:
self.assertEqual(member["instance_count"], 1)
self.assertEqual(member["health"], mfi.HEALTH_RUNNING)
def test_healthy_fleet_satisfies_the_mutation_gate(self):
result = classify(healthy_rows())
self.assertTrue(result["inventory_complete"])
self.assertTrue(result["mutation_gate_satisfied"])
self.assertIsNone(result["blocked_reason"])
self.assertTrue(result["single_cohort"])
self.assertFalse(result["mixed_cohort"])
self.assertFalse(result["mixed_revision"])
def test_every_expected_namespace_and_profile_is_present(self):
result = classify(healthy_rows())
self.assertEqual(
sorted(m["profile"] for m in result["configured_members"]),
[
"prgs-author",
"prgs-controller",
"prgs-merger",
"prgs-reconciler",
"prgs-reviewer",
],
)
class TestDuplicateMember(_AliveMixin, unittest.TestCase):
"""AC2: two processes on one profile are a duplicate and fail the gate."""
def test_duplicate_is_reported_with_every_pid(self):
rows = healthy_rows()
rows.append(row("gitea-author", "prgs-author", "author", 49999))
result = classify(rows)
self.assertEqual(len(result["duplicate_members"]), 1)
duplicate = result["duplicate_members"][0]
self.assertEqual(duplicate["profile"], "prgs-author")
self.assertEqual(duplicate["instance_count"], 2)
self.assertEqual(duplicate["pids"], [_BASE_PID, 49999])
def test_duplicate_fails_the_mutation_gate(self):
rows = healthy_rows()
rows.append(row("gitea-author", "prgs-author", "author", 49999))
result = classify(rows)
self.assertFalse(result["exactly_one_per_profile"])
self.assertFalse(result["mutation_gate_satisfied"])
self.assertIn("duplicate", result["blocked_reason"])
def test_duplicate_is_not_reported_as_missing_or_unexpected(self):
rows = healthy_rows()
rows.append(row("gitea-author", "prgs-author", "author", 49999))
result = classify(rows)
self.assertEqual(result["missing_members"], [])
self.assertEqual(result["unexpected_members"], [])
class TestMissingMember(_AliveMixin, unittest.TestCase):
"""AC3: a missing expected server is identified and fails the gate."""
def test_missing_member_is_named(self):
rows = [r for r in healthy_rows() if r["profile"] != "prgs-merger"]
result = classify(rows)
self.assertEqual(len(result["missing_members"]), 1)
self.assertEqual(result["missing_members"][0]["profile"], "prgs-merger")
self.assertEqual(result["missing_members"][0]["health"], mfi.HEALTH_MISSING)
def test_missing_member_fails_the_mutation_gate(self):
rows = [r for r in healthy_rows() if r["profile"] != "prgs-merger"]
result = classify(rows)
self.assertFalse(result["exactly_one_per_profile"])
self.assertFalse(result["mutation_gate_satisfied"])
self.assertIn("prgs-merger", result["blocked_reason"])
def test_missing_member_still_lists_all_configured_members(self):
rows = [r for r in healthy_rows() if r["profile"] != "prgs-merger"]
result = classify(rows)
self.assertEqual(len(result["configured_members"]), 5)
merger = [
m for m in result["configured_members"] if m["profile"] == "prgs-merger"
][0]
self.assertEqual(merger["instance_count"], 0)
self.assertEqual(merger["health"], mfi.HEALTH_MISSING)
class TestUnexpectedMember(_AliveMixin, unittest.TestCase):
"""AC4: an unexpected PRGS server is reported explicitly."""
def test_unexpected_member_is_its_own_category(self):
rows = healthy_rows()
rows.append(row("gitea-shadow", "prgs-shadow", "author", 47777))
result = classify(rows)
self.assertEqual(len(result["unexpected_members"]), 1)
self.assertEqual(result["unexpected_members"][0]["profile"], "prgs-shadow")
self.assertEqual(
result["unexpected_members"][0]["health"], mfi.HEALTH_UNEXPECTED
)
def test_unexpected_member_is_not_collapsed_into_duplicates_or_missing(self):
rows = healthy_rows()
rows.append(row("gitea-shadow", "prgs-shadow", "author", 47777))
result = classify(rows)
self.assertEqual(result["duplicate_members"], [])
self.assertEqual(result["missing_members"], [])
self.assertTrue(result["exactly_one_per_profile"])
self.assertFalse(result["no_unexpected_members"])
self.assertFalse(result["mutation_gate_satisfied"])
def test_unexpected_member_blocks_with_its_own_reason(self):
rows = healthy_rows()
rows.append(row("gitea-shadow", "prgs-shadow", "author", 47777))
result = classify(rows)
self.assertTrue(
any("unexpected" in reason for reason in result["blocked_reasons"])
)
class TestMixedCohort(_AliveMixin, unittest.TestCase):
"""AC5: members from different client cohorts are detected."""
def test_two_cohorts_are_detected(self):
rows = healthy_rows()
rows[0]["cohort_id"] = COHORT_B
result = classify(rows)
self.assertTrue(result["mixed_cohort"])
self.assertFalse(result["single_cohort"])
self.assertEqual(result["cohort_ids"], sorted([COHORT_A, COHORT_B]))
def test_mixed_cohort_fails_the_mutation_gate(self):
rows = healthy_rows()
rows[0]["cohort_id"] = COHORT_B
result = classify(rows)
self.assertFalse(result["mutation_gate_satisfied"])
self.assertTrue(any("cohort" in reason for reason in result["blocked_reasons"]))
def test_mixed_cohort_is_not_a_duplicate_or_missing_report(self):
rows = healthy_rows()
rows[0]["cohort_id"] = COHORT_B
result = classify(rows)
self.assertEqual(result["duplicate_members"], [])
self.assertEqual(result["missing_members"], [])
class TestMixedRevision(_AliveMixin, unittest.TestCase):
"""AC6: members running different startup revisions are detected."""
def test_two_revisions_are_detected(self):
rows = healthy_rows()
rows[0]["startup_head"] = HEAD_B
result = classify(rows)
self.assertTrue(result["mixed_revision"])
self.assertEqual(result["startup_revisions"], sorted([HEAD_A, HEAD_B]))
def test_mixed_revision_fails_the_mutation_gate(self):
rows = healthy_rows()
rows[0]["startup_head"] = HEAD_B
result = classify(rows)
self.assertFalse(result["mutation_gate_satisfied"])
self.assertTrue(
any("revision" in reason for reason in result["blocked_reasons"])
)
def test_mixed_revision_does_not_by_itself_imply_mixed_cohort(self):
rows = healthy_rows()
rows[0]["startup_head"] = HEAD_B
result = classify(rows)
self.assertTrue(result["single_cohort"])
self.assertFalse(result["mixed_cohort"])
class TestRevisionIsNotCohort(_AliveMixin, unittest.TestCase):
"""AC7: matching Git revisions alone do not establish a single cohort."""
def test_identical_revisions_with_unknown_cohort_stay_unknown(self):
rows = healthy_rows()
for r in rows:
r["cohort_id"] = None
result = classify(rows)
self.assertEqual(len({r["startup_head"] for r in rows}), 1)
self.assertIsNone(result["single_cohort"])
self.assertIsNone(result["mixed_cohort"])
def test_identical_revisions_with_unknown_cohort_fail_closed(self):
rows = healthy_rows()
for r in rows:
r["cohort_id"] = None
result = classify(rows)
self.assertFalse(result["inventory_complete"])
self.assertFalse(result["mutation_gate_satisfied"])
self.assertTrue(
any("cohort" in reason for reason in result["incomplete_reasons"])
)
def test_one_unknown_cohort_among_known_ones_still_fails_closed(self):
rows = healthy_rows()
rows[0]["cohort_id"] = None
result = classify(rows)
self.assertIsNone(result["single_cohort"])
self.assertFalse(result["mutation_gate_satisfied"])
def test_cohort_derivation_never_consults_revisions(self):
"""The cohort helper takes no revision input at all."""
identity = mfi.derive_cohort_identity({mfi.COHORT_ID_ENV: "cohort-x"})
self.assertEqual(identity["cohort_id"], "cohort-x")
self.assertEqual(identity["cohort_source"], "explicit_env")
def test_orphaned_process_reports_unknown_cohort(self):
with mock.patch("os.getppid", return_value=1):
identity = mfi.derive_cohort_identity({})
self.assertIsNone(identity["cohort_id"])
self.assertEqual(identity["cohort_source"], "unknown")
def test_parent_process_is_the_cohort_when_no_env_is_set(self):
with mock.patch("os.getppid", return_value=40990):
identity = mfi.derive_cohort_identity({})
self.assertEqual(identity["cohort_id"], COHORT_A)
self.assertEqual(identity["cohort_source"], "parent_process")
class TestConfigurationIsNotRunning(_AliveMixin, unittest.TestCase):
"""AC8: configuration without a live worker is not a running member."""
def test_no_registry_rows_means_every_member_is_missing(self):
result = classify([], scan=scan_for([]))
self.assertEqual(len(result["missing_members"]), 5)
self.assertEqual(result["running_members"], [])
self.assertFalse(result["mutation_gate_satisfied"])
def test_dead_pid_is_stale_not_running(self):
rows = healthy_rows()
self.dead_pids = {rows[0]["pid"]}
result = classify(rows)
self.assertEqual(len(result["running_members"]), 4)
self.assertEqual(len(result["stale_members"]), 1)
self.assertEqual(result["stale_members"][0]["liveness"], mfi.LIVENESS_DEAD)
self.assertEqual(result["stale_members"][0]["health"], mfi.HEALTH_STALE)
self.assertEqual(len(result["missing_members"]), 1)
def test_registry_row_without_a_matching_process_is_not_running(self):
rows = healthy_rows()
scan = scan_for(rows[1:]) # first member's process is absent
result = classify(rows, scan=scan)
self.assertEqual(len(result["running_members"]), 4)
self.assertEqual(
result["stale_members"][0]["liveness"], mfi.LIVENESS_UNOBSERVED
)
self.assertFalse(result["mutation_gate_satisfied"])
def test_recycled_pid_does_not_impersonate_a_dead_server(self):
rows = healthy_rows()
scan = scan_for(rows)
# The process now holding the first PID started *after* registration.
scan["processes"][0]["started_at"] = "2026-07-28T02:30:00Z"
result = classify(rows, scan=scan)
stale = [
m
for m in result["stale_members"]
if m["liveness"] == mfi.LIVENESS_PID_RECYCLED
]
self.assertEqual(len(stale), 1)
self.assertEqual(len(result["running_members"]), 4)
self.assertFalse(result["mutation_gate_satisfied"])
class TestIncompleteEvidenceFailsClosed(_AliveMixin, unittest.TestCase):
"""AC9: unknown or unavailable evidence produces a fail-closed result."""
def test_unavailable_process_listing_fails_closed(self):
rows = healthy_rows()
result = classify(
rows,
scan={"available": False, "processes": [], "reason": "ps unavailable"},
)
self.assertFalse(result["inventory_complete"])
self.assertFalse(result["mutation_gate_satisfied"])
self.assertIn("ps unavailable", result["incomplete_reasons"])
self.assertEqual(result["running_members"], [])
def test_unreadable_registry_fails_closed(self):
result = classify(
[],
scan=scan_for([]),
registry_available=False,
registry_error="registry unreadable",
)
self.assertFalse(result["inventory_complete"])
self.assertFalse(result["mutation_gate_satisfied"])
self.assertIn("registry unreadable", result["incomplete_reasons"])
def test_unregistered_running_process_fails_closed(self):
rows = healthy_rows()
result = classify(rows, scan=scan_for(rows, extra_pids=[59999]))
self.assertEqual([p["pid"] for p in result["unregistered_processes"]], [59999])
self.assertFalse(result["inventory_complete"])
self.assertFalse(result["mutation_gate_satisfied"])
self.assertTrue(
any("59999" in reason for reason in result["incomplete_reasons"])
)
def test_undeterminable_pid_liveness_is_unknown_not_healthy(self):
rows = healthy_rows()
self.unknown_pids = {rows[0]["pid"]}
result = classify(rows)
self.assertEqual(result["stale_members"][0]["liveness"], mfi.LIVENESS_UNKNOWN)
self.assertEqual(result["stale_members"][0]["health"], mfi.HEALTH_UNKNOWN)
self.assertFalse(result["inventory_complete"])
self.assertFalse(result["mutation_gate_satisfied"])
def test_unknown_startup_revision_fails_closed(self):
rows = healthy_rows()
rows[0]["startup_head"] = None
result = classify(rows)
self.assertFalse(result["inventory_complete"])
self.assertFalse(result["mutation_gate_satisfied"])
self.assertTrue(
any("startup revision" in r for r in result["incomplete_reasons"])
)
def test_unknown_is_distinguished_from_healthy(self):
rows = healthy_rows()
for r in rows:
r["cohort_id"] = None
result = classify(rows)
# Not "unhealthy" in the sense of a named defect: nothing is missing,
# duplicated or unexpected. It is *unknown*, and that still fails closed.
self.assertEqual(result["missing_members"], [])
self.assertEqual(result["duplicate_members"], [])
self.assertEqual(result["unexpected_members"], [])
self.assertIsNone(result["single_cohort"])
self.assertFalse(result["mutation_gate_satisfied"])
class TestBindingAndRoleConsistency(_AliveMixin, unittest.TestCase):
"""Repository-binding and role/profile mismatches fail closed."""
def test_repository_binding_mismatch_is_reported(self):
rows = healthy_rows()
rows[0]["repo"] = "Some-Other-Repo"
result = classify(rows)
self.assertEqual(len(result["repository_binding_mismatches"]), 1)
self.assertEqual(
result["repository_binding_mismatches"][0]["repo"], "Some-Other-Repo"
)
self.assertFalse(result["mutation_gate_satisfied"])
def test_incomplete_repository_binding_is_reported(self):
rows = healthy_rows()
rows[0]["org"] = None
result = classify(rows)
self.assertEqual(len(result["repository_binding_mismatches"]), 1)
self.assertIn("complete", result["repository_binding_mismatches"][0]["reason"])
self.assertFalse(result["mutation_gate_satisfied"])
def test_role_mismatch_is_reported(self):
rows = healthy_rows()
rows[0]["role"] = "merger" # prgs-author is configured as author
result = classify(rows)
self.assertEqual(len(result["role_mismatches"]), 1)
self.assertEqual(result["role_mismatches"][0]["expected_role"], "author")
self.assertEqual(result["role_mismatches"][0]["declared_role"], "merger")
self.assertFalse(result["mutation_gate_satisfied"])
def test_profile_served_from_the_wrong_namespace_is_reported(self):
rows = healthy_rows()
rows[0]["namespace"] = "gitea-reviewer"
result = classify(rows)
self.assertTrue(
any(
m.get("expected_namespace") == "gitea-author"
for m in result["role_mismatches"]
)
)
self.assertFalse(result["mutation_gate_satisfied"])
def test_matching_binding_produces_no_mismatch(self):
result = classify(healthy_rows())
self.assertEqual(result["repository_binding_mismatches"], [])
self.assertEqual(result["role_mismatches"], [])
class TestDeterminism(_AliveMixin, unittest.TestCase):
"""AC10 / stable ordering and deterministic structured output."""
def test_identical_snapshots_produce_identical_results(self):
rows = healthy_rows()
first = classify(rows, scan=scan_for(rows))
second = classify(healthy_rows(), scan=scan_for(healthy_rows()))
self.assertEqual(first, second)
def test_row_order_does_not_change_the_result(self):
rows = healthy_rows()
shuffled = list(reversed(healthy_rows()))
forward = classify(rows, scan=scan_for(rows))
backward = classify(shuffled, scan=scan_for(shuffled))
self.assertEqual(forward["running_members"], backward["running_members"])
self.assertEqual(forward["configured_members"], backward["configured_members"])
self.assertEqual(
forward["mutation_gate_satisfied"], backward["mutation_gate_satisfied"]
)
def test_members_are_sorted_by_namespace_then_profile_then_pid(self):
rows = healthy_rows()
rows.append(row("gitea-author", "prgs-author", "author", 40001))
result = classify(rows)
keys = [
(m["namespace"], m["profile"], m["pid"]) for m in result["running_members"]
]
self.assertEqual(keys, sorted(keys))
def test_answering_namespace_does_not_change_the_verdict(self):
"""AC10: controller and reconciler agree for one fleet snapshot."""
rows = healthy_rows()
controller = classify(
rows, scan=scan_for(rows), answering_namespace="gitea-controller"
)
reconciler = classify(
healthy_rows(),
scan=scan_for(healthy_rows()),
answering_namespace="gitea-reconciler",
)
self.assertEqual(controller["answering_namespace"], "gitea-controller")
self.assertEqual(reconciler["answering_namespace"], "gitea-reconciler")
for key in sorted(set(controller) - {"answering_namespace"}):
self.assertEqual(controller[key], reconciler[key], f"{key} disagreed")
def test_controller_and_reconciler_agree_on_an_unhealthy_fleet(self):
rows = [r for r in healthy_rows() if r["profile"] != "prgs-merger"]
controller = classify(
rows, scan=scan_for(rows), answering_namespace="gitea-controller"
)
reconciler = classify(
rows, scan=scan_for(rows), answering_namespace="gitea-reconciler"
)
self.assertEqual(controller["blocked_reason"], reconciler["blocked_reason"])
self.assertEqual(
controller["mutation_gate_satisfied"],
reconciler["mutation_gate_satisfied"],
)
class TestReadOnly(_AliveMixin, unittest.TestCase):
"""AC11: the capability mutates nothing."""
def test_classification_reports_no_mutations(self):
result = classify(healthy_rows())
self.assertEqual(result["mutations_performed"], [])
self.assertTrue(result["read_only"])
def test_classification_does_not_write_the_input_rows_back(self):
rows = healthy_rows()
snapshot = [dict(r) for r in rows]
classify(rows)
self.assertEqual(rows, snapshot)
class TestLivenessProbe(unittest.TestCase):
"""AC11: the real probe never sends a terminating signal.
Deliberately *not* using ``_AliveMixin`` these tests exercise
``probe_pid_alive`` itself, which the mixin replaces.
"""
def test_only_signal_zero_is_ever_sent(self):
with mock.patch("os.kill") as killer:
mfi.probe_pid_alive(4242)
killer.assert_called_once_with(4242, 0)
def test_classifying_a_duplicate_never_terminates_it(self):
rows = healthy_rows()
rows.append(row("gitea-author", "prgs-author", "author", 49999))
with mock.patch("os.kill") as killer:
result = mfi.classify_fleet_inventory(
runtime_rows=rows,
process_scan=scan_for(rows),
expected_binding=BINDING,
now=NOW,
)
self.assertTrue(killer.call_args_list, "liveness must actually be probed")
for call in killer.call_args_list:
self.assertEqual(call.args[1], 0, "only signal 0 may ever be sent")
self.assertEqual(result["mutations_performed"], [])
def test_liveness_probe_tolerates_a_missing_process(self):
with mock.patch("os.kill", side_effect=ProcessLookupError):
self.assertFalse(mfi.probe_pid_alive(4242))
def test_liveness_probe_treats_permission_error_as_alive(self):
with mock.patch("os.kill", side_effect=PermissionError):
self.assertTrue(mfi.probe_pid_alive(4242))
def test_liveness_probe_returns_unknown_on_other_os_errors(self):
with mock.patch("os.kill", side_effect=OSError):
self.assertIsNone(mfi.probe_pid_alive(4242))
def test_invalid_pid_is_unknown_and_probes_nothing(self):
with mock.patch("os.kill") as killer:
self.assertIsNone(mfi.probe_pid_alive(None))
self.assertIsNone(mfi.probe_pid_alive(0))
killer.assert_not_called()
class TestMultiClientRegression(_AliveMixin, unittest.TestCase):
"""The multi-LLM duplicate-server scenario that motivated #949."""
@staticmethod
def _two_client_rows():
first = healthy_rows()
second = [
row(
entry["namespace"],
entry["profile"],
entry["role"],
50000 + index,
cohort_id=COHORT_B,
)
for index, entry in enumerate(mfi.EXPECTED_PRGS_FLEET)
]
return first + second
def test_second_client_running_the_same_five_profiles_is_caught(self):
"""Two clients, ten servers, same profiles, same revision.
Every member self-reports the same parity, which is exactly why the old
per-process surfaces reported success. The fleet inventory must report
five duplicates, two cohorts, and a closed gate.
"""
rows = self._two_client_rows()
result = classify(rows, scan=scan_for(rows))
self.assertEqual(len(result["duplicate_members"]), 5)
self.assertFalse(result["exactly_one_per_profile"])
self.assertTrue(result["mixed_cohort"])
self.assertFalse(result["single_cohort"])
self.assertFalse(result["mixed_revision"], "both clients share a revision")
self.assertFalse(result["mutation_gate_satisfied"])
self.assertEqual(result["missing_members"], [])
def test_identical_parity_across_ten_servers_is_not_health(self):
rows = self._two_client_rows()
result = classify(rows, scan=scan_for(rows))
self.assertEqual(result["startup_revisions"], [HEAD_A])
self.assertFalse(result["mutation_gate_satisfied"])
def test_every_duplicate_pid_is_named_for_the_operator(self):
rows = self._two_client_rows()
result = classify(rows, scan=scan_for(rows))
for duplicate in result["duplicate_members"]:
self.assertEqual(len(duplicate["pids"]), 2)
class TestProcessScan(unittest.TestCase):
"""The process observation is server-side and fails closed."""
def test_scan_parses_mcp_server_processes(self):
stdout = (
" PID STARTED COMMAND\n"
" 41000 Mon Jul 27 20:00:00 2026 python /path/mcp_server.py\n"
" 41001 Mon Jul 27 20:00:01 2026 python /path/other_server.py\n"
)
result = mfi.scan_mcp_server_processes(
runner=lambda *a, **k: mock.Mock(stdout=stdout)
)
self.assertTrue(result["available"])
self.assertEqual([p["pid"] for p in result["processes"]], [41000])
def test_scan_failure_reports_unavailable_rather_than_empty(self):
def boom(*args, **kwargs):
raise OSError("ps missing")
result = mfi.scan_mcp_server_processes(runner=boom)
self.assertFalse(result["available"])
self.assertEqual(result["processes"], [])
self.assertIn("ps missing", result["reason"])
def test_scan_results_are_sorted_by_pid(self):
stdout = (
" PID STARTED COMMAND\n"
" 41005 Mon Jul 27 20:00:00 2026 python /path/mcp_server.py\n"
" 41001 Mon Jul 27 20:00:01 2026 python /path/mcp_server.py\n"
)
result = mfi.scan_mcp_server_processes(
runner=lambda *a, **k: mock.Mock(stdout=stdout)
)
self.assertEqual([p["pid"] for p in result["processes"]], [41001, 41005])
class TestRuntimeRecord(unittest.TestCase):
"""The row a server writes about itself."""
def test_record_describes_the_calling_process(self):
record = mfi.build_process_runtime_record(
namespace="gitea-controller",
profile="prgs-controller",
role="controller",
remote="prgs",
org=BINDING["org"],
repo=BINDING["repo"],
pid=41022,
startup_head=HEAD_A,
transport="stdio",
client_provenance="client_managed",
env={mfi.COHORT_ID_ENV: COHORT_A},
boot_id="deadbeefcafe0001",
)
self.assertEqual(
record["runtime_id"], "gitea-controller:41022:deadbeefcafe0001"
)
self.assertEqual(record["namespace"], "gitea-controller")
self.assertEqual(record["cohort_id"], COHORT_A)
self.assertEqual(record["cohort_source"], "explicit_env")
self.assertEqual(record["startup_head"], HEAD_A)
self.assertEqual(record["status"], "running")
def test_record_defaults_pid_to_the_current_process(self):
record = mfi.build_process_runtime_record(
namespace="gitea-author", profile="prgs-author", role="author"
)
self.assertEqual(record["pid"], os.getpid())
def test_each_boot_gets_a_distinct_runtime_id(self):
first = mfi.build_process_runtime_record(
namespace="gitea-author", profile="prgs-author", role="author", pid=1
)
second = mfi.build_process_runtime_record(
namespace="gitea-author", profile="prgs-author", role="author", pid=1
)
self.assertNotEqual(first["runtime_id"], second["runtime_id"])
def test_registered_at_uses_the_control_plane_timestamp_format(self):
record = mfi.build_process_runtime_record(
namespace="gitea-author", profile="prgs-author", role="author", pid=1
)
datetime.strptime(record["registered_at"], "%Y-%m-%dT%H:%M:%SZ")
class TestSummary(_AliveMixin, unittest.TestCase):
def test_healthy_summary_names_the_counts(self):
result = classify(healthy_rows())
self.assertIn("5 of 5", mfi.summarize(result))
def test_blocked_summary_repeats_the_blocked_reason(self):
rows = [r for r in healthy_rows() if r["profile"] != "prgs-merger"]
result = classify(rows)
self.assertIn(result["blocked_reason"], mfi.summarize(result))
class TestHeartbeatAge(_AliveMixin, unittest.TestCase):
def test_heartbeat_age_is_reported_for_diagnosis(self):
rows = healthy_rows()
stamp = (NOW - timedelta(minutes=30)).strftime("%Y-%m-%dT%H:%M:%SZ")
rows[0]["last_heartbeat_at"] = stamp
result = classify(rows)
member = [m for m in result["running_members"] if m["pid"] == rows[0]["pid"]][0]
self.assertEqual(member["heartbeat_age_seconds"], 1800)
def test_missing_heartbeat_is_reported_as_unknown_age(self):
rows = healthy_rows()
rows[0]["last_heartbeat_at"] = None
result = classify(rows)
member = [m for m in result["running_members"] if m["pid"] == rows[0]["pid"]][0]
self.assertIsNone(member["heartbeat_age_seconds"])
if __name__ == "__main__":
unittest.main()
+190
View File
@@ -0,0 +1,190 @@
"""Tests for Sentry/GlitchTip observability console (#649, Phase 4)."""
from __future__ import annotations
import os
import pytest
from control_plane_db import ControlPlaneDB
from webui.app import create_app
from webui.console_authz import authorize, resolve_principal
from webui.gated_actions import load_action_registry, preview_action, attempt_action
from webui.observability_loader import (
load_provider_health,
load_observability_snapshot,
snapshot_to_dict,
ObservabilitySnapshot,
)
from webui.observability_views import render_observability_page
from tests.webui_testclient import TestClient
@pytest.fixture
def test_db(tmp_path):
db_path = str(tmp_path / "test_control_plane.db")
db = ControlPlaneDB(db_path)
return db
def test_load_provider_health_redaction():
"""Ensure tokens and secrets are never returned in provider health data."""
env = {
"SENTRY_BASE_URL": "https://sentry.prgs.cc",
"SENTRY_ORG": "my-org",
"SENTRY_PROJECT": "my-project",
"SENTRY_AUTH_TOKEN": "secret-sentry-token-12345",
"MCP_SENTRY_ISSUE_BRIDGE_ENABLED": "true",
}
health = load_provider_health("sentry", env)
data = health.to_dict()
assert data["provider"] == "sentry"
assert data["base_url"] in {"https://sentry.prgs.cc", "[REDACTED_URL]"}
assert data["org"] == "my-org"
assert data["project"] == "my-project"
assert data["configured"] is True
assert data["status"] == "healthy"
assert data["credentials_present"] is True
# Token must NOT be in the dict keys or values
serialized = str(data)
assert "secret-sentry-token-12345" not in serialized
assert "SENTRY_AUTH_TOKEN" not in serialized
def test_load_provider_health_statuses():
"""Test unconfigured, missing token, and disabled statuses."""
# Not configured
h1 = load_provider_health("sentry", {})
d1 = h1.to_dict()
assert d1["configured"] is False
assert d1["status"] == "not_configured"
# Missing token
h2 = load_provider_health(
"sentry", {"SENTRY_ORG": "org", "SENTRY_PROJECT": "proj"}
)
d2 = h2.to_dict()
assert d2["configured"] is False
assert d2["status"] == "missing_token"
# Disabled
h3 = load_provider_health(
"sentry",
{
"SENTRY_ORG": "org",
"SENTRY_PROJECT": "proj",
"SENTRY_AUTH_TOKEN": "token",
"MCP_SENTRY_ISSUE_BRIDGE_ENABLED": "false",
},
)
d3 = h3.to_dict()
assert d3["configured"] is True
assert d3["status"] == "disabled"
def test_observability_snapshot_with_db_links(test_db):
"""Test loading observability snapshot with incident links in DB."""
test_db.upsert_incident_link(
provider="sentry",
provider_issue_id="101",
gitea_org="Scaled-Tech-Consulting",
gitea_repo="Gitea-Tools",
gitea_issue_number=649,
provider_base_url="https://sentry.prgs.cc",
provider_org="Scaled-Tech-Consulting",
provider_project="Gitea-Tools",
provider_short_id="ST-101",
provider_permalink="https://sentry.prgs.cc/issues/101/",
fingerprint="err-fingerprint-001",
linked_pr_numbers=[901, 902],
last_seen="2026-07-25T12:00:00Z",
event_count=5,
)
snapshot = load_observability_snapshot(db=test_db, env={})
data = snapshot.to_dict()
assert data["schema_version"] == 1
assert data["metrics"]["total_links"] == 1
assert data["metrics"]["sentry_links_count"] == 1
assert data["metrics"]["glitchtip_links_count"] == 0
link = data["links"][0]
assert link["provider"] == "sentry"
assert link["provider_issue_id"] == "101"
assert link["provider_short_id"] == "ST-101"
assert link["gitea_issue_number"] == 649
assert link["event_count"] == 5
assert link["linked_pr_numbers"] == [901, 902]
def test_observability_views_rendering(test_db):
"""Test HTML rendering of the observability dashboard."""
snapshot = load_observability_snapshot(db=test_db, env={})
html_output = render_observability_page(snapshot)
assert "Observability &amp; Incident Bridge (#649)" in html_output or "Observability & Incident Bridge (#649)" in html_output or "Observability" in html_output
assert "ADR Authority Model:" in html_output
assert "Provider Connections" in html_output
assert "Correlated Incidents" in html_output
def test_webui_observability_routes():
"""Test Starlette HTTP routes for /observability and /api/v1/observability."""
client = TestClient(create_app())
# HTML page route
res_html = client.get("/observability")
assert res_html.status_code == 200
assert "text/html" in res_html.headers["content-type"]
assert "Observability" in res_html.text
# Versioned API route
res_api_v1 = client.get("/api/v1/observability")
assert res_api_v1.status_code == 200
assert "application/json" in res_api_v1.headers["content-type"]
data_v1 = res_api_v1.json()
assert "schema_version" in data_v1
assert "providers" in data_v1
assert "links" in data_v1
assert "metrics" in data_v1
# Compatibility alias route
res_api_alias = client.get("/api/observability")
assert res_api_alias.status_code == 200
assert res_api_alias.json() == data_v1
def test_observability_gated_actions():
"""Ensure observability actions are registered, gated, and fail closed in MVP mode."""
registry = load_action_registry()
action_reconcile = registry.get("observability_reconcile_incident")
assert action_reconcile is not None
assert action_reconcile.task_key == "observability_reconcile_incident"
assert action_reconcile.mcp_tool == "gitea_observability_reconcile_incident"
action_link = registry.get("observability_link_issue")
assert action_link is not None
assert action_link.task_key == "observability_link_issue"
# Preview returns mutation ledger
prev = preview_action("observability_reconcile_incident", provider="sentry", issue_id="101")
assert prev["action_id"] == "observability_reconcile_incident"
assert prev["enabled"] is False
# Execution fails closed in MVP mode
att = attempt_action("observability_reconcile_incident", provider="sentry", issue_id="101")
assert att["success"] is False
assert att["error"] == "action_disabled"
def test_observability_authz_rbac():
"""Test RBAC authorization for observability actions."""
principal = resolve_principal({})
# Check authorize decision
decision = authorize("observability_reconcile_incident", principal)
assert decision.action_id == "observability_reconcile_incident"
# Phase 4 action denies in Phase 1 runtime by default
assert decision.allowed is False
+21
View File
@@ -87,6 +87,11 @@ from webui.notifications import (
from webui.notification_views import render_notifications_page
from webui import request_service
from webui.request_views import render_requests_page
from webui.observability_loader import (
load_observability_snapshot,
snapshot_to_dict as observability_snapshot_to_dict,
)
from webui.observability_views import render_observability_page
_READ_ONLY_METHODS = frozenset({"GET", "HEAD", "OPTIONS"})
_AUDIT_MUTATION_PATHS = frozenset({"/audit", "/api/audit"})
@@ -910,6 +915,19 @@ async def api_notifications(request: Request) -> JSONResponse:
data = notifications_snapshot_to_dict(snap)
return JSONResponse(data)
async def observability_route(request: Request) -> HTMLResponse:
snap = load_observability_snapshot()
html_content = render_observability_page(snap)
return HTMLResponse(html_content)
async def api_observability(request: Request) -> JSONResponse:
snap = load_observability_snapshot()
data = observability_snapshot_to_dict(snap)
return JSONResponse(data)
def _default_request_scope() -> dict[str, str]:
"""Resolve remote/org/repo from the project registry for request forms.
@@ -1075,6 +1093,9 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
Route("/api/analytics", api_v1_analytics, methods=["GET"]),
Route("/api/v1/analytics", api_v1_analytics, methods=["GET"]),
Route("/api/v1/analytics/usage", api_v1_analytics_ingest, methods=["POST"]),
Route("/observability", observability_route, methods=["GET"]),
Route("/api/observability", api_observability, methods=["GET"]),
Route("/api/v1/observability", api_observability, methods=["GET"]),
Route("/audit", audit, methods=["GET", "POST"]),
Route("/api/audit", api_audit, methods=["GET", "POST"]),
Route("/worktrees", worktrees, methods=["GET"]),
+23
View File
@@ -317,6 +317,29 @@ _ACTION_SPECS: tuple[ConsoleAction, ...] = (
phase=2,
summary="Run reconciler cleanup for merged or superseded PR branches.",
),
# #649: Phase 4 observability & incident bridge actions.
ConsoleAction(
action_id="observability_reconcile_incident",
task_key="observability_reconcile_incident",
action_class=CLASS_WRITE,
minimum_role=OPERATOR,
requires_confirmation=True,
dual_control=False,
break_glass=False,
phase=4,
summary="Trigger/reconcile durable Gitea issue creation from a provider incident.",
),
ConsoleAction(
action_id="observability_link_issue",
task_key="observability_link_issue",
action_class=CLASS_WRITE,
minimum_role=OPERATOR,
requires_confirmation=True,
dual_control=False,
break_glass=False,
phase=4,
summary="Link a provider incident to an existing Gitea issue.",
),
# #643: submit a work request — desired role, issue/PR, intent — and let
# the allocator reserve it. This is the one Phase 2 action whose execution
# path is actually implemented (``webui.request_service``), so it carries
+5
View File
@@ -185,6 +185,11 @@ def build_action_registry() -> ActionRegistry:
"console.rebind_session_worktree", "Rebind session worktree to verified lease."),
("system.reconcile_cleanups", "Reconcile cleanups", "reconcile_cleanups",
"console.reconcile_cleanups", "Run reconciler cleanup for merged or superseded PRs."),
# #649: Phase 4 observability & incident bridge actions.
("observability_reconcile_incident", "Reconcile incident", "observability_reconcile_incident",
"gitea_observability_reconcile_incident", "Trigger or dry-run durable issue reconciliation for a provider incident."),
("observability_link_issue", "Link incident issue", "observability_link_issue",
"gitea_observability_link_issue", "Link a provider incident to a Gitea tracking issue."),
)
actions = tuple(
GatedAction(
+1
View File
@@ -73,6 +73,7 @@ NAV_GROUPS: tuple[NavGroup, ...] = (
)),
NavGroup("Insights", (
NavItem("/insights", "Insights", "stub"),
NavItem("/observability", "Observability"),
NavItem("/analytics", "Analytics"),
NavItem("/audit", "Audit"),
)),
+275
View File
@@ -0,0 +1,275 @@
"""Sentry/GlitchTip observability and incident correlation loader for the console (#649, Phase 4).
Operators need to inspect provider connection status (Sentry/GlitchTip), error
correlations, and durable Gitea issue linkage without treating raw incidents
as allocator work.
ADR authority model:
* Gitea owns work.
* Providers (Sentry/GlitchTip) observe incidents.
* Control-plane DB coordinates incident links.
* The #612 bridge reconciles observations into durable Gitea issues.
* The web console projects read-only state and gates mutations.
Redaction boundary:
* Provider auth tokens, DSNs, Authorization headers, and sensitive local file
paths are ALWAYS redacted before leaving this module.
"""
from __future__ import annotations
import os
from dataclasses import dataclass
from typing import Any
from control_plane_db import ControlPlaneDB
import sentry_incident_bridge
from webui import console_redaction
OBSERVABILITY_SCHEMA_VERSION = 1
@dataclass(frozen=True)
class ProviderHealth:
"""Connection and health status of an observability provider."""
provider: str
base_url: str
org: str
project: str
configured: bool
status: str
bridge_enabled: bool
lookback: str
min_events_for_issue: int
self_hosted: bool
environment: str | None = None
credentials_present: bool = False
def to_dict(self) -> dict[str, Any]:
data = {
"provider": self.provider,
"base_url": self.base_url,
"org": self.org,
"project": self.project,
"configured": self.configured,
"status": self.status,
"bridge_enabled": self.bridge_enabled,
"lookback": self.lookback,
"min_events_for_issue": self.min_events_for_issue,
"self_hosted": self.self_hosted,
"environment": self.environment,
"credentials_present": self.credentials_present,
}
return console_redaction.redact_payload(data)
@dataclass(frozen=True)
class CorrelatedIncidentLink:
"""One linked provider incident ↔ Gitea issue correlation record."""
link_id: int
provider: str
provider_base_url: str
provider_org: str
provider_project: str
provider_issue_id: str
provider_short_id: str | None
provider_permalink: str | None
fingerprint: str | None
gitea_org: str
gitea_repo: str
gitea_issue_number: int
linked_pr_numbers: list[int]
last_seen: str | None
event_count: int
created_at: str | None
updated_at: str | None
def to_dict(self) -> dict[str, Any]:
data = {
"link_id": self.link_id,
"provider": self.provider,
"provider_base_url": self.provider_base_url,
"provider_org": self.provider_org,
"provider_project": self.provider_project,
"provider_issue_id": self.provider_issue_id,
"provider_short_id": self.provider_short_id,
"provider_permalink": self.provider_permalink,
"fingerprint": self.fingerprint,
"gitea_org": self.gitea_org,
"gitea_repo": self.gitea_repo,
"gitea_issue_number": self.gitea_issue_number,
"linked_pr_numbers": self.linked_pr_numbers,
"last_seen": self.last_seen,
"event_count": self.event_count,
"created_at": self.created_at,
"updated_at": self.updated_at,
}
return console_redaction.redact_payload(data)
@dataclass(frozen=True)
class ObservabilitySnapshot:
"""Read-only snapshot of observability provider status and incident correlations."""
schema_version: int
providers: list[ProviderHealth]
links: list[CorrelatedIncidentLink]
total_links: int
sentry_links_count: int
glitchtip_links_count: int
bridge_active: bool
def to_dict(self) -> dict[str, Any]:
return {
"schema_version": self.schema_version,
"providers": [p.to_dict() for p in self.providers],
"links": [link.to_dict() for link in self.links],
"metrics": {
"total_links": self.total_links,
"sentry_links_count": self.sentry_links_count,
"glitchtip_links_count": self.glitchtip_links_count,
"bridge_active": self.bridge_active,
},
}
def load_provider_health(
provider_name: str = "sentry",
env: dict[str, str] | None = None,
) -> ProviderHealth:
"""Inspect configuration and connection health for an observability provider."""
source_env = dict(env if env is not None else os.environ)
if provider_name.lower() == "sentry":
config = sentry_incident_bridge.load_bridge_config(source_env)
token = sentry_incident_bridge.resolve_token(source_env)
has_token = bool(token)
configured = bool(config.org and config.project and has_token)
if not config.org or not config.project:
status = "not_configured"
elif not has_token:
status = "missing_token"
elif not config.bridge_enabled:
status = "disabled"
else:
status = "healthy"
return ProviderHealth(
provider="sentry",
base_url=config.base_url,
org=config.org or "unconfigured",
project=config.project or "unconfigured",
configured=configured,
status=status,
bridge_enabled=config.bridge_enabled,
lookback=config.lookback,
min_events_for_issue=config.min_events_for_issue,
self_hosted=not config.base_url.rstrip("/").endswith("sentry.io"),
environment=config.environment,
credentials_present=has_token,
)
# GlitchTip or fallback provider configuration
glitchtip_url = (source_env.get("GLITCHTIP_BASE_URL") or "https://glitchtip.prgs.cc").strip()
glitchtip_org = (source_env.get("GLITCHTIP_ORG") or "").strip()
glitchtip_proj = (source_env.get("GLITCHTIP_PROJECT") or "").strip()
glitchtip_token = (source_env.get("GLITCHTIP_AUTH_TOKEN") or "").strip()
has_token = bool(glitchtip_token)
configured = bool(glitchtip_org and glitchtip_proj and has_token)
status = "healthy" if configured else ("missing_token" if glitchtip_org and glitchtip_proj else "not_configured")
return ProviderHealth(
provider="glitchtip",
base_url=glitchtip_url,
org=glitchtip_org or "unconfigured",
project=glitchtip_proj or "unconfigured",
configured=configured,
status=status,
bridge_enabled=configured,
lookback="24h",
min_events_for_issue=2,
self_hosted=True,
environment=source_env.get("GLITCHTIP_ENVIRONMENT"),
credentials_present=has_token,
)
def _parse_pr_numbers(raw: Any) -> list[int]:
if isinstance(raw, list):
return [int(x) for x in raw if str(x).isdigit()]
if isinstance(raw, str) and raw.strip():
import json
try:
parsed = json.loads(raw)
if isinstance(parsed, list):
return [int(x) for x in parsed if str(x).isdigit()]
except Exception:
pass
return []
def load_observability_snapshot(
db: ControlPlaneDB | None = None,
env: dict[str, str] | None = None,
) -> ObservabilitySnapshot:
"""Build a read-only snapshot of observability connection health and incident links."""
sentry_health = load_provider_health("sentry", env)
glitchtip_health = load_provider_health("glitchtip", env)
providers = [sentry_health, glitchtip_health]
target_db = db or ControlPlaneDB()
raw_links = target_db.list_incident_links(limit=100)
links: list[CorrelatedIncidentLink] = []
sentry_cnt = 0
glitchtip_cnt = 0
for r in raw_links:
prov = (r.get("provider") or "sentry").lower()
if prov == "sentry":
sentry_cnt += 1
elif prov == "glitchtip":
glitchtip_cnt += 1
pr_nums = _parse_pr_numbers(r.get("linked_pr_numbers"))
links.append(
CorrelatedIncidentLink(
link_id=int(r.get("link_id", 0)),
provider=prov,
provider_base_url=r.get("provider_base_url") or "",
provider_org=r.get("provider_org") or "",
provider_project=r.get("provider_project") or "",
provider_issue_id=str(r.get("provider_issue_id") or ""),
provider_short_id=r.get("provider_short_id"),
provider_permalink=r.get("provider_permalink"),
fingerprint=r.get("fingerprint"),
gitea_org=r.get("gitea_org") or "Scaled-Tech-Consulting",
gitea_repo=r.get("gitea_repo") or "Gitea-Tools",
gitea_issue_number=int(r.get("gitea_issue_number", 0)),
linked_pr_numbers=pr_nums,
last_seen=r.get("last_seen"),
event_count=int(r.get("event_count", 1)),
created_at=r.get("created_at"),
updated_at=r.get("updated_at"),
)
)
bridge_active = any(p.bridge_enabled for p in providers)
return ObservabilitySnapshot(
schema_version=OBSERVABILITY_SCHEMA_VERSION,
providers=providers,
links=links,
total_links=len(links),
sentry_links_count=sentry_cnt,
glitchtip_links_count=glitchtip_cnt,
bridge_active=bridge_active,
)
def snapshot_to_dict(snapshot: ObservabilitySnapshot) -> dict[str, Any]:
return snapshot.to_dict()
+143
View File
@@ -0,0 +1,143 @@
"""HTML view renderer for the Sentry/GlitchTip observability console (#649, Phase 4).
Renders connection status widgets, error correlation links, and gated issue creation
affordances over the read-only observability snapshot.
"""
from __future__ import annotations
import html
from typing import Any
from webui.layout import render_page
from webui.observability_loader import ObservabilitySnapshot, snapshot_to_dict
def _badge(status: str) -> str:
st = (status or "").lower()
if st == "healthy":
return '<span class="badge badge-success">healthy</span>'
if st == "disabled":
return '<span class="badge badge-warning">disabled (dry-run)</span>'
if st in {"missing_token", "not_configured"}:
return f'<span class="badge badge-muted">{html.escape(st)}</span>'
return f'<span class="badge">{html.escape(st)}</span>'
def _provider_card(p: dict[str, Any]) -> str:
name = html.escape(str(p.get("provider", "provider")).upper())
base_url = html.escape(str(p.get("base_url", "")))
org = html.escape(str(p.get("org", "")))
proj = html.escape(str(p.get("project", "")))
status_badge = _badge(str(p.get("status", "")))
min_events = p.get("min_events_for_issue", 2)
lookback = html.escape(str(p.get("lookback", "24h")))
bridge_enabled = "yes" if p.get("bridge_enabled") else "no"
return f"""
<div class="card" style="margin-bottom: 1rem; padding: 1rem; border: 1px solid #ccc; border-radius: 6px;">
<div style="display: flex; justify-content: space-between; align-items: center;">
<h3 style="margin: 0;">{name} Connection</h3>
<div>{status_badge}</div>
</div>
<table style="width: 100%; margin-top: 0.5rem; border-collapse: collapse;">
<tr><td><strong>Base URL:</strong></td><td><code>{base_url}</code></td></tr>
<tr><td><strong>Scope:</strong></td><td><code>{org} / {proj}</code></td></tr>
<tr><td><strong>Bridge Enabled:</strong></td><td><code>{bridge_enabled}</code></td></tr>
<tr><td><strong>Min Events for Issue:</strong></td><td><code>{min_events}</code></td></tr>
<tr><td><strong>Lookback Window:</strong></td><td><code>{lookback}</code></td></tr>
</table>
</div>
"""
def render_observability_page(snapshot: ObservabilitySnapshot | dict[str, Any]) -> str:
"""Render the observability dashboard HTML page."""
data = snapshot.to_dict() if isinstance(snapshot, ObservabilitySnapshot) else dict(snapshot)
providers_raw = data.get("providers", [])
provider_cards = "".join(_provider_card(p) for p in providers_raw) if providers_raw else "<p>No providers configured.</p>"
links = data.get("links", [])
link_rows = []
for l in links:
prov = html.escape(str(l.get("provider", "")))
p_issue_id = html.escape(str(l.get("provider_issue_id", "")))
fingerprint = html.escape(str(l.get("fingerprint") or ""))
g_issue_num = int(l.get("gitea_issue_number", 0))
g_org = html.escape(str(l.get("gitea_org", "")))
g_repo = html.escape(str(l.get("gitea_repo", "")))
g_issue_link = f'<strong>#{g_issue_num}</strong> ({g_org}/{g_repo})'
event_cnt = int(l.get("event_count", 1))
last_seen = html.escape(str(l.get("last_seen") or ""))
short_id = html.escape(str(l.get("provider_short_id") or p_issue_id))
link_rows.append(f"""
<tr>
<td><code>{prov}</code></td>
<td><strong>{short_id}</strong><br><small style="color: #666;">id: {p_issue_id}</small></td>
<td><code>{fingerprint}</code></td>
<td>{g_issue_link}</td>
<td>{event_cnt}</td>
<td><small>{last_seen}</small></td>
</tr>
""")
table_body = "".join(link_rows) if link_rows else '<tr><td colspan="6" style="text-align: center; padding: 1.5rem; color: #666;">No correlated incident links stored. Bridge operates under dry-run default.</td></tr>'
metrics = data.get("metrics", {})
total_links = metrics.get("total_links", 0)
sentry_cnt = metrics.get("sentry_links_count", 0)
glitchtip_cnt = metrics.get("glitchtip_links_count", 0)
body_html = f"""
<h2>Observability & Incident Bridge (#649)</h2>
<p>Read-only console surface for Sentry/GlitchTip provider connections, error correlation,
and durable Gitea issue linkage.</p>
<div class="alert alert-info" style="background: #f0f4f8; padding: 1rem; border-left: 4px solid #0052cc; margin-bottom: 1.5rem;">
<strong>ADR Authority Model:</strong> Gitea records durable issue history. Control-plane DB coordinates incident links.
Sentry/GlitchTip observe errors. Raw monitoring incidents are <em>never</em> assignable control-plane work items.
Durable issue creation is gated and dry-runable via the <code>#612</code> bridge APIs.
</div>
<h3>Provider Connections</h3>
<div style="display: grid; grid-template-columns: repeat(auto-fit, minmax(300px, 1fr)); gap: 1rem; margin-bottom: 2rem;">
{provider_cards}
</div>
<div style="display: flex; justify-content: space-between; align-items: center; margin-bottom: 1rem;">
<h3 style="margin: 0;">Correlated Incidents ({total_links})</h3>
<div>
<span class="badge" style="margin-right: 0.5rem;">Sentry: {sentry_cnt}</span>
<span class="badge">GlitchTip: {glitchtip_cnt}</span>
</div>
</div>
<table class="table" style="width: 100%; border-collapse: collapse; border: 1px solid #ddd;">
<thead>
<tr style="background: #f9f9f9; text-align: left;">
<th style="padding: 0.5rem; border-bottom: 2px solid #ddd;">Provider</th>
<th style="padding: 0.5rem; border-bottom: 2px solid #ddd;">Incident ID</th>
<th style="padding: 0.5rem; border-bottom: 2px solid #ddd;">Fingerprint</th>
<th style="padding: 0.5rem; border-bottom: 2px solid #ddd;">Gitea Issue Link</th>
<th style="padding: 0.5rem; border-bottom: 2px solid #ddd;">Events</th>
<th style="padding: 0.5rem; border-bottom: 2px solid #ddd;">Last Seen</th>
</tr>
</thead>
<tbody>
{table_body}
</tbody>
</table>
<div style="margin-top: 2rem; padding: 1rem; background: #fafafa; border: 1px solid #eee; border-radius: 4px;">
<h4 style="margin-top: 0;">Reconcile & Link Controls (Gated)</h4>
<p style="margin-bottom: 0.5rem; color: #555;">
Create or reconcile durable Gitea issues from provider observations using the <code>#612</code> incident bridge:
</p>
<code>mcp call gitea_observability_reconcile_incident --provider sentry --apply false</code>
</div>
"""
return render_page(title="Observability", body_html=body_html)