From 4a1cc63e94ad9df20c5e0935234c914e54d84b11 Mon Sep 17 00:00:00 2001 From: Jason Walker <913443@dadeschools.net> Date: Thu, 30 Jul 2026 03:14:15 -0400 Subject: [PATCH] feat(controller): expose instance-level fleet identity and health snapshots Add a pure fleet snapshot assessor and a read-only controller/reconciler MCP tool so multi-instance fleets can be enumerated by client_instance_id without treating shared client_type as duplication. Trusted launcher instance IDs, classification (missing/unmanifested/collisions/historical), registry schema extensions, docs, and regression tests for #978. Closes #978 Co-Authored-By: Claude Opus 4.8 (1M context) --- docs/instance-fleet-identity.md | 171 +++ docs/mcp-tool-inventory.md | 1 + gitea_config.py | 3 + gitea_mcp_server.py | 223 +++- mcp_fleet_snapshot.py | 847 +++++++++++++ mcp_worker_identity.py | 99 +- task_capability_map.py | 12 + .../test_issue_978_instance_fleet_snapshot.py | 1054 +++++++++++++++++ 8 files changed, 2402 insertions(+), 8 deletions(-) create mode 100644 docs/instance-fleet-identity.md create mode 100644 mcp_fleet_snapshot.py create mode 100644 tests/test_issue_978_instance_fleet_snapshot.py diff --git a/docs/instance-fleet-identity.md b/docs/instance-fleet-identity.md new file mode 100644 index 0000000..d7bbd47 --- /dev/null +++ b/docs/instance-fleet-identity.md @@ -0,0 +1,171 @@ +# Instance-level fleet identity and health snapshots + +Issue **#978**. Companion primitives: **#948** (worker ownership), **#975** +(heartbeat lifecycle), **#951** (restart receipts). + +## Why this exists + +The multi-instance fleet gate (#963) needs a *native* answer to: + +* which application launches are live; +* which of the five namespace workers belong to each launch; +* whether two Codex (or Claude, Grok, …) launches are distinct; +* whether any live collision is real rather than a shared client type. + +Process tables, PID proximity, configuration files, and direct SQLite access +are **not** production evidence. The sanctioned surface is the read-only MCP +tool `gitea_snapshot_instance_fleet` on **controller** and **reconciler** +namespaces. + +## Identity hierarchy + +| Identity | Scope | Who generates it | Lifetime | +| --- | --- | --- | --- | +| `client_type` | Application family (`codex`, `claude_code`, `gemini`, `grok`, …) | Launcher sets `GITEA_MCP_CLIENT` | Stable for the product | +| `client_instance_id` | One running application launch | **Trusted launcher**, once per launch, as `GITEA_MCP_CLIENT_INSTANCE` | Fresh launch → new ID; reconnect of same launch → same ID; full app restart → new ID | +| `fleet_run_id` | Operator-approved enrollment / canary cohort | Operator / controller sets `GITEA_MCP_FLEET_RUN_ID` | Duration of the approved rollout | +| `namespace` | `author` \| `reviewer` \| `merger` \| `controller` \| `reconciler` | Profile / MCP server binding | Process lifetime | +| `worker_id` / `worker_identity` | One namespace worker process | Runtime registry at first registration | New on worker restart; not reused | +| `session_id` | Client session ownership | Launcher `GITEA_MCP_CLIENT_SESSION` or runtime | Session lifetime | +| `generation_id` | One daemon launch | Runtime at process boot | New on process restart | +| `process_identity` / PID | OS process | Runtime | Process lifetime | + +### Rules + +1. Multiple active instances **may** share the same `client_type`. +2. Every application launch receives a **distinct** `client_instance_id`. +3. All five namespace workers of one launch report the **same** + `client_instance_id`. +4. Each namespace worker has a **distinct** `worker_identity`, process + identity, generation, and PID. +5. Instance identity is **never** inferred from PID proximity, timestamps, or + client type alone. +6. Live reuse of one `client_instance_id` with conflicting workers fails closed. +7. Sharing only a profile or `client_type` is **not** a duplicate. + +This deliberately **replaces** any permanent `exactly_one_per_profile` fleet +model (#949 assumption) as the operating rule for multi-instance fleets. + +## How five workers join one instance + +1. The host starts one application instance (for example one Codex session). +2. The trusted launcher mints `client_instance_id` (see + `mcp_fleet_snapshot.generate_client_instance_id`) and injects: + + ```text + GITEA_MCP_CLIENT= + GITEA_MCP_CLIENT_INSTANCE= + GITEA_MCP_CLIENT_SESSION= # optional but recommended + GITEA_MCP_FLEET_RUN_ID= # when on an approved canary + ``` + +3. The launcher starts the five MCP namespace processes (or attaches five + role profiles) with **that same environment**. +4. Each worker registers once into the worker registry with its own + `worker_identity`, `namespace`, `generation_id`, and PID, but the shared + `client_instance_id`. + +Worker reconnect (same process, same registration) keeps the instance ID. +Worker restart (new process) mints a new worker identity and generation but +must still receive the same `GITEA_MCP_CLIENT_INSTANCE` from the parent +application if it is the same launch. Full application restart mints a new +`client_instance_id`. + +## Approved fleet manifest + +An operator (or controller enrollment step) obtains an approved manifest as a +list of expected instances, for example: + +```json +{ + "instances": [ + { + "client_type": "codex", + "client_instance_id": "inst-codex-20260730T120000Z-abc123def456", + "fleet_run_id": "canary-963-2026-07-30", + "namespaces": ["author", "reviewer", "merger", "controller", "reconciler"] + }, + { + "client_type": "codex", + "client_instance_id": "inst-codex-20260730T120100Z-fed654cba321", + "fleet_run_id": "canary-963-2026-07-30", + "namespaces": ["author", "reviewer", "merger", "controller", "reconciler"] + } + ] +} +``` + +Pass that list as `expected_manifest` to `gitea_snapshot_instance_fleet`. +When a manifest is supplied: + +* instances on the manifest but not live → `missing_expected`; +* live instances not on the manifest → `unmanifested`; +* both are **active blockers** for fleet-gate safety. + +Without a manifest, the snapshot still enumerates the live fleet and classifies +identity collisions; it does not invent enrollment policy. + +## Snapshot consistency + +Each snapshot includes: + +* `snapshot_at` — UTC timestamp; +* `consistency_token` / `registry_revision` — content digest over worker + identity, instance id, status, heartbeat, fencing, and generation. + +Two successive snapshots can prove heartbeat continuity via +`mcp_fleet_snapshot.compare_snapshot_heartbeats`. + +## Classification (active vs historical) + +**Active blockers** (make `live_fleet_safe=false`): + +* missing expected instance; +* unmanifested extra instance; +* duplicate namespace worker within one instance; +* live `client_instance_id` collision; +* reused worker / session / generation / process / PID / fencing identity; +* orphaned or unowned workers; +* unknown client; +* foreign-repository workers; +* old-revision workers; +* stale workers still marked active; +* legacy incomplete instance identity. + +**Historical dead rows** (`status` released/superseded) are reported as +`historical_dead` findings with `active_blocker=false`. They never +automatically make the live fleet unsafe. + +## Fail-closed behaviour and recovery + +| Condition | Mutation safety | Diagnostic reads | Recovery | +| --- | --- | --- | --- | +| Two instances, same `client_type`, distinct IDs | Safe (if otherwise healthy) | Available | None needed | +| Live reuse of one `client_instance_id` | Unsafe | Available | Stop the colliding launch or re-issue a distinct ID | +| Two workers, same namespace, one instance | Unsafe | Available | Stop the extra worker | +| Legacy registration without trusted instance ID | Unsafe for fleet mutation | Available | Relaunch with launcher-issued `GITEA_MCP_CLIENT_INSTANCE` | +| Historical dead row only | Does not block alone | Available | No action required for fleet safety | +| Registry unavailable | Snapshot fails closed | N/A | Repair registry path / reconnect namespaces | + +Incomplete or untrusted instance identity **cannot** authorize unsafe +mutation. Diagnostic reads remain available where `gitea.read` allows. + +## Permissions + +* **Allowed:** `controller`, `reconciler` with `gitea.read`. +* **Denied:** author, reviewer, merger (even with `gitea.read` for other tools). +* **No new mutation permissions** are granted to any role. + +## Related surfaces + +* `mcp_worker_identity` — registry, worker identity, heartbeats (#948, #975). +* `mcp_fleet_snapshot` — pure snapshot + classification (#978). +* `gitea_snapshot_instance_fleet` — sanctioned MCP tool (#978). +* `gitea_get_runtime_context` — single-process view (not fleet-wide). + +## Non-goals + +* Starting the #963 canary. +* Purging historical registry rows. +* Preserving “exactly one process per profile” as the permanent model. +* Using shell / process-table / SQLite inspection as production fleet evidence. diff --git a/docs/mcp-tool-inventory.md b/docs/mcp-tool-inventory.md index ae61faa..7264f6b 100644 --- a/docs/mcp-tool-inventory.md +++ b/docs/mcp-tool-inventory.md @@ -159,6 +159,7 @@ that gates each call, not which tools exist. - `gitea_sentry_reconcile_issue` - `gitea_sentry_watchdog` - `gitea_set_issue_labels` +- `gitea_snapshot_instance_fleet` - `gitea_submit_pr_review` - `gitea_update_pr_branch_by_merge` - `gitea_validate_review_final_report` diff --git a/gitea_config.py b/gitea_config.py index f99362d..c220012 100644 --- a/gitea_config.py +++ b/gitea_config.py @@ -1203,6 +1203,9 @@ RECOGNIZED_GITEA_ENV_KEYS = frozenset({ "GITEA_MCP_CLIENT", "GITEA_MCP_CLIENT_INSTANCE", "GITEA_MCP_CLIENT_SESSION", + # #978: operator-approved fleet enrollment identity shared by one launch. + "GITEA_MCP_FLEET_RUN_ID", + "GITEA_MCP_PROCESS_IDENTITY", # #975 review 652 B1: production also consumes HEARTBEAT_INTERVAL_ENV from # mcp_worker_identity via gitea_mcp_server._start_worker_heartbeat. Omitting # it reproduced the same unsupported-env → runtime_reconnect_required diff --git a/gitea_mcp_server.py b/gitea_mcp_server.py index 9f864d9..698c5ce 100644 --- a/gitea_mcp_server.py +++ b/gitea_mcp_server.py @@ -15733,6 +15733,7 @@ _WORKER_HEARTBEAT_SUPERVISOR = None CLIENT_NAME_ENV = "GITEA_MCP_CLIENT" CLIENT_INSTANCE_ENV = "GITEA_MCP_CLIENT_INSTANCE" CLIENT_SESSION_ENV = "GITEA_MCP_CLIENT_SESSION" +FLEET_RUN_ENV = "GITEA_MCP_FLEET_RUN_ID" def _worker_registry(): @@ -15756,13 +15757,34 @@ def _worker_registry(): def _client_identity_hints() -> dict: - """What the launcher told us about itself. Unset fields stay unset.""" + """What the launcher told us about itself (#948 / #975 / #978). + + ``client_instance_id`` is established by the trusted application launcher + once per launch and shared by all five namespace workers. When the launcher + key is absent we still register under an explicit *legacy* placeholder so + diagnostic reads work, but the registration is marked incomplete and cannot + authorize multi-instance fleet mutation safety. + """ + import mcp_fleet_snapshot + + client_name = (os.environ.get(CLIENT_NAME_ENV) or "").strip() or None + raw_instance = (os.environ.get(CLIENT_INSTANCE_ENV) or "").strip() or None + instance = mcp_fleet_snapshot.assess_instance_identity(raw_instance) + # Legacy placeholder only when the launcher omitted the key — never invent a + # trusted ID from PID proximity. + if not instance["client_instance_id"]: + placeholder = f"legacy-pid-{os.getpid()}" + instance = mcp_fleet_snapshot.assess_instance_identity( + placeholder, source="pid_fallback" + ) return { - "client_name": (os.environ.get(CLIENT_NAME_ENV) or "").strip() or None, - "client_instance_id": (os.environ.get(CLIENT_INSTANCE_ENV) or "").strip() - or f"pid-{os.getpid()}", + "client_name": client_name, + "client_instance_id": instance["client_instance_id"], + "instance_id_provenance": instance["provenance"], + "instance_identity_trusted": instance["trusted"], "session_id": (os.environ.get(CLIENT_SESSION_ENV) or "").strip() or f"proc-{os.getpid()}-{_process_boot_head_sha or 'nohead'}", + "fleet_run_id": (os.environ.get(FLEET_RUN_ENV) or "").strip() or None, } @@ -15794,20 +15816,43 @@ def _active_worker_identity() -> str | None: identity = mcp_worker_identity.generate_worker_identity( hints["client_name"], hints["session_id"] ) + role = _active_role_kind_safe() + profile_name = ( + (os.environ.get(gitea_config.ENV_PROFILE) or "").strip() or None + ) + # Namespace is the role namespace (author/reviewer/…) when known. + namespace = role if role in { + "author", "reviewer", "merger", "controller", "reconciler" + } else None + parity = None + try: + parity = _current_master_parity() + except Exception: + parity = None outcome = registry.register( worker_identity=identity, client_name=hints["client_name"], client_instance_id=hints["client_instance_id"], session_id=hints["session_id"], generation_id=generation, - role=_active_role_kind_safe(), - profile=(os.environ.get(gitea_config.ENV_PROFILE) or "").strip() or None, + role=role, + profile=profile_name, + namespace=namespace, remote=(os.environ.get("GITEA_MCP_REMOTE") or "").strip() or None, repository_binding=PROJECT_ROOT, pid=os.getpid(), transport=native.get("bound_transport"), token_fingerprint=native.get("token_fingerprint"), pid_alive_probe=issue_lock_store.is_process_alive, + fleet_run_id=hints.get("fleet_run_id"), + process_identity=f"pid-{os.getpid()}", + startup_revision=(parity or {}).get("startup_head") + or _process_boot_head_sha, + loaded_revision=(parity or {}).get("local_head") + or (parity or {}).get("current_head"), + parity_revision=(parity or {}).get("current_head"), + live_revision=(parity or {}).get("live_remote_head"), + instance_id_provenance=hints.get("instance_id_provenance"), ) if outcome.get("registered"): _WORKER_IDENTITY = identity @@ -19373,6 +19418,172 @@ def gitea_validate_review_final_report( ) +@mcp.tool() +def gitea_snapshot_instance_fleet( + remote: str = "dadeschools", + host: str | None = None, + org: str | None = None, + repo: str | None = None, + expected_manifest: list | dict | str | None = None, + include_historical: bool = True, + canonical_repository: str | None = None, + expected_live_revision: str | None = None, +) -> dict: + """Read-only: authoritative instance-level fleet identity and health snapshot (#978). + + Controller and reconciler only. Returns a point-in-time snapshot of every + registered namespace worker with instance attribution, heartbeat freshness, + and structured classification (expected/missing/unmanifested, duplicate + namespace workers, identity collisions, foreign repository, old revision, + historical dead rows). Sharing only a client type is never a duplicate. + + Does not grant any mutation capability. Does not scan process tables or + open foreign databases for production evidence — the worker registry is + the sole authority. + + Args: + remote: Known instance — 'dadeschools' or 'prgs'. + host: Optional host override. + org: Optional org override (audit context only). + repo: Optional repo override (audit context only). + expected_manifest: Optional approved fleet enrollment list + (list of dicts with client_instance_id / client_type / fleet_run_id). + When supplied, missing and unmanifested instances are classified. + include_historical: Include released/superseded rows (default True); + historical rows never make the live fleet unsafe by themselves. + canonical_repository: Expected repository binding path/slug for + foreign-repository classification. + expected_live_revision: When set, workers whose recorded revisions + differ are classified as old-revision. + """ + import json as _json + + import mcp_fleet_snapshot + + read_block = _profile_operation_gate("gitea.read") + if read_block: + return { + "success": False, + "read_only": True, + "reasons": read_block, + "permission_report": _permission_block_report("gitea.read"), + } + + profile = get_profile() + role = _profile_role_kind(profile) + if role not in {"controller", "reconciler"}: + return { + "success": False, + "read_only": True, + "allowed": False, + "denied_role": role, + "required_roles": ["controller", "reconciler"], + "mutation_performed": False, + "reasons": [ + f"gitea_snapshot_instance_fleet is restricted to controller and " + f"reconciler roles; active role_kind is {role!r}. Author, " + "reviewer, and merger profiles keep gitea.read for diagnosis " + "elsewhere but do not receive this fleet surface, and no " + "unrelated mutation permission is granted." + ], + "exact_next_action": ( + "Re-run from a prgs-controller or prgs-reconciler namespace." + ), + } + + manifest: list | None = None + if expected_manifest is not None: + raw = expected_manifest + if isinstance(raw, str): + try: + raw = _json.loads(raw) + except Exception as exc: + return { + "success": False, + "read_only": True, + "reasons": [ + f"expected_manifest is not valid JSON: {type(exc).__name__}" + ], + } + if isinstance(raw, dict): + # Accept {"instances": [...]} or a single instance object. + if "instances" in raw and isinstance(raw["instances"], list): + raw = raw["instances"] + else: + raw = [raw] + if not isinstance(raw, list): + return { + "success": False, + "read_only": True, + "reasons": ["expected_manifest must be a list of instance objects"], + } + manifest = raw + + registry = _worker_registry() + if registry is None: + return { + "success": False, + "read_only": True, + "reasons": [ + "worker registry is unavailable; cannot produce an authoritative " + "fleet snapshot (fail closed)" + ], + "exact_next_action": ( + "Ensure GITEA_WORKER_REGISTRY_DB is writable and re-run after " + "workers have registered." + ), + } + + try: + if include_historical: + records = registry.list_workers(status=None) + else: + records = registry.list_workers(status=mcp_worker_identity.STATUS_ACTIVE) + except Exception as exc: + return { + "success": False, + "read_only": True, + "reasons": [ + f"failed to read worker registry: {type(exc).__name__}: {_redact(str(exc))}" + ], + } + + parity = None + try: + parity = _current_master_parity() + except Exception: + parity = None + live_rev = expected_live_revision or (parity or {}).get("live_remote_head") + canon = canonical_repository or PROJECT_ROOT + + snapshot = mcp_fleet_snapshot.snapshot_instance_fleet( + records, + expected_manifest=manifest, + pid_alive_probe=issue_lock_store.is_process_alive, + canonical_repository=canon, + expected_live_revision=live_rev, + ) + snapshot["role_kind"] = role + snapshot["profile"] = profile.get("profile_name") + snapshot["remote"] = _effective_remote(remote) + snapshot["repository"] = { + "org": org, + "repo": repo, + "canonical_repository": canon, + } + snapshot["mutation_performed"] = False + snapshot["permission_scope"] = { + "read_only": True, + "granted_operations": ["gitea.read"], + "denied_unrelated_mutations": True, + "note": ( + "This capability is strictly observational. It does not authorize " + "branch, issue, PR, review, merge, or restart mutations." + ), + } + return snapshot + + @mcp.tool() def gitea_get_runtime_context( remote: str = "dadeschools", diff --git a/mcp_fleet_snapshot.py b/mcp_fleet_snapshot.py new file mode 100644 index 0000000..a5b6bb8 --- /dev/null +++ b/mcp_fleet_snapshot.py @@ -0,0 +1,847 @@ +"""Instance-level fleet identity and health snapshots (#978). + +#948 established per-worker ownership; #975 made heartbeats keep those rows +live. Neither surface could enumerate the fleet at *instance* granularity: +which application launch owns which five namespace workers, whether two +Codex launches are distinct, or whether a live collision is real rather than +a shared client type. + +This module is pure. Callers supply registry rows (and optional enrichments); +nothing here opens SQLite, scans process tables, or mutates state. Production +evidence for the fleet gate is the snapshot returned by the sanctioned +controller/reconciler tool that wraps this assessor. + +Identity hierarchy (highest → lowest): + +* ``client_type`` — application family (``codex``, ``claude_code``, …) +* ``client_instance_id`` — one running application launch (trusted launcher) +* ``fleet_run_id`` — operator-approved enrollment / canary cohort +* ``worker_id`` / ``worker_identity`` — one namespace worker process +* ``namespace`` — author | reviewer | merger | controller | reconciler + +Multiple simultaneous instances of the same ``client_type`` are first-class. +Sharing only a profile or client type is never a duplicate. +""" + +from __future__ import annotations + +import hashlib +import secrets +from collections import defaultdict +from datetime import datetime, timezone +from typing import Any, Callable, Iterable, Mapping + +import mcp_worker_identity as mwi + +# --- Classification labels ------------------------------------------------ + +CLASS_EXPECTED = "expected_enrolled" +CLASS_MISSING = "missing_expected" +CLASS_UNMANIFESTED = "unmanifested" +CLASS_DUPLICATE_NAMESPACE = "duplicate_namespace_worker" +CLASS_INSTANCE_ID_COLLISION = "instance_id_collision" +CLASS_WORKER_ID_COLLISION = "worker_identity_collision" +CLASS_SESSION_COLLISION = "session_identity_collision" +CLASS_GENERATION_COLLISION = "generation_identity_collision" +CLASS_PROCESS_COLLISION = "process_identity_collision" +CLASS_PID_COLLISION = "pid_collision" +CLASS_OWNERSHIP_COLLISION = "ownership_fencing_collision" +CLASS_ORPHANED = "orphaned_unowned" +CLASS_UNKNOWN_CLIENT = "unknown_client" +CLASS_FOREIGN_REPOSITORY = "foreign_repository" +CLASS_OLD_REVISION = "old_revision" +CLASS_STALE_WORKER = "stale_orphaned_worker" +CLASS_LEGACY_INCOMPLETE = "legacy_incomplete_identity" +CLASS_HISTORICAL = "historical_dead" +CLASS_HEALTHY = "healthy" + +#: Active blockers that make the live fleet unsafe for mutation-gated work. +ACTIVE_BLOCKER_CLASSES = frozenset( + { + CLASS_MISSING, + CLASS_UNMANIFESTED, + CLASS_DUPLICATE_NAMESPACE, + CLASS_INSTANCE_ID_COLLISION, + CLASS_WORKER_ID_COLLISION, + CLASS_SESSION_COLLISION, + CLASS_GENERATION_COLLISION, + CLASS_PROCESS_COLLISION, + CLASS_PID_COLLISION, + CLASS_OWNERSHIP_COLLISION, + CLASS_ORPHANED, + CLASS_UNKNOWN_CLIENT, + CLASS_FOREIGN_REPOSITORY, + CLASS_OLD_REVISION, + CLASS_STALE_WORKER, + CLASS_LEGACY_INCOMPLETE, + } +) + +SANCTIONED_NAMESPACES = frozenset( + {"author", "reviewer", "merger", "controller", "reconciler"} +) + +INSTANCE_ID_PROVENANCE_TRUSTED = "trusted_launcher" +INSTANCE_ID_PROVENANCE_LEGACY = "legacy_incomplete" +INSTANCE_ID_PROVENANCE_MISSING = "missing" + +CLIENT_INSTANCE_ENV = "GITEA_MCP_CLIENT_INSTANCE" +FLEET_RUN_ENV = "GITEA_MCP_FLEET_RUN_ID" +PROCESS_IDENTITY_ENV = "GITEA_MCP_PROCESS_IDENTITY" + +_LEGACY_INSTANCE_PREFIXES = ("pid-", "proc-", "legacy-") + + +def _utc_now() -> datetime: + return datetime.now(timezone.utc) + + +def _ts(value: datetime) -> str: + return value.astimezone(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") + + +def generate_client_instance_id( + client_type: str | None, + *, + launch_nonce: str | None = None, + now: datetime | None = None, +) -> str: + """Mint a distinct instance ID for one application launch (#978). + + The trusted launcher (or host that starts all five namespace workers) + generates this once per launch and injects it as ``GITEA_MCP_CLIENT_INSTANCE`` + into every worker environment. Workers never invent their own instance ID + from PID proximity or timestamps. + """ + client = mwi.normalize_client_name(client_type) + stamp = (now or _utc_now()).astimezone(timezone.utc).strftime( + mwi.IDENTITY_TIMESTAMP_FORMAT + ) + nonce = launch_nonce if launch_nonce is not None else secrets.token_hex(16) + digest = hashlib.sha256( + f"{client}\x1f{stamp}\x1f{nonce}".encode("utf-8") + ).hexdigest()[:12] + return f"inst-{client}-{stamp}-{digest}" + + +def assess_instance_identity( + raw_instance_id: str | None, + *, + source: str | None = None, +) -> dict[str, Any]: + """Classify whether a client_instance_id is trusted enough for mutation. + + Trusted instance IDs are non-empty and not the pre-#978 PID/proc fallbacks. + Incomplete identities remain visible for diagnosis but cannot authorize + unsafe mutation. + """ + text = (raw_instance_id or "").strip() + if not text: + return { + "client_instance_id": None, + "complete": False, + "trusted": False, + "provenance": INSTANCE_ID_PROVENANCE_MISSING, + "reasons": [ + "client_instance_id is missing; the trusted launcher must set " + f"{CLIENT_INSTANCE_ENV} once per application launch" + ], + } + lowered = text.lower() + if lowered.startswith(_LEGACY_INSTANCE_PREFIXES) or source == "pid_fallback": + return { + "client_instance_id": text, + "complete": False, + "trusted": False, + "provenance": INSTANCE_ID_PROVENANCE_LEGACY, + "reasons": [ + f"client_instance_id {text!r} is a legacy PID/process fallback, " + "not a trusted launcher-issued instance identity" + ], + } + return { + "client_instance_id": text, + "complete": True, + "trusted": True, + "provenance": INSTANCE_ID_PROVENANCE_TRUSTED, + "reasons": [], + } + + +def resolve_client_instance_from_env( + env: Mapping[str, str] | None = None, + *, + pid: int | None = None, +) -> dict[str, Any]: + """Resolve instance identity from launcher env without inventing one. + + When the trusted key is absent, return incomplete evidence rather than a + silent ``pid-`` identity. Callers that still need a non-empty registry + key may choose a legacy placeholder deliberately; they must not treat it as + trusted. + """ + source = dict(env or {}) + raw = (source.get(CLIENT_INSTANCE_ENV) or "").strip() + assessment = assess_instance_identity(raw or None) + assessment["fleet_run_id"] = (source.get(FLEET_RUN_ENV) or "").strip() or None + assessment["process_identity"] = ( + (source.get(PROCESS_IDENTITY_ENV) or "").strip() + or (f"pid-{pid}" if pid is not None else None) + ) + return assessment + + +def _public_worker(record: Mapping[str, Any]) -> dict[str, Any]: + return { + "worker_id": record.get("worker_identity") or record.get("worker_id"), + "worker_identity": record.get("worker_identity") or record.get("worker_id"), + "client_type": mwi.normalize_client_name( + record.get("client_name") or record.get("client_type") + ), + "client_instance_id": record.get("client_instance_id"), + "fleet_run_id": record.get("fleet_run_id"), + "namespace": record.get("namespace"), + "profile": record.get("profile"), + "declared_role": record.get("role") or record.get("declared_role"), + "authenticated_account": record.get("authenticated_account"), + "session_id": record.get("session_id"), + "generation_id": record.get("generation_id"), + "process_identity": record.get("process_identity") + or ( + f"pid-{record['pid']}" + if record.get("pid") is not None + else None + ), + "pid": record.get("pid"), + "repository_binding": record.get("repository_binding"), + "remote": record.get("remote"), + "startup_revision": record.get("startup_revision"), + "loaded_revision": record.get("loaded_revision"), + "parity_revision": record.get("parity_revision"), + "live_revision": record.get("live_revision"), + "runtime_provenance": record.get("runtime_provenance") + or record.get("transport"), + "transport": record.get("transport"), + "started_at": record.get("started_at"), + "last_heartbeat_at": record.get("last_heartbeat_at"), + "heartbeat_ttl_seconds": record.get("heartbeat_ttl_seconds"), + "fencing_epoch": record.get("fencing_epoch"), + "status": record.get("status"), + "instance_id_provenance": record.get("instance_id_provenance"), + } + + +def _liveness( + record: Mapping[str, Any], + *, + now: datetime | None, + pid_alive_probe: Callable[[int | None], bool | None] | None, +) -> dict[str, Any]: + pid_alive = None + if pid_alive_probe is not None and record.get("pid") is not None: + try: + pid_alive = pid_alive_probe(record.get("pid")) + except Exception: + pid_alive = None + return mwi.WorkerRegistry.is_live(record, now=now, pid_alive=pid_alive) + + +def _consistency_token(rows: Iterable[Mapping[str, Any]], snapshot_at: str) -> str: + material = [snapshot_at] + for row in sorted( + rows, + key=lambda r: ( + str(r.get("worker_identity") or ""), + str(r.get("last_heartbeat_at") or ""), + str(r.get("fencing_epoch") or ""), + ), + ): + material.append( + "|".join( + [ + str(row.get("worker_identity") or ""), + str(row.get("client_instance_id") or ""), + str(row.get("status") or ""), + str(row.get("last_heartbeat_at") or ""), + str(row.get("fencing_epoch") or ""), + str(row.get("generation_id") or ""), + ] + ) + ) + digest = hashlib.sha256("\n".join(material).encode("utf-8")).hexdigest()[:16] + return f"fleetrev-{digest}" + + +def build_worker_snapshot_row( + record: Mapping[str, Any], + *, + now: datetime | None = None, + pid_alive_probe: Callable[[int | None], bool | None] | None = None, + canonical_repository: str | None = None, + expected_live_revision: str | None = None, + heartbeat_supervised: bool | None = None, +) -> dict[str, Any]: + """One point-in-time worker row for the fleet snapshot.""" + stamp = now or _utc_now() + base = _public_worker(record) + identity = assess_instance_identity( + base.get("client_instance_id"), + source=record.get("instance_id_source"), + ) + liveness = _liveness(record, now=stamp, pid_alive_probe=pid_alive_probe) + is_historical = str(record.get("status") or "") != mwi.STATUS_ACTIVE + live = bool(liveness.get("live")) and not is_historical + + repo = (base.get("repository_binding") or "").strip() or None + foreign_repo = bool( + canonical_repository + and repo + and repo.rstrip("/") != str(canonical_repository).rstrip("/") + ) + old_revision = False + if expected_live_revision: + for key in ("startup_revision", "loaded_revision", "parity_revision", "live_revision"): + rev = (base.get(key) or "").strip() + if rev and rev != expected_live_revision: + old_revision = True + break + + ownership_state = "historical" if is_historical else ( + "live" if live else "stale" + ) + if live and not identity["trusted"]: + ownership_state = "live_untrusted_identity" + if live and not base.get("session_id"): + ownership_state = "orphaned" + + mutation_safe = bool( + live + and identity["trusted"] + and not foreign_repo + and not old_revision + and ownership_state == "live" + and base.get("client_type") != mwi.UNKNOWN_CLIENT + ) + + restart_required = bool( + old_revision + or (live and not liveness.get("heartbeat_fresh", True)) + ) + + return { + **base, + "instance_identity": identity, + "client_instance_id": identity["client_instance_id"] or base.get("client_instance_id"), + "instance_id_provenance": identity["provenance"], + "instance_identity_trusted": identity["trusted"], + "live": live, + "historical": is_historical, + "liveness": liveness, + "heartbeat": { + "registered": bool(base.get("last_heartbeat_at")), + "supervised": heartbeat_supervised, + "age_seconds": liveness.get("heartbeat_age_seconds"), + "ttl_seconds": liveness.get("heartbeat_ttl_seconds"), + "fresh": liveness.get("heartbeat_fresh"), + "last_heartbeat_at": base.get("last_heartbeat_at"), + }, + "fencing": { + "fencing_epoch": base.get("fencing_epoch"), + "generation_id": base.get("generation_id"), + }, + "ownership_state": ownership_state, + "foreign_repository": foreign_repo, + "old_revision": old_revision, + "stale": not live and not is_historical, + "restart_required": restart_required, + "mutation_safe": mutation_safe, + "conflicting_live_sessions": [], + } + + +def _collision_groups( + live_rows: list[dict[str, Any]], + key_fn, +) -> dict[str, list[dict[str, Any]]]: + groups: dict[str, list[dict[str, Any]]] = defaultdict(list) + for row in live_rows: + key = key_fn(row) + if key is None or key == "" or key == "None": + continue + groups[str(key)].append(row) + return {k: v for k, v in groups.items() if len(v) > 1} + + +def snapshot_instance_fleet( + records: Iterable[Mapping[str, Any]], + *, + expected_manifest: list[Mapping[str, Any]] | None = None, + now: datetime | None = None, + pid_alive_probe: Callable[[int | None], bool | None] | None = None, + canonical_repository: str | None = None, + expected_live_revision: str | None = None, + registry_revision: str | None = None, + known_client_types: Iterable[str] | None = None, +) -> dict[str, Any]: + """Authoritative point-in-time fleet snapshot with classification (#978). + + Historical dead rows are reported separately and never automatically make + the live fleet unsafe. + """ + stamp = now or _utc_now() + snapshot_at = _ts(stamp) + known = { + mwi.normalize_client_name(c) + for c in (known_client_types or mwi.CLIENT_ALIASES.values()) + } + known.discard(mwi.UNKNOWN_CLIENT) + + all_rows: list[dict[str, Any]] = [] + for record in records: + all_rows.append( + build_worker_snapshot_row( + record, + now=stamp, + pid_alive_probe=pid_alive_probe, + canonical_repository=canonical_repository, + expected_live_revision=expected_live_revision, + ) + ) + + live_rows = [r for r in all_rows if r["live"]] + historical_rows = [r for r in all_rows if r["historical"]] + stale_rows = [r for r in all_rows if r["stale"]] + + # --- identity collisions among live workers --- + findings: list[dict[str, Any]] = [] + + def _finding( + classification: str, + *, + severity: str, + workers: list[dict[str, Any]] | None = None, + instance_ids: list[str] | None = None, + detail: str, + active_blocker: bool, + ) -> None: + findings.append( + { + "classification": classification, + "severity": severity, + "active_blocker": active_blocker, + "detail": detail, + "client_instance_ids": instance_ids or sorted( + { + str(w.get("client_instance_id")) + for w in (workers or []) + if w.get("client_instance_id") + } + ), + "worker_identities": [ + w.get("worker_identity") for w in (workers or []) + ], + } + ) + + # Duplicate worker identity (should not happen with PK, still detect) + for wid, group in _collision_groups( + live_rows, lambda r: r.get("worker_identity") + ).items(): + _finding( + CLASS_WORKER_ID_COLLISION, + severity="blocker", + workers=group, + detail=f"worker identity {wid!r} is claimed by {len(group)} live workers", + active_blocker=True, + ) + + # Reused session identity across live workers + for sid, group in _collision_groups(live_rows, lambda r: r.get("session_id")).items(): + # Same session may appear once; collision only when multiple workers share it + # across different worker identities (always true for group size > 1). + _finding( + CLASS_SESSION_COLLISION, + severity="blocker", + workers=group, + detail=f"session identity {sid!r} is reused by {len(group)} live workers", + active_blocker=True, + ) + + # Generation claimed by multiple live sessions/workers is a conflict when + # the workers are not the five sanctioned namespaces of one instance. + for gen, group in _collision_groups( + live_rows, lambda r: r.get("generation_id") + ).items(): + namespaces = {g.get("namespace") for g in group if g.get("namespace")} + instance_ids = {g.get("client_instance_id") for g in group} + # Multiple workers under one generation is only valid if they share one + # instance and distinct namespaces. Same generation + same namespace = bad. + by_ns: dict[str, list] = defaultdict(list) + for g in group: + by_ns[str(g.get("namespace") or "")].append(g) + ns_dups = {ns: rows for ns, rows in by_ns.items() if ns and len(rows) > 1} + if ns_dups or len(instance_ids) > 1: + _finding( + CLASS_GENERATION_COLLISION, + severity="blocker", + workers=group, + detail=( + f"generation {gen!r} is contested across namespaces/instances " + f"(namespaces={sorted(namespaces)}, " + f"instances={sorted(str(i) for i in instance_ids if i)})" + ), + active_blocker=True, + ) + + # Process identity / PID collisions across distinct workers + for proc, group in _collision_groups( + live_rows, lambda r: r.get("process_identity") + ).items(): + if len({r.get("worker_identity") for r in group}) > 1: + _finding( + CLASS_PROCESS_COLLISION, + severity="blocker", + workers=group, + detail=f"process identity {proc!r} is shared by distinct live workers", + active_blocker=True, + ) + + for pid, group in _collision_groups(live_rows, lambda r: r.get("pid")).items(): + if len({r.get("worker_identity") for r in group}) > 1: + _finding( + CLASS_PID_COLLISION, + severity="blocker", + workers=group, + detail=f"PID {pid} is shared by distinct live workers", + active_blocker=True, + ) + + # Fencing/ownership: same fencing epoch on different workers of different instances + for epoch, group in _collision_groups( + live_rows, + lambda r: ( + f"{r.get('generation_id')}:{r.get('fencing_epoch')}" + if r.get("generation_id") is not None and r.get("fencing_epoch") is not None + else None + ), + ).items(): + if len({r.get("client_instance_id") for r in group}) > 1: + _finding( + CLASS_OWNERSHIP_COLLISION, + severity="blocker", + workers=group, + detail=( + f"fencing token {epoch!r} spans more than one client_instance_id" + ), + active_blocker=True, + ) + + # Per-instance grouping + by_instance: dict[str, list[dict[str, Any]]] = defaultdict(list) + unkeyed_live: list[dict[str, Any]] = [] + for row in live_rows: + iid = row.get("client_instance_id") + if not iid: + unkeyed_live.append(row) + continue + by_instance[str(iid)].append(row) + + instances: list[dict[str, Any]] = [] + for iid, workers in sorted(by_instance.items()): + client_types = sorted({w.get("client_type") for w in workers if w.get("client_type")}) + trusted = all(w.get("instance_identity_trusted") for w in workers) + namespaces = [w.get("namespace") for w in workers] + ns_counts: dict[str, int] = defaultdict(int) + for ns in namespaces: + if ns: + ns_counts[str(ns)] += 1 + dup_ns = sorted(ns for ns, n in ns_counts.items() if n > 1) + if dup_ns: + _finding( + CLASS_DUPLICATE_NAMESPACE, + severity="blocker", + workers=[w for w in workers if w.get("namespace") in dup_ns], + instance_ids=[iid], + detail=( + f"instance {iid!r} has more than one live worker for " + f"namespace(s) {dup_ns}" + ), + active_blocker=True, + ) + + # Live reuse of one instance ID with incompatible client types + if len(client_types) > 1: + _finding( + CLASS_INSTANCE_ID_COLLISION, + severity="blocker", + workers=workers, + instance_ids=[iid], + detail=( + f"client_instance_id {iid!r} is live under multiple client " + f"types {client_types}" + ), + active_blocker=True, + ) + + if not trusted: + _finding( + CLASS_LEGACY_INCOMPLETE, + severity="blocker", + workers=workers, + instance_ids=[iid], + detail=( + f"instance {iid!r} lacks trusted launcher-issued instance " + "identity; diagnostic reads remain available" + ), + active_blocker=True, + ) + + unknown = [w for w in workers if w.get("client_type") == mwi.UNKNOWN_CLIENT] + if unknown: + _finding( + CLASS_UNKNOWN_CLIENT, + severity="blocker", + workers=unknown, + instance_ids=[iid], + detail=f"instance {iid!r} has worker(s) with unknown client_type", + active_blocker=True, + ) + + foreign = [w for w in workers if w.get("foreign_repository")] + if foreign: + _finding( + CLASS_FOREIGN_REPOSITORY, + severity="blocker", + workers=foreign, + instance_ids=[iid], + detail=f"instance {iid!r} has foreign-repository workers", + active_blocker=True, + ) + + old = [w for w in workers if w.get("old_revision")] + if old: + _finding( + CLASS_OLD_REVISION, + severity="blocker", + workers=old, + instance_ids=[iid], + detail=f"instance {iid!r} has old-revision workers", + active_blocker=True, + ) + + orphans = [w for w in workers if w.get("ownership_state") == "orphaned"] + if orphans: + _finding( + CLASS_ORPHANED, + severity="blocker", + workers=orphans, + instance_ids=[iid], + detail=f"instance {iid!r} has orphaned/unowned workers", + active_blocker=True, + ) + + instances.append( + { + "client_instance_id": iid, + "client_types": client_types, + "client_type": client_types[0] if len(client_types) == 1 else None, + "fleet_run_ids": sorted( + {w.get("fleet_run_id") for w in workers if w.get("fleet_run_id")} + ), + "worker_count": len(workers), + "namespaces": sorted({n for n in namespaces if n}), + "namespace_counts": dict(ns_counts), + "duplicate_namespaces": dup_ns, + "trusted_instance_identity": trusted, + "workers": workers, + "mutation_safe": all(w.get("mutation_safe") for w in workers) + and not dup_ns + and trusted, + } + ) + + for row in unkeyed_live: + _finding( + CLASS_LEGACY_INCOMPLETE, + severity="blocker", + workers=[row], + detail="live worker has no client_instance_id", + active_blocker=True, + ) + if row.get("client_type") == mwi.UNKNOWN_CLIENT: + _finding( + CLASS_UNKNOWN_CLIENT, + severity="blocker", + workers=[row], + detail="live worker has unknown client_type and no instance id", + active_blocker=True, + ) + + for row in stale_rows: + _finding( + CLASS_STALE_WORKER, + severity="warning", + workers=[row], + detail=( + f"worker {row.get('worker_identity')!r} is active in the registry " + "but not live (stale heartbeat or dead pid)" + ), + active_blocker=True, + ) + + for row in historical_rows: + _finding( + CLASS_HISTORICAL, + severity="info", + workers=[row], + detail=( + f"historical registration {row.get('worker_identity')!r} " + f"(status={row.get('status')!r}) is not an active blocker" + ), + active_blocker=False, + ) + + # Manifest comparison + expected = list(expected_manifest or []) + expected_ids = { + str(item.get("client_instance_id")).strip() + for item in expected + if (item.get("client_instance_id") or "").strip() + } + live_ids = set(by_instance.keys()) + missing_ids = sorted(expected_ids - live_ids) + unmanifested_ids = sorted(live_ids - expected_ids) if expected_ids else [] + + for iid in missing_ids: + _finding( + CLASS_MISSING, + severity="blocker", + instance_ids=[iid], + detail=f"expected enrolled instance {iid!r} is missing from the live fleet", + active_blocker=True, + ) + for iid in unmanifested_ids: + _finding( + CLASS_UNMANIFESTED, + severity="blocker", + instance_ids=[iid], + workers=by_instance.get(iid, []), + detail=( + f"live instance {iid!r} is not on the approved fleet manifest " + "(unmanifested)" + ), + active_blocker=True, + ) + + # Same client_type multi-instance is healthy when each has distinct instance IDs + by_type: dict[str, list[str]] = defaultdict(list) + for inst in instances: + for ct in inst.get("client_types") or []: + by_type[str(ct)].append(inst["client_instance_id"]) + multi_instance_same_type = { + ct: ids for ct, ids in by_type.items() if len(ids) > 1 + } + + active_blockers = [f for f in findings if f.get("active_blocker")] + historical_only = [f for f in findings if f.get("classification") == CLASS_HISTORICAL] + live_safe = not active_blockers + + consistency = registry_revision or _consistency_token(all_rows, snapshot_at) + + return { + "success": True, + "read_only": True, + "snapshot_at": snapshot_at, + "consistency_token": consistency, + "registry_revision": consistency, + "live_worker_count": len(live_rows), + "historical_worker_count": len(historical_rows), + "stale_worker_count": len(stale_rows), + "instance_count": len(instances), + "workers": all_rows, + "live_workers": live_rows, + "historical_workers": historical_rows, + "stale_workers": stale_rows, + "instances": instances, + "multi_instance_same_client_type": multi_instance_same_type, + "same_client_type_not_duplicate": True, + "expected_manifest": [ + { + "client_instance_id": item.get("client_instance_id"), + "client_type": item.get("client_type"), + "fleet_run_id": item.get("fleet_run_id"), + "namespaces": item.get("namespaces"), + } + for item in expected + ], + "missing_expected_instance_ids": missing_ids, + "unmanifested_instance_ids": unmanifested_ids, + "findings": findings, + "active_blockers": active_blockers, + "historical_findings": historical_only, + "live_fleet_safe": live_safe, + "mutation_safe": live_safe and all( + inst.get("mutation_safe") for inst in instances + ) + if instances + else live_safe, + "classification_model": { + "exactly_one_process_per_profile": False, + "exactly_one_instance_per_client_type": False, + "multiple_instances_per_client_type": True, + "duplicate_requires": [ + "live client_instance_id collision", + "duplicate namespace worker within one instance", + "reused worker/session/generation/process/pid/fencing identity", + "cross-instance ownership collision", + "unmanifested instance when a manifest is required", + ], + }, + "reasons": [f["detail"] for f in active_blockers], + } + + +def compare_snapshot_heartbeats( + earlier: Mapping[str, Any], + later: Mapping[str, Any], +) -> dict[str, Any]: + """Prove heartbeat continuity and stable ownership across two snapshots.""" + earlier_live = { + w.get("worker_identity"): w for w in earlier.get("live_workers") or [] + } + later_live = { + w.get("worker_identity"): w for w in later.get("live_workers") or [] + } + shared = sorted(set(earlier_live) & set(later_live)) + continuity: list[dict[str, Any]] = [] + stable_ownership = True + for wid in shared: + a = earlier_live[wid] + b = later_live[wid] + same_instance = a.get("client_instance_id") == b.get("client_instance_id") + same_session = a.get("session_id") == b.get("session_id") + same_generation = a.get("generation_id") == b.get("generation_id") + hb_advanced_or_equal = True + if a.get("last_heartbeat_at") and b.get("last_heartbeat_at"): + hb_advanced_or_equal = b["last_heartbeat_at"] >= a["last_heartbeat_at"] + if not (same_instance and same_session and same_generation): + stable_ownership = False + continuity.append( + { + "worker_identity": wid, + "same_client_instance_id": same_instance, + "same_session_id": same_session, + "same_generation_id": same_generation, + "heartbeat_non_decreasing": hb_advanced_or_equal, + "earlier_heartbeat": a.get("last_heartbeat_at"), + "later_heartbeat": b.get("last_heartbeat_at"), + } + ) + return { + "shared_live_workers": shared, + "continuity": continuity, + "stable_ownership": stable_ownership + and all(c["heartbeat_non_decreasing"] for c in continuity), + "dropped_workers": sorted(set(earlier_live) - set(later_live)), + "new_workers": sorted(set(later_live) - set(earlier_live)), + } diff --git a/mcp_worker_identity.py b/mcp_worker_identity.py index 491d4db..f9363e9 100644 --- a/mcp_worker_identity.py +++ b/mcp_worker_identity.py @@ -301,8 +301,22 @@ CREATE INDEX IF NOT EXISTS idx_worker_session ON worker_registrations(session_id, status); CREATE INDEX IF NOT EXISTS idx_worker_profile ON worker_registrations(profile, status); +CREATE INDEX IF NOT EXISTS idx_worker_instance + ON worker_registrations(client_instance_id, status); """ +#: #978 optional columns added without rewriting historical rows. +_SCHEMA_OPTIONAL_COLUMNS: tuple[tuple[str, str], ...] = ( + ("fleet_run_id", "TEXT"), + ("authenticated_account", "TEXT"), + ("process_identity", "TEXT"), + ("startup_revision", "TEXT"), + ("loaded_revision", "TEXT"), + ("parity_revision", "TEXT"), + ("live_revision", "TEXT"), + ("instance_id_provenance", "TEXT"), +) + class WorkerRegistryError(RuntimeError): """Raised for registry misuse that is a programming error, not a refusal.""" @@ -634,6 +648,17 @@ class WorkerRegistry: conn = self._connect() try: conn.executescript(_SCHEMA_SQL) + existing = { + row[1] + for row in conn.execute( + "PRAGMA table_info(worker_registrations)" + ).fetchall() + } + for name, decl in _SCHEMA_OPTIONAL_COLUMNS: + if name not in existing: + conn.execute( + f"ALTER TABLE worker_registrations ADD COLUMN {name} {decl}" + ) conn.commit() finally: conn.close() @@ -758,6 +783,14 @@ class WorkerRegistry: heartbeat_ttl_seconds: float = DEFAULT_HEARTBEAT_TTL_SECONDS, now: datetime | None = None, pid_alive_probe=None, + fleet_run_id: str | None = None, + authenticated_account: str | None = None, + process_identity: str | None = None, + startup_revision: str | None = None, + loaded_revision: str | None = None, + parity_revision: str | None = None, + live_revision: str | None = None, + instance_id_provenance: str | None = None, ) -> dict[str, Any]: """Atomically register one worker identity. @@ -766,6 +799,10 @@ class WorkerRegistry: to mint a different identity and register that instead (AC32). The existing registration is returned untouched so the caller can see what it collided with. + + #978 also fails closed when a *live* registration already holds the same + ``client_instance_id`` for the same namespace: two workers of one + namespace cannot share one instance. """ parsed = parse_worker_identity(worker_identity) if not parsed["valid"]: @@ -783,6 +820,9 @@ class WorkerRegistry: } stamp = _ts(now or _utc_now()) + proc_id = process_identity or ( + f"pid-{int(pid)}" if pid is not None else None + ) with self._tx() as conn: existing = conn.execute( "SELECT * FROM worker_registrations WHERE worker_identity = ?", @@ -819,6 +859,41 @@ class WorkerRegistry: ), } + # #978: refuse a second live worker for the same (instance, namespace). + if namespace and client_instance_id: + peers = conn.execute( + "SELECT * FROM worker_registrations " + "WHERE client_instance_id = ? AND namespace = ? AND status = ?", + (client_instance_id, namespace, STATUS_ACTIVE), + ).fetchall() + for peer_row in peers: + peer = self._row_to_record(peer_row) + pid_alive = ( + pid_alive_probe(peer.get("pid")) + if pid_alive_probe is not None + else None + ) + if self.is_live(peer, now=now, pid_alive=pid_alive)["live"]: + return { + "success": False, + "registered": False, + "mutation_performed": False, + "blocker_kind": BLOCKER_CONFLICTING_SESSIONS, + "collision": True, + "collision_kind": "duplicate_namespace_worker", + "existing_registration": _public_record(peer), + "reasons": [ + f"client_instance_id {client_instance_id!r} already " + f"has a live worker for namespace {namespace!r} " + f"({peer.get('worker_identity')!r}); #978 fails closed " + "rather than registering a second worker" + ], + "exact_next_action": ( + "Stop the extra namespace worker, or use a distinct " + "client_instance_id for a separate application launch." + ), + } + epoch = self._next_epoch(conn, generation_id) conn.execute( """ @@ -827,8 +902,11 @@ class WorkerRegistry: generation_id, role, profile, namespace, remote, repository_binding, pid, transport, token_fingerprint, started_at, last_heartbeat_at, heartbeat_ttl_seconds, - fencing_epoch, status - ) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?) + fencing_epoch, status, + fleet_run_id, authenticated_account, process_identity, + startup_revision, loaded_revision, parity_revision, + live_revision, instance_id_provenance + ) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?) """, ( worker_identity, @@ -849,6 +927,14 @@ class WorkerRegistry: float(heartbeat_ttl_seconds), epoch, STATUS_ACTIVE, + (fleet_run_id or "").strip() or None, + (authenticated_account or "").strip() or None, + proc_id, + (startup_revision or "").strip() or None, + (loaded_revision or "").strip() or None, + (parity_revision or "").strip() or None, + (live_revision or "").strip() or None, + (instance_id_provenance or "").strip() or None, ), ) row = conn.execute( @@ -1149,8 +1235,17 @@ def _public_record(record: dict[str, Any]) -> dict[str, Any]: "transport": record.get("transport"), "started_at": record.get("started_at"), "last_heartbeat_at": record.get("last_heartbeat_at"), + "heartbeat_ttl_seconds": record.get("heartbeat_ttl_seconds"), "fencing_epoch": record.get("fencing_epoch"), "status": record.get("status"), + "fleet_run_id": record.get("fleet_run_id"), + "authenticated_account": record.get("authenticated_account"), + "process_identity": record.get("process_identity"), + "startup_revision": record.get("startup_revision"), + "loaded_revision": record.get("loaded_revision"), + "parity_revision": record.get("parity_revision"), + "live_revision": record.get("live_revision"), + "instance_id_provenance": record.get("instance_id_provenance"), } diff --git a/task_capability_map.py b/task_capability_map.py index 3766fb9..d306e03 100644 --- a/task_capability_map.py +++ b/task_capability_map.py @@ -163,6 +163,18 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = { "permission": "gitea.read", "role": "author", }, + # #978: instance-level fleet identity/health snapshot. gitea.read is the + # operation gate; the tool additionally restricts role_kind to + # controller|reconciler so author/reviewer/merger cannot use it as a + # mutation surface and no unrelated write permission is introduced. + "snapshot_instance_fleet": { + "permission": "gitea.read", + "role": "controller", + }, + "gitea_snapshot_instance_fleet": { + "permission": "gitea.read", + "role": "controller", + }, # #644: Phase 2 Web Console recovery tasks. "clear_stale_binding": { "permission": "gitea.read", diff --git a/tests/test_issue_978_instance_fleet_snapshot.py b/tests/test_issue_978_instance_fleet_snapshot.py new file mode 100644 index 0000000..8ce41e0 --- /dev/null +++ b/tests/test_issue_978_instance_fleet_snapshot.py @@ -0,0 +1,1054 @@ +"""Instance-level fleet identity and health snapshots (#978). + +Covers the acceptance criteria for multi-instance fleets: two Codex launches +remain distinct, ten workers attribute correctly, same client_type is not a +duplicate, collisions fail closed, historical rows stay non-blocking, and the +controller/reconciler read surface grants no mutation authority. + +All identifiers are synthetic. +""" + +from __future__ import annotations + +import inspect +import json +import os +import tempfile +import unittest +from datetime import datetime, timedelta, timezone +from typing import Any +from unittest import mock + +import mcp_fleet_snapshot as fleet +import mcp_worker_identity as mwi + + +NOW = datetime(2026, 7, 30, 7, 0, 0, tzinfo=timezone.utc) +TTL = 900.0 +NAMESPACES = ("author", "reviewer", "merger", "controller", "reconciler") + + +def _registry() -> mwi.WorkerRegistry: + handle, path = tempfile.mkstemp(suffix=".sqlite3") + os.close(handle) + os.unlink(path) + return mwi.WorkerRegistry(path) + + +def _alive(_pid: int | None) -> bool: + return True + + +def _dead(_pid: int | None) -> bool: + return False + + +def _register( + registry: mwi.WorkerRegistry, + *, + client: str, + instance: str, + session: str, + generation: str, + namespace: str, + pid: int, + now: datetime = NOW, + ttl: float = TTL, + fleet_run_id: str | None = "run-canary", + status_after: str | None = None, + worker_identity: str | None = None, + process_identity: str | None = None, + repository_binding: str = "/repo/Gitea-Tools", + startup_revision: str = "rev-live", + instance_id_provenance: str = fleet.INSTANCE_ID_PROVENANCE_TRUSTED, + role: str | None = None, +) -> dict[str, Any]: + identity = worker_identity or mwi.generate_worker_identity( + client, f"{session}-{namespace}", now=now, nonce=f"{pid}-{namespace}" + ) + outcome = registry.register( + worker_identity=identity, + client_name=client, + client_instance_id=instance, + session_id=session, + generation_id=generation, + role=role or namespace, + profile=f"prgs-{namespace}", + namespace=namespace, + remote="prgs", + repository_binding=repository_binding, + pid=pid, + heartbeat_ttl_seconds=ttl, + now=now, + fleet_run_id=fleet_run_id, + process_identity=process_identity or f"pid-{pid}", + startup_revision=startup_revision, + loaded_revision=startup_revision, + parity_revision=startup_revision, + live_revision=startup_revision, + instance_id_provenance=instance_id_provenance, + pid_alive_probe=_alive, + ) + outcome["identity"] = identity + if status_after and outcome.get("registered"): + if status_after == mwi.STATUS_RELEASED: + registry.release(worker_identity=identity, now=now) + elif status_after == mwi.STATUS_SUPERSEDED: + # Direct status write for historical rows. + with registry._tx() as conn: # noqa: SLF001 — test-only fixture + conn.execute( + "UPDATE worker_registrations SET status = ? WHERE worker_identity = ?", + (mwi.STATUS_SUPERSEDED, identity), + ) + return outcome + + +def _five_workers( + registry: mwi.WorkerRegistry, + *, + client: str, + instance: str, + session_prefix: str, + gen_prefix: str, + pid_base: int, + now: datetime = NOW, +) -> list[dict[str, Any]]: + out = [] + for i, ns in enumerate(NAMESPACES): + out.append( + _register( + registry, + client=client, + instance=instance, + session=f"{session_prefix}-{ns}", + generation=f"{gen_prefix}-{ns}", + namespace=ns, + pid=pid_base + i, + now=now, + ) + ) + return out + + +class InstanceIdentityTests(unittest.TestCase): + def test_generate_distinct_instance_ids(self): + a = fleet.generate_client_instance_id("codex", launch_nonce="a", now=NOW) + b = fleet.generate_client_instance_id("codex", launch_nonce="b", now=NOW) + self.assertNotEqual(a, b) + self.assertTrue(a.startswith("inst-codex-")) + self.assertTrue(fleet.assess_instance_identity(a)["trusted"]) + + def test_legacy_pid_fallback_not_trusted(self): + a = fleet.assess_instance_identity("pid-1234") + self.assertFalse(a["trusted"]) + self.assertEqual(a["provenance"], fleet.INSTANCE_ID_PROVENANCE_LEGACY) + + def test_missing_instance_id(self): + a = fleet.assess_instance_identity(None) + self.assertFalse(a["complete"]) + self.assertEqual(a["provenance"], fleet.INSTANCE_ID_PROVENANCE_MISSING) + + +class TwoCodexInstancesTests(unittest.TestCase): + """AC1–AC3: two Codex instances, ten workers, same type not duplicate.""" + + def test_two_codex_instances_ten_workers(self): + registry = _registry() + inst_a = fleet.generate_client_instance_id("codex", launch_nonce="A", now=NOW) + inst_b = fleet.generate_client_instance_id("codex", launch_nonce="B", now=NOW) + _five_workers( + registry, + client="codex", + instance=inst_a, + session_prefix="sessA", + gen_prefix="genA", + pid_base=1000, + ) + _five_workers( + registry, + client="codex", + instance=inst_b, + session_prefix="sessB", + gen_prefix="genB", + pid_base=2000, + ) + snap = fleet.snapshot_instance_fleet( + registry.list_workers(status=None), + now=NOW + timedelta(seconds=10), + pid_alive_probe=_alive, + canonical_repository="/repo/Gitea-Tools", + expected_live_revision="rev-live", + ) + self.assertEqual(snap["live_worker_count"], 10) + self.assertEqual(snap["instance_count"], 2) + self.assertIn("codex", snap["multi_instance_same_client_type"]) + self.assertEqual(len(snap["multi_instance_same_client_type"]["codex"]), 2) + self.assertTrue(snap["same_client_type_not_duplicate"]) + # No finding that treats same client_type alone as a duplicate. + for f in snap["active_blockers"]: + self.assertNotIn("client_type alone", f["detail"].lower()) + by_inst = {i["client_instance_id"]: i for i in snap["instances"]} + self.assertEqual(by_inst[inst_a]["worker_count"], 5) + self.assertEqual(by_inst[inst_b]["worker_count"], 5) + self.assertEqual(set(by_inst[inst_a]["namespaces"]), set(NAMESPACES)) + self.assertEqual(set(by_inst[inst_b]["namespaces"]), set(NAMESPACES)) + self.assertTrue(snap["live_fleet_safe"]) + + def test_mixed_client_types(self): + registry = _registry() + _five_workers( + registry, + client="codex", + instance="inst-codex-1", + session_prefix="c", + gen_prefix="gc", + pid_base=10, + ) + _five_workers( + registry, + client="claude_code", + instance="inst-claude-1", + session_prefix="l", + gen_prefix="gl", + pid_base=20, + ) + snap = fleet.snapshot_instance_fleet( + registry.list_workers(), + now=NOW + timedelta(seconds=5), + pid_alive_probe=_alive, + ) + self.assertEqual(snap["instance_count"], 2) + self.assertTrue(snap["live_fleet_safe"]) + + +class CollisionTests(unittest.TestCase): + def test_live_reuse_of_instance_id_different_client_types(self): + registry = _registry() + shared = "inst-shared-collision" + _register( + registry, + client="codex", + instance=shared, + session="s1", + generation="g1", + namespace="author", + pid=1, + ) + # Second client type forced onto same instance id via direct SQL to + # bypass register's same-namespace guard and exercise classification. + with registry._tx() as conn: # noqa: SLF001 + conn.execute( + """ + INSERT INTO worker_registrations ( + worker_identity, client_name, client_instance_id, session_id, + generation_id, role, profile, namespace, remote, + repository_binding, pid, transport, token_fingerprint, + started_at, last_heartbeat_at, heartbeat_ttl_seconds, + fencing_epoch, status, process_identity, instance_id_provenance + ) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?) + """, + ( + "claude_code-20260730T070000Z-aaaaaaaaaaaa", + "claude_code", + shared, + "s2", + "g2", + "author", + "prgs-author", + "reviewer", + "prgs", + "/repo/Gitea-Tools", + 2, + None, + None, + mwi._ts(NOW), # noqa: SLF001 + mwi._ts(NOW), # noqa: SLF001 + TTL, + 1, + mwi.STATUS_ACTIVE, + "pid-2", + fleet.INSTANCE_ID_PROVENANCE_TRUSTED, + ), + ) + snap = fleet.snapshot_instance_fleet( + registry.list_workers(), + now=NOW + timedelta(seconds=1), + pid_alive_probe=_alive, + ) + classes = {f["classification"] for f in snap["active_blockers"]} + self.assertIn(fleet.CLASS_INSTANCE_ID_COLLISION, classes) + self.assertFalse(snap["live_fleet_safe"]) + + def test_duplicate_namespace_within_instance(self): + registry = _registry() + inst = "inst-dup-ns" + first = _register( + registry, + client="codex", + instance=inst, + session="s-a", + generation="g-a", + namespace="author", + pid=11, + ) + self.assertTrue(first["registered"], first) + # Direct second author under same instance (bypass register guard). + with registry._tx() as conn: # noqa: SLF001 + conn.execute( + """ + INSERT INTO worker_registrations ( + worker_identity, client_name, client_instance_id, session_id, + generation_id, role, profile, namespace, remote, + repository_binding, pid, transport, token_fingerprint, + started_at, last_heartbeat_at, heartbeat_ttl_seconds, + fencing_epoch, status, process_identity, instance_id_provenance + ) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?) + """, + ( + "codex-20260730T070000Z-bbbbbbbbbbbb", + "codex", + inst, + "s-b", + "g-b", + "author", + "prgs-author", + "author", + "prgs", + "/repo/Gitea-Tools", + 12, + None, + None, + mwi._ts(NOW), # noqa: SLF001 + mwi._ts(NOW), # noqa: SLF001 + TTL, + 2, + mwi.STATUS_ACTIVE, + "pid-12", + fleet.INSTANCE_ID_PROVENANCE_TRUSTED, + ), + ) + snap = fleet.snapshot_instance_fleet( + registry.list_workers(), + now=NOW + timedelta(seconds=1), + pid_alive_probe=_alive, + ) + classes = {f["classification"] for f in snap["active_blockers"]} + self.assertIn(fleet.CLASS_DUPLICATE_NAMESPACE, classes) + + def test_register_refuses_duplicate_namespace(self): + registry = _registry() + inst = "inst-reg-dup" + a = _register( + registry, + client="codex", + instance=inst, + session="s1", + generation="g1", + namespace="author", + pid=21, + ) + self.assertTrue(a["registered"]) + b = _register( + registry, + client="codex", + instance=inst, + session="s2", + generation="g2", + namespace="author", + pid=22, + ) + self.assertFalse(b["registered"]) + self.assertEqual(b.get("collision_kind"), "duplicate_namespace_worker") + + def test_reused_worker_identity(self): + registry = _registry() + identity = mwi.generate_worker_identity("codex", "sess", now=NOW, nonce="fixed") + first = _register( + registry, + client="codex", + instance="inst-1", + session="s1", + generation="g1", + namespace="author", + pid=31, + worker_identity=identity, + ) + self.assertTrue(first["registered"]) + second = registry.register( + worker_identity=identity, + client_name="codex", + client_instance_id="inst-2", + session_id="s2", + generation_id="g2", + namespace="author", + pid=32, + now=NOW, + pid_alive_probe=_alive, + ) + self.assertFalse(second["registered"]) + self.assertTrue(second["collision"]) + + def test_reused_session_identity(self): + registry = _registry() + _register( + registry, + client="codex", + instance="inst-1", + session="shared-session", + generation="g1", + namespace="author", + pid=41, + ) + _register( + registry, + client="codex", + instance="inst-2", + session="shared-session", + generation="g2", + namespace="reviewer", + pid=42, + ) + snap = fleet.snapshot_instance_fleet( + registry.list_workers(), + now=NOW + timedelta(seconds=1), + pid_alive_probe=_alive, + ) + classes = {f["classification"] for f in snap["active_blockers"]} + self.assertIn(fleet.CLASS_SESSION_COLLISION, classes) + + def test_reused_generation_across_instances(self): + registry = _registry() + _register( + registry, + client="codex", + instance="inst-1", + session="s1", + generation="gen-shared", + namespace="author", + pid=51, + ) + _register( + registry, + client="codex", + instance="inst-2", + session="s2", + generation="gen-shared", + namespace="author", + pid=52, + ) + snap = fleet.snapshot_instance_fleet( + registry.list_workers(), + now=NOW + timedelta(seconds=1), + pid_alive_probe=_alive, + ) + classes = {f["classification"] for f in snap["active_blockers"]} + self.assertIn(fleet.CLASS_GENERATION_COLLISION, classes) + + def test_reused_pid_and_process_identity(self): + registry = _registry() + _register( + registry, + client="codex", + instance="inst-1", + session="s1", + generation="g1", + namespace="author", + pid=61, + process_identity="proc-same", + ) + _register( + registry, + client="codex", + instance="inst-2", + session="s2", + generation="g2", + namespace="reviewer", + pid=61, + process_identity="proc-same", + ) + snap = fleet.snapshot_instance_fleet( + registry.list_workers(), + now=NOW + timedelta(seconds=1), + pid_alive_probe=_alive, + ) + classes = {f["classification"] for f in snap["active_blockers"]} + self.assertIn(fleet.CLASS_PID_COLLISION, classes) + self.assertIn(fleet.CLASS_PROCESS_COLLISION, classes) + + def test_ownership_fencing_collision(self): + registry = _registry() + _register( + registry, + client="codex", + instance="inst-1", + session="s1", + generation="gen-fence", + namespace="author", + pid=71, + ) + # Second instance forced onto same generation+epoch via SQL. + with registry._tx() as conn: # noqa: SLF001 + conn.execute( + """ + INSERT INTO worker_registrations ( + worker_identity, client_name, client_instance_id, session_id, + generation_id, role, profile, namespace, remote, + repository_binding, pid, transport, token_fingerprint, + started_at, last_heartbeat_at, heartbeat_ttl_seconds, + fencing_epoch, status, process_identity, instance_id_provenance + ) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?) + """, + ( + "codex-20260730T070000Z-cccccccccccc", + "codex", + "inst-2", + "s2", + "gen-fence", + "reviewer", + "prgs-reviewer", + "reviewer", + "prgs", + "/repo/Gitea-Tools", + 72, + None, + None, + mwi._ts(NOW), # noqa: SLF001 + mwi._ts(NOW), # noqa: SLF001 + TTL, + 1, # same epoch as first (first got epoch 1) + mwi.STATUS_ACTIVE, + "pid-72", + fleet.INSTANCE_ID_PROVENANCE_TRUSTED, + ), + ) + snap = fleet.snapshot_instance_fleet( + registry.list_workers(), + now=NOW + timedelta(seconds=1), + pid_alive_probe=_alive, + ) + classes = {f["classification"] for f in snap["active_blockers"]} + self.assertTrue( + {fleet.CLASS_OWNERSHIP_COLLISION, fleet.CLASS_GENERATION_COLLISION} + & classes + ) + + +class ManifestAndEdgeTests(unittest.TestCase): + def test_missing_and_unmanifested(self): + registry = _registry() + live = "inst-live" + expected = "inst-expected-missing" + _register( + registry, + client="codex", + instance=live, + session="s", + generation="g", + namespace="author", + pid=81, + ) + snap = fleet.snapshot_instance_fleet( + registry.list_workers(), + expected_manifest=[ + {"client_instance_id": expected, "client_type": "codex"}, + ], + now=NOW + timedelta(seconds=1), + pid_alive_probe=_alive, + ) + self.assertEqual(snap["missing_expected_instance_ids"], [expected]) + self.assertEqual(snap["unmanifested_instance_ids"], [live]) + classes = {f["classification"] for f in snap["active_blockers"]} + self.assertIn(fleet.CLASS_MISSING, classes) + self.assertIn(fleet.CLASS_UNMANIFESTED, classes) + + def test_unknown_client(self): + registry = _registry() + with registry._tx() as conn: # noqa: SLF001 + conn.execute( + """ + INSERT INTO worker_registrations ( + worker_identity, client_name, client_instance_id, session_id, + generation_id, role, profile, namespace, remote, + repository_binding, pid, transport, token_fingerprint, + started_at, last_heartbeat_at, heartbeat_ttl_seconds, + fencing_epoch, status, process_identity, instance_id_provenance + ) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?) + """, + ( + "unknown_client-20260730T070000Z-dddddddddddd", + "unknown_client", + "inst-unk", + "s", + "g", + "author", + "prgs-author", + "author", + "prgs", + "/repo/Gitea-Tools", + 91, + None, + None, + mwi._ts(NOW), # noqa: SLF001 + mwi._ts(NOW), # noqa: SLF001 + TTL, + 1, + mwi.STATUS_ACTIVE, + "pid-91", + fleet.INSTANCE_ID_PROVENANCE_TRUSTED, + ), + ) + snap = fleet.snapshot_instance_fleet( + registry.list_workers(), + now=NOW + timedelta(seconds=1), + pid_alive_probe=_alive, + ) + classes = {f["classification"] for f in snap["active_blockers"]} + self.assertIn(fleet.CLASS_UNKNOWN_CLIENT, classes) + + def test_foreign_repository(self): + registry = _registry() + _register( + registry, + client="codex", + instance="inst-fr", + session="s", + generation="g", + namespace="author", + pid=101, + repository_binding="/other/repo", + ) + snap = fleet.snapshot_instance_fleet( + registry.list_workers(), + now=NOW + timedelta(seconds=1), + pid_alive_probe=_alive, + canonical_repository="/repo/Gitea-Tools", + ) + classes = {f["classification"] for f in snap["active_blockers"]} + self.assertIn(fleet.CLASS_FOREIGN_REPOSITORY, classes) + + def test_old_revision(self): + registry = _registry() + _register( + registry, + client="codex", + instance="inst-old", + session="s", + generation="g", + namespace="author", + pid=111, + startup_revision="rev-old", + ) + snap = fleet.snapshot_instance_fleet( + registry.list_workers(), + now=NOW + timedelta(seconds=1), + pid_alive_probe=_alive, + expected_live_revision="rev-live", + ) + classes = {f["classification"] for f in snap["active_blockers"]} + self.assertIn(fleet.CLASS_OLD_REVISION, classes) + + def test_stale_worker(self): + registry = _registry() + _register( + registry, + client="codex", + instance="inst-stale", + session="s", + generation="g", + namespace="author", + pid=121, + now=NOW, + ttl=60.0, + ) + snap = fleet.snapshot_instance_fleet( + registry.list_workers(), + now=NOW + timedelta(seconds=120), + pid_alive_probe=_dead, + ) + classes = {f["classification"] for f in snap["active_blockers"]} + self.assertIn(fleet.CLASS_STALE_WORKER, classes) + + def test_historical_not_active_blocker(self): + registry = _registry() + _register( + registry, + client="codex", + instance="inst-hist", + session="s", + generation="g", + namespace="author", + pid=131, + status_after=mwi.STATUS_RELEASED, + ) + _register( + registry, + client="codex", + instance="inst-live", + session="s2", + generation="g2", + namespace="author", + pid=132, + ) + snap = fleet.snapshot_instance_fleet( + registry.list_workers(status=None), + now=NOW + timedelta(seconds=1), + pid_alive_probe=_alive, + ) + self.assertEqual(snap["historical_worker_count"], 1) + self.assertTrue( + all(f["classification"] == fleet.CLASS_HISTORICAL for f in snap["historical_findings"]) + ) + # Only historical finding for the dead row — live fleet otherwise safe. + hist_blockers = [ + f + for f in snap["active_blockers"] + if f["classification"] == fleet.CLASS_HISTORICAL + ] + self.assertEqual(hist_blockers, []) + self.assertTrue(snap["live_fleet_safe"]) + + def test_legacy_incomplete_identity(self): + registry = _registry() + _register( + registry, + client="codex", + instance="pid-999", + session="s", + generation="g", + namespace="author", + pid=141, + instance_id_provenance=fleet.INSTANCE_ID_PROVENANCE_LEGACY, + ) + snap = fleet.snapshot_instance_fleet( + registry.list_workers(), + now=NOW + timedelta(seconds=1), + pid_alive_probe=_alive, + ) + classes = {f["classification"] for f in snap["active_blockers"]} + self.assertIn(fleet.CLASS_LEGACY_INCOMPLETE, classes) + self.assertFalse(snap["mutation_safe"]) + + +class LifecycleTests(unittest.TestCase): + def test_heartbeat_continuity_across_snapshots(self): + registry = _registry() + outcome = _register( + registry, + client="codex", + instance="inst-hb", + session="s", + generation="g", + namespace="author", + pid=151, + ) + identity = outcome["identity"] + earlier = fleet.snapshot_instance_fleet( + registry.list_workers(), + now=NOW + timedelta(seconds=1), + pid_alive_probe=_alive, + ) + later_time = NOW + timedelta(seconds=30) + registry.heartbeat( + worker_identity=identity, + fencing_epoch=outcome["fencing_epoch"], + now=later_time, + expected_session_id="s", + expected_generation_id="g", + expected_client_name="codex", + expected_pid=151, + ) + later = fleet.snapshot_instance_fleet( + registry.list_workers(), + now=later_time + timedelta(seconds=1), + pid_alive_probe=_alive, + ) + cmp = fleet.compare_snapshot_heartbeats(earlier, later) + self.assertTrue(cmp["stable_ownership"]) + self.assertEqual(cmp["shared_live_workers"], [identity]) + self.assertTrue(cmp["continuity"][0]["heartbeat_non_decreasing"]) + + def test_worker_reconnect_preserves_instance(self): + """Reconnect = heartbeat renew; same instance/worker/session.""" + registry = _registry() + o = _register( + registry, + client="codex", + instance="inst-rc", + session="s-rc", + generation="g-rc", + namespace="author", + pid=161, + ) + before = registry.get(o["identity"]) + registry.heartbeat( + worker_identity=o["identity"], + fencing_epoch=o["fencing_epoch"], + now=NOW + timedelta(seconds=20), + expected_session_id="s-rc", + expected_generation_id="g-rc", + expected_client_name="codex", + expected_pid=161, + ) + after = registry.get(o["identity"]) + self.assertEqual(before["client_instance_id"], after["client_instance_id"]) + self.assertEqual(before["session_id"], after["session_id"]) + self.assertEqual(before["worker_identity"], after["worker_identity"]) + + def test_worker_restart_new_worker_same_instance(self): + registry = _registry() + inst = "inst-restart" + old = _register( + registry, + client="codex", + instance=inst, + session="s-old", + generation="g-old", + namespace="author", + pid=171, + ) + registry.release(worker_identity=old["identity"], now=NOW + timedelta(seconds=5)) + new = _register( + registry, + client="codex", + instance=inst, + session="s-new", + generation="g-new", + namespace="author", + pid=172, + now=NOW + timedelta(seconds=10), + ) + self.assertTrue(new["registered"]) + self.assertNotEqual(old["identity"], new["identity"]) + snap = fleet.snapshot_instance_fleet( + registry.list_workers(status=None), + now=NOW + timedelta(seconds=15), + pid_alive_probe=_alive, + ) + live = [w for w in snap["live_workers"] if w["client_instance_id"] == inst] + self.assertEqual(len(live), 1) + self.assertEqual(live[0]["worker_identity"], new["identity"]) + self.assertEqual(snap["historical_worker_count"], 1) + + def test_full_application_restart_new_instance(self): + registry = _registry() + old_inst = fleet.generate_client_instance_id("codex", launch_nonce="old", now=NOW) + new_inst = fleet.generate_client_instance_id( + "codex", launch_nonce="new", now=NOW + timedelta(hours=1) + ) + for ns, pid in zip(NAMESPACES, range(181, 186)): + o = _register( + registry, + client="codex", + instance=old_inst, + session=f"old-{ns}", + generation=f"oldg-{ns}", + namespace=ns, + pid=pid, + ) + registry.release(worker_identity=o["identity"], now=NOW + timedelta(minutes=1)) + for ns, pid in zip(NAMESPACES, range(191, 196)): + _register( + registry, + client="codex", + instance=new_inst, + session=f"new-{ns}", + generation=f"newg-{ns}", + namespace=ns, + pid=pid, + now=NOW + timedelta(minutes=2), + ) + snap = fleet.snapshot_instance_fleet( + registry.list_workers(status=None), + now=NOW + timedelta(minutes=3), + pid_alive_probe=_alive, + ) + live_ids = {i["client_instance_id"] for i in snap["instances"]} + self.assertEqual(live_ids, {new_inst}) + self.assertEqual(snap["historical_worker_count"], 5) + + +class PermissionAndSingleClientTests(unittest.TestCase): + def test_single_client_still_supported(self): + registry = _registry() + _five_workers( + registry, + client="codex", + instance="inst-single", + session_prefix="solo", + gen_prefix="solo-g", + pid_base=300, + ) + snap = fleet.snapshot_instance_fleet( + registry.list_workers(), + now=NOW + timedelta(seconds=1), + pid_alive_probe=_alive, + ) + self.assertEqual(snap["instance_count"], 1) + self.assertEqual(snap["live_worker_count"], 5) + self.assertTrue(snap["live_fleet_safe"]) + self.assertTrue(snap["mutation_safe"]) + + def test_snapshot_has_consistency_token(self): + registry = _registry() + _register( + registry, + client="codex", + instance="inst-tok", + session="s", + generation="g", + namespace="author", + pid=400, + ) + snap = fleet.snapshot_instance_fleet( + registry.list_workers(), + now=NOW + timedelta(seconds=1), + pid_alive_probe=_alive, + ) + self.assertTrue(snap["consistency_token"].startswith("fleetrev-")) + self.assertEqual(snap["consistency_token"], snap["registry_revision"]) + self.assertIn("snapshot_at", snap) + + def test_tool_role_gate_denies_author(self): + import gitea_mcp_server as server + + with mock.patch.object(server, "_profile_operation_gate", return_value=[]): + with mock.patch.object( + server, + "get_profile", + return_value={ + "profile_name": "prgs-author", + "role": "author", + "allowed_operations": ["gitea.read"], + "forbidden_operations": [], + }, + ): + result = server.gitea_snapshot_instance_fleet(remote="prgs") + self.assertFalse(result.get("success")) + self.assertEqual(result.get("denied_role"), "author") + self.assertTrue(result.get("read_only")) + self.assertFalse(result.get("mutation_performed", True)) + + def test_tool_allows_controller(self): + import gitea_mcp_server as server + + registry = _registry() + _register( + registry, + client="codex", + instance="inst-ctrl", + session="s", + generation="g", + namespace="controller", + pid=410, + ) + with mock.patch.object(server, "_profile_operation_gate", return_value=[]): + with mock.patch.object( + server, + "get_profile", + return_value={ + "profile_name": "prgs-controller", + "role": "controller", + "allowed_operations": ["gitea.read"], + "forbidden_operations": [], + }, + ): + with mock.patch.object(server, "_worker_registry", return_value=registry): + with mock.patch.object( + server, "_current_master_parity", return_value={} + ): + result = server.gitea_snapshot_instance_fleet(remote="prgs") + self.assertTrue(result.get("success"), result) + self.assertTrue(result.get("read_only")) + self.assertTrue(result["permission_scope"]["denied_unrelated_mutations"]) + self.assertEqual(result["permission_scope"]["granted_operations"], ["gitea.read"]) + + def test_tool_allows_reconciler(self): + import gitea_mcp_server as server + + registry = _registry() + with mock.patch.object(server, "_profile_operation_gate", return_value=[]): + with mock.patch.object( + server, + "get_profile", + return_value={ + "profile_name": "prgs-reconciler", + "role": "reconciler", + "allowed_operations": ["gitea.read"], + "forbidden_operations": [], + }, + ): + with mock.patch.object(server, "_worker_registry", return_value=registry): + with mock.patch.object( + server, "_current_master_parity", return_value={} + ): + result = server.gitea_snapshot_instance_fleet(remote="prgs") + self.assertTrue(result.get("success"), result) + + def test_capability_map_entry(self): + import task_capability_map as tcm + + entry = tcm.TASK_CAPABILITY_MAP["snapshot_instance_fleet"] + self.assertEqual(entry["permission"], "gitea.read") + self.assertEqual(entry["role"], "controller") + self.assertIn("gitea_snapshot_instance_fleet", tcm.TASK_CAPABILITY_MAP) + + def test_fleet_run_env_allowlisted(self): + import gitea_config + + self.assertIn("GITEA_MCP_FLEET_RUN_ID", gitea_config.RECOGNIZED_GITEA_ENV_KEYS) + + def test_docs_exist(self): + root = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) + path = os.path.join(root, "docs", "instance-fleet-identity.md") + self.assertTrue(os.path.isfile(path)) + with open(path, encoding="utf-8") as fh: + text = fh.read() + self.assertIn("client_instance_id", text) + self.assertIn("exactly_one_per_profile", text) + self.assertIn("gitea_snapshot_instance_fleet", text) + + def test_tool_is_registered(self): + import gitea_mcp_server as server + + self.assertTrue(callable(server.gitea_snapshot_instance_fleet)) + src = inspect.getsource(server.gitea_snapshot_instance_fleet) + self.assertIn("controller", src) + self.assertIn("reconciler", src) + + +class ClientHintsTests(unittest.TestCase): + def test_client_hints_legacy_when_unset(self): + import gitea_mcp_server as server + + env = {"GITEA_MCP_CLIENT": "codex"} + with mock.patch.dict(os.environ, env, clear=False): + # Ensure instance key is absent. + os.environ.pop("GITEA_MCP_CLIENT_INSTANCE", None) + hints = server._client_identity_hints() + self.assertFalse(hints["instance_identity_trusted"]) + self.assertTrue( + str(hints["client_instance_id"]).startswith("legacy-pid-") + or hints["instance_id_provenance"] == fleet.INSTANCE_ID_PROVENANCE_LEGACY + ) + + def test_client_hints_trusted_when_set(self): + import gitea_mcp_server as server + + inst = fleet.generate_client_instance_id("codex", launch_nonce="t", now=NOW) + with mock.patch.dict( + os.environ, + { + "GITEA_MCP_CLIENT": "codex", + "GITEA_MCP_CLIENT_INSTANCE": inst, + "GITEA_MCP_FLEET_RUN_ID": "run-1", + }, + clear=False, + ): + hints = server._client_identity_hints() + self.assertTrue(hints["instance_identity_trusted"]) + self.assertEqual(hints["client_instance_id"], inst) + self.assertEqual(hints["fleet_run_id"], "run-1") + + +if __name__ == "__main__": + unittest.main()