Compare commits

...
Author SHA1 Message Date
sysadmin e344e68a13 fix(fleet): dedicated retirement capability, identity proof, external fencing
Addresses the three blocking findings in review 657 on PR #982 and nothing
else.

B1 — apply was authorized by gitea.read
---------------------------------------
Plan and apply shared the observational permission class, so any profile that
could look could also destroy; the mutation landing in the local control-plane
registry rather than in Gitea makes it no less a mutation.

Apply now requires its own capability, gitea.worker_registry.retire, checked at
entry and re-resolved immediately before the registry mutation. Plan stays on
gitea.read. No profile holds the new permission by default, so author,
reviewer, merger, and read-only profiles fail closed on the permission itself
rather than on the role check alone; the controller/reconciler role restriction
remains as defence in depth. Granting it is a deliberate operator edit to
profiles.json, and removing it revokes apply completely. No new Gitea write
permission is introduced and no author permission is broadened.

B2 — complete legacy rows could be retired without identity proof
-----------------------------------------------------------------
"Incomplete identity" previously meant only that pre-existing columns were
null, which a legacy-pid row satisfies trivially. Retirement now requires an
affirmative two-part proof: launcher-minted inst- attribution, plus fencing
evidence (host_id, boot_id, process_start_time) that turns a bare pid into a
statement about one process incarnation.

New mcp_process_fencing supplies those probes; every one returns None rather
than guessing, and None always preserves. Rows written before these columns
existed, and rows on legacy instance identities, are preserved permanently —
they are retired only after re-registering under a trusted identity. That
legacy rows would otherwise remain outstanding is explicitly not treated as
grounds for a weaker proof. A live pid stays an absolute block even across a
boot boundary, and pid reuse is reported distinctly from a live worker.

B3 — lease and OS liveness sat outside the registry transaction
---------------------------------------------------------------
BEGIN IMMEDIATE locks the worker registry only, and the per-target loop runs
after revalidation, so a lease acquired or a pid revived in between would go
unnoticed — the registry-column guard cannot catch it because no registry
column changed.

Two mechanisms now close that window, both applied per target immediately
before its own write: an external_fence_fn version token over active leases,
captured inside the transaction before the authoritative read and re-compared
before every guarded UPDATE (movement aborts the whole transaction; an
unreadable lease store raises rather than comparing equal), and a liveness_fn
re-probe that must affirmatively re-establish that this exact process is gone,
comparing process_start_time so a reused pid is refused. Omitting the re-probe
retires nothing rather than proceeding unfenced. The guarded UPDATE also
asserts the fencing triple is unchanged, and all three columns participate in
the CAS token.

Preserved from the accepted work: registry_fingerprint still excludes
observation time and is order- and numeric-typing stable, plan still mutates
nothing, drift still retires zero, retired rows stay historical, and ordinary
author operations remain ungated by fleet state.

Tests: tests/test_issue_980_stale_worker_retirement.py 87 passed, 46 subtests
(was 40 passed, 23 subtests), covering all three blockers' required
regressions. Adjacent suites (#980, #978, #975, #948, task-capability role
invariants) 242 passed, 156 subtests. Full suite from the branch worktree
28 failed, 6253 passed, 6 skipped, 1152 subtests; master baseline at
108cbfa173 from branches/baseline-980-108cbfa 28 failed, 6166 passed, 6
skipped, 1106 subtests. The failing sets are identical under comm, so zero
regressions and zero masked pre-existing failures; the +87 passing delta is
this branch's tests.

No live worker-registry row was retired — every test uses a throwaway SQLite
database. BAA, issue #981, PR #906, issue #650, review 624, unrelated branches
and worktrees, profiles, configuration, credentials, sessions, and running
processes were not touched.

