Compare commits

..
2 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
9 changed files with 4104 additions and 3 deletions
+5
View File
@@ -189,6 +189,11 @@ mutation. Diagnostic reads remain available where `gitea.read` allows.
* `mcp_fleet_snapshot` — pure snapshot + classification (#978). * `mcp_fleet_snapshot` — pure snapshot + classification (#978).
* `gitea_snapshot_instance_fleet` — sanctioned MCP tool (#978). * `gitea_snapshot_instance_fleet` — sanctioned MCP tool (#978).
* `gitea_get_runtime_context` — single-process view (not fleet-wide). * `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 ## Non-goals
+2
View File
@@ -51,6 +51,7 @@ that gates each call, not which tools exist.
- `gitea_adopt_merger_pr_lease` - `gitea_adopt_merger_pr_lease`
- `gitea_adopt_workflow_lease` - `gitea_adopt_workflow_lease`
- `gitea_allocate_next_work` - `gitea_allocate_next_work`
- `gitea_apply_stale_worker_retirement`
- `gitea_assess_already_landed_reconciliation` - `gitea_assess_already_landed_reconciliation`
- `gitea_assess_conflict_fix_classification` - `gitea_assess_conflict_fix_classification`
- `gitea_assess_conflict_fix_push` - `gitea_assess_conflict_fix_push`
@@ -124,6 +125,7 @@ that gates each call, not which tools exist.
- `gitea_observability_link_issue` - `gitea_observability_link_issue`
- `gitea_observability_list_projects` - `gitea_observability_list_projects`
- `gitea_observability_reconcile_incident` - `gitea_observability_reconcile_incident`
- `gitea_plan_stale_worker_retirement`
- `gitea_post_heartbeat` - `gitea_post_heartbeat`
- `gitea_publish_unpublished_issue_branch` - `gitea_publish_unpublished_issue_branch`
- `gitea_quarantine_contaminated_review` - `gitea_quarantine_contaminated_review`
+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.
+695
View File
@@ -19623,6 +19623,701 @@ def gitea_snapshot_instance_fleet(
return snapshot return snapshot
# --- #980 CAS-protected stale worker retirement ---------------------------
#: #980 review 657 B1. Apply is a mutation and must not be authorized by the
#: observational ``gitea.read`` class. Declared once here and consumed by both
#: the entry gate and the immediately-pre-mutation re-check so the two can
#: never drift apart; ``task_capability_map`` maps the apply tasks to the same
#: literal.
RETIREMENT_MUTATION_PERMISSION = "gitea.worker_registry.retire"
def _retirement_role_block(task: str) -> dict | None:
"""Refuse the retirement surface to any non-controller/reconciler role.
Role is a *secondary* control. The primary authority for apply is the
dedicated ``gitea.worker_registry.retire`` capability (#980 review 657 B1);
this check additionally pins the surface to the two roles that own fleet
reconciliation, so a profile mistakenly granted the permission still cannot
reach it from an author, reviewer, or merger role. Plan remains
``gitea.read`` and is observational.
"""
profile = get_profile()
role = _profile_role_kind(profile)
if role in {"controller", "reconciler"}:
return None
return {
"success": False,
"allowed": False,
"mutation_performed": False,
"retired_count": 0,
"denied_role": role,
"required_roles": ["controller", "reconciler"],
"requested_task": task,
"reasons": [
f"{task} is restricted to controller and reconciler roles; active "
f"role_kind is {role!r}. Author, reviewer, and merger profiles keep "
"gitea.read for diagnosis elsewhere but never receive worker-"
"registry retirement, and no unrelated mutation permission is "
"granted."
],
"exact_next_action": (
"Re-run from a prgs-controller or prgs-reconciler namespace."
),
}
def _retirement_runtime_block() -> list[str]:
"""Fail-closed runtime reasons that must stop a retirement apply (#980).
``gitea.read`` deliberately bypasses the #420 parity gate and the #615
stable-runtime gate, because a stale server may still be *inspected*. An
apply is a mutation, so both gates are re-asserted explicitly here rather
than inherited.
"""
reasons: list[str] = []
try:
parity = _current_master_parity()
except Exception as exc:
return [
f"master parity could not be assessed (fail closed): {_redact(str(exc))}"
]
if not parity.get("mutation_safe"):
reasons.append(
"runtime parity is not mutation-safe: "
f"{parity.get('summary') or 'stale runtime'}"
)
try:
reasons.extend(
stable_control_runtime.runtime_block_reasons(
_current_runtime_mode_report()
)
)
except Exception as exc:
reasons.append(
f"runtime mode could not be assessed (fail closed): {_redact(str(exc))}"
)
return reasons
def _retirement_protected_owners() -> tuple[dict | None, dict]:
"""Workflow owners that must never be retired as stale workers.
A registration whose process still owns an active control-plane lease is a
live workflow participant needing its own reconciliation, not a stale
orphan. Failure to enumerate leases is ambiguity, so it aborts rather than
proceeding with an empty protection set.
"""
db, errs = _control_plane_db_or_error()
if db is None:
return (
{
"success": False,
"mutation_performed": False,
"retired_count": 0,
"reasons": [
"active workflow leases could not be enumerated, so "
"protected owners are unknown (fail closed)",
*errs,
],
},
{},
)
try:
leases = db.list_leases(statuses=["active"], limit=1000)
except Exception as exc:
return (
{
"success": False,
"mutation_performed": False,
"retired_count": 0,
"reasons": [
"active workflow leases could not be enumerated, so "
"protected owners are unknown (fail closed): "
f"{_redact(str(exc))}"
],
},
{},
)
session_ids: set[str] = set()
pids: set[int] = set()
for lease in leases:
owner = lease.get("session_id") or lease.get("owner_session_id")
if owner:
session_ids.add(str(owner))
for key in ("owner_pid", "session_pid"):
value = lease.get(key)
if value is None:
continue
try:
pids.add(int(value))
except (TypeError, ValueError):
continue
return None, {
"session_ids": sorted(session_ids),
"pids": sorted(pids),
"active_lease_count": len(leases),
}
def _retirement_plan(records, canonical_repository: str, protected: dict) -> dict:
"""The single planning path shared by dry run and in-transaction revalidation."""
import mcp_fleet_retirement
import mcp_process_fencing
return mcp_fleet_retirement.plan_stale_worker_retirement(
records,
pid_alive_probe=issue_lock_store.is_process_alive,
canonical_repository=canonical_repository,
protected_session_ids=(protected or {}).get("session_ids"),
protected_pids=(protected or {}).get("pids"),
# #980 review 657 B2: the fencing evidence that makes a recorded pid
# interpretable. Probed live so a plan produced on one host can never
# authorize a retirement carried out on another.
current_host_id=mcp_process_fencing.current_host_id(),
current_boot_id=mcp_process_fencing.current_boot_id(),
start_time_probe=mcp_process_fencing.process_start_time,
)
def _retirement_external_fence() -> str:
"""Version token over lease + liveness state consumed by a retirement (#980 B3).
Recomputed inside the retirement transaction immediately before every
guarded write. Because neither the control-plane lease database nor the OS
process table is covered by the registry's ``BEGIN IMMEDIATE``, this token
is the only thing that can detect either of them moving mid-transaction.
A failure to read leases raises rather than returning a token: an
unreadable external world is indistinguishable from a changed one, so it
must abort the transaction rather than silently compare equal.
"""
import mcp_fleet_retirement
db, errs = _control_plane_db_or_error()
if db is None:
raise RuntimeError(
"active workflow leases could not be read while fencing the "
f"retirement transaction: {errs}"
)
leases = db.list_leases(statuses=["active"], limit=1000)
return mcp_fleet_retirement.external_state_fingerprint(leases)
def _retirement_liveness_reprobe(row: dict) -> dict:
"""Immediate pre-write re-establishment that one specific process is gone.
#980 review 657 B3: revalidation happens once per transaction, but the
per-target loop runs afterwards, so this re-probe closes the remaining
window between "this row was judged safe" and "this row is written". It
re-reads OS state rather than trusting the plan, and it compares the
fencing triple so a pid number reused since the plan cannot pass.
"""
import mcp_fleet_retirement
import mcp_process_fencing
pid = row.get("pid")
try:
alive = issue_lock_store.is_process_alive(pid)
except Exception as exc:
return {
"safe": False,
"reason_code": mcp_fleet_retirement.REASON_PID_UNKNOWN,
"detail": (
f"pid {pid!r} could not be re-probed immediately before the "
f"write: {_redact(str(exc))}"
),
}
if alive:
live_start = mcp_process_fencing.process_start_time(pid)
recorded_start = (row.get("process_start_time") or "").strip()
reused = bool(live_start) and live_start != recorded_start
return {
"safe": False,
"reason_code": (
mcp_fleet_retirement.REASON_PID_REUSED
if reused
else mcp_fleet_retirement.REASON_PID_ALIVE
),
"detail": (
f"pid {pid!r} is alive at write time"
+ (
f" but is a different incarnation ({live_start!r} != "
f"{recorded_start!r})"
if reused
else ""
)
),
"evidence": {
"pid": pid,
"pid_alive": True,
"live_process_start_time": live_start,
"recorded_process_start_time": recorded_start,
},
}
host_now = mcp_process_fencing.current_host_id()
boot_now = mcp_process_fencing.current_boot_id()
if not host_now or (row.get("host_id") or "").strip() != host_now:
return {
"safe": False,
"reason_code": mcp_fleet_retirement.REASON_HOST_UNPROVEN,
"detail": (
"host identity no longer agrees with the registration at write "
f"time (row {row.get('host_id')!r}, now {host_now!r})"
),
}
if not boot_now:
return {
"safe": False,
"reason_code": mcp_fleet_retirement.REASON_BOOT_UNKNOWN,
"detail": "boot identity became unobtainable before the write",
}
return {
"safe": True,
"evidence": {
"pid": pid,
"pid_alive": False,
"host_id": host_now,
"boot_id": boot_now,
},
}
def _retirement_revalidation_plan(records, canonical_repository: str) -> dict:
"""Revalidation planner used *inside* the retirement transaction.
Deliberately re-reads the active workflow leases rather than reusing the
set captured before the transaction opened: a lease acquired after planning
must still preserve its worker. An enumeration failure raises, which rolls
the transaction back and retires nothing.
"""
block, protected = _retirement_protected_owners()
if block:
raise RuntimeError(
"active workflow leases could not be re-read inside the retirement "
"transaction; refusing to retire anything"
)
return _retirement_plan(records, canonical_repository, protected)
@mcp.tool()
def gitea_plan_stale_worker_retirement(
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
canonical_repository: str | None = None,
) -> dict:
"""Read-only: authoritative retirement plan for stale worker rows (#980).
Controller and reconciler only. Reads the worker registry, classifies every
registration with the same assessor the #978 fleet snapshot uses, and
returns the exact set of registrations that are conclusively stale orphans
together with a **stable** ``registry_fingerprint`` and an exact
``candidate_fingerprint``.
The fingerprint is derived only from canonical retirement-relevant registry
content never from ``snapshot_at``, wall-clock, request, or report time
so two plans over an unchanged registry agree and the apply compare-and-swap
can actually pass. Row order never affects it.
Retires nothing. Live workers, workers whose PID cannot be probed, workers
with unparsable heartbeats, workers sharing identity evidence with a live or
unprobeable worker, foreign or unbound repositories, and workers that still
own an active workflow lease are all preserved with a reason code.
Args:
remote: Known instance 'dadeschools' or 'prgs'.
host: Optional host override.
org: Optional org override (audit context only).
repo: Optional repo override (audit context only).
canonical_repository: Expected repository binding for
foreign-repository classification (defaults to the process root).
"""
read_block = _profile_operation_gate("gitea.read")
if read_block:
return {
"success": False,
"read_only": True,
"mutation_performed": False,
"reasons": read_block,
"permission_report": _permission_block_report("gitea.read"),
}
role_block = _retirement_role_block("gitea_plan_stale_worker_retirement")
if role_block:
role_block["read_only"] = True
return role_block
registry = _worker_registry()
if registry is None:
return {
"success": False,
"read_only": True,
"mutation_performed": False,
"reasons": [
"worker registry is unavailable; cannot produce an authoritative "
"retirement plan (fail closed)"
],
"exact_next_action": (
"Ensure GITEA_WORKER_REGISTRY_DB is writable and re-run after "
"workers have registered."
),
}
protected_block, protected = _retirement_protected_owners()
if protected_block:
protected_block["read_only"] = True
return protected_block
try:
records = registry.list_workers(status=None)
except Exception as exc:
return {
"success": False,
"read_only": True,
"mutation_performed": False,
"reasons": [
f"failed to read worker registry: {type(exc).__name__}: "
f"{_redact(str(exc))}"
],
}
canon = canonical_repository or PROJECT_ROOT
plan = _retirement_plan(records, canon, protected)
profile = get_profile()
plan["role_kind"] = _profile_role_kind(profile)
plan["profile"] = profile.get("profile_name")
plan["remote"] = _effective_remote(remote)
plan["repository"] = {"org": org, "repo": repo, "canonical_repository": canon}
plan["protected_active_workflow_owners"] = protected
plan["apply_tool"] = "gitea_apply_stale_worker_retirement"
plan["permission_scope"] = {
"read_only": True,
"granted_operations": ["gitea.read"],
"denied_unrelated_mutations": True,
"note": (
"Planning is strictly observational. It does not authorize branch, "
"issue, PR, review, merge, or restart mutations, and it retires "
"nothing."
),
}
plan["exact_next_action"] = (
"Pass registry_fingerprint, candidate_fingerprint, and the exact "
"candidate_worker_identities to gitea_apply_stale_worker_retirement."
if plan.get("candidate_count")
else "No registration is conclusively stale; nothing to apply."
)
return plan
@mcp.tool()
def gitea_apply_stale_worker_retirement(
registry_fingerprint: str,
candidate_fingerprint: str,
worker_identities: list | str,
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
canonical_repository: str | None = None,
) -> dict:
"""Retire conclusively stale worker registrations under CAS (#980).
Controller and reconciler only. Requires the exact ``registry_fingerprint``
and ``candidate_fingerprint`` returned by
``gitea_plan_stale_worker_retirement`` plus the exact candidate identity
list. A matching token is necessary but never sufficient: inside a single
``BEGIN IMMEDIATE`` transaction the registry is re-read, the fingerprint is
recomputed from those rows, the eligibility plan is recomputed from those
rows, and every target is independently revalidated immediately before its
own guarded ``UPDATE``. Any drift retires zero workers and reports
``registry_revision_moved`` or ``candidate_set_moved``.
Retiring stale rows does not repair untrusted live identity. Live workers
registered under legacy ``pid-``/``proc-`` instance identities remain
untouched and their ``legacy_incomplete_identity`` blockers remain
outstanding, so the result never claims the fleet became safe.
Args:
registry_fingerprint: Stable token from the plan (CAS expectation).
candidate_fingerprint: Exact candidate-set token from the plan.
worker_identities: The exact candidate identities the plan returned
(list, or a JSON / comma-separated string).
remote: Known instance 'dadeschools' or 'prgs'.
host: Optional host override.
org: Optional org override (audit context only).
repo: Optional repo override (audit context only).
canonical_repository: Expected repository binding (defaults to the
process root); must match the value the plan used.
"""
import json as _json
import mcp_fleet_retirement
# #980 review 657 B1: the dedicated mutation capability, not gitea.read.
# A profile holding only the observational permission cannot get past this
# line, and the role check below is defence in depth rather than the sole
# authority.
mutation_block = _profile_operation_gate(RETIREMENT_MUTATION_PERMISSION)
if mutation_block:
return {
"success": False,
"mutation_performed": False,
"retired_count": 0,
"requested_task": "apply_stale_worker_retirement",
"required_operation_permission": RETIREMENT_MUTATION_PERMISSION,
"reasons": mutation_block,
"permission_report": _permission_block_report(
RETIREMENT_MUTATION_PERMISSION
),
}
role_block = _retirement_role_block("gitea_apply_stale_worker_retirement")
if role_block:
return role_block
runtime_reasons = _retirement_runtime_block()
if runtime_reasons:
return {
"success": False,
"mutation_performed": False,
"retired_count": 0,
"blocker_kind": "runtime_not_mutation_safe",
"reasons": runtime_reasons,
"exact_next_action": (
"Restore runtime parity on the stable control checkout, then "
"re-plan and re-apply."
),
}
# Fresh identity + capability resolution immediately before mutation.
profile = get_profile()
role_kind = _profile_role_kind(profile)
required_permission = task_capability_map.required_permission(
"apply_stale_worker_retirement"
)
required_role = task_capability_map.required_role(
"apply_stale_worker_retirement"
)
permission_ok, permission_reason = gitea_config.check_operation(
required_permission,
profile.get("allowed_operations") or [],
profile.get("forbidden_operations") or [],
)
if not permission_ok:
return {
"success": False,
"mutation_performed": False,
"retired_count": 0,
"requested_task": "apply_stale_worker_retirement",
"required_operation_permission": required_permission,
"required_role_kind": required_role,
"reasons": [
"capability resolution immediately before apply refused this "
f"session: {permission_reason}"
],
}
try:
resolved_host = host or REMOTES[_effective_remote(remote)]["host"]
authenticated_username = _authenticated_username(resolved_host)
except Exception:
authenticated_username = None
registry = _worker_registry()
if registry is None:
return {
"success": False,
"mutation_performed": False,
"retired_count": 0,
"reasons": [
"worker registry is unavailable; refusing to retire anything "
"(fail closed)"
],
}
raw = worker_identities
if isinstance(raw, str):
text = raw.strip()
try:
raw = _json.loads(text)
except Exception:
raw = [part.strip() for part in text.split(",") if part.strip()]
if isinstance(raw, str):
raw = [raw]
if not isinstance(raw, list):
return {
"success": False,
"mutation_performed": False,
"retired_count": 0,
"reasons": ["worker_identities must be a list of worker identities"],
}
targets = [str(item).strip() for item in raw if str(item).strip()]
protected_block, protected = _retirement_protected_owners()
if protected_block:
return protected_block
# #948 daemon-cohort uniqueness: a contested generation or a reused worker
# identity means ownership is ambiguous fleet-wide, so retire nothing.
try:
cohort = mcp_worker_identity.classify_cohort(
registry.list_workers(status=mcp_worker_identity.STATUS_ACTIVE),
pid_alive_probe=issue_lock_store.is_process_alive,
)
except Exception as exc:
return {
"success": False,
"mutation_performed": False,
"retired_count": 0,
"reasons": [
"daemon-cohort uniqueness could not be assessed (fail closed): "
f"{_redact(str(exc))}"
],
}
if cohort.get("blocked"):
return {
"success": False,
"mutation_performed": False,
"retired_count": 0,
"blocker_kind": cohort.get("blocker_kind"),
"reasons": [
"daemon-cohort uniqueness failed; worker ownership is contested",
*(cohort.get("reasons") or []),
],
"cohort": {
"blocked_worker_identities": cohort.get("blocked_worker_identities"),
"duplicate_identities": cohort.get("duplicate_identities"),
"contested_generations": cohort.get("contested_generations"),
},
}
canon = canonical_repository or PROJECT_ROOT
acting = "/".join(
part
for part in (authenticated_username, profile.get("profile_name"))
if part
) or "unknown"
result = registry.retire_stale_workers(
expected_registry_fingerprint=registry_fingerprint,
expected_candidate_fingerprint=candidate_fingerprint,
worker_identities=targets,
fingerprint_fn=mcp_fleet_retirement.registry_fingerprint,
plan_fn=lambda rows: _retirement_revalidation_plan(rows, canon),
retired_by=acting,
retirement_reason=mcp_fleet_retirement.REASON_ELIGIBLE,
external_fence_fn=_retirement_external_fence,
liveness_fn=_retirement_liveness_reprobe,
)
result["role_kind"] = role_kind
result["profile"] = profile.get("profile_name")
result["remote"] = _effective_remote(remote)
result["repository"] = {"org": org, "repo": repo, "canonical_repository": canon}
result["protected_active_workflow_owners"] = protected
result["permission_scope"] = {
"granted_operations": [RETIREMENT_MUTATION_PERMISSION],
"mutation_capability": RETIREMENT_MUTATION_PERMISSION,
"plan_capability": "gitea.read",
"control_plane_mutation": "worker_registrations.status -> retired",
"denied_unrelated_mutations": True,
"note": (
"This capability retires local worker-registry rows only. It grants "
"no branch, issue, PR, review, merge, or restart authority, and it "
"never kills or restarts a process."
),
}
result["post_apply"] = _retirement_post_apply(
registry, result, canonical_repository=canon
)
try:
gitea_audit.write_event(
gitea_audit.build_event(
action="gitea_apply_stale_worker_retirement",
result=(
gitea_audit.SUCCEEDED
if result.get("mutation_performed")
else gitea_audit.BLOCKED
if not result.get("success")
else gitea_audit.ALLOWED
),
remote=_effective_remote(remote),
repository=canon,
profile_name=profile.get("profile_name"),
audit_label=profile.get("audit_label"),
authenticated_username=authenticated_username,
task_role=role_kind,
operation="worker_registry.retire_stale_workers",
reason=result.get("outcome"),
request_metadata=mcp_fleet_retirement.summarize_plan(result),
)
)
except Exception:
pass
return result
def _retirement_post_apply(
registry, result: dict, *, canonical_repository: str
) -> dict:
"""Fresh fleet verification after a retirement attempt (#980 requirement 6)."""
import mcp_fleet_retirement
import mcp_fleet_snapshot as _fleet_after
try:
after_records = registry.list_workers(status=None)
after = _fleet_after.snapshot_instance_fleet(
after_records,
pid_alive_probe=issue_lock_store.is_process_alive,
canonical_repository=canonical_repository,
)
except Exception as exc:
return {
"available": False,
"reasons": [
f"post-apply verification could not be produced: {_redact(str(exc))}"
],
}
retired_ids = {
str(r.get("worker_identity")) for r in result.get("retired") or []
}
return {
"available": True,
"registry_fingerprint": mcp_fleet_retirement.registry_fingerprint(
after_records
),
"live_worker_count": after.get("live_worker_count"),
"stale_worker_count": after.get("stale_worker_count"),
"historical_worker_count": after.get("historical_worker_count"),
"retired_still_counted_live": sorted(
str(w.get("worker_identity"))
for w in after.get("live_workers") or []
if str(w.get("worker_identity")) in retired_ids
),
"retired_still_counted_stale": sorted(
str(w.get("worker_identity"))
for w in after.get("stale_workers") or []
if str(w.get("worker_identity")) in retired_ids
),
"live_fleet_safe": after.get("live_fleet_safe"),
"remaining_blockers": [
{"classification": f.get("classification"), "detail": f.get("detail")}
for f in after.get("active_blockers") or []
],
"note": (
"Stale retirement does not repair untrusted live identity; residual "
"legacy_incomplete_identity blockers keep live_fleet_safe false and "
"that is a truthful result."
),
}
@mcp.tool() @mcp.tool()
def gitea_get_runtime_context( def gitea_get_runtime_context(
remote: str = "dadeschools", remote: str = "dadeschools",
+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"),
}
+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),
}
+428 -3
View File
@@ -39,7 +39,7 @@ import threading
import time import time
from contextlib import contextmanager from contextlib import contextmanager
from datetime import datetime, timezone from datetime import datetime, timezone
from typing import Any, Iterator from typing import Any, Callable, Iterator, Sequence
# --- Provenance verdicts ------------------------------------------------- # --- Provenance verdicts -------------------------------------------------
@@ -224,6 +224,13 @@ def _heartbeat_expectation_drift(
STATUS_ACTIVE = "active" STATUS_ACTIVE = "active"
STATUS_SUPERSEDED = "superseded" STATUS_SUPERSEDED = "superseded"
STATUS_RELEASED = "released" STATUS_RELEASED = "released"
#: #980 terminal state for a registration whose owning process is conclusively
#: gone. Distinct from ``released`` (the worker said goodbye) and
#: ``superseded`` (a newer generation took over): ``retired`` records that the
#: *control plane* concluded the row was a stale orphan and retired it under a
#: compare-and-swap. Like every non-active status it is not live, so a retired
#: row counts as historical rather than stale in the #978 fleet snapshot.
STATUS_RETIRED = "retired"
_TRUE_VALUES = frozenset({"1", "true", "yes", "client_managed"}) _TRUE_VALUES = frozenset({"1", "true", "yes", "client_managed"})
_FALSE_VALUES = frozenset({"0", "false", "no", "manual", "manual_launch"}) _FALSE_VALUES = frozenset({"0", "false", "no", "manual", "manual_launch"})
@@ -315,13 +322,53 @@ _SCHEMA_OPTIONAL_COLUMNS: tuple[tuple[str, str], ...] = (
("parity_revision", "TEXT"), ("parity_revision", "TEXT"),
("live_revision", "TEXT"), ("live_revision", "TEXT"),
("instance_id_provenance", "TEXT"), ("instance_id_provenance", "TEXT"),
# #980 review 657 B2/B3: fencing evidence that turns a bare pid into a
# statement about one specific process. A registration written before these
# columns existed carries NULL and can never satisfy the retirement identity
# proof, which is the intended fail-closed outcome — absence of evidence is
# not evidence of staleness.
("host_id", "TEXT"),
("boot_id", "TEXT"),
("process_start_time", "TEXT"),
# #980 retirement bookkeeping. Deliberately outside the CAS fingerprint
# field set: they record *that* a retirement happened, and are written only
# by the retirement transaction itself.
("retired_at", "TEXT"),
("retired_by", "TEXT"),
("retirement_reason", "TEXT"),
) )
def _default_process_fencing(pid: int | None) -> dict[str, str | None]:
"""Fencing triple for *pid*, degrading to unknowns rather than raising.
Imported lazily so this storage module keeps no import-time dependency on
the probe layer, and so a platform where the probes are unavailable still
registers workers — it simply records no fencing evidence, and those rows
are then permanently ineligible for retirement.
"""
try:
import mcp_process_fencing
return mcp_process_fencing.current_process_fencing(pid)
except Exception:
return {"host_id": None, "boot_id": None, "process_start_time": None}
class WorkerRegistryError(RuntimeError): class WorkerRegistryError(RuntimeError):
"""Raised for registry misuse that is a programming error, not a refusal.""" """Raised for registry misuse that is a programming error, not a refusal."""
class _ExternalStateMoved(RuntimeError):
"""External safety evidence changed inside the retirement transaction.
Raised so the surrounding ``with self._tx()`` rolls back: once lease state
or process liveness has moved, every remaining per-row decision was computed
against a world that no longer exists, so the whole attempt is abandoned
rather than partially applied.
"""
def _utc_now() -> datetime: def _utc_now() -> datetime:
return datetime.now(timezone.utc) return datetime.now(timezone.utc)
@@ -791,6 +838,9 @@ class WorkerRegistry:
parity_revision: str | None = None, parity_revision: str | None = None,
live_revision: str | None = None, live_revision: str | None = None,
instance_id_provenance: str | None = None, instance_id_provenance: str | None = None,
host_id: str | None = None,
boot_id: str | None = None,
process_start_time: str | None = None,
) -> dict[str, Any]: ) -> dict[str, Any]:
"""Atomically register one worker identity. """Atomically register one worker identity.
@@ -823,6 +873,20 @@ class WorkerRegistry:
proc_id = process_identity or ( proc_id = process_identity or (
f"pid-{int(pid)}" if pid is not None else None f"pid-{int(pid)}" if pid is not None else None
) )
# #980 B2/B3: capture the fencing triple for the pid being registered.
# A caller may supply it (tests, or a launcher that already probed);
# otherwise it is probed here, at the only moment the process is known
# to be the one that owns this registration. Any probe that cannot
# answer stores NULL, which permanently withholds retirement eligibility
# from the row rather than granting it on absent evidence.
fencing = _default_process_fencing(pid)
host_id = host_id if host_id is not None else fencing["host_id"]
boot_id = boot_id if boot_id is not None else fencing["boot_id"]
process_start_time = (
process_start_time
if process_start_time is not None
else fencing["process_start_time"]
)
with self._tx() as conn: with self._tx() as conn:
existing = conn.execute( existing = conn.execute(
"SELECT * FROM worker_registrations WHERE worker_identity = ?", "SELECT * FROM worker_registrations WHERE worker_identity = ?",
@@ -905,8 +969,9 @@ class WorkerRegistry:
fencing_epoch, status, fencing_epoch, status,
fleet_run_id, authenticated_account, process_identity, fleet_run_id, authenticated_account, process_identity,
startup_revision, loaded_revision, parity_revision, startup_revision, loaded_revision, parity_revision,
live_revision, instance_id_provenance live_revision, instance_id_provenance,
) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?) host_id, boot_id, process_start_time
) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)
""", """,
( (
worker_identity, worker_identity,
@@ -935,6 +1000,9 @@ class WorkerRegistry:
(parity_revision or "").strip() or None, (parity_revision or "").strip() or None,
(live_revision or "").strip() or None, (live_revision or "").strip() or None,
(instance_id_provenance or "").strip() or None, (instance_id_provenance or "").strip() or None,
(host_id or "").strip() or None,
(boot_id or "").strip() or None,
(process_start_time or "").strip() or None,
), ),
) )
row = conn.execute( row = conn.execute(
@@ -1217,6 +1285,363 @@ class WorkerRegistry:
"reasons": [], "reasons": [],
} }
def retire_stale_workers(
self,
*,
expected_registry_fingerprint: str,
expected_candidate_fingerprint: str,
worker_identities: Sequence[str],
fingerprint_fn: Callable[[list[dict[str, Any]]], str],
plan_fn: Callable[[list[dict[str, Any]]], dict[str, Any]],
retired_by: str | None = None,
retirement_reason: str = "",
now: datetime | None = None,
external_fence_fn: Callable[[], str] | None = None,
liveness_fn: Callable[[dict[str, Any]], dict[str, Any]] | None = None,
) -> dict[str, Any]:
"""Compare-and-swap retirement of conclusively stale registrations (#980).
The whole decision happens inside one ``BEGIN IMMEDIATE`` transaction:
the authoritative rows are re-read, the stable registry fingerprint is
recomputed from *those* rows, the eligibility plan is recomputed from
*those* rows, and only then are the approved targets retired — each
with a per-row guarded ``UPDATE`` that also asserts the row's identity,
liveness, and ownership columns are byte-identical to what the
revalidation just read. There is no window in which a safety check and
its matching write are separated by another statement, so a worker that
comes back to life, changes ownership, or is retired concurrently
cannot be deleted on the strength of a stale observation.
``fingerprint_fn`` and ``plan_fn`` are injected rather than imported so
the storage layer never depends on the decision layer; production wires
in :func:`mcp_fleet_retirement.registry_fingerprint` and
:func:`mcp_fleet_retirement.plan_stale_worker_retirement`, which is
exactly what the plan surface used.
Any exception — including a failure to commit — rolls the transaction
back and is reported as ``transaction_failed`` with zero retirements
and ``mutation_performed`` false, so a partial write can never be
reported as success.
**External state (#980 review 657 B3).** ``BEGIN IMMEDIATE`` locks this
registry and nothing else, so two critical inputs live 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 enough — the per-target loop takes time, so a
lease acquired (or a pid revived) after ``plan_fn`` returned but before
*this* row's ``UPDATE`` would go unnoticed, and the registry-column guard
cannot catch it because no registry column changed.
Two mechanisms close that gap, both applied per target and immediately
before its own write:
``external_fence_fn``
A version token over all external state the decision consumed. It is
captured inside the transaction before revalidation and re-read
before every guarded ``UPDATE``; any movement aborts the whole
transaction rather than retiring against evidence that has changed.
``liveness_fn``
A per-row re-probe of process liveness and fencing identity (host,
boot, start time). It runs immediately before the row's write and
must affirmatively re-establish that this exact process is gone.
Both default to ``None`` only so the storage layer stays independent of
the decision and control-plane layers; production always supplies them,
and a caller that omits ``liveness_fn`` gets no retirement at all rather
than an unfenced one.
"""
requested = [str(w) for w in (worker_identities or []) if str(w).strip()]
stamp = _ts(now or _utc_now())
result: dict[str, Any] | None = None
def _base(outcome: str) -> dict[str, Any]:
return {
"success": True,
"outcome": outcome,
"mutation_performed": False,
"retired": [],
"retired_count": 0,
"preserved": [],
"preserved_count": 0,
"requested_count": len(requested),
"expected_registry_fingerprint": expected_registry_fingerprint,
"expected_candidate_fingerprint": expected_candidate_fingerprint,
"acting_identity": retired_by,
"reasons": [],
}
if not requested:
outcome = _base("nothing_requested")
outcome["reasons"] = [
"no worker identities were supplied; nothing to retire"
]
return outcome
try:
with self._tx() as conn:
# Captured *after* BEGIN IMMEDIATE and *before* the
# authoritative read, so every later comparison is against the
# external state this decision was actually built on.
fence_at_plan = external_fence_fn() if external_fence_fn else None
rows = [
self._row_to_record(r)
for r in conn.execute(
"SELECT * FROM worker_registrations"
).fetchall()
]
by_identity = {
str(row.get("worker_identity")): row for row in rows
}
current_registry_fingerprint = fingerprint_fn(rows)
if current_registry_fingerprint != expected_registry_fingerprint:
already = [
wid
for wid in requested
if str(
(by_identity.get(wid) or {}).get("status") or ""
)
== STATUS_RETIRED
]
idempotent = len(already) == len(requested)
result = _base(
"already_retired" if idempotent else "registry_revision_moved"
)
result["idempotent"] = idempotent
result["current_registry_fingerprint"] = (
current_registry_fingerprint
)
result["reasons"] = [
"the worker registry changed between plan and apply; "
"retiring zero workers"
if not idempotent
else "every requested registration is already retired; "
"safe no-op"
]
result["preserved"] = [
{
"worker_identity": wid,
"reason_code": "registry_revision_moved",
"detail": (
"aborted before any retirement: registry "
"fingerprint moved"
),
}
for wid in requested
]
result["preserved_count"] = len(requested)
return result
fresh_plan = plan_fn(rows)
current_candidate_fingerprint = fresh_plan.get(
"candidate_fingerprint"
)
if current_candidate_fingerprint != expected_candidate_fingerprint:
result = _base("candidate_set_moved")
result["current_registry_fingerprint"] = (
current_registry_fingerprint
)
result["current_candidate_fingerprint"] = (
current_candidate_fingerprint
)
result["reasons"] = [
"the retirement candidate set changed between plan and "
"apply; retiring zero workers"
]
result["preserved"] = [
{
"worker_identity": wid,
"reason_code": "candidate_set_moved",
"detail": (
"aborted before any retirement: candidate "
"fingerprint moved"
),
}
for wid in requested
]
result["preserved_count"] = len(requested)
return result
eligible = {
str(c.get("worker_identity")): c
for c in fresh_plan.get("candidates") or []
}
preserved_index = {
str(p.get("worker_identity")): p
for p in fresh_plan.get("preserved") or []
}
retired: list[dict[str, Any]] = []
preserved: list[dict[str, Any]] = []
for wid in requested:
row = by_identity.get(wid)
if row is None:
preserved.append(
{
"worker_identity": wid,
"reason_code": "registration_missing",
"detail": "no registration row with this identity",
}
)
continue
if wid not in eligible:
blocked = preserved_index.get(wid) or {}
preserved.append(
{
"worker_identity": wid,
"reason_code": blocked.get("reason_code")
or "not_in_current_plan",
"detail": blocked.get("detail")
or (
"revalidation immediately before retirement "
"no longer finds this worker eligible"
),
"evidence": blocked.get("evidence"),
}
)
continue
# --- External-state fence, immediately before this write ---
#
# Lease state and OS liveness are outside this transaction,
# so they are re-checked here rather than trusted from
# revalidation. Movement aborts the whole transaction: a
# changed world invalidates every remaining decision, not
# only this row's.
if external_fence_fn is not None:
fence_now = external_fence_fn()
if fence_now != fence_at_plan:
raise _ExternalStateMoved(
"external safety state (workflow leases or "
"process liveness) changed inside the retirement "
"transaction; rolling back and retiring nothing"
)
if liveness_fn is None:
preserved.append(
{
"worker_identity": wid,
"reason_code": "liveness_reprobe_unavailable",
"detail": (
"no immediate pre-write liveness re-probe was "
"supplied; refusing to retire on revalidation "
"evidence alone (fail closed)"
),
}
)
continue
verdict = liveness_fn(dict(row))
if not verdict.get("safe"):
preserved.append(
{
"worker_identity": wid,
"reason_code": verdict.get("reason_code")
or "liveness_reprobe_refused",
"detail": verdict.get("detail")
or (
"immediate pre-write re-probe could not "
"re-establish that this process is gone"
),
"evidence": verdict.get("evidence"),
}
)
continue
cursor = conn.execute(
"UPDATE worker_registrations "
"SET status = ?, retired_at = ?, retired_by = ?, "
" retirement_reason = ? "
"WHERE worker_identity = ? "
" AND status = ? "
" AND last_heartbeat_at = ? "
" AND generation_id = ? "
" AND session_id = ? "
" AND fencing_epoch = ? "
" AND IFNULL(pid, -1) = IFNULL(?, -1) "
" AND IFNULL(host_id, '') = IFNULL(?, '') "
" AND IFNULL(boot_id, '') = IFNULL(?, '') "
" AND IFNULL(process_start_time, '') = IFNULL(?, '')",
(
STATUS_RETIRED,
stamp,
retired_by,
retirement_reason
or (eligible[wid].get("reason_code") or ""),
wid,
STATUS_ACTIVE,
row.get("last_heartbeat_at"),
row.get("generation_id"),
row.get("session_id"),
row.get("fencing_epoch"),
row.get("pid"),
row.get("host_id"),
row.get("boot_id"),
row.get("process_start_time"),
),
)
if cursor.rowcount == 1:
retired.append(
{
"worker_identity": wid,
"reason_code": eligible[wid].get("reason_code"),
"retired_at": stamp,
"retired_by": retired_by,
"evidence": eligible[wid].get("evidence"),
}
)
else:
preserved.append(
{
"worker_identity": wid,
"reason_code": "row_changed_since_plan",
"detail": (
"guarded update matched no row; the "
"registration changed inside the retirement "
"transaction"
),
}
)
result = _base("applied")
result["mutation_performed"] = bool(retired)
result["retired"] = retired
result["retired_count"] = len(retired)
result["preserved"] = preserved
result["preserved_count"] = len(preserved)
result["current_registry_fingerprint"] = (
current_registry_fingerprint
)
result["current_candidate_fingerprint"] = (
current_candidate_fingerprint
)
result["retired_at"] = stamp if retired else None
result["external_fence"] = fence_at_plan
except Exception as exc: # rolled back by _tx; report, never half-claim
moved = isinstance(exc, _ExternalStateMoved)
failure = _base("external_state_moved" if moved else "transaction_failed")
failure["success"] = False
failure["reasons"] = [
str(exc)
if moved
else "retirement transaction failed and was rolled back; zero "
f"registrations were retired: {type(exc).__name__}: {exc}"
]
failure["preserved"] = [
{
"worker_identity": wid,
"reason_code": (
"external_state_moved" if moved else "transaction_failed"
),
"detail": "transaction rolled back before any commit",
}
for wid in requested
]
failure["preserved_count"] = len(requested)
return failure
return result
def _public_record(record: dict[str, Any]) -> dict[str, Any]: def _public_record(record: dict[str, Any]) -> dict[str, Any]:
"""Registry row minus anything that should not travel to an LLM surface.""" """Registry row minus anything that should not travel to an LLM surface."""
+38
View File
@@ -175,6 +175,44 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"permission": "gitea.read", "permission": "gitea.read",
"role": "controller", "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. # #644: Phase 2 Web Console recovery tasks.
"clear_stale_binding": { "clear_stale_binding": {
"permission": "gitea.read", "permission": "gitea.read",
File diff suppressed because it is too large Load Diff