Compare commits

...
Author SHA1 Message Date
sysadminandClaude Opus 4.8 4a1cc63e94 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) <[email protected]>
2026-07-30 03:14:15 -04:00
sysadmin 6596b259fc Merge pull request 'fix(runtime): recognize client identity environment and refresh worker registrations' (#976) from fix/issue-975-client-identity-heartbeat into master 2026-07-29 23:26:34 -05:00
8 changed files with 2402 additions and 8 deletions
+171
View File
@@ -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=<client_type>
GITEA_MCP_CLIENT_INSTANCE=<client_instance_id>
GITEA_MCP_CLIENT_SESSION=<session_id> # optional but recommended
GITEA_MCP_FLEET_RUN_ID=<enrollment 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.
+1
View File
@@ -159,6 +159,7 @@ that gates each call, not which tools exist.
- `gitea_sentry_reconcile_issue` - `gitea_sentry_reconcile_issue`
- `gitea_sentry_watchdog` - `gitea_sentry_watchdog`
- `gitea_set_issue_labels` - `gitea_set_issue_labels`
- `gitea_snapshot_instance_fleet`
- `gitea_submit_pr_review` - `gitea_submit_pr_review`
- `gitea_update_pr_branch_by_merge` - `gitea_update_pr_branch_by_merge`
- `gitea_validate_review_final_report` - `gitea_validate_review_final_report`
+3
View File
@@ -1203,6 +1203,9 @@ RECOGNIZED_GITEA_ENV_KEYS = frozenset({
"GITEA_MCP_CLIENT", "GITEA_MCP_CLIENT",
"GITEA_MCP_CLIENT_INSTANCE", "GITEA_MCP_CLIENT_INSTANCE",
"GITEA_MCP_CLIENT_SESSION", "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 # #975 review 652 B1: production also consumes HEARTBEAT_INTERVAL_ENV from
# mcp_worker_identity via gitea_mcp_server._start_worker_heartbeat. Omitting # mcp_worker_identity via gitea_mcp_server._start_worker_heartbeat. Omitting
# it reproduced the same unsupported-env → runtime_reconnect_required # it reproduced the same unsupported-env → runtime_reconnect_required
+217 -6
View File
@@ -15733,6 +15733,7 @@ _WORKER_HEARTBEAT_SUPERVISOR = None
CLIENT_NAME_ENV = "GITEA_MCP_CLIENT" CLIENT_NAME_ENV = "GITEA_MCP_CLIENT"
CLIENT_INSTANCE_ENV = "GITEA_MCP_CLIENT_INSTANCE" CLIENT_INSTANCE_ENV = "GITEA_MCP_CLIENT_INSTANCE"
CLIENT_SESSION_ENV = "GITEA_MCP_CLIENT_SESSION" CLIENT_SESSION_ENV = "GITEA_MCP_CLIENT_SESSION"
FLEET_RUN_ENV = "GITEA_MCP_FLEET_RUN_ID"
def _worker_registry(): def _worker_registry():
@@ -15756,13 +15757,34 @@ def _worker_registry():
def _client_identity_hints() -> dict: 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 { return {
"client_name": (os.environ.get(CLIENT_NAME_ENV) or "").strip() or None, "client_name": client_name,
"client_instance_id": (os.environ.get(CLIENT_INSTANCE_ENV) or "").strip() "client_instance_id": instance["client_instance_id"],
or f"pid-{os.getpid()}", "instance_id_provenance": instance["provenance"],
"instance_identity_trusted": instance["trusted"],
"session_id": (os.environ.get(CLIENT_SESSION_ENV) or "").strip() "session_id": (os.environ.get(CLIENT_SESSION_ENV) or "").strip()
or f"proc-{os.getpid()}-{_process_boot_head_sha or 'nohead'}", 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( identity = mcp_worker_identity.generate_worker_identity(
hints["client_name"], hints["session_id"] 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( outcome = registry.register(
worker_identity=identity, worker_identity=identity,
client_name=hints["client_name"], client_name=hints["client_name"],
client_instance_id=hints["client_instance_id"], client_instance_id=hints["client_instance_id"],
session_id=hints["session_id"], session_id=hints["session_id"],
generation_id=generation, generation_id=generation,
role=_active_role_kind_safe(), role=role,
profile=(os.environ.get(gitea_config.ENV_PROFILE) or "").strip() or None, profile=profile_name,
namespace=namespace,
remote=(os.environ.get("GITEA_MCP_REMOTE") or "").strip() or None, remote=(os.environ.get("GITEA_MCP_REMOTE") or "").strip() or None,
repository_binding=PROJECT_ROOT, repository_binding=PROJECT_ROOT,
pid=os.getpid(), pid=os.getpid(),
transport=native.get("bound_transport"), transport=native.get("bound_transport"),
token_fingerprint=native.get("token_fingerprint"), token_fingerprint=native.get("token_fingerprint"),
pid_alive_probe=issue_lock_store.is_process_alive, 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"): if outcome.get("registered"):
_WORKER_IDENTITY = identity _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() @mcp.tool()
def gitea_get_runtime_context( def gitea_get_runtime_context(
remote: str = "dadeschools", remote: str = "dadeschools",
+847
View File
@@ -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-<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)),
}
+97 -2
View File
@@ -301,8 +301,22 @@ CREATE INDEX IF NOT EXISTS idx_worker_session
ON worker_registrations(session_id, status); ON worker_registrations(session_id, status);
CREATE INDEX IF NOT EXISTS idx_worker_profile CREATE INDEX IF NOT EXISTS idx_worker_profile
ON worker_registrations(profile, status); 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): class WorkerRegistryError(RuntimeError):
"""Raised for registry misuse that is a programming error, not a refusal.""" """Raised for registry misuse that is a programming error, not a refusal."""
@@ -634,6 +648,17 @@ class WorkerRegistry:
conn = self._connect() conn = self._connect()
try: try:
conn.executescript(_SCHEMA_SQL) 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() conn.commit()
finally: finally:
conn.close() conn.close()
@@ -758,6 +783,14 @@ class WorkerRegistry:
heartbeat_ttl_seconds: float = DEFAULT_HEARTBEAT_TTL_SECONDS, heartbeat_ttl_seconds: float = DEFAULT_HEARTBEAT_TTL_SECONDS,
now: datetime | None = None, now: datetime | None = None,
pid_alive_probe=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]: ) -> dict[str, Any]:
"""Atomically register one worker identity. """Atomically register one worker identity.
@@ -766,6 +799,10 @@ class WorkerRegistry:
to mint a different identity and register that instead (AC32). The to mint a different identity and register that instead (AC32). The
existing registration is returned untouched so the caller can see what existing registration is returned untouched so the caller can see what
it collided with. 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) parsed = parse_worker_identity(worker_identity)
if not parsed["valid"]: if not parsed["valid"]:
@@ -783,6 +820,9 @@ class WorkerRegistry:
} }
stamp = _ts(now or _utc_now()) 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: with self._tx() as conn:
existing = conn.execute( existing = conn.execute(
"SELECT * FROM worker_registrations WHERE worker_identity = ?", "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) epoch = self._next_epoch(conn, generation_id)
conn.execute( conn.execute(
""" """
@@ -827,8 +902,11 @@ class WorkerRegistry:
generation_id, role, profile, namespace, remote, generation_id, role, profile, namespace, remote,
repository_binding, pid, transport, token_fingerprint, repository_binding, pid, transport, token_fingerprint,
started_at, last_heartbeat_at, heartbeat_ttl_seconds, started_at, last_heartbeat_at, heartbeat_ttl_seconds,
fencing_epoch, status fencing_epoch, status,
) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?) fleet_run_id, authenticated_account, process_identity,
startup_revision, loaded_revision, parity_revision,
live_revision, instance_id_provenance
) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)
""", """,
( (
worker_identity, worker_identity,
@@ -849,6 +927,14 @@ class WorkerRegistry:
float(heartbeat_ttl_seconds), float(heartbeat_ttl_seconds),
epoch, epoch,
STATUS_ACTIVE, 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( row = conn.execute(
@@ -1149,8 +1235,17 @@ def _public_record(record: dict[str, Any]) -> dict[str, Any]:
"transport": record.get("transport"), "transport": record.get("transport"),
"started_at": record.get("started_at"), "started_at": record.get("started_at"),
"last_heartbeat_at": record.get("last_heartbeat_at"), "last_heartbeat_at": record.get("last_heartbeat_at"),
"heartbeat_ttl_seconds": record.get("heartbeat_ttl_seconds"),
"fencing_epoch": record.get("fencing_epoch"), "fencing_epoch": record.get("fencing_epoch"),
"status": record.get("status"), "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"),
} }
+12
View File
@@ -163,6 +163,18 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"permission": "gitea.read", "permission": "gitea.read",
"role": "author", "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. # #644: Phase 2 Web Console recovery tasks.
"clear_stale_binding": { "clear_stale_binding": {
"permission": "gitea.read", "permission": "gitea.read",
File diff suppressed because it is too large Load Diff