Compare commits

...
1 Commits
Author SHA1 Message Date
sysadminandClaude Opus 4.8 c0c6d14b73 feat(fleet): add CAS-protected stale worker retirement capability
Adds a sanctioned controller/reconciler capability that retires conclusively
stale MCP worker-registry rows through a dry-run-first, compare-and-swap
protected workflow (#980).

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

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

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

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

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

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

Closes #980

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-30 15:22:30 -04:00
8 changed files with 2362 additions and 1 deletions
+5
View File
@@ -189,6 +189,11 @@ mutation. Diagnostic reads remain available where `gitea.read` allows.
* `mcp_fleet_snapshot` — pure snapshot + classification (#978).
* `gitea_snapshot_instance_fleet` — sanctioned MCP tool (#978).
* `gitea_get_runtime_context` — single-process view (not fleet-wide).
* [stale-worker-retirement.md](stale-worker-retirement.md) — CAS-protected
retirement of conclusively stale registrations (#980). Note that the
`registry_revision` this snapshot returns is **time-seeded** and unsuitable
for compare-and-swap; retirement derives its own stable
`registry_fingerprint` from registry content alone.
## Non-goals
+2
View File
@@ -51,6 +51,7 @@ that gates each call, not which tools exist.
- `gitea_adopt_merger_pr_lease`
- `gitea_adopt_workflow_lease`
- `gitea_allocate_next_work`
- `gitea_apply_stale_worker_retirement`
- `gitea_assess_already_landed_reconciliation`
- `gitea_assess_conflict_fix_classification`
- `gitea_assess_conflict_fix_push`
@@ -124,6 +125,7 @@ that gates each call, not which tools exist.
- `gitea_observability_link_issue`
- `gitea_observability_list_projects`
- `gitea_observability_reconcile_incident`
- `gitea_plan_stale_worker_retirement`
- `gitea_post_heartbeat`
- `gitea_publish_unpublished_issue_branch`
- `gitea_quarantine_contaminated_review`
+203
View File
@@ -0,0 +1,203 @@
# 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` |
| Not an active workflow-lease owner | `protected_active_workflow_owner` |
Eligible rows carry `eligible_stale_orphan`.
Two 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.
* **Trusted launcher identity is not required.** #980 places trusted
`client_instance_id` propagation out of scope and lists backfilling trusted
identity for legacy workers as a non-goal. Requiring `inst-…` provenance here
would preserve every legacy row forever and make the feature inert. What is
required is that the registry *fields* the conjunction reads are present.
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
* **Allowed:** `controller`, `reconciler`.
* **Denied:** author, reviewer, merger — they keep `gitea.read` for diagnosis
elsewhere and are refused this surface by role.
* The Gitea operation gate stays `gitea.read` because the mutation lands in the
**local control-plane worker registry**, not in Gitea — the same model the
#601 lease lifecycle uses. **No new Gitea write permission is introduced for
any profile**, and no author permission is broadened.
* The fleet snapshot remains observational: nothing here turns it into a gate on
ordinary author work.
## 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.
+562
View File
@@ -19623,6 +19623,568 @@ def gitea_snapshot_instance_fleet(
return snapshot
# --- #980 CAS-protected stale worker retirement ---------------------------
def _retirement_role_block(task: str) -> dict | None:
"""Refuse the retirement surface to any non-controller/reconciler role.
Mirrors ``gitea_snapshot_instance_fleet``: ``gitea.read`` is the operation
gate (the mutation lands in the local worker registry, not in Gitea), and
the role restriction is what actually keeps author, reviewer, and merger
profiles out. No unrelated permission is granted to anyone.
"""
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
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"),
)
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
read_block = _profile_operation_gate("gitea.read")
if read_block:
return {
"success": False,
"mutation_performed": False,
"retired_count": 0,
"reasons": read_block,
"permission_report": _permission_block_report("gitea.read"),
}
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,
)
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": ["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()
def gitea_get_runtime_context(
remote: str = "dadeschools",
+479
View File
@@ -0,0 +1,479 @@
"""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"
#: Registry columns that must carry a usable value before the eligibility
#: conjunction can even be evaluated. Absence is ambiguity, not permission.
REQUIRED_IDENTITY_FIELDS: tuple[str, ...] = (
"worker_identity",
"client_instance_id",
"session_id",
"generation_id",
"status",
"started_at",
"last_heartbeat_at",
"heartbeat_ttl_seconds",
"pid",
)
#: 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",
)
_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 _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,
"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 _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,
) -> 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] = {}
for index, row in enumerate(all_rows):
pid_alive = _probe_pid(row.get("pid"), pid_alive_probe)
pid_alive_by_index[index] = pid_alive
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))
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:
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}",
)
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,
}
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"),
}
+279 -1
View File
@@ -39,7 +39,7 @@ import threading
import time
from contextlib import contextmanager
from datetime import datetime, timezone
from typing import Any, Iterator
from typing import Any, Callable, Iterator, Sequence
# --- Provenance verdicts -------------------------------------------------
@@ -224,6 +224,13 @@ def _heartbeat_expectation_drift(
STATUS_ACTIVE = "active"
STATUS_SUPERSEDED = "superseded"
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"})
_FALSE_VALUES = frozenset({"0", "false", "no", "manual", "manual_launch"})
@@ -315,6 +322,12 @@ _SCHEMA_OPTIONAL_COLUMNS: tuple[tuple[str, str], ...] = (
("parity_revision", "TEXT"),
("live_revision", "TEXT"),
("instance_id_provenance", "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"),
)
@@ -1217,6 +1230,271 @@ class WorkerRegistry:
"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,
) -> 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.
"""
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:
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
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)",
(
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"),
),
)
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
except Exception as exc: # rolled back by _tx; report, never half-claim
failure = _base("transaction_failed")
failure["success"] = False
failure["reasons"] = [
"retirement transaction failed and was rolled back; zero "
f"registrations were retired: {type(exc).__name__}: {exc}"
]
failure["preserved"] = [
{
"worker_identity": wid,
"reason_code": "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]:
"""Registry row minus anything that should not travel to an LLM surface."""
+26
View File
@@ -175,6 +175,32 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"permission": "gitea.read",
"role": "controller",
},
# #980: CAS-protected retirement of conclusively stale worker
# registrations. The mutation lands in the local control-plane worker
# registry, not in Gitea, so — exactly like the #601 lease lifecycle — the
# Gitea operation gate stays ``gitea.read`` and no new Gitea write
# permission is introduced for any profile. The real authority is enforced
# in the tools themselves: role_kind must be controller or reconciler, the
# runtime must be parity-clean and cohort-unique, and apply additionally
# requires the exact stable registry + candidate fingerprints returned by
# the plan. Author, reviewer, and merger profiles keep gitea.read for
# diagnosis elsewhere and are refused this surface.
"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.read",
"role": "controller",
},
"gitea_apply_stale_worker_retirement": {
"permission": "gitea.read",
"role": "controller",
},
# #644: Phase 2 Web Console recovery tasks.
"clear_stale_binding": {
"permission": "gitea.read",
@@ -0,0 +1,806 @@
"""CAS-protected stale worker retirement (#980).
Covers the acceptance criteria: the registry CAS token is stable across time
and row order but moves on any retirement-relevant change, plan mutates
nothing, apply fails closed on drift, every target is revalidated inside the
retirement transaction, and live / ambiguous / incomplete rows are preserved.
The regression test that matters most is
``test_snapshot_at_would_have_moved_the_token``: it demonstrates the exact
defect the R3-C review found — ``mcp_fleet_snapshot._consistency_token``
changes when only the observation time changed — and proves the new token does
not.
All identifiers are synthetic.
"""
from __future__ import annotations
import os
import tempfile
import unittest
from datetime import datetime, timedelta, timezone
from typing import Any
import mcp_fleet_retirement as retire
import mcp_fleet_snapshot as fleet
import mcp_worker_identity as mwi
NOW = datetime(2026, 7, 30, 7, 0, 0, tzinfo=timezone.utc)
TTL = 900.0
REPO = "/synthetic/repo/Gitea-Tools"
_ROW_COLUMNS = (
"worker_identity",
"client_name",
"client_instance_id",
"session_id",
"generation_id",
"role",
"profile",
"namespace",
"remote",
"repository_binding",
"pid",
"process_identity",
"transport",
"token_fingerprint",
"started_at",
"last_heartbeat_at",
"heartbeat_ttl_seconds",
"fencing_epoch",
"status",
"fleet_run_id",
"authenticated_account",
"instance_id_provenance",
)
def _registry() -> mwi.WorkerRegistry:
handle, path = tempfile.mkstemp(suffix=".sqlite3")
os.close(handle)
os.unlink(path)
return mwi.WorkerRegistry(path)
def _row(
*,
identity: str,
pid: int | None,
instance: str = "inst-codex-20260730T070000Z-0123456789ab",
session: str = "sess-a",
generation: str = "gen-a",
namespace: str = "author",
heartbeat: datetime | None = None,
status: str = mwi.STATUS_ACTIVE,
repository_binding: str | None = REPO,
ttl: float = TTL,
**overrides: Any,
) -> dict[str, Any]:
"""One synthetic ``worker_registrations`` row."""
record = {
"worker_identity": identity,
"client_name": "codex",
"client_instance_id": instance,
"session_id": session,
"generation_id": generation,
"role": namespace,
"profile": f"prgs-{namespace}",
"namespace": namespace,
"remote": "prgs",
"repository_binding": repository_binding,
"pid": pid,
"process_identity": f"pid-{pid}" if pid is not None else None,
"transport": "stdio",
"token_fingerprint": None,
"started_at": "2026-07-30T06:00:00Z",
"last_heartbeat_at": mwi._ts(heartbeat or (NOW - timedelta(seconds=7200))),
"heartbeat_ttl_seconds": ttl,
"fencing_epoch": 1,
"status": status,
"fleet_run_id": "run-canary",
"authenticated_account": "synthetic-user",
"instance_id_provenance": fleet.INSTANCE_ID_PROVENANCE_TRUSTED,
}
record.update(overrides)
return record
def _dead(_pid: int | None) -> bool:
return False
def _alive(_pid: int | None) -> bool:
return True
def _selective(alive_pids: set[int]):
def probe(pid: int | None) -> bool:
return pid in alive_pids
return probe
def _plan(rows, *, probe=_dead, now=NOW, **kwargs):
return retire.plan_stale_worker_retirement(
rows,
now=now,
pid_alive_probe=probe,
canonical_repository=REPO,
**kwargs,
)
def _register(registry: mwi.WorkerRegistry, row: dict[str, Any]) -> None:
"""Insert a synthetic row directly, bypassing register()'s live-now stamps."""
columns = [name for name in _ROW_COLUMNS if name in row]
placeholders = ", ".join("?" for _ in columns)
with registry._tx() as conn:
conn.execute(
f"INSERT INTO worker_registrations ({', '.join(columns)}) "
f"VALUES ({placeholders})",
[row[name] for name in columns],
)
class RegistryFingerprintStabilityTests(unittest.TestCase):
"""AC: identical contents observed at different times produce one token."""
def test_same_contents_different_observation_times_same_token(self):
rows = [_row(identity="w-1", pid=101), _row(identity="w-2", pid=102)]
self.assertEqual(
retire.registry_fingerprint(rows), retire.registry_fingerprint(rows)
)
# Recompute after the clock has moved a full hour: the token is derived
# from content only, so it cannot notice.
plan_early = _plan(rows, now=NOW)
plan_late = _plan(rows, now=NOW + timedelta(hours=1))
self.assertEqual(
plan_early["registry_fingerprint"], plan_late["registry_fingerprint"]
)
def test_snapshot_at_would_have_moved_the_token(self):
"""Regression: the old time-seeded derivation moved, the new one does not."""
rows = [_row(identity="w-1", pid=101)]
old_early = fleet._consistency_token(rows, "2026-07-30T07:00:00Z")
old_late = fleet._consistency_token(rows, "2026-07-30T07:00:01Z")
self.assertNotEqual(
old_early,
old_late,
"the #980 defect: one second of observation drift changed the token",
)
self.assertEqual(
retire.registry_fingerprint(rows), retire.registry_fingerprint(rows)
)
def test_row_order_does_not_change_the_token(self):
rows = [
_row(identity="w-1", pid=101),
_row(identity="w-2", pid=102),
_row(identity="w-3", pid=103),
]
self.assertEqual(
retire.registry_fingerprint(rows),
retire.registry_fingerprint(list(reversed(rows))),
)
def test_numeric_typing_does_not_change_the_token(self):
as_float = [_row(identity="w-1", pid=101, heartbeat_ttl_seconds=900.0)]
as_int = [_row(identity="w-1", pid=101, heartbeat_ttl_seconds=900)]
self.assertEqual(
retire.registry_fingerprint(as_float),
retire.registry_fingerprint(as_int),
)
def test_retirement_relevant_changes_move_the_token(self):
base = [_row(identity="w-1", pid=101)]
baseline = retire.registry_fingerprint(base)
mutations = {
"row added": base + [_row(identity="w-2", pid=102)],
"row removed": [],
"status changed": [_row(identity="w-1", pid=101, status="released")],
"identity changed": [
_row(
identity="w-1",
pid=101,
instance="inst-codex-20260730T070000Z-ffffffffffff",
)
],
"ownership changed": [_row(identity="w-1", pid=101, session="sess-other")],
"generation changed": [
_row(identity="w-1", pid=101, generation="gen-other")
],
"pid changed": [_row(identity="w-1", pid=999)],
"repository binding changed": [
_row(identity="w-1", pid=101, repository_binding="/elsewhere")
],
}
for label, rows in mutations.items():
with self.subTest(change=label):
self.assertNotEqual(baseline, retire.registry_fingerprint(rows))
def test_heartbeat_change_moves_the_token(self):
base = [_row(identity="w-1", pid=101)]
moved = [_row(identity="w-1", pid=101, heartbeat=NOW - timedelta(seconds=30))]
self.assertNotEqual(
retire.registry_fingerprint(base), retire.registry_fingerprint(moved)
)
def test_ttl_change_moves_the_token(self):
base = [_row(identity="w-1", pid=101)]
moved = [_row(identity="w-1", pid=101, ttl=60.0)]
self.assertNotEqual(
retire.registry_fingerprint(base), retire.registry_fingerprint(moved)
)
class PlanEligibilityTests(unittest.TestCase):
"""AC: only conclusively stale orphans are selected; everything else stays."""
def test_plan_selects_dead_stale_orphan(self):
plan = _plan([_row(identity="w-1", pid=101)])
self.assertEqual(plan["candidate_count"], 1)
self.assertEqual(plan["candidate_worker_identities"], ["w-1"])
self.assertEqual(plan["candidates"][0]["reason_code"], retire.REASON_ELIGIBLE)
self.assertFalse(plan["mutation_performed"])
self.assertTrue(plan["read_only"])
def test_plan_performs_no_mutation(self):
registry = _registry()
_register(registry, _row(identity="w-1", pid=101))
before = registry.list_workers(status=None)
plan = _plan(before)
after = registry.list_workers(status=None)
self.assertEqual(plan["candidate_count"], 1)
self.assertEqual(before, after)
self.assertEqual([r["status"] for r in after], [mwi.STATUS_ACTIVE])
self.assertFalse(plan["mutation_performed"])
def test_live_worker_is_preserved(self):
row = _row(identity="w-live", pid=101, heartbeat=NOW - timedelta(seconds=10))
plan = _plan([row], probe=_alive)
self.assertEqual(plan["candidate_count"], 0)
self.assertEqual(
plan["preserved"][0]["reason_code"], retire.REASON_WORKER_LIVE
)
def test_fresh_heartbeat_with_dead_pid_still_fails_closed(self):
"""A dead pid withdraws liveness; the unexpired heartbeat still preserves."""
row = _row(identity="w-fresh", pid=101, heartbeat=NOW - timedelta(seconds=10))
plan = _plan([row], probe=_dead)
self.assertEqual(plan["candidate_count"], 0)
self.assertEqual(
plan["preserved"][0]["reason_code"], retire.REASON_HEARTBEAT_FRESH
)
def test_unprobeable_pid_is_preserved(self):
plan = _plan([_row(identity="w-1", pid=101)], probe=lambda _pid: None)
self.assertEqual(plan["candidate_count"], 0)
self.assertEqual(
plan["preserved"][0]["reason_code"], retire.REASON_PID_UNKNOWN
)
def test_incomplete_legacy_identity_is_preserved(self):
"""AC: incomplete legacy identities remain fail-closed."""
rows = [
_row(identity="w-nopid", pid=None, instance="legacy-pid-27833"),
_row(identity="w-nosession", pid=102, session=""),
]
plan = _plan(rows)
self.assertEqual(plan["candidate_count"], 0)
self.assertEqual(
{p["reason_code"] for p in plan["preserved"]},
{retire.REASON_INCOMPLETE_IDENTITY},
)
detail = next(
p["detail"] for p in plan["preserved"] if p["worker_identity"] == "w-nopid"
)
self.assertIn("pid", detail)
def test_unparsable_heartbeat_is_preserved(self):
rows = [_row(identity="w-1", pid=101, last_heartbeat_at="not-a-stamp")]
plan = _plan(rows)
self.assertEqual(plan["candidate_count"], 0)
self.assertEqual(
plan["preserved"][0]["reason_code"], retire.REASON_UNPARSABLE_HEARTBEAT
)
def test_identity_shared_with_live_worker_is_preserved(self):
"""Two rows, one live: the dead one shares session evidence, so it stays."""
rows = [
_row(
identity="w-live",
pid=101,
session="sess-shared",
heartbeat=NOW - timedelta(seconds=5),
),
_row(identity="w-dead", pid=102, session="sess-shared"),
]
plan = _plan(rows, probe=_selective({101}))
self.assertEqual(plan["candidate_count"], 0)
codes = {p["reason_code"] for p in plan["preserved"]}
self.assertIn(retire.REASON_CONFLICTING_IDENTITY, codes)
def test_foreign_or_missing_repository_binding_is_preserved(self):
rows = [
_row(identity="w-foreign", pid=101, repository_binding="/other/repo"),
_row(identity="w-unbound", pid=102, repository_binding=None),
]
plan = _plan(rows)
self.assertEqual(plan["candidate_count"], 0)
self.assertEqual(
{p["reason_code"] for p in plan["preserved"]},
{retire.REASON_FOREIGN_REPOSITORY},
)
def test_active_workflow_owner_is_preserved(self):
rows = [_row(identity="w-1", pid=101, session="sess-leased")]
plan = _plan(rows, protected_session_ids=["sess-leased"])
self.assertEqual(plan["candidate_count"], 0)
self.assertEqual(
plan["preserved"][0]["reason_code"], retire.REASON_PROTECTED_OWNER
)
def test_terminal_rows_are_not_retired_again(self):
rows = [_row(identity="w-1", pid=101, status=mwi.STATUS_RELEASED)]
plan = _plan(rows)
self.assertEqual(plan["candidate_count"], 0)
self.assertEqual(
plan["preserved"][0]["reason_code"], retire.REASON_ALREADY_TERMINAL
)
def test_mixed_fleet_produces_correct_per_worker_outcomes(self):
rows = [
_row(identity="w-stale-1", pid=101, session="s1", generation="g1"),
_row(identity="w-stale-2", pid=102, session="s2", generation="g2"),
_row(
identity="w-live",
pid=103,
session="s3",
generation="g3",
heartbeat=NOW - timedelta(seconds=5),
),
_row(identity="w-nopid", pid=None, session="s4", generation="g4"),
_row(
identity="w-foreign",
pid=105,
session="s5",
generation="g5",
repository_binding="/other",
),
_row(
identity="w-terminal",
pid=106,
session="s6",
generation="g6",
status=mwi.STATUS_SUPERSEDED,
),
]
plan = _plan(rows, probe=_selective({103}))
self.assertEqual(
sorted(plan["candidate_worker_identities"]), ["w-stale-1", "w-stale-2"]
)
by_identity = {
p["worker_identity"]: p["reason_code"] for p in plan["preserved"]
}
self.assertEqual(by_identity["w-live"], retire.REASON_WORKER_LIVE)
self.assertEqual(by_identity["w-nopid"], retire.REASON_INCOMPLETE_IDENTITY)
self.assertEqual(by_identity["w-foreign"], retire.REASON_FOREIGN_REPOSITORY)
self.assertEqual(by_identity["w-terminal"], retire.REASON_ALREADY_TERMINAL)
self.assertEqual(plan["assessed_count"], 6)
self.assertEqual(plan["preserved_count"], 4)
def test_plan_is_deterministic_for_identical_contents(self):
rows = [_row(identity="w-1", pid=101), _row(identity="w-2", pid=102)]
first = _plan(rows)
second = _plan(list(reversed(rows)), now=NOW + timedelta(minutes=5))
self.assertEqual(first["registry_fingerprint"], second["registry_fingerprint"])
self.assertEqual(
first["candidate_fingerprint"], second["candidate_fingerprint"]
)
self.assertEqual(
sorted(first["candidate_worker_identities"]),
sorted(second["candidate_worker_identities"]),
)
class ApplyCasTests(unittest.TestCase):
"""AC: apply is compare-and-swap protected and revalidates every target."""
def setUp(self) -> None:
self.registry = _registry()
self.probe = _dead
def _rows(self):
return self.registry.list_workers(status=None)
def _apply(self, plan, *, probe=None, identities=None, **kwargs):
chosen = probe or self.probe
return self.registry.retire_stale_workers(
expected_registry_fingerprint=plan["registry_fingerprint"],
expected_candidate_fingerprint=plan["candidate_fingerprint"],
worker_identities=(
identities
if identities is not None
else plan["candidate_worker_identities"]
),
fingerprint_fn=retire.registry_fingerprint,
plan_fn=lambda rows: _plan(rows, probe=chosen, **kwargs),
retired_by="synthetic-user/prgs-reconciler",
now=NOW,
)
def test_matching_token_retires_the_planned_set(self):
_register(self.registry, _row(identity="w-1", pid=101, session="s1"))
_register(self.registry, _row(identity="w-2", pid=102, session="s2"))
plan = _plan(self._rows(), probe=self.probe)
result = self._apply(plan)
self.assertEqual(result["outcome"], "applied")
self.assertTrue(result["mutation_performed"])
self.assertEqual(result["retired_count"], 2)
statuses = {r["worker_identity"]: r["status"] for r in self._rows()}
self.assertEqual(
statuses, {"w-1": mwi.STATUS_RETIRED, "w-2": mwi.STATUS_RETIRED}
)
retired_row = self.registry.get("w-1")
self.assertEqual(retired_row["retired_at"], mwi._ts(NOW))
self.assertEqual(retired_row["retired_by"], "synthetic-user/prgs-reconciler")
def test_moved_registry_token_retires_zero_workers(self):
_register(self.registry, _row(identity="w-1", pid=101))
plan = _plan(self._rows(), probe=self.probe)
# An unrelated registration lands between plan and apply.
_register(self.registry, _row(identity="w-2", pid=102, session="s2"))
result = self._apply(plan)
self.assertEqual(result["outcome"], retire.OUTCOME_REGISTRY_MOVED)
self.assertEqual(result["retired_count"], 0)
self.assertFalse(result["mutation_performed"])
self.assertTrue(all(r["status"] == mwi.STATUS_ACTIVE for r in self._rows()))
def test_moved_candidate_set_retires_zero_workers(self):
"""Registry unchanged, but the eligibility verdict is no longer the same."""
_register(self.registry, _row(identity="w-1", pid=101))
plan = _plan(self._rows(), probe=self.probe)
# The row did not change; the process came back (pid probes alive), so
# revalidation inside the transaction finds no candidates at all.
result = self._apply(plan, probe=_alive)
self.assertEqual(result["outcome"], retire.OUTCOME_CANDIDATES_MOVED)
self.assertEqual(result["retired_count"], 0)
self.assertFalse(result["mutation_performed"])
self.assertEqual(self.registry.get("w-1")["status"], mwi.STATUS_ACTIVE)
def test_worker_that_becomes_live_between_plan_and_apply_is_preserved(self):
_register(self.registry, _row(identity="w-dead", pid=101, session="s1"))
_register(self.registry, _row(identity="w-back", pid=102, session="s2"))
plan = _plan(self._rows(), probe=self.probe)
self.assertEqual(len(plan["candidate_worker_identities"]), 2)
# w-back's process is alive by the time apply runs, so the candidate
# fingerprint moves and nothing at all is retired.
result = self._apply(plan, probe=_selective({102}))
self.assertEqual(result["outcome"], retire.OUTCOME_CANDIDATES_MOVED)
self.assertEqual(result["retired_count"], 0)
self.assertEqual(self.registry.get("w-back")["status"], mwi.STATUS_ACTIVE)
self.assertEqual(self.registry.get("w-dead")["status"], mwi.STATUS_ACTIVE)
def test_worker_that_becomes_ambiguous_between_plan_and_apply_is_preserved(self):
_register(self.registry, _row(identity="w-1", pid=101, session="s1"))
plan = _plan(self._rows(), probe=self.probe)
result = self._apply(plan, probe=lambda _pid: None)
self.assertEqual(result["outcome"], retire.OUTCOME_CANDIDATES_MOVED)
self.assertEqual(result["retired_count"], 0)
self.assertEqual(self.registry.get("w-1")["status"], mwi.STATUS_ACTIVE)
def test_apply_revalidates_and_will_not_narrow_the_approved_set(self):
"""A protection appearing after plan moves the CAS, so nothing is retired."""
_register(self.registry, _row(identity="w-1", pid=101, session="s1"))
_register(
self.registry, _row(identity="w-protected", pid=102, session="sess-leased")
)
unprotected_plan = _plan(self._rows(), probe=self.probe)
self.assertEqual(unprotected_plan["candidate_count"], 2)
result = self._apply(
unprotected_plan, protected_session_ids=["sess-leased"]
)
self.assertEqual(result["outcome"], retire.OUTCOME_CANDIDATES_MOVED)
self.assertEqual(result["retired_count"], 0)
self.assertEqual(
self.registry.get("w-protected")["status"], mwi.STATUS_ACTIVE
)
def test_target_outside_the_current_plan_is_preserved(self):
_register(
self.registry,
_row(identity="w-1", pid=101, session="s1", generation="g1"),
)
_register(
self.registry,
_row(
identity="w-live",
pid=102,
session="s2",
generation="g2",
heartbeat=NOW - timedelta(seconds=5),
),
)
probe = _selective({102})
plan = _plan(self._rows(), probe=probe)
self.assertEqual(plan["candidate_worker_identities"], ["w-1"])
# A caller that appends a live identity to the approved list gets it
# preserved with the live reason code; only the real candidate retires.
result = self._apply(plan, probe=probe, identities=["w-1", "w-live"])
self.assertEqual(result["outcome"], "applied")
self.assertEqual(result["retired_count"], 1)
self.assertEqual(result["retired"][0]["worker_identity"], "w-1")
preserved = {
p["worker_identity"]: p["reason_code"] for p in result["preserved"]
}
self.assertEqual(preserved["w-live"], retire.REASON_WORKER_LIVE)
self.assertEqual(self.registry.get("w-live")["status"], mwi.STATUS_ACTIVE)
def test_missing_registration_is_reported_not_invented(self):
_register(self.registry, _row(identity="w-1", pid=101))
plan = _plan(self._rows(), probe=self.probe)
result = self._apply(plan, identities=["w-1", "w-ghost"])
self.assertEqual(result["retired_count"], 1)
preserved = {
p["worker_identity"]: p["reason_code"] for p in result["preserved"]
}
self.assertEqual(preserved["w-ghost"], "registration_missing")
def test_reapplying_a_completed_plan_is_a_safe_no_op(self):
_register(self.registry, _row(identity="w-1", pid=101))
plan = _plan(self._rows(), probe=self.probe)
first = self._apply(plan)
self.assertEqual(first["retired_count"], 1)
second = self._apply(plan)
self.assertEqual(second["outcome"], retire.OUTCOME_ALREADY_RETIRED)
self.assertTrue(second["idempotent"])
self.assertEqual(second["retired_count"], 0)
self.assertFalse(second["mutation_performed"])
self.assertEqual(self.registry.get("w-1")["status"], mwi.STATUS_RETIRED)
def test_empty_target_list_reports_nothing_requested(self):
_register(self.registry, _row(identity="w-1", pid=101))
plan = _plan(self._rows(), probe=self.probe)
result = self._apply(plan, identities=[])
self.assertEqual(result["outcome"], retire.OUTCOME_NOTHING_REQUESTED)
self.assertFalse(result["mutation_performed"])
self.assertEqual(self.registry.get("w-1")["status"], mwi.STATUS_ACTIVE)
def test_transaction_failure_cannot_report_partial_success(self):
_register(self.registry, _row(identity="w-1", pid=101, session="s1"))
_register(self.registry, _row(identity="w-2", pid=102, session="s2"))
plan = _plan(self._rows(), probe=self.probe)
calls = {"n": 0}
def exploding_plan(rows):
calls["n"] += 1
raise RuntimeError("synthetic revalidation failure")
result = self.registry.retire_stale_workers(
expected_registry_fingerprint=plan["registry_fingerprint"],
expected_candidate_fingerprint=plan["candidate_fingerprint"],
worker_identities=plan["candidate_worker_identities"],
fingerprint_fn=retire.registry_fingerprint,
plan_fn=exploding_plan,
retired_by="synthetic-user/prgs-reconciler",
now=NOW,
)
self.assertEqual(calls["n"], 1)
self.assertFalse(result["success"])
self.assertEqual(result["outcome"], "transaction_failed")
self.assertEqual(result["retired_count"], 0)
self.assertFalse(result["mutation_performed"])
self.assertTrue(
all(r["status"] == mwi.STATUS_ACTIVE for r in self._rows()),
"a failed transaction must roll back every retirement",
)
def test_structured_result_fields_are_accurate(self):
_register(self.registry, _row(identity="w-1", pid=101))
plan = _plan(self._rows(), probe=self.probe)
blocked = self._apply(plan, probe=_alive)
self.assertFalse(blocked["mutation_performed"])
self.assertEqual(blocked["retired"], [])
self.assertEqual(blocked["requested_count"], 1)
self.assertEqual(
blocked["expected_registry_fingerprint"], plan["registry_fingerprint"]
)
applied = self._apply(plan)
self.assertTrue(applied["mutation_performed"])
self.assertEqual(applied["acting_identity"], "synthetic-user/prgs-reconciler")
self.assertEqual(applied["retired"][0]["reason_code"], retire.REASON_ELIGIBLE)
self.assertIn("evidence", applied["retired"][0])
class PostRetirementFleetCompatibilityTests(unittest.TestCase):
"""AC: retired rows leave the stale count and never fake fleet safety."""
def test_retired_rows_are_historical_not_stale(self):
registry = _registry()
_register(
registry, _row(identity="w-1", pid=101, session="s1", generation="g1")
)
_register(
registry,
_row(
identity="w-live",
pid=102,
session="s2",
generation="g2",
instance="legacy-pid-27833",
instance_id_provenance=fleet.INSTANCE_ID_PROVENANCE_LEGACY,
heartbeat=NOW - timedelta(seconds=5),
),
)
probe = _selective({102})
before = fleet.snapshot_instance_fleet(
registry.list_workers(status=None),
now=NOW,
pid_alive_probe=probe,
canonical_repository=REPO,
)
self.assertEqual(before["stale_worker_count"], 1)
plan = _plan(registry.list_workers(status=None), probe=probe)
result = registry.retire_stale_workers(
expected_registry_fingerprint=plan["registry_fingerprint"],
expected_candidate_fingerprint=plan["candidate_fingerprint"],
worker_identities=plan["candidate_worker_identities"],
fingerprint_fn=retire.registry_fingerprint,
plan_fn=lambda rows: _plan(rows, probe=probe),
retired_by="synthetic-user/prgs-reconciler",
now=NOW,
)
self.assertEqual(result["retired_count"], 1)
after = fleet.snapshot_instance_fleet(
registry.list_workers(status=None),
now=NOW,
pid_alive_probe=probe,
canonical_repository=REPO,
)
self.assertEqual(after["stale_worker_count"], 0)
self.assertEqual(after["live_worker_count"], 1)
self.assertEqual(after["historical_worker_count"], 1)
# The surviving live worker still carries a legacy instance identity, so
# the fleet must not be declared safe.
self.assertFalse(after["live_fleet_safe"])
self.assertIn(
fleet.CLASS_LEGACY_INCOMPLETE,
{f["classification"] for f in after["active_blockers"]},
)
def test_existing_snapshot_shape_is_unchanged(self):
rows = [_row(identity="w-1", pid=101)]
snapshot = fleet.snapshot_instance_fleet(
rows, now=NOW, pid_alive_probe=_dead, canonical_repository=REPO
)
for key in (
"consistency_token",
"registry_revision",
"snapshot_at",
"live_workers",
"stale_workers",
"historical_workers",
"findings",
"live_fleet_safe",
):
self.assertIn(key, snapshot)
self.assertTrue(snapshot["registry_revision"].startswith("fleetrev-"))
self.assertTrue(snapshot["read_only"])
self.assertFalse(snapshot.get("mutation_performed", False))
class CapabilityExposureTests(unittest.TestCase):
"""AC: only controller/reconciler reach the surface; author policy unchanged."""
def test_capability_map_declares_controller_role(self):
import task_capability_map as tcm
for task in (
"plan_stale_worker_retirement",
"gitea_plan_stale_worker_retirement",
"apply_stale_worker_retirement",
"gitea_apply_stale_worker_retirement",
):
with self.subTest(task=task):
self.assertEqual(tcm.required_permission(task), "gitea.read")
self.assertEqual(tcm.required_role(task), "controller")
def test_ordinary_author_capability_resolution_is_unaffected(self):
import task_capability_map as tcm
self.assertEqual(tcm.required_permission("work_issue"), "gitea.pr.create")
self.assertEqual(tcm.required_role("work_issue"), "author")
self.assertEqual(tcm.required_permission("create_pr"), "gitea.pr.create")
self.assertEqual(tcm.required_role("create_pr"), "author")
self.assertEqual(tcm.required_permission("commit_files"), "gitea.repo.commit")
self.assertEqual(tcm.required_permission("push_branch"), "gitea.branch.push")
self.assertEqual(
tcm.required_permission("snapshot_instance_fleet"), "gitea.read"
)
def test_no_new_gitea_write_permission_is_introduced(self):
import task_capability_map as tcm
for task in ("plan_stale_worker_retirement", "apply_stale_worker_retirement"):
with self.subTest(task=task):
self.assertNotIn(
tcm.required_permission(task),
{
"gitea.pr.approve",
"gitea.pr.merge",
"gitea.pr.review",
"gitea.branch.push",
"gitea.repo.commit",
},
)
class SurfaceRegistrationTests(unittest.TestCase):
"""AC: the dry-run and apply modes are exposed as sanctioned native tools."""
def test_tools_are_registered_and_role_restricted(self):
import inspect
import gitea_mcp_server as server
for name in (
"gitea_plan_stale_worker_retirement",
"gitea_apply_stale_worker_retirement",
):
with self.subTest(tool=name):
tool = getattr(server, name)
self.assertTrue(callable(tool))
source = inspect.getsource(tool)
# Both tools route their role check through the shared guard.
self.assertIn(f'_retirement_role_block("{name}")', source)
guard = inspect.getsource(server._retirement_role_block)
self.assertIn('{"controller", "reconciler"}', guard)
self.assertIn("required_roles", guard)
apply_source = inspect.getsource(server.gitea_apply_stale_worker_retirement)
self.assertIn("registry_fingerprint", apply_source)
self.assertIn("candidate_fingerprint", apply_source)
self.assertIn("_retirement_runtime_block", apply_source)
self.assertIn("classify_cohort", apply_source)
# Revalidation inside the transaction re-reads lease protection rather
# than reusing the pre-transaction snapshot.
self.assertIn("_retirement_revalidation_plan", apply_source)
revalidation = inspect.getsource(server._retirement_revalidation_plan)
self.assertIn("_retirement_protected_owners", revalidation)
self.assertIn("raise RuntimeError", revalidation)
def test_plan_tool_declares_itself_read_only(self):
import inspect
import gitea_mcp_server as server
source = inspect.getsource(server.gitea_plan_stale_worker_retirement)
self.assertIn("read_only", source)
self.assertNotIn("retire_stale_workers", source)
def test_docs_exist(self):
root = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
path = os.path.join(root, "docs", "stale-worker-retirement.md")
self.assertTrue(os.path.isfile(path))
with open(path, encoding="utf-8") as handle:
text = handle.read()
for token in (
"registry_fingerprint",
"candidate_fingerprint",
"gitea_plan_stale_worker_retirement",
"gitea_apply_stale_worker_retirement",
"eligible_stale_orphan",
"registry_revision_moved",
"does not repair untrusted live identity",
):
with self.subTest(token=token):
self.assertIn(token, text)
if __name__ == "__main__": # pragma: no cover
unittest.main()