diff --git a/control_plane_db.py b/control_plane_db.py index 422003c..4d159ea 100644 --- a/control_plane_db.py +++ b/control_plane_db.py @@ -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( diff --git a/docs/mcp-fleet-inventory.md b/docs/mcp-fleet-inventory.md new file mode 100644 index 0000000..44469c5 --- /dev/null +++ b/docs/mcp-fleet-inventory.md @@ -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. diff --git a/docs/mcp-tool-inventory.md b/docs/mcp-tool-inventory.md index 799d0d0..7d4b7a1 100644 --- a/docs/mcp-tool-inventory.md +++ b/docs/mcp-tool-inventory.md @@ -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` diff --git a/gitea_mcp_server.py b/gitea_mcp_server.py index 3721039..36b7b2e 100644 --- a/gitea_mcp_server.py +++ b/gitea_mcp_server.py @@ -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 @@ -18811,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. @@ -24344,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__": @@ -24353,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). diff --git a/mcp_fleet_inventory.py b/mcp_fleet_inventory.py new file mode 100644 index 0000000..a73da1c --- /dev/null +++ b/mcp_fleet_inventory.py @@ -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')}" diff --git a/mcp_tool_inventory.py b/mcp_tool_inventory.py index 68d7982..5e2cdd9 100644 --- a/mcp_tool_inventory.py +++ b/mcp_tool_inventory.py @@ -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", diff --git a/skills/llm-project-workflow/SKILL.md b/skills/llm-project-workflow/SKILL.md index 58762c0..cd2ba0f 100644 --- a/skills/llm-project-workflow/SKILL.md +++ b/skills/llm-project-workflow/SKILL.md @@ -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 diff --git a/task_capability_map.py b/task_capability_map.py index a30b7f9..7f344c1 100644 --- a/task_capability_map.py +++ b/task_capability_map.py @@ -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": { diff --git a/tests/test_control_plane_db.py b/tests/test_control_plane_db.py index 6065f38..8ac1c51 100644 --- a/tests/test_control_plane_db.py +++ b/tests/test_control_plane_db.py @@ -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: diff --git a/tests/test_control_plane_db_server_runtimes.py b/tests/test_control_plane_db_server_runtimes.py new file mode 100644 index 0000000..c408c50 --- /dev/null +++ b/tests/test_control_plane_db_server_runtimes.py @@ -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() diff --git a/tests/test_fleet_inventory_tool.py b/tests/test_fleet_inventory_tool.py new file mode 100644 index 0000000..21539d2 --- /dev/null +++ b/tests/test_fleet_inventory_tool.py @@ -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() diff --git a/tests/test_mcp_fleet_inventory.py b/tests/test_mcp_fleet_inventory.py new file mode 100644 index 0000000..3800f08 --- /dev/null +++ b/tests/test_mcp_fleet_inventory.py @@ -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()