Refs #980
Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-30 18:01:20 -04:00
sysadminandClaude Opus 4.8 c0c6d14b73 feat(fleet): add CAS-protected stale worker retirement capability
Adds a sanctioned controller/reconciler capability that retires conclusively
stale MCP worker-registry rows through a dry-run-first, compare-and-swap
protected workflow (#980).

The #978 fleet snapshot is read-only, and its registry_revision is seeded with
snapshot_at at second precision, so a CAS gated on it can never pass. That
token is left unchanged; retirement gets its own stable registry_fingerprint
derived exclusively from canonical retirement-relevant registry content, plus a
candidate_fingerprint pinning the approved set. Identical contents observed at
different times produce the same token; row order never affects it; any
retirement-relevant change moves it.

WorkerRegistry.retire_stale_workers performs the whole decision inside one
BEGIN IMMEDIATE transaction: re-read rows, recompute the fingerprint from those
rows, recompute the eligibility plan from those rows (re-reading active workflow
leases), compare the candidate fingerprint, revalidate every target, then retire
each survivor with a guarded UPDATE asserting its status, heartbeat, generation,
session, fencing epoch, and pid are unchanged. Drift retires zero workers and
reports registry_revision_moved or candidate_set_moved; any exception rolls back
and reports transaction_failed, so a partial write is never reported as success.

Retirement requires a full conjunction: active status, complete registry fields,
parsable heartbeat, not live, pid_alive false, expired heartbeat, stale
ownership, canonical repository binding, no identity evidence shared with a live
or unprobeable worker, and no active workflow-lease ownership. Everything else
is preserved with a structured reason code. Retired rows become historical
rather than stale and keep their history; nothing is deleted.

Retirement does not repair untrusted live identity: live workers on legacy
instance identities are preserved and keep their legacy_incomplete_identity
blockers, so the result never claims the fleet became safe.

Exposed as gitea_plan_stale_worker_retirement and
gitea_apply_stale_worker_retirement, restricted to controller/reconciler role
kinds. The Gitea operation gate stays gitea.read because the mutation lands in
the local control-plane registry, matching the #601 lease lifecycle; no new
Gitea write permission is introduced and no author permission is broadened.

Tests: tests/test_issue_980_stale_worker_retirement.py (40 passed, 23 subtests),
including the regression test proving the old snapshot_at derivation moved the
token one second apart while the new registry CAS token does not. Full suite
from the branch worktree matches the master baseline exactly: 28 failed / 6206
passed vs 28 failed / 6166 passed, identical failure set.

Closes #980

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-30 15:22:30 -04:00
sysadmin 108cbfa173 Merge pull request 'feat(controller): expose instance-level fleet identity and health snapshots (Closes #978)' (#979) from feat/issue-978-instance-fleet-snapshot into master 2026-07-30 09:34:33 -05:00
sysadminandClaude Opus 4.8 0dbff6dcd5 fix(controller): production client_instance_id launch path and instance-aware mutation gate
Remediate review 655 blockers on PR #979 (issue #978):

B1 — Production launcher (mcp_application_launcher) mints one trusted
client_instance_id per application launch and propagates it to all five
namespace workers via GITEA_MCP_CLIENT_INSTANCE. launcher_entry and
multi_namespace_launcher_entries use that path. Workers never invent a
trusted ID; missing/malformed/user-supplied values fail closed.

B2 — _check_mcp_runtimes_diagnostics is instance-aware: two legitimate
instances sharing a profile are allowed when each has a distinct trusted
client_instance_id; duplicate workers for the same (instance, profile)
still fail closed. Worker identity/generation exported for peer scans.

Tests cover shared ID across five namespaces, distinct launches, multi-
instance same profile, same-instance duplicates, untrusted attribution,
and the production launcher path.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-30 10:12:16 -04:00
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
sysadminandClaude Opus 4.8 b0868be6b3 fix(runtime): recognize heartbeat interval env and wire production lifecycle tests
Address REQUEST_CHANGES review #652 on PR #976 (issue #975).

Blocker 1: name GITEA_WORKER_HEARTBEAT_INTERVAL_SECONDS individually in
RECOGNIZED_GITEA_ENV_KEYS so the production-consumed operator override is no
longer classified as unsupported-env / runtime_reconnect_required. No prefix
broadening; unknown overrides remain fail-closed.

Blocker 2: add focused production-path tests that call _active_worker_identity
and _start_worker_heartbeat, proving supervisor attachment, identity/session/
generation/pid/epoch agreement, no duplicates, failed/fenced paths, orderly
shutdown, configured interval, and off-loop sqlite heartbeats.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-29 23:41:37 -04:00
jcwalker3andClaude Opus 4.8 1a38ef95e3 fix(runtime): recognize client identity environment and refresh worker registrations
Two defects left behind by #948 made sanctioned multi-client operation
impossible. They are inseparable: fixing either alone still leaves the
multi-client canary unable to run.

1. Client-identity environment keys were not recognized.

   gitea_mcp_server reads GITEA_MCP_CLIENT, GITEA_MCP_CLIENT_INSTANCE and
   GITEA_MCP_CLIENT_SESSION as the authoritative client-identity inputs for
   worker registration, but none of the three appeared in
   RECOGNIZED_GITEA_ENV_KEYS or matched a recognized prefix. The runtime
   diagnostic scans every peer mcp_server.py process environment and classifies
   any unlisted GITEA_* key as an unsupported override, which is raised as a
   runtime blocker, so gitea_resolve_task_capability returned
   blocker_kind=runtime_reconnect_required with stop_required=true. Because that
   resolver is the mandatory preflight for every author, reviewer and merger
   mutation, setting the very variable #948 requires closed the mutation gate
   for the whole fleet, and reconnecting could not clear it: the variable is
   re-exported from the client's server definition on every launch.

   The three keys are now named individually in the recognized-key set. No
   prefix is added, so an unrecognized GITEA_* override is still refused
   exactly as before.

   A related inconsistency in the same path is also fixed. The diagnostic
   reasons are raised as one RuntimeError, but the preflight re-raise
   recognized only "stale-runtime:", so an "unsupported-env:" reason was
   silently swallowed there while still failing the resolver. Both reason
   families now live in RUNTIME_DIAGNOSTIC_HARD_PREFIXES beside the function
   that produces them, and both propagate identically. This only widens what is
   refused, never what is permitted.

2. WorkerRegistry.heartbeat() had no production caller.

   #948 delivered heartbeat() but only tests called it. The single production
   writer registers once per process behind an attempted-once flag, and
   register() stamps the same timestamp into both started_at and
   last_heartbeat_at. Nothing advanced it afterwards: no lifespan hook, no
   background task, no atexit handler in a process that blocks in mcp.run().
   Since liveness is age against heartbeat_ttl_seconds, that TTL was not a
   liveness window at all but a hard cap on how long any client could stay
   attached; at 900 seconds a healthy, connected, client-managed process became
   session_ownership=unowned with blocker_kind=session_attachment_missing.

   WorkerHeartbeatSupervisor in mcp_worker_identity is the missing caller,
   started from _active_worker_identity() at the moment register() succeeds,
   because that is the only point where identity and fencing_epoch are both
   known. It is a daemon thread rather than an asyncio task or a request-driven
   refresh because renewal must survive an idle session, and because the
   registry performs blocking BEGIN IMMEDIATE sqlite writes that must not run on
   the server's event loop. daemon=True is deliberate: a hard kill takes the
   thread with it, so a dead worker still goes stale on the normal TTL.

   heartbeat_interval_for() returns one third of the TTL, hard-capped at one
   half, so two consecutive beats can be lost without the row expiring and no
   override can produce an interval that outlives the registration it renews.
   heartbeat() gains optional keyword-only expectations (session, generation,
   client name, pid); each supplied one must match the recorded row or the
   renewal is refused with the existing BLOCKER_FENCED literal rather than a new
   blocker_kind, since consumers switch on that value. Omitting them preserves
   the pre-existing behavior exactly. Client names are compared normalized, so
   several namespaces of one application stay one client while separate
   applications stay distinct.

   A terminal refusal stops the supervisor permanently and records why, so a
   fenced session can never beat its way back into ownership. A transient
   failure is counted and beating continues. An atexit hook stops it on orderly
   shutdown. status() is surfaced read-only as worker_heartbeat on
   gitea_get_runtime_context so a stopped heartbeat is diagnosable before the
   TTL turns it into session_attachment_missing; it grants nothing.

   claim_generation() still has no production caller. It bumps fencing_epoch,
   which would fence the supervisor's cached epoch, and the strict refusal is
   left in place deliberately: auto-re-adopting a bumped epoch would defeat
   fencing.

No lock or lease TTL is changed, including the author issue-lock TTL, and no
mutation refused today becomes permitted.

Tests: tests/test_issue_975_client_identity_heartbeat.py adds 40 focused tests
covering all 13 acceptance criteria. Every TTL assertion uses an injected
clock; no test waits for a real TTL. The thread-loop tests use a
millisecond-scale interval with bounded polling.

Focused: 40 passed, 15 subtests passed.
Full suite from inside the branches worktree: 28 failed, 6105 passed, 6 skipped,
1105 subtests — an identical failure set to the 324a0c8a baseline measured in a
sibling branches worktree (28 failed, 6065 passed, 1090 subtests). Zero new
failures; the delta is exactly the added tests.

Closes #975

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-29 20:06:10 -05:00
sysadmin 324a0c8a8d Merge pull request 'fix(runtime): support validated cross-repository canonical roots (#973)' (#974) from fix/issue-973-cross-repo-canonical-roots into master 2026-07-29 17:45:44 -05:00
15 changed files with 9208 additions and 62 deletions
+203
View File
@@ -0,0 +1,203 @@
# 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 **production application launcher**
(`mcp_application_launcher.build_application_mcp_servers` /
`gitea_config.multi_namespace_launcher_entries`) mints exactly one
`client_instance_id` via `mcp_fleet_snapshot.generate_client_instance_id`
and injects the same env into every namespace worker:
```text
GITEA_MCP_CLIENT=<client_type>
GITEA_MCP_CLIENT_INSTANCE=<client_instance_id>
GITEA_MCP_INSTANCE_PROVENANCE=trusted_launcher
GITEA_MCP_CLIENT_SESSION=<session_id> # optional but recommended
GITEA_MCP_FLEET_RUN_ID=<enrollment id> # when on an approved canary
GITEA_CLIENT_MANAGED=1
GITEA_MCP_PROFILE=<role-profile>
```
3. The MCP client starts the five namespace processes (`gitea-author`,
`gitea-reviewer`, `gitea-merger`, `gitea-controller`, `gitea-reconciler`)
from that generated `mcpServers` map — each entry carries the **same**
instance ID.
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`. Workers never mint a trusted instance ID themselves.
Only IDs matching the launcher format `inst-<client>-<timestamp>-<digest>`
are trusted. Missing, `legacy-pid-*` / `pid-*` placeholders, and ordinary
user-supplied strings fail closed as untrusted and cannot authorize
multi-instance fleet mutation safety.
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` (call the launcher without a prior ID).
### Mutation gate (instance-aware)
The capability / runtime diagnostic gate
(`_check_mcp_runtimes_diagnostics`) is **instance-aware**:
* Two legitimate application instances that share a profile (same
`GITEA_MCP_PROFILE`) are allowed when each has a distinct trusted
`client_instance_id`.
* More than one live worker for the same
`(client_instance_id, profile/namespace)` fails closed as a duplicate
namespace worker.
* Multiple processes sharing a profile **without** trusted instance
evidence still fail closed (indistinguishable from a duplicate).
## 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).
* [stale-worker-retirement.md](stale-worker-retirement.md) — CAS-protected
retirement of conclusively stale registrations (#980). Note that the
`registry_revision` this snapshot returns is **time-seeded** and unsuitable
for compare-and-swap; retirement derives its own stable
`registry_fingerprint` from registry content alone.
## 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.
+3
View File
@@ -51,6 +51,7 @@ that gates each call, not which tools exist.
- `gitea_adopt_merger_pr_lease`
- `gitea_adopt_workflow_lease`
- `gitea_allocate_next_work`
- `gitea_apply_stale_worker_retirement`
- `gitea_assess_already_landed_reconciliation`
- `gitea_assess_conflict_fix_classification`
- `gitea_assess_conflict_fix_push`
@@ -124,6 +125,7 @@ that gates each call, not which tools exist.
- `gitea_observability_link_issue`
- `gitea_observability_list_projects`
- `gitea_observability_reconcile_incident`
- `gitea_plan_stale_worker_retirement`
- `gitea_post_heartbeat`
- `gitea_publish_unpublished_issue_branch`
- `gitea_quarantine_contaminated_review`
@@ -159,6 +161,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`
+282
View File
@@ -0,0 +1,282 @@
# Stale worker retirement (#980)
The #978 instance-fleet snapshot made registry accuracy observable but
deliberately read-only: a registry full of rows whose owning processes are long
gone stays full. This document describes the sanctioned way to retire those
rows — a dry-run-first, compare-and-swap-protected workflow available only to
controller and reconciler namespaces.
Related: [instance-fleet-identity.md](instance-fleet-identity.md) (#978),
[post-restart-reconcile.md](post-restart-reconcile.md) (#662).
## Why a dedicated registry token
`mcp_fleet_snapshot._consistency_token` seeds its digest with `snapshot_at`,
formatted at second precision. Its output — surfaced as `registry_revision` and
`consistency_token` on the snapshot — therefore changes on **every call**, even
when no registry row changed. Any compare-and-swap gated on it can never pass:
a dry-run/apply cycle spanning more than one second aborts unconditionally.
That token remains useful as an observation stamp and is unchanged. #980 adds a
separate, *stable* token instead:
| Token | Module | Derived from | Stable across time? |
| --- | --- | --- | --- |
| `registry_revision` / `consistency_token` | `mcp_fleet_snapshot` | `snapshot_at` + a subset of row fields | **No** — moves every second |
| `registry_fingerprint` | `mcp_fleet_retirement` | canonical retirement-relevant row content only | **Yes** |
| `candidate_fingerprint` | `mcp_fleet_retirement` | canonical content of the selected candidate rows | **Yes** |
`registry_fingerprint` guarantees:
* identical canonical registry contents always produce the same token, whenever
they are observed;
* row iteration order never affects the token (serialized rows are sorted);
* any retirement-relevant change moves it — row creation or deletion, identity
change, heartbeat or TTL change, ownership change, registration-state change,
PID change, or repository-binding change.
The exact field set is `mcp_fleet_retirement.FINGERPRINT_FIELDS`. Deliberately
excluded: `token_fingerprint` (credential-adjacent, never a retirement input),
the four `*_revision` columns (revision drift is an independent restart concern
and is not part of the eligibility conjunction), and the `retired_*` bookkeeping
columns this feature adds. Numeric values are canonicalised, so a TTL that
round-trips through SQLite as `900.0` hashes identically to `900`.
## Eligibility — the conjunction
A registration is retired only when **every** one of these holds. Any missing
or contradictory evidence preserves the row.
| Requirement | Preserve reason code when it fails |
| --- | --- |
| `status` is `active` | `already_terminal_registration` |
| Every field the conjunction reads is present (`REQUIRED_IDENTITY_FIELDS`) | `incomplete_registry_identity` |
| `last_heartbeat_at` parses as a UTC stamp | `unparsable_heartbeat` |
| Worker is not live | `worker_live` |
| PID probe returns a definite answer | `pid_liveness_unknown` |
| PID probe says the process is gone | `pid_alive` |
| Heartbeat has expired under the canonical TTL | `heartbeat_not_expired` |
| `ownership_state` is exactly `stale` | `ambiguous_ownership_state` |
| Repository binding present and canonical | `repository_binding_ambiguous` |
| No identity evidence shared with a live or unprobeable worker | `conflicting_identity_evidence` |
| No other active row claims the same (instance, namespace) while one may be live | `client_instance_conflict` |
| Not an active workflow-lease owner | `protected_active_workflow_owner` |
| Instance identity is launcher-minted (`inst-…`) | `untrusted_identity_provenance` |
| Row's `host_id` matches the host running retirement | `host_binding_unproven` |
| Boot identity is known on both sides | `boot_identity_unknown` |
| The pid number is not occupied by a different incarnation | `pid_reuse_detected` |
Eligible rows carry `eligible_stale_orphan`.
Three properties are worth stating explicitly:
* **`pid_alive` can only withdraw liveness, never grant it**
(`WorkerRegistry.is_live`, #948 AC7). A heartbeat-lapsed but still-running
process therefore classifies as `stale` in the snapshot, yet #980's added
`pid_alive is False` requirement preserves it. An unprobeable PID (`None`)
also fails closed.
* **Affirmative identity proof is required (review 657 B2).** An earlier
revision required only that the pre-existing registry columns were non-null —
which a legacy `legacy-pid-…` row satisfies trivially, so a row that proved
nothing about *which* process it described was retireable. Retirement now
needs both halves of a positive proof:
* **Attribution** — a launcher-minted `inst-…` `client_instance_id`, so the
row is known to belong to one specific application launch rather than
having been inferred from pid proximity.
* **Fencing** — `host_id`, `boot_id`, and `process_start_time`, which turn a
bare pid into a statement about one process: which machine it ran on, which
boot of that machine, and which incarnation of that pid number.
**Consequence, stated plainly:** registrations written before these columns
existed, and any row on a legacy instance identity, are preserved
*permanently*. They are retired only after their worker re-registers under a
trusted identity — never on weaker evidence. That the alternative would leave
legacy rows outstanding indefinitely is not a reason to relax the proof.
* **A live pid is an absolute block.** Even across a boot boundary, where the
number provably cannot belong to the registered process, an occupied pid
preserves the row rather than being argued away by the fencing proof.
Multiple processes belonging to one legitimate worker cohort are not treated as
multiple independent workers: the fleet model from #948/#978 is preserved
unchanged, and sharing a role or profile is never a duplicate.
## Tools
### `gitea_plan_stale_worker_retirement`
Read-only. Controller and reconciler only.
| Parameter | Meaning |
| --- | --- |
| `remote` | `dadeschools` or `prgs` |
| `host`, `org`, `repo` | Optional overrides (audit context) |
| `canonical_repository` | Expected repository binding; defaults to the process root |
Returns `registry_fingerprint`, `candidate_fingerprint`,
`candidate_worker_identities`, per-worker `candidates` and `preserved` entries
(each with `reason_code`, `detail`, and structured `evidence`),
`preserved_reason_counts`, `assessed_count`, `candidate_count`,
`preserved_count`, and `protected_active_workflow_owners`.
`mutation_performed` is always `false` and `read_only` is always `true`.
Planning is deterministic: the same authoritative registry contents produce the
same plan and the same tokens regardless of when they are observed.
### `gitea_apply_stale_worker_retirement`
Mutating. Controller and reconciler only.
| Parameter | Meaning |
| --- | --- |
| `registry_fingerprint` | The exact stable token the plan returned |
| `candidate_fingerprint` | The exact candidate-set token the plan returned |
| `worker_identities` | The exact candidate identities (list, or JSON / comma-separated string) |
| `remote`, `host`, `org`, `repo` | As above |
| `canonical_repository` | Must match the value the plan used |
Before touching the registry, apply fails closed on: profile permission, role
kind, master parity (`mutation_safe`), stable-runtime mode, capability
resolution refreshed immediately before mutation, worker-registry availability,
workflow-lease enumeration failure, and daemon-cohort uniqueness
(`classify_cohort`).
A matching token is necessary but never sufficient. Inside one
`BEGIN IMMEDIATE` transaction (`WorkerRegistry.retire_stale_workers`) the
server:
1. re-reads the authoritative rows;
2. recomputes `registry_fingerprint` from *those* rows and compares — a mismatch
returns `registry_revision_moved` with `retired_count: 0` and no write;
3. recomputes the eligibility plan from *those* rows — re-reading the active
workflow leases rather than reusing the set captured before the transaction
opened, so a lease acquired after planning still preserves its worker — and
compares `candidate_fingerprint`. A mismatch returns `candidate_set_moved`
with `retired_count: 0` and no write; a lease-enumeration failure raises and
rolls the transaction back;
4. revalidates every requested identity against that fresh plan;
5. retires each survivor with a guarded `UPDATE` that additionally asserts
`status`, `last_heartbeat_at`, `generation_id`, `session_id`,
`fencing_epoch`, and `pid` are unchanged. A guard that matches no row
preserves the worker with `row_changed_since_plan`.
There is no window between a safety check and its matching write, so a worker
that comes back to life, changes ownership, or is retired concurrently cannot be
removed on the strength of a stale observation. Any exception — including a
commit failure — rolls the whole transaction back and returns
`transaction_failed` with `success: false`, `retired_count: 0`, and
`mutation_performed: false`; a partial write can never be reported as success.
### Outcomes
| `outcome` | Meaning | `mutation_performed` |
| --- | --- | --- |
| `planned` | Dry-run result | `false` |
| `applied` | Transaction ran; see `retired` / `preserved` | `true` only if something was retired |
| `registry_revision_moved` | Registry changed between plan and apply | `false` |
| `candidate_set_moved` | Eligibility verdict changed between plan and apply | `false` |
| `already_retired` | Every requested row is already retired (idempotent replay) | `false` |
| `nothing_requested` | Empty target list | `false` |
| `transaction_failed` | Rolled back; nothing retired | `false` |
## What retirement does to the fleet snapshot
A retired row keeps its history: `status` moves to `retired` and `retired_at`,
`retired_by`, `retirement_reason` are recorded. Nothing is deleted. Because
`retired` is not `active`, the #978 snapshot counts the row as **historical**,
not stale, so `stale_worker_count` falls and historical rows never make the live
fleet unsafe by themselves.
**Retirement does not repair untrusted live identity.** Live workers registered
under legacy `pid-`/`proc-` instance identities are preserved untouched and
keep their `legacy_incomplete_identity` blockers. Retiring every stale row can
therefore legitimately produce:
* `stale_worker_count: 0`
* residual live `legacy_incomplete_identity` blockers
* `live_fleet_safe: false`
That is a truthful result, and the `post_apply` block reports the remaining
blockers rather than claiming the fleet became safe. Trusted
`client_instance_id` propagation through launchers is a separate enrollment
problem.
## Permissions
Plan and apply are authorized differently, and deliberately so (review 657 B1).
| | Plan | Apply |
| --- | --- | --- |
| Capability | `gitea.read` | `gitea.worker_registry.retire` |
| Nature | observational; opens no transaction, writes nothing | mutation |
| Roles | `controller`, `reconciler` | `controller`, `reconciler` |
An earlier revision authorized apply with `gitea.read` alone, reasoning that
the mutation lands in the local control-plane registry rather than in Gitea.
That the write is local makes it **no less a mutation**: sharing an
observational permission class with plan meant any profile that could *look*
could also *destroy*. Apply now requires its own capability.
* **Denied:** author, reviewer, merger, and every ordinary read-only profile —
they lack the capability, so they fail closed on the permission itself rather
than on the role check alone. The role restriction remains as defence in
depth: a profile mistakenly granted the capability still cannot reach apply
from an author, reviewer, or merger role.
* The capability is checked at entry **and** re-resolved immediately before the
registry mutation, so a profile change mid-call cannot be outrun.
* **No new Gitea write permission is introduced.**
`gitea.worker_registry.retire` authorizes exactly one local control-plane
transition (`worker_registrations.status -> retired`) and grants no branch,
issue, PR, review, merge, or restart authority. No author permission is
broadened.
* The fleet snapshot remains observational: nothing here turns it into a gate on
ordinary author work.
### Operator step
No profile holds `gitea.worker_registry.retire` by default, so apply is inert
until an operator adds it to the `allowed_operations` of the controller or
reconciler profile in `profiles.json`. Removing it again immediately and
completely revokes apply, while leaving plan and every other capability
untouched. That grant is a configuration change and is outside the scope of the
code that implements this feature.
## External-state fencing
`BEGIN IMMEDIATE` locks the worker registry and nothing else, so two inputs the
decision depends on sit outside the transaction's isolation domain: the
workflow-lease table in a separate control-plane database, and OS process
liveness. Re-reading them once during revalidation is not sufficient — the
per-target loop runs afterwards, so a lease acquired (or a pid revived) after
revalidation but before a given row's `UPDATE` would go unnoticed, and the
registry-column guard cannot catch it because no registry column changed.
Two mechanisms close that window, both applied per target immediately before
its own write:
* **`external_fence_fn`** — a version token over active leases
(`external_state_fingerprint`), captured inside the transaction *before* the
authoritative read and re-compared before every guarded `UPDATE`. Any movement
raises, rolling back the whole transaction: once the world has changed, every
remaining per-row decision was computed against a world that no longer exists.
An unreadable lease store raises rather than returning a token, because
"unreadable" must not silently compare equal to "unchanged".
* **`liveness_fn`** — a re-probe of process liveness and fencing identity that
must affirmatively re-establish that this exact process is gone. It compares
`process_start_time`, so a pid number reused since the plan is refused rather
than accepted.
A caller that supplies no `liveness_fn` retires nothing
(`liveness_reprobe_unavailable`) rather than proceeding unfenced. The guarded
`UPDATE` additionally asserts `host_id`, `boot_id`, and `process_start_time` are
unchanged, and all three participate in the CAS token, so fencing movement
alone is enough to abort.
## Non-goals
* Killing or restarting processes.
* Editing session files or configuration.
* Direct database cleanup outside the sanctioned transaction.
* Retiring live workers.
* Backfilling trusted identity for legacy workers.
* Rewriting worker ownership.
* Cleaning unrelated workflow-lease or issue-claim registries.
+86 -15
View File
@@ -1193,6 +1193,31 @@ RECOGNIZED_GITEA_ENV_KEYS = frozenset({
"GITEA_IRRECOVERABLE_HMAC_SECRET",
"GITEA_FORCE_MCP_RUNTIME_CHECK",
"GITEA_FORCE_CLIENT_MANAGED",
# #975: the client-identity inputs the server actually consumes at startup
# (CLIENT_NAME_ENV / CLIENT_INSTANCE_ENV / CLIENT_SESSION_ENV in
# gitea_mcp_server). Production read them while this allowlist omitted them,
# so the peer-env scan classified them as unsupported overrides and the
# capability resolver refused every mutation fleet-wide. Named individually
# on purpose: no prefix is added, so an unrecognised GITEA_* override is
# still refused exactly as it was before.
"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",
# #978 B1/B2: launcher-sealed provenance + peer-visible worker identity so
# the instance-aware mutation gate can distinguish two legitimate
# application instances that share a profile from a true duplicate worker.
"GITEA_MCP_INSTANCE_PROVENANCE",
"GITEA_MCP_WORKER_IDENTITY",
"GITEA_MCP_GENERATION_ID",
# #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
# failure mode for the documented operator override. Named individually;
# no GITEA_* / GITEA_WORKER_* prefix is added.
"GITEA_WORKER_HEARTBEAT_INTERVAL_SECONDS",
})
RECOGNIZED_GITEA_ENV_PREFIXES = (
@@ -1220,24 +1245,70 @@ def get_unconsumed_gitea_env_overrides(env=None) -> dict[str, str]:
return unconsumed
def launcher_entry(profile_name, config_path=None):
def launcher_entry(
profile_name,
config_path=None,
*,
client_type=None,
client_instance_id=None,
fleet_run_id=None,
session_id=None,
launch_nonce=None,
):
"""Return a thin MCP launcher entry for *profile_name*.
Contains command/args and the GITEA_MCP_* / GITEA_CLIENT_MANAGED env vars — never a token
or password. Suitable for Claude / Gemini / Codex ``mcpServers`` blocks.
Contains command/args and the GITEA_MCP_* / GITEA_CLIENT_MANAGED env vars —
never a token or password. Suitable for Claude / Gemini / Codex
``mcpServers`` blocks.
#978 B1: production launches mint (or reuse) a trusted
``GITEA_MCP_CLIENT_INSTANCE`` so workers never fall back to the legacy
placeholder identity on the real serve path. Pass *client_instance_id* to
resume the same launch; omit it for a fresh launch (new ID).
"""
command, args = server_command()
return {
"gitea-tools": {
"command": command,
"args": args,
"env": {
"GITEA_MCP_CONFIG": config_path or DEFAULT_CONFIG_PATH,
"GITEA_MCP_PROFILE": profile_name,
"GITEA_CLIENT_MANAGED": "1",
},
}
}
import mcp_application_launcher as app_launcher
built = app_launcher.launcher_entry_for_profile(
profile_name,
client_type=client_type or "unknown",
config_path=config_path or DEFAULT_CONFIG_PATH,
client_instance_id=client_instance_id,
fleet_run_id=fleet_run_id,
session_id=session_id,
launch_nonce=launch_nonce,
server_key="gitea-tools",
)
# Public shape stays {server_key: {command, args, env}} for existing callers.
return {"gitea-tools": built["gitea-tools"]}
def multi_namespace_launcher_entries(
profile_by_namespace,
*,
client_type,
config_path=None,
client_instance_id=None,
fleet_run_id=None,
session_id=None,
launch_nonce=None,
):
"""Build production mcpServers for all five namespaces of one application launch.
One shared trusted ``client_instance_id`` is minted (or reused) and
propagated to every namespace worker env. See
:func:`mcp_application_launcher.build_application_mcp_servers`.
"""
import mcp_application_launcher as app_launcher
return app_launcher.build_application_mcp_servers(
profile_by_namespace,
client_type=client_type,
config_path=config_path or DEFAULT_CONFIG_PATH,
client_instance_id=client_instance_id,
fleet_run_id=fleet_run_id,
session_id=session_id,
launch_nonce=launch_nonce,
)
+1157 -43
View File
File diff suppressed because it is too large Load Diff
+393
View File
@@ -0,0 +1,393 @@
"""Production application launcher for multi-namespace MCP fleets (#978 B1).
One real LLM application launch mints exactly one trusted
``client_instance_id`` and propagates it to every Gitea MCP namespace worker
started for that launch. Workers never invent a trusted instance identity from
PID proximity, timestamps, or ordinary untrusted environment values.
This module is the production serve-path authority for instance identity.
Tests and fixtures may call the same functions, but production registration
receives the identity from the env this launcher builds — not from a hand-set
test-only helper that bypasses it.
"""
from __future__ import annotations
import os
import secrets
from typing import Any, Mapping
import mcp_fleet_snapshot as fleet
import mcp_worker_identity as mwi
#: Canonical MCP server names for the five role namespaces.
NAMESPACE_SERVER_NAMES: tuple[str, ...] = (
"gitea-author",
"gitea-reviewer",
"gitea-merger",
"gitea-controller",
"gitea-reconciler",
)
#: Map MCP server name → role/namespace kind.
SERVER_TO_NAMESPACE: dict[str, str] = {
"gitea-author": "author",
"gitea-reviewer": "reviewer",
"gitea-merger": "merger",
"gitea-controller": "controller",
"gitea-reconciler": "reconciler",
}
SANCTIONED_NAMESPACES: tuple[str, ...] = (
"author",
"reviewer",
"merger",
"controller",
"reconciler",
)
# Env the launcher may set on every worker of one application launch.
CLIENT_NAME_ENV = "GITEA_MCP_CLIENT"
CLIENT_INSTANCE_ENV = fleet.CLIENT_INSTANCE_ENV
FLEET_RUN_ENV = fleet.FLEET_RUN_ENV
CLIENT_SESSION_ENV = "GITEA_MCP_CLIENT_SESSION"
CLIENT_MANAGED_ENV = "GITEA_CLIENT_MANAGED"
PROFILE_ENV = "GITEA_MCP_PROFILE"
CONFIG_ENV = "GITEA_MCP_CONFIG"
INSTANCE_PROVENANCE_ENV = "GITEA_MCP_INSTANCE_PROVENANCE"
WORKER_IDENTITY_ENV = "GITEA_MCP_WORKER_IDENTITY"
GENERATION_ID_ENV = "GITEA_MCP_GENERATION_ID"
# Marker the trusted launcher alone writes; ordinary user env without this
# marker is never classified as launcher-trusted provenance.
LAUNCHER_PROVENANCE_VALUE = fleet.INSTANCE_ID_PROVENANCE_TRUSTED
def mint_application_launch(
client_type: str | None,
*,
launch_nonce: str | None = None,
fleet_run_id: str | None = None,
session_id: str | None = None,
now=None,
) -> dict[str, Any]:
"""Mint one trusted application-instance identity for a production launch.
Called exactly once per real application launch. The returned
``client_instance_id`` is injected into every namespace worker environment
for that launch. A second call (separate launch) yields a different ID.
"""
client = mwi.normalize_client_name(client_type)
instance_id = fleet.generate_client_instance_id(
client, launch_nonce=launch_nonce, now=now
)
assessment = fleet.assess_instance_identity(instance_id)
if not assessment["trusted"]:
# generate_client_instance_id always produces a trusted format; fail
# closed if that invariant ever breaks rather than shipping untrusted.
raise RuntimeError(
f"launcher produced untrusted client_instance_id {instance_id!r}: "
f"{assessment.get('reasons')}"
)
session = (session_id or "").strip() or f"launch-{secrets.token_hex(12)}"
return {
"client_type": client,
"client_instance_id": instance_id,
"instance_id_provenance": LAUNCHER_PROVENANCE_VALUE,
"instance_identity_trusted": True,
"fleet_run_id": (fleet_run_id or "").strip() or None,
"session_id": session,
"namespaces": list(SANCTIONED_NAMESPACES),
"namespace_server_names": list(NAMESPACE_SERVER_NAMES),
}
def namespace_worker_env(
*,
profile_name: str,
client_type: str | None,
client_instance_id: str,
config_path: str | None = None,
fleet_run_id: str | None = None,
session_id: str | None = None,
extra_env: Mapping[str, str] | None = None,
) -> dict[str, str]:
"""Build the environment for one namespace worker of a trusted launch.
Never invents a client_instance_id. The caller must supply the launch-minted
identity so all five workers receive the same value.
"""
assessment = fleet.assess_instance_identity(client_instance_id)
if not assessment["trusted"]:
raise ValueError(
"namespace_worker_env refuses untrusted client_instance_id "
f"{client_instance_id!r}: {assessment.get('reasons')}"
)
client = mwi.normalize_client_name(client_type)
env: dict[str, str] = {
PROFILE_ENV: str(profile_name),
CLIENT_MANAGED_ENV: "1",
"GITEA_MCP_CLIENT": client,
CLIENT_INSTANCE_ENV: assessment["client_instance_id"],
INSTANCE_PROVENANCE_ENV: LAUNCHER_PROVENANCE_VALUE,
}
if config_path:
env[CONFIG_ENV] = str(config_path)
if fleet_run_id:
env[FLEET_RUN_ENV] = str(fleet_run_id)
if session_id:
env[CLIENT_SESSION_ENV] = str(session_id)
if extra_env:
# Never let untrusted callers override the trusted instance keys.
protected = {
CLIENT_INSTANCE_ENV,
INSTANCE_PROVENANCE_ENV,
"GITEA_MCP_CLIENT",
CLIENT_MANAGED_ENV,
}
for key, value in extra_env.items():
if key in protected:
continue
env[str(key)] = str(value)
return env
def build_application_mcp_servers(
profile_by_namespace: Mapping[str, str],
*,
client_type: str | None,
config_path: str | None = None,
client_instance_id: str | None = None,
fleet_run_id: str | None = None,
session_id: str | None = None,
launch_nonce: str | None = None,
command: str | None = None,
args: list[str] | None = None,
now=None,
) -> dict[str, Any]:
"""Build a production ``mcpServers`` map for one application launch.
Mints one ``client_instance_id`` when *client_instance_id* is omitted (fresh
launch). When the caller supplies a previously minted trusted ID (resume of
the same launch / config rewrite), that ID is reused so reconnect keeps
attribution. A full new application restart omits the ID and receives a
fresh mint.
Every namespace server entry receives the **same** instance ID. Separate
calls with no supplied ID receive distinct IDs.
"""
missing = [
ns for ns in SANCTIONED_NAMESPACES if not profile_by_namespace.get(ns)
]
if missing:
raise ValueError(
"build_application_mcp_servers requires a profile for every "
f"sanctioned namespace; missing: {missing}"
)
if client_instance_id is None:
launch = mint_application_launch(
client_type,
launch_nonce=launch_nonce,
fleet_run_id=fleet_run_id,
session_id=session_id,
now=now,
)
else:
assessment = fleet.assess_instance_identity(client_instance_id)
if not assessment["trusted"]:
raise ValueError(
"refusing to propagate untrusted client_instance_id "
f"{client_instance_id!r} into production launch envs: "
f"{assessment.get('reasons')}"
)
launch = {
"client_type": mwi.normalize_client_name(client_type),
"client_instance_id": assessment["client_instance_id"],
"instance_id_provenance": LAUNCHER_PROVENANCE_VALUE,
"instance_identity_trusted": True,
"fleet_run_id": (fleet_run_id or "").strip() or None,
"session_id": (session_id or "").strip()
or f"launch-{secrets.token_hex(12)}",
"namespaces": list(SANCTIONED_NAMESPACES),
"namespace_server_names": list(NAMESPACE_SERVER_NAMES),
}
# Resolve command/args from the production server entry when not provided.
if command is None or args is None:
import gitea_config
cmd, cmd_args = gitea_config.server_command()
command = command or cmd
args = args if args is not None else list(cmd_args)
servers: dict[str, Any] = {}
shared_id = launch["client_instance_id"]
for namespace in SANCTIONED_NAMESPACES:
server_name = f"gitea-{namespace}"
profile = profile_by_namespace[namespace]
env = namespace_worker_env(
profile_name=profile,
client_type=launch["client_type"],
client_instance_id=shared_id,
config_path=config_path,
fleet_run_id=launch.get("fleet_run_id"),
session_id=launch.get("session_id"),
)
servers[server_name] = {
"command": command,
"args": list(args),
"env": env,
}
return {
"mcpServers": servers,
"launch": launch,
"client_instance_id": shared_id,
"client_type": launch["client_type"],
"namespaces": list(SANCTIONED_NAMESPACES),
"shared_instance_id_across_namespaces": True,
"namespace_count": len(SANCTIONED_NAMESPACES),
}
def launcher_entry_for_profile(
profile_name: str,
*,
client_type: str | None = None,
config_path: str | None = None,
client_instance_id: str | None = None,
fleet_run_id: str | None = None,
session_id: str | None = None,
server_key: str = "gitea-tools",
launch_nonce: str | None = None,
now=None,
) -> dict[str, Any]:
"""Thin single-server production launcher entry with trusted instance ID.
Used when only one namespace is being configured. Still mints (or reuses)
a trusted ``client_instance_id`` so production never relies on the legacy
placeholder identity for normal launches.
"""
import gitea_config
if client_instance_id is None:
launch = mint_application_launch(
client_type or "unknown",
launch_nonce=launch_nonce,
fleet_run_id=fleet_run_id,
session_id=session_id,
now=now,
)
client_instance_id = launch["client_instance_id"]
client = launch["client_type"]
fleet_run = launch.get("fleet_run_id")
session = launch.get("session_id")
else:
assessment = fleet.assess_instance_identity(client_instance_id)
if not assessment["trusted"]:
raise ValueError(
f"untrusted client_instance_id {client_instance_id!r}"
)
client = mwi.normalize_client_name(client_type)
fleet_run = (fleet_run_id or "").strip() or None
session = (session_id or "").strip() or None
client_instance_id = assessment["client_instance_id"]
command, args = gitea_config.server_command()
env = namespace_worker_env(
profile_name=profile_name,
client_type=client,
client_instance_id=client_instance_id,
config_path=config_path or gitea_config.DEFAULT_CONFIG_PATH,
fleet_run_id=fleet_run,
session_id=session,
)
return {
server_key: {
"command": command,
"args": args,
"env": env,
},
"client_instance_id": client_instance_id,
"client_type": client,
}
def collect_instance_ids_from_mcp_servers(
mcp_servers: Mapping[str, Any],
) -> dict[str, Any]:
"""Inspect a production mcpServers map for shared instance attribution.
Returns the unique set of client_instance_id values across gitea-* servers
and whether all five namespaces share exactly one trusted ID.
"""
ids: list[str] = []
by_server: dict[str, str | None] = {}
for name in NAMESPACE_SERVER_NAMES:
entry = mcp_servers.get(name) or {}
env = entry.get("env") or {}
raw = (env.get(CLIENT_INSTANCE_ENV) or "").strip() or None
by_server[name] = raw
if raw:
ids.append(raw)
unique = sorted(set(ids))
trusted = [
i
for i in unique
if fleet.assess_instance_identity(i)["trusted"]
]
return {
"server_instance_ids": by_server,
"unique_instance_ids": unique,
"trusted_instance_ids": trusted,
"shared_single_trusted_id": (
len(unique) == 1
and len(trusted) == 1
and all(by_server.get(n) == unique[0] for n in NAMESPACE_SERVER_NAMES)
),
"namespace_server_count": sum(
1 for n in NAMESPACE_SERVER_NAMES if n in mcp_servers
),
}
def inherit_or_refuse_client_instance(
env: Mapping[str, str] | None = None,
) -> dict[str, Any]:
"""Resolve instance identity for a worker process at serve time.
Production workers inherit the launcher-issued ID. They never mint a
trusted ID themselves. Missing / legacy / malformed values fail soft into
an untrusted assessment so registration can still record a diagnostic row
without authorizing multi-instance fleet mutation.
"""
source = dict(env if env is not None else os.environ)
raw = (source.get(CLIENT_INSTANCE_ENV) or "").strip() or None
provenance_marker = (source.get(INSTANCE_PROVENANCE_ENV) or "").strip()
assessment = fleet.assess_instance_identity(raw)
# Ordinary user-supplied values without launcher provenance marker are
# still format-checked by assess_instance_identity. When the format is
# trusted but the launcher marker is absent, keep the ID but note that
# provenance is not launcher-sealed (operator hand-set or legacy config).
if assessment["trusted"] and provenance_marker != LAUNCHER_PROVENANCE_VALUE:
assessment = dict(assessment)
assessment["launcher_sealed"] = False
assessment["reasons"] = list(assessment.get("reasons") or []) + [
f"{INSTANCE_PROVENANCE_ENV} is not {LAUNCHER_PROVENANCE_VALUE!r}; "
"identity format is valid but not sealed by the production launcher"
]
else:
assessment = dict(assessment)
assessment["launcher_sealed"] = bool(
assessment["trusted"]
and provenance_marker == LAUNCHER_PROVENANCE_VALUE
)
assessment["fleet_run_id"] = (source.get(FLEET_RUN_ENV) or "").strip() or None
assessment["session_id"] = (
(source.get(CLIENT_SESSION_ENV) or "").strip() or None
)
assessment["client_type"] = mwi.normalize_client_name(
(source.get("GITEA_MCP_CLIENT") or "").strip() or None
)
return assessment
+776
View File
@@ -0,0 +1,776 @@
"""CAS-protected retirement planning for stale worker registrations (#980).
#978 (merged PR #979) made the fleet observable: every registered namespace
worker, its instance attribution, its heartbeat freshness, and a structured
classification. It deliberately stopped there — the snapshot is read-only and
the control plane still had no sanctioned way to retire registry rows whose
owning process is conclusively gone.
This module is the *decision layer* for that retirement. It is pure: callers
supply registry rows, a clock, and a PID probe; nothing here opens SQLite,
scans process tables, or mutates state. The transactional apply lives in
:meth:`mcp_worker_identity.WorkerRegistry.retire_stale_workers`, which calls
back into these same pure functions so plan and apply can never disagree about
what "the registry looks like" or "which rows are eligible".
Why a separate token
--------------------
``mcp_fleet_snapshot._consistency_token`` seeds its digest with ``snapshot_at``
at second precision, so ``registry_revision`` changes on every call even when
no registry row changed. A compare-and-swap gated on it can never pass — a
dry-run/apply cycle spanning more than one second aborts unconditionally. That
token is still useful as an observation stamp, so it is left exactly as it is;
#980 gets its own :func:`registry_fingerprint`, derived *only* from canonical
retirement-relevant row content:
* identical registry contents observed at any two times produce the same token,
* row order never affects the token (serialized rows are sorted),
* any create/delete/identity/liveness/ownership/registration-state change to a
retirement-relevant field changes the token.
Fail-closed posture
-------------------
A worker is retired only when the control plane *conclusively* establishes it
is a stale orphan. Missing evidence is never read as permission: an unprobeable
PID, an unparsable heartbeat, a row that shares identity evidence with a live
or unprobeable worker, a foreign or absent repository binding, or a worker that
still owns an active workflow lease all preserve the row.
Trusted launcher identity (``inst-…`` provenance) is deliberately *not* part of
the conjunction. #980 places trusted ``client_instance_id`` propagation out of
scope and lists "backfilling trusted identity for legacy workers" as a non-goal;
requiring it here would preserve every legacy row forever and make the feature
inert. What *is* required is that the registry fields the conjunction reads are
actually present — see :data:`REQUIRED_IDENTITY_FIELDS`.
"""
from __future__ import annotations
import hashlib
from datetime import datetime
from typing import Any, Callable, Iterable, Mapping, Sequence
import mcp_fleet_snapshot as fleet
import mcp_worker_identity as mwi
# --- Outcomes -------------------------------------------------------------
OUTCOME_PLANNED = "planned"
OUTCOME_APPLIED = "applied"
OUTCOME_REGISTRY_MOVED = "registry_revision_moved"
OUTCOME_CANDIDATES_MOVED = "candidate_set_moved"
OUTCOME_ALREADY_RETIRED = "already_retired"
OUTCOME_NOTHING_REQUESTED = "nothing_requested"
# --- Reason codes ---------------------------------------------------------
#: The only reason code that authorizes retirement.
REASON_ELIGIBLE = "eligible_stale_orphan"
REASON_ALREADY_TERMINAL = "already_terminal_registration"
REASON_AMBIGUOUS_OWNERSHIP = "ambiguous_ownership_state"
REASON_CONFLICTING_IDENTITY = "conflicting_identity_evidence"
REASON_FOREIGN_REPOSITORY = "repository_binding_ambiguous"
REASON_HEARTBEAT_FRESH = "heartbeat_not_expired"
REASON_INCOMPLETE_IDENTITY = "incomplete_registry_identity"
REASON_NOT_IN_PLAN = "not_in_current_plan"
REASON_PID_ALIVE = "pid_alive"
REASON_PID_UNKNOWN = "pid_liveness_unknown"
REASON_PROTECTED_OWNER = "protected_active_workflow_owner"
REASON_ROW_CHANGED = "row_changed_since_plan"
REASON_ROW_MISSING = "registration_missing"
REASON_UNPARSABLE_HEARTBEAT = "unparsable_heartbeat"
REASON_WORKER_LIVE = "worker_live"
# --- #980 review 657 B2: affirmative identity/liveness proof ---------------
#: The registration's instance identity is not launcher-minted (``inst-…``),
#: so nothing proves which application launch this row belongs to.
REASON_UNTRUSTED_PROVENANCE = "untrusted_identity_provenance"
#: The row does not record which host its pid belongs to, or records a
#: different host than the one probing. A local pid probe cannot speak for a
#: process on another machine.
REASON_HOST_UNPROVEN = "host_binding_unproven"
#: Boot identity is missing on the row or unobtainable here, so a recorded pid
#: cannot be compared against a live pid at all.
REASON_BOOT_UNKNOWN = "boot_identity_unknown"
#: The pid is alive but belongs to a different process incarnation than the one
#: registered — reported distinctly from a plain live worker.
REASON_PID_REUSED = "pid_reuse_detected"
#: Two active registrations claim one client instance within one namespace.
REASON_INSTANCE_CONFLICT = "client_instance_conflict"
#: The immediate pre-write re-probe could not re-establish death (#980 B3).
REASON_LIVENESS_REPROBE = "liveness_reprobe_refused"
#: Instance-identity prefix minted by the trusted launcher. Kept in sync with
#: ``mcp_fleet_snapshot._TRUSTED_INSTANCE_PREFIX`` through
#: :func:`mcp_fleet_snapshot.assess_instance_identity`, which stays the single
#: authority on what "trusted" means — this module never re-implements it.
TRUSTED_INSTANCE_PREFIX = "inst-"
#: Registry columns that must carry a usable value before the eligibility
#: conjunction can even be evaluated. Absence is ambiguity, not permission.
#:
#: #980 review 657 B2 added the fencing triple. Before it, "complete identity"
#: meant only that the pre-existing columns were non-null, which a legacy
#: ``legacy-pid-…`` row satisfies trivially — so a row that proved nothing about
#: *which* process it described was retireable. The triple is what makes a
#: recorded pid interpretable: which machine it ran on, which boot of that
#: machine, and which incarnation of that pid number. A registration written
#: before these columns existed carries NULL and is therefore preserved
#: permanently, which is the intended fail-closed outcome.
REQUIRED_IDENTITY_FIELDS: tuple[str, ...] = (
"worker_identity",
"client_instance_id",
"session_id",
"generation_id",
"status",
"started_at",
"last_heartbeat_at",
"heartbeat_ttl_seconds",
"pid",
"host_id",
"boot_id",
"process_start_time",
)
#: Canonical retirement-relevant content. Ordering here is fixed and part of
#: the token contract; adding a field changes every fingerprint, so a change
#: here is a deliberate contract revision.
#:
#: Deliberately excluded: ``token_fingerprint`` (credential-adjacent, never a
#: retirement input), the four ``*_revision`` columns (revision drift is an
#: independent restart concern and is not part of the eligibility conjunction),
#: and the ``retired_*`` bookkeeping columns this feature adds.
FINGERPRINT_FIELDS: tuple[str, ...] = (
"worker_identity",
"client_name",
"client_instance_id",
"session_id",
"generation_id",
"role",
"profile",
"namespace",
"remote",
"repository_binding",
"pid",
"process_identity",
"transport",
"started_at",
"last_heartbeat_at",
"heartbeat_ttl_seconds",
"fencing_epoch",
"status",
"fleet_run_id",
"authenticated_account",
"instance_id_provenance",
# #980 review 657 B3: fencing evidence is a retirement input, so moving it
# must move the CAS token. Without these, a row whose host, boot, or
# process incarnation changed would hash identically to the row the plan
# approved.
"host_id",
"boot_id",
"process_start_time",
)
_FINGERPRINT_VERSION = "registryfp-v1"
_CANDIDATE_VERSION = "candidatefp-v1"
_UNIT = "\x1f"
_RECORD = "\x1e"
def _canon(value: Any) -> str:
"""Stable text for one field value, independent of Python/SQLite typing.
``900`` and ``900.0`` are the same TTL and must hash the same; a value that
round-trips through SQLite as REAL must not produce a different token than
the same value supplied by a caller as ``int``.
"""
if value is None:
return ""
if isinstance(value, bool):
return "true" if value else "false"
if isinstance(value, float):
if value != value or value in (float("inf"), float("-inf")):
return repr(value)
if value.is_integer():
return str(int(value))
return repr(value)
if isinstance(value, int):
return str(value)
return str(value)
def _serialize_row(row: Mapping[str, Any]) -> str:
return _UNIT.join(f"{name}={_canon(row.get(name))}" for name in FINGERPRINT_FIELDS)
def _digest(version: str, serialized: Sequence[str], prefix: str) -> str:
ordered = sorted(serialized)
material = _RECORD.join([version, str(len(ordered)), *ordered])
return f"{prefix}-{hashlib.sha256(material.encode('utf-8')).hexdigest()[:32]}"
def registry_fingerprint(rows: Iterable[Mapping[str, Any]]) -> str:
"""Content-derived compare-and-swap token for the worker registry (#980).
Derived exclusively from :data:`FINGERPRINT_FIELDS` across every row. It
contains no ``snapshot_at``, wall-clock, request, or report-generation
time, so two observations of an unchanged registry always agree, and the
serialized rows are sorted so iteration order cannot perturb the digest.
"""
return _digest(
_FINGERPRINT_VERSION,
[_serialize_row(row) for row in rows],
"registryfp",
)
def candidate_fingerprint(candidate_rows: Iterable[Mapping[str, Any]]) -> str:
"""Exact-candidate-set token over the selected rows' canonical content.
A matching :func:`registry_fingerprint` already implies these rows are
unchanged; this second token additionally pins *which* rows the operator
approved, so an apply can never widen or narrow the approved set.
"""
return _digest(
_CANDIDATE_VERSION,
[_serialize_row(row) for row in candidate_rows],
"candidatefp",
)
def _probe_pid(
pid: Any, pid_alive_probe: Callable[[int | None], bool | None] | None
) -> bool | None:
if pid_alive_probe is None or pid is None:
return None
try:
result = pid_alive_probe(pid)
except Exception:
return None
return None if result is None else bool(result)
def _reuse_detected(
row: Mapping[str, Any],
live_start_time: str | None,
current_host_id: str | None,
current_boot_id: str | None,
) -> bool:
"""Is the pid occupied by a *different* incarnation than the one recorded?
Only meaningful when the recorded pid is comparable to the live one — same
machine, same boot. Across hosts or boots the number is unrelated by
construction and reuse is not the interesting question.
"""
recorded_host = (row.get("host_id") or "").strip()
recorded_boot = (row.get("boot_id") or "").strip()
if not current_host_id or recorded_host != current_host_id:
return False
if not current_boot_id or recorded_boot != current_boot_id:
return False
recorded_start = (row.get("process_start_time") or "").strip()
return bool(live_start_time) and live_start_time != recorded_start
def _instance_key(row: Mapping[str, Any]) -> tuple[str, str] | None:
"""The (instance, namespace) pair #978 requires to be unique among live rows."""
instance = _canon(row.get("client_instance_id"))
namespace = _canon(row.get("namespace"))
if not instance or not namespace:
return None
return (instance, namespace)
def _probe_start_time(
pid: Any, start_time_probe: Callable[[Any], str | None] | None
) -> str | None:
if start_time_probe is None or pid is None:
return None
try:
return start_time_probe(pid)
except Exception:
return None
def _evidence(
row: Mapping[str, Any],
snapshot_row: Mapping[str, Any],
pid_alive: bool | None,
) -> dict[str, Any]:
liveness = snapshot_row.get("liveness") or {}
return {
"worker_identity": row.get("worker_identity"),
"client_type": snapshot_row.get("client_type"),
"client_instance_id": row.get("client_instance_id"),
"fleet_run_id": row.get("fleet_run_id"),
"namespace": row.get("namespace"),
"profile": row.get("profile"),
"declared_role": row.get("role"),
"session_id": row.get("session_id"),
"generation_id": row.get("generation_id"),
"fencing_epoch": row.get("fencing_epoch"),
"process_identity": snapshot_row.get("process_identity"),
"pid": row.get("pid"),
"pid_alive": pid_alive,
"host_id": row.get("host_id"),
"boot_id": row.get("boot_id"),
"process_start_time": row.get("process_start_time"),
"repository_binding": row.get("repository_binding"),
"foreign_repository": bool(snapshot_row.get("foreign_repository")),
"status": row.get("status"),
"started_at": row.get("started_at"),
"last_heartbeat_at": row.get("last_heartbeat_at"),
"heartbeat_ttl_seconds": row.get("heartbeat_ttl_seconds"),
"heartbeat_age_seconds": liveness.get("heartbeat_age_seconds"),
"heartbeat_fresh": liveness.get("heartbeat_fresh"),
"live": bool(snapshot_row.get("live")),
"ownership_state": snapshot_row.get("ownership_state"),
"instance_id_provenance": snapshot_row.get("instance_id_provenance"),
"instance_identity_trusted": bool(
snapshot_row.get("instance_identity_trusted")
),
}
def _conflict_keys(
row: Mapping[str, Any], snapshot_row: Mapping[str, Any]
) -> list[tuple[str, str]]:
keys: list[tuple[str, str]] = []
for name, value in (
("session_id", row.get("session_id")),
("generation_id", row.get("generation_id")),
("process_identity", snapshot_row.get("process_identity")),
("pid", row.get("pid")),
):
text = _canon(value)
if text:
keys.append((name, text))
return keys
def external_state_fingerprint(
leases: Iterable[Mapping[str, Any]],
*,
liveness: Iterable[tuple[Any, Any]] = (),
) -> str:
"""Version token over the external state a retirement decision consumed.
#980 review 657 B3: ``BEGIN IMMEDIATE`` on the worker registry does not
cover the control-plane lease table or the OS process table, so those
inputs can move while the transaction is open. This token lets the
transaction detect that movement: it is captured before the authoritative
read and re-compared immediately before every guarded write, and any
difference aborts rather than retiring against evidence that has changed.
Only ownership-relevant lease fields participate, so unrelated churn (a
heartbeat timestamp advancing on an unrelated lease) does not cause
spurious aborts, while an acquire, release, or owner change always does.
"""
lease_units: list[str] = []
for lease in leases:
lease_units.append(
_UNIT.join(
f"{name}={_canon(lease.get(name))}"
for name in (
"lease_id",
"role",
"target",
"status",
"session_id",
"owner_session_id",
"owner_pid",
"session_pid",
"generation",
)
)
)
for pid, alive in liveness:
lease_units.append(f"liveness{_UNIT}pid={_canon(pid)}{_UNIT}alive={_canon(alive)}")
return _digest("externalfp-v1", lease_units, "externalfp")
def assess_retirement_identity_proof(
row: Mapping[str, Any],
snapshot_row: Mapping[str, Any],
*,
current_host_id: str | None,
current_boot_id: str | None,
live_start_time: str | None,
pid_alive: bool | None,
) -> tuple[str, str] | None:
"""Affirmative proof that this row names one specific, now-dead process.
Returns ``None`` when the proof holds, or ``(reason_code, detail)`` naming
the first thing that could not be established. #980 review 657 B2: absence
of evidence is never read as staleness, so every branch here refuses on
*missing* information exactly as firmly as on contradictory information.
The proof has two independent halves and needs both:
* **Attribution** — a launcher-minted ``inst-…`` instance identity, so the
row is known to belong to one specific application launch rather than
having been inferred from pid proximity.
* **Fencing** — the row's host matches the host doing the probing, boot
identity is known on both sides, and the recorded process incarnation
agrees with whatever currently occupies that pid number.
Requiring trusted attribution means pre-#978 ``legacy-pid-…`` rows are
preserved permanently. That is deliberate. The reviewer specifically
rejected the argument that legacy rows "would remain forever" as grounds
for a weaker proof, and #980 lists backfilling trusted identity for legacy
workers as a non-goal — so those rows are retired only after their worker
re-registers under a trusted identity, never on weaker evidence.
"""
if not snapshot_row.get("instance_identity_trusted"):
return (
REASON_UNTRUSTED_PROVENANCE,
"client_instance_id "
f"{row.get('client_instance_id')!r} is not launcher-minted "
f"({snapshot_row.get('instance_id_provenance')!r}); nothing proves "
"which application launch this registration belongs to",
)
recorded_host = (row.get("host_id") or "").strip()
if not current_host_id:
return (
REASON_HOST_UNPROVEN,
"this process cannot establish its own host identity, so a local "
"pid probe cannot be attributed to any machine",
)
if recorded_host != current_host_id:
return (
REASON_HOST_UNPROVEN,
f"registration is bound to host {recorded_host!r} but retirement is "
f"running on {current_host_id!r}; a local pid probe says nothing "
"about a process on another machine",
)
recorded_boot = (row.get("boot_id") or "").strip()
if not current_boot_id:
return (
REASON_BOOT_UNKNOWN,
"the current boot identity could not be determined, so a recorded "
"pid cannot be compared against a live pid",
)
recorded_start = (row.get("process_start_time") or "").strip()
if recorded_boot != current_boot_id:
# A different boot is the strongest possible death evidence: every pid
# from a previous boot is gone, and pid numbers restart, so whatever
# occupies this number now is unrelated by construction.
return None
# Same boot: the pid number is directly comparable, so the recorded
# incarnation must still agree with whatever holds that number.
if pid_alive and live_start_time and live_start_time != recorded_start:
return (
REASON_PID_REUSED,
f"pid {row.get('pid')!r} is alive but started at "
f"{live_start_time!r}, not the registered {recorded_start!r}; the "
"number was reused by an unrelated process and this registration's "
"own liveness is therefore unproven",
)
if pid_alive:
return (
REASON_PID_ALIVE,
f"recorded pid {row.get('pid')!r} is still running on this host and "
"boot",
)
return None
def _missing_identity_fields(row: Mapping[str, Any]) -> list[str]:
missing: list[str] = []
for name in REQUIRED_IDENTITY_FIELDS:
value = row.get(name)
if value is None or (isinstance(value, str) and not value.strip()):
missing.append(name)
return missing
def plan_stale_worker_retirement(
rows: Iterable[Mapping[str, Any]],
*,
now: datetime | None = None,
pid_alive_probe: Callable[[int | None], bool | None] | None = None,
canonical_repository: str | None = None,
protected_worker_identities: Iterable[str] | None = None,
protected_session_ids: Iterable[str] | None = None,
protected_pids: Iterable[Any] | None = None,
current_host_id: str | None = None,
current_boot_id: str | None = None,
start_time_probe: Callable[[Any], str | None] | None = None,
) -> dict[str, Any]:
"""Decide, without mutating anything, which registrations may be retired.
Every row lands in exactly one of ``candidates`` (eligible) or
``preserved`` (with the reason code that stopped it), so the output
explains the whole registry rather than only the interesting part.
"""
all_rows = [dict(row) for row in rows]
protected_ids = {str(w) for w in (protected_worker_identities or []) if w}
protected_sessions = {str(s) for s in (protected_session_ids or []) if s}
protected_pid_set = {_canon(p) for p in (protected_pids or []) if p is not None}
snapshots: dict[int, dict[str, Any]] = {}
pid_alive_by_index: dict[int, bool | None] = {}
start_time_by_index: dict[int, str | None] = {}
for index, row in enumerate(all_rows):
pid_alive = _probe_pid(row.get("pid"), pid_alive_probe)
pid_alive_by_index[index] = pid_alive
start_time_by_index[index] = _probe_start_time(
row.get("pid"), start_time_probe
)
snapshots[index] = fleet.build_worker_snapshot_row(
row,
now=now,
pid_alive_probe=(lambda _pid, _value=pid_alive: _value),
canonical_repository=canonical_repository,
)
# Identity evidence owned by a worker that is live, or whose liveness could
# not be established, is ambiguous: anything sharing it is preserved.
ambiguous_keys: set[tuple[str, str]] = set()
identity_counts: dict[str, int] = {}
for index, row in enumerate(all_rows):
identity = _canon(row.get("worker_identity"))
if identity:
identity_counts[identity] = identity_counts.get(identity, 0) + 1
snapshot_row = snapshots[index]
liveness = snapshot_row.get("liveness") or {}
unresolved = (
pid_alive_by_index[index] is None
or liveness.get("heartbeat_fresh") is None
)
if snapshot_row.get("live") or (
str(row.get("status") or "") == mwi.STATUS_ACTIVE and unresolved
):
ambiguous_keys.update(_conflict_keys(row, snapshot_row))
# Two active registrations claiming one (client_instance_id, namespace)
# violate the #978 uniqueness invariant — but only when one of them might
# still be running. Several *dead* rows accumulating on one slot across
# restarts is ordinary history and every one of them is safely retirable;
# a slot shared with a live or unprobeable worker is genuinely ambiguous,
# because which row that process belongs to cannot be settled from the
# registry alone.
#
# The key is deliberately the (instance, namespace) pair, not the instance
# alone: one legitimate cohort is exactly one instance spread across
# distinct namespaces, so keying on the instance would make every cohort
# look self-conflicting and preserve the whole fleet forever.
instance_members: dict[tuple[str, str], list[int]] = {}
for index, row in enumerate(all_rows):
if str(row.get("status") or "") != mwi.STATUS_ACTIVE:
continue
key = _instance_key(row)
if key is None:
continue
instance_members.setdefault(key, []).append(index)
instance_conflicts: dict[tuple[str, str], bool] = {}
for key, members in instance_members.items():
if len(members) < 2:
continue
contested = any(
snapshots[i].get("live") or pid_alive_by_index[i] is None
for i in members
)
if contested:
instance_conflicts[key] = True
candidates: list[dict[str, Any]] = []
candidate_rows: list[Mapping[str, Any]] = []
preserved: list[dict[str, Any]] = []
for index, row in enumerate(all_rows):
snapshot_row = snapshots[index]
pid_alive = pid_alive_by_index[index]
liveness = snapshot_row.get("liveness") or {}
evidence = _evidence(row, snapshot_row, pid_alive)
blocked: tuple[str, str] | None = None
identity = _canon(row.get("worker_identity"))
missing = _missing_identity_fields(row)
shared = sorted(
f"{name}={value}"
for name, value in _conflict_keys(row, snapshot_row)
if (name, value) in ambiguous_keys
)
binding = (row.get("repository_binding") or "").strip()
protected_hits: list[str] = []
if identity and identity in protected_ids:
protected_hits.append(f"worker_identity={identity}")
if _canon(row.get("session_id")) in protected_sessions:
protected_hits.append(f"session_id={_canon(row.get('session_id'))}")
if _canon(row.get("pid")) in protected_pid_set:
protected_hits.append(f"pid={_canon(row.get('pid'))}")
if identity and identity_counts.get(identity, 0) > 1:
blocked = (
REASON_CONFLICTING_IDENTITY,
f"worker identity {identity!r} appears on more than one registry row",
)
elif str(row.get("status") or "") != mwi.STATUS_ACTIVE:
blocked = (
REASON_ALREADY_TERMINAL,
f"registration status is {row.get('status')!r}; nothing to retire",
)
elif missing:
blocked = (
REASON_INCOMPLETE_IDENTITY,
"registry row is missing field(s) the retirement conjunction "
f"reads: {missing}",
)
elif mwi._parse_ts(row.get("last_heartbeat_at")) is None:
blocked = (
REASON_UNPARSABLE_HEARTBEAT,
"last_heartbeat_at is not a parsable UTC stamp; liveness is unknown",
)
elif snapshot_row.get("live"):
blocked = (REASON_WORKER_LIVE, "worker is live and must not be retired")
elif pid_alive is None:
blocked = (
REASON_PID_UNKNOWN,
f"pid {row.get('pid')!r} could not be probed; liveness is unproven",
)
elif pid_alive and _reuse_detected(
row, start_time_by_index[index], current_host_id, current_boot_id
):
# Reported before the generic live-pid branch so the operator sees
# *why* the number is occupied: an unrelated process inherited it,
# which means this registration's own liveness is unproven rather
# than positively established.
blocked = (
REASON_PID_REUSED,
f"pid {row.get('pid')!r} is alive but started at "
f"{start_time_by_index[index]!r}, not the registered "
f"{row.get('process_start_time')!r}; the number was reused",
)
elif pid_alive:
blocked = (
REASON_PID_ALIVE,
f"recorded pid {row.get('pid')!r} is still running",
)
elif liveness.get("heartbeat_fresh") is not False:
blocked = (
REASON_HEARTBEAT_FRESH,
"heartbeat has not expired under the canonical TTL policy",
)
elif snapshot_row.get("ownership_state") != "stale":
blocked = (
REASON_AMBIGUOUS_OWNERSHIP,
"ownership_state is "
f"{snapshot_row.get('ownership_state')!r}, not 'stale'",
)
elif not binding or snapshot_row.get("foreign_repository"):
blocked = (
REASON_FOREIGN_REPOSITORY,
"repository binding is absent or does not match the canonical "
"repository; retirement scope is ambiguous",
)
elif shared:
blocked = (
REASON_CONFLICTING_IDENTITY,
"identity evidence is shared with a live or unprobeable worker: "
f"{shared}",
)
elif protected_hits:
blocked = (
REASON_PROTECTED_OWNER,
"worker still owns active workflow state requiring separate "
f"reconciliation: {protected_hits}",
)
elif instance_conflicts.get(_instance_key(row)):
blocked = (
REASON_INSTANCE_CONFLICT,
"another active registration claims client_instance_id "
f"{row.get('client_instance_id')!r} in namespace "
f"{row.get('namespace')!r}; instance ownership is ambiguous",
)
else:
# Affirmative identity + fencing proof runs last: everything above
# establishes the row is *inert*, and this establishes it is
# unambiguously *this* worker (#980 review 657 B2).
blocked = assess_retirement_identity_proof(
row,
snapshot_row,
current_host_id=current_host_id,
current_boot_id=current_boot_id,
live_start_time=start_time_by_index[index],
pid_alive=pid_alive,
)
if blocked is not None:
preserved.append(
{
"worker_identity": row.get("worker_identity"),
"reason_code": blocked[0],
"detail": blocked[1],
"evidence": evidence,
}
)
continue
candidates.append(
{
"worker_identity": row.get("worker_identity"),
"reason_code": REASON_ELIGIBLE,
"detail": (
"dead pid, expired heartbeat, stale ownership, unambiguous "
"identity, canonical repository binding, no active workflow "
"ownership"
),
"evidence": evidence,
}
)
candidate_rows.append(row)
counts: dict[str, int] = {}
for entry in preserved:
counts[entry["reason_code"]] = counts.get(entry["reason_code"], 0) + 1
return {
"success": True,
"read_only": True,
"mutation_performed": False,
"outcome": OUTCOME_PLANNED,
"registry_fingerprint": registry_fingerprint(all_rows),
"candidate_fingerprint": candidate_fingerprint(candidate_rows),
"assessed_count": len(all_rows),
"candidate_count": len(candidates),
"preserved_count": len(preserved),
"candidates": candidates,
"candidate_worker_identities": [c["worker_identity"] for c in candidates],
"preserved": preserved,
"preserved_reason_counts": counts,
"protected_inputs": {
"worker_identities": sorted(protected_ids),
"session_ids": sorted(protected_sessions),
"pids": sorted(protected_pid_set),
},
"canonical_repository": canonical_repository,
"fencing_context": {
"current_host_id": current_host_id,
"current_boot_id": current_boot_id,
"start_time_probe_available": start_time_probe is not None,
},
}
def summarize_plan(plan: Mapping[str, Any]) -> dict[str, Any]:
"""Compact, log-safe view of a plan or apply result."""
return {
"outcome": plan.get("outcome"),
"registry_fingerprint": plan.get("registry_fingerprint"),
"candidate_fingerprint": plan.get("candidate_fingerprint"),
"assessed_count": plan.get("assessed_count"),
"candidate_count": plan.get("candidate_count"),
"retired_count": plan.get("retired_count"),
"preserved_count": plan.get("preserved_count"),
"mutation_performed": plan.get("mutation_performed"),
}
+869
View File
@@ -0,0 +1,869 @@
"""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"
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)),
}
+205
View File
@@ -0,0 +1,205 @@
"""Host, boot, and process-start fencing evidence for retirement safety (#980).
Review 657 B2/B3 established that a local ``os.kill(pid, 0)`` probe is not, on
its own, evidence that a *particular registered worker* is gone:
* **Host ambiguity.** A registry row written on host A records pid 1234. Probing
pid 1234 on host B answers a question nobody asked. "Not running here" is not
"not running".
* **PID reuse.** Pid 1234 may be alive and belong to an unrelated process that
the kernel handed the number to after the original exited. The naive probe
reads that as "the worker is live" (safe, over-preserving) — but the converse
matters too: evidence captured about pid 1234 at time T must not be honoured
at time T+n if the process behind that number changed in between.
* **Boot boundaries.** Every pid from a previous boot is conclusively gone, and
pid numbers restart, so a recorded pid is only comparable to a live pid when
both belong to the same boot.
This module supplies the three pieces of evidence that turn a bare pid into a
statement about one specific process:
``host_id``
Which machine the pid belongs to.
``boot_id``
Which boot of that machine the pid belongs to. Pids are only comparable
within a single boot.
``process_start_time``
Which *incarnation* of that pid number. Two processes on the same host and
boot sharing a pid number cannot share a start time, so comparing start
times defeats reuse.
Every probe here is read-only, never raises, and returns ``None`` when the
evidence cannot be established. ``None`` means *unknown*, and the retirement
conjunction is required to treat unknown as "preserve", never as "safe".
"""
from __future__ import annotations
import os
import platform
import subprocess
# --- Host identity --------------------------------------------------------
def current_host_id() -> str | None:
"""Stable identifier for the machine this process runs on.
Deliberately the kernel node name rather than anything network-derived: it
does not change when an interface goes down or a VPN reassigns an address,
and retirement must not become unsafe because DNS moved.
"""
try:
node = (platform.node() or "").strip()
except Exception:
return None
return node or None
# --- Boot identity --------------------------------------------------------
def _linux_boot_id() -> str | None:
try:
with open("/proc/sys/kernel/random/boot_id", encoding="utf-8") as handle:
value = handle.read().strip()
except Exception:
return None
return value or None
def _darwin_boot_id() -> str | None:
"""macOS boot identity, derived from ``kern.boottime``.
``sysctl`` prints e.g. ``{ sec = 1785400000, usec = 123456 } Wed Jul 30 ...``.
Only the integer seconds are kept: the trailing human-readable date is
locale-dependent and would make the identifier unstable across environments
for the very same boot.
"""
try:
completed = subprocess.run(
["/usr/sbin/sysctl", "-n", "kern.boottime"],
capture_output=True,
text=True,
timeout=5,
check=False,
)
except Exception:
return None
if completed.returncode != 0:
return None
text = (completed.stdout or "").strip()
marker = "sec = "
start = text.find(marker)
if start < 0:
return None
tail = text[start + len(marker) :]
digits = ""
for char in tail:
if char.isdigit():
digits += char
else:
break
return f"boot-{digits}" if digits else None
def current_boot_id() -> str | None:
"""Identifier for the current boot, or ``None`` when it cannot be proven."""
system = ""
try:
system = (platform.system() or "").strip().lower()
except Exception:
system = ""
if system == "linux":
return _linux_boot_id()
if system == "darwin":
return _darwin_boot_id()
# An unrecognised platform yields no boot evidence rather than a guess.
return None
# --- Process start time ---------------------------------------------------
def _linux_process_start_time(pid: int) -> str | None:
"""Field 22 of ``/proc/<pid>/stat`` — start time in clock ticks since boot.
The executable name in field 2 is parenthesised and may itself contain
spaces and parentheses, so the fields are located from the *last* ``)``
rather than by splitting the whole line.
"""
try:
with open(f"/proc/{pid}/stat", encoding="utf-8") as handle:
raw = handle.read()
except Exception:
return None
close = raw.rfind(")")
if close < 0:
return None
fields = raw[close + 1 :].split()
# After the ')' the next field is state (field 3), so field 22 is index 19.
if len(fields) < 20:
return None
value = fields[19].strip()
return f"ticks-{value}" if value else None
def _darwin_process_start_time(pid: int) -> str | None:
"""macOS process start, from ``ps -o lstart=``.
``lstart`` is the absolute wall-clock start of that pid's current
incarnation. Two processes reusing one pid number report different values,
which is exactly the discrimination reuse detection needs.
"""
try:
completed = subprocess.run(
["/bin/ps", "-o", "lstart=", "-p", str(int(pid))],
capture_output=True,
text=True,
timeout=5,
check=False,
)
except Exception:
return None
if completed.returncode != 0:
return None
value = " ".join((completed.stdout or "").split())
return f"lstart-{value}" if value else None
def process_start_time(pid: int | None) -> str | None:
"""Start-time token for *pid*'s current incarnation, or ``None``.
``None`` is returned both when the pid does not exist and when the platform
cannot answer. Callers must not read either case as evidence of death — a
dead pid is established by the liveness probe, and this value only ever
*withdraws* a retirement that pid-level evidence would otherwise allow.
"""
if pid is None:
return None
try:
numeric = int(pid)
except (TypeError, ValueError):
return None
if numeric <= 0:
return None
system = ""
try:
system = (platform.system() or "").strip().lower()
except Exception:
system = ""
if system == "linux":
return _linux_process_start_time(numeric)
if system == "darwin":
return _darwin_process_start_time(numeric)
return None
def current_process_fencing(pid: int | None = None) -> dict[str, str | None]:
"""The full fencing triple for *pid* (defaults to this process)."""
target = os.getpid() if pid is None else pid
return {
"host_id": current_host_id(),
"boot_id": current_boot_id(),
"process_start_time": process_start_time(target),
}
+909 -3
View File
File diff suppressed because it is too large Load Diff
+50
View File
@@ -163,6 +163,56 @@ 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",
},
# #980: CAS-protected retirement of conclusively stale worker
# registrations.
#
# Planning is observational and stays on ``gitea.read``: it opens no
# transaction, writes nothing, and returns only what a fleet snapshot
# already exposes to the same roles.
#
# Applying is a mutation and review 657 B1 established that ``gitea.read``
# cannot authorize it. The mutation landing in the local control-plane
# registry rather than in Gitea makes it *no less* a mutation, and sharing
# an observational permission class with plan meant any profile that could
# look could also destroy. It now requires its own permission,
# ``gitea.worker_registry.retire``, which no profile holds by default — so
# author, reviewer, merger, and ordinary read-only profiles fail closed on
# the permission itself rather than relying on the role check alone. The
# role restriction (controller/reconciler), runtime parity, cohort
# uniqueness, and the exact registry + candidate fingerprints all remain,
# and are now defence in depth behind the capability rather than a
# substitute for it.
#
# Granting the permission is a deliberate operator act in profiles.json;
# removing it from a profile immediately and completely revokes apply.
"plan_stale_worker_retirement": {
"permission": "gitea.read",
"role": "controller",
},
"gitea_plan_stale_worker_retirement": {
"permission": "gitea.read",
"role": "controller",
},
"apply_stale_worker_retirement": {
"permission": "gitea.worker_registry.retire",
"role": "controller",
},
"gitea_apply_stale_worker_retirement": {
"permission": "gitea.worker_registry.retire",
"role": "controller",
},
# #644: Phase 2 Web Console recovery tasks.
"clear_stale_binding": {
"permission": "gitea.read",
+14 -1
View File
@@ -175,8 +175,21 @@ class TestLauncherSnippets(unittest.TestCase):
def test_only_safe_keys_no_secrets(self):
entry = gitea_config.launcher_entry("prgs", "/cfg/profiles.json")["gitea-tools"]
self.assertEqual(set(entry), {"command", "args", "env"})
self.assertEqual(set(entry["env"]), {"GITEA_MCP_CONFIG", "GITEA_MCP_PROFILE", "GITEA_CLIENT_MANAGED"})
# #978 B1: production launcher also injects trusted client identity.
required = {
"GITEA_MCP_CONFIG",
"GITEA_MCP_PROFILE",
"GITEA_CLIENT_MANAGED",
"GITEA_MCP_CLIENT",
"GITEA_MCP_CLIENT_INSTANCE",
"GITEA_MCP_INSTANCE_PROVENANCE",
}
self.assertTrue(required.issubset(set(entry["env"])), entry["env"])
self.assertEqual(entry["env"]["GITEA_MCP_PROFILE"], "prgs")
self.assertTrue(
entry["env"]["GITEA_MCP_CLIENT_INSTANCE"].startswith("inst-"),
entry["env"]["GITEA_MCP_CLIENT_INSTANCE"],
)
blob = json.dumps(entry).lower()
for word in ("token", "password", "secret"):
self.assertNotIn(word, blob)
File diff suppressed because it is too large Load Diff
File diff suppressed because it is too large Load Diff
File diff suppressed because it is too large Load Diff