Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e344e68a13 | ||
|
|
c0c6d14b73 |
@@ -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
|
||||||
|
|
||||||
|
|||||||
@@ -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`
|
||||||
|
|||||||
@@ -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.
|
||||||
@@ -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",
|
||||||
|
|||||||
@@ -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"),
|
||||||
|
}
|
||||||
@@ -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
@@ -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."""
|
||||||
|
|||||||
@@ -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
Reference in New Issue
Block a user