From c0c6d14b73add36ac157d501c603fd9ea51de1da Mon Sep 17 00:00:00 2001 From: Jason Walker <913443@dadeschools.net> Date: Thu, 30 Jul 2026 15:22:30 -0400 Subject: [PATCH] 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) --- docs/instance-fleet-identity.md | 5 + docs/mcp-tool-inventory.md | 2 + docs/stale-worker-retirement.md | 203 +++++ gitea_mcp_server.py | 562 ++++++++++++ mcp_fleet_retirement.py | 479 +++++++++++ mcp_worker_identity.py | 280 +++++- task_capability_map.py | 26 + .../test_issue_980_stale_worker_retirement.py | 806 ++++++++++++++++++ 8 files changed, 2362 insertions(+), 1 deletion(-) create mode 100644 docs/stale-worker-retirement.md create mode 100644 mcp_fleet_retirement.py create mode 100644 tests/test_issue_980_stale_worker_retirement.py diff --git a/docs/instance-fleet-identity.md b/docs/instance-fleet-identity.md index d4bc396..4a4b8c2 100644 --- a/docs/instance-fleet-identity.md +++ b/docs/instance-fleet-identity.md @@ -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 diff --git a/docs/mcp-tool-inventory.md b/docs/mcp-tool-inventory.md index 7264f6b..ac4c87f 100644 --- a/docs/mcp-tool-inventory.md +++ b/docs/mcp-tool-inventory.md @@ -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` diff --git a/docs/stale-worker-retirement.md b/docs/stale-worker-retirement.md new file mode 100644 index 0000000..6c7b90d --- /dev/null +++ b/docs/stale-worker-retirement.md @@ -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. diff --git a/gitea_mcp_server.py b/gitea_mcp_server.py index e365da7..8da17e7 100644 --- a/gitea_mcp_server.py +++ b/gitea_mcp_server.py @@ -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", diff --git a/mcp_fleet_retirement.py b/mcp_fleet_retirement.py new file mode 100644 index 0000000..133f2fc --- /dev/null +++ b/mcp_fleet_retirement.py @@ -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"), + } diff --git a/mcp_worker_identity.py b/mcp_worker_identity.py index f9363e9..7c6f83e 100644 --- a/mcp_worker_identity.py +++ b/mcp_worker_identity.py @@ -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.""" diff --git a/task_capability_map.py b/task_capability_map.py index d306e03..10723b5 100644 --- a/task_capability_map.py +++ b/task_capability_map.py @@ -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", diff --git a/tests/test_issue_980_stale_worker_retirement.py b/tests/test_issue_980_stale_worker_retirement.py new file mode 100644 index 0000000..b35be53 --- /dev/null +++ b/tests/test_issue_980_stale_worker_retirement.py @@ -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()