Files
Gitea-Tools/mcp_fleet_snapshot.py
T
jcwalker3andClaude Opus 4.8 7e19079b5c feat(launcher): trusted client-instance identity for project-scoped launches
Issue #985. The production launcher could mint one trusted
GITEA_MCP_CLIENT_INSTANCE per launch, but two gaps kept real launches on
untrusted legacy-pid-* identities:

1. build_application_mcp_servers required a profile for all five sanctioned
   namespaces, so a project-scoped configuration exposing only author,
   reviewer, and merger could not use it without inventing controller and
   reconciler profiles that must not exist.
2. The module produced configuration data but had no runnable entry point, so
   every real launch bypassed it entirely.

Changes:

- resolve_launch_namespaces() validates an explicit namespace subset as an
  allow-list; unknown, duplicate, and empty selections are refused rather than
  silently narrowing a launch. Omitting it preserves five-namespace behaviour.
- build_application_mcp_servers() accepts that subset, requires profiles only
  for the launched namespaces, starts only those workers, and reports
  excluded_namespaces / project_scoped.
- collect_instance_ids_from_mcp_servers() inspects the launch's own namespaces
  instead of an assumed five, and reports missing_servers, so a three-namespace
  launch can prove shared attribution without reading as two absent workers.
- Runnable entry point: python3 -m mcp_application_launcher mints one trusted
  identity, writes a per-launch 0600 mcpServers config, and execs the client.
  CLIENT_LAUNCH_SPECS is a data-driven registry so other supported clients use
  the same mint-once/propagate-to-all mechanism. --dry-run prints the plan.
- Provenance sealing now fails closed. The inst- format is public and
  reproducible, so format alone could previously let anyone who set one
  environment variable manufacture a trusted identity. Trust now additionally
  requires GITEA_MCP_INSTANCE_PROVENANCE=trusted_launcher, which only the
  launcher writes; a well-formed but unsealed value is classified
  unsealed_launcher and refused, while still being reported for diagnosis.

Deliberate behaviour change: tests/test_issue_978_instance_fleet_snapshot.py
test_client_hints_trusted_when_set previously asserted that a well-formed ID
alone was trusted. It now supplies the launcher seal, and a new companion test
asserts the unsealed case fails closed. This tightens the contract; no
assertion was weakened.

No static or persistent per-project instance IDs are introduced, duplicate
worker and cohort detection are untouched, and no fleet or mutation gate is
relaxed.

Tests: tests/test_issue_985_project_scoped_launcher.py, 47 passed, 3 subtests.
Full suite from a branches/ worktree: 28 failed, 6252 passed, 6 skipped against
a master baseline at 32ab8392 of 28 failed, 6204 passed, 6 skipped; the failing
sets are byte-identical, so zero regressions and zero masked failures.

Closes #985

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-31 03:15:53 -05:00

874 lines
32 KiB
Python

"""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"
#: Well-formed ``inst-…`` identity presented without the launcher's own
#: provenance seal (#985). The format alone is public and reproducible, so an
#: unsealed value is treated as manually asserted trust and fails closed.
INSTANCE_ID_PROVENANCE_UNSEALED = "unsealed_launcher"
CLIENT_INSTANCE_ENV = "GITEA_MCP_CLIENT_INSTANCE"
FLEET_RUN_ENV = "GITEA_MCP_FLEET_RUN_ID"
PROCESS_IDENTITY_ENV = "GITEA_MCP_PROCESS_IDENTITY"
INSTANCE_PROVENANCE_ENV = "GITEA_MCP_INSTANCE_PROVENANCE"
_LEGACY_INSTANCE_PREFIXES = ("pid-", "proc-", "legacy-")
#: Trusted launcher-issued IDs use the reserved ``inst-`` prefix
#: (see :func:`generate_client_instance_id`). Ordinary user-supplied strings
#: without that prefix never count as trusted attribution.
_TRUSTED_INSTANCE_PREFIX = "inst-"
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, match the launcher-minted ``inst-…``
format from :func:`generate_client_instance_id`, and are not pre-#978
PID/proc/legacy placeholders. Ordinary user-supplied or malformed values
fail closed as untrusted so they cannot spoof multi-instance attribution.
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"
],
}
if not text.startswith(_TRUSTED_INSTANCE_PREFIX) or len(text) <= len(
_TRUSTED_INSTANCE_PREFIX
):
return {
"client_instance_id": text,
"complete": False,
"trusted": False,
"provenance": INSTANCE_ID_PROVENANCE_LEGACY,
"reasons": [
f"client_instance_id {text!r} is malformed or user-supplied and "
f"does not use the trusted launcher prefix "
f"{_TRUSTED_INSTANCE_PREFIX!r}; refusing trusted attribution"
],
}
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-<n>`` 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)),
}