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..043034d --- /dev/null +++ b/docs/stale-worker-retirement.md @@ -0,0 +1,282 @@ +# Stale worker retirement (#980) + +The #978 instance-fleet snapshot made registry accuracy observable but +deliberately read-only: a registry full of rows whose owning processes are long +gone stays full. This document describes the sanctioned way to retire those +rows — a dry-run-first, compare-and-swap-protected workflow available only to +controller and reconciler namespaces. + +Related: [instance-fleet-identity.md](instance-fleet-identity.md) (#978), +[post-restart-reconcile.md](post-restart-reconcile.md) (#662). + +## Why a dedicated registry token + +`mcp_fleet_snapshot._consistency_token` seeds its digest with `snapshot_at`, +formatted at second precision. Its output — surfaced as `registry_revision` and +`consistency_token` on the snapshot — therefore changes on **every call**, even +when no registry row changed. Any compare-and-swap gated on it can never pass: +a dry-run/apply cycle spanning more than one second aborts unconditionally. + +That token remains useful as an observation stamp and is unchanged. #980 adds a +separate, *stable* token instead: + +| Token | Module | Derived from | Stable across time? | +| --- | --- | --- | --- | +| `registry_revision` / `consistency_token` | `mcp_fleet_snapshot` | `snapshot_at` + a subset of row fields | **No** — moves every second | +| `registry_fingerprint` | `mcp_fleet_retirement` | canonical retirement-relevant row content only | **Yes** | +| `candidate_fingerprint` | `mcp_fleet_retirement` | canonical content of the selected candidate rows | **Yes** | + +`registry_fingerprint` guarantees: + +* identical canonical registry contents always produce the same token, whenever + they are observed; +* row iteration order never affects the token (serialized rows are sorted); +* any retirement-relevant change moves it — row creation or deletion, identity + change, heartbeat or TTL change, ownership change, registration-state change, + PID change, or repository-binding change. + +The exact field set is `mcp_fleet_retirement.FINGERPRINT_FIELDS`. Deliberately +excluded: `token_fingerprint` (credential-adjacent, never a retirement input), +the four `*_revision` columns (revision drift is an independent restart concern +and is not part of the eligibility conjunction), and the `retired_*` bookkeeping +columns this feature adds. Numeric values are canonicalised, so a TTL that +round-trips through SQLite as `900.0` hashes identically to `900`. + +## Eligibility — the conjunction + +A registration is retired only when **every** one of these holds. Any missing +or contradictory evidence preserves the row. + +| Requirement | Preserve reason code when it fails | +| --- | --- | +| `status` is `active` | `already_terminal_registration` | +| Every field the conjunction reads is present (`REQUIRED_IDENTITY_FIELDS`) | `incomplete_registry_identity` | +| `last_heartbeat_at` parses as a UTC stamp | `unparsable_heartbeat` | +| Worker is not live | `worker_live` | +| PID probe returns a definite answer | `pid_liveness_unknown` | +| PID probe says the process is gone | `pid_alive` | +| Heartbeat has expired under the canonical TTL | `heartbeat_not_expired` | +| `ownership_state` is exactly `stale` | `ambiguous_ownership_state` | +| Repository binding present and canonical | `repository_binding_ambiguous` | +| No identity evidence shared with a live or unprobeable worker | `conflicting_identity_evidence` | +| No other active row claims the same (instance, namespace) while one may be live | `client_instance_conflict` | +| Not an active workflow-lease owner | `protected_active_workflow_owner` | +| Instance identity is launcher-minted (`inst-…`) | `untrusted_identity_provenance` | +| Row's `host_id` matches the host running retirement | `host_binding_unproven` | +| Boot identity is known on both sides | `boot_identity_unknown` | +| The pid number is not occupied by a different incarnation | `pid_reuse_detected` | + +Eligible rows carry `eligible_stale_orphan`. + +Three properties are worth stating explicitly: + +* **`pid_alive` can only withdraw liveness, never grant it** + (`WorkerRegistry.is_live`, #948 AC7). A heartbeat-lapsed but still-running + process therefore classifies as `stale` in the snapshot, yet #980's added + `pid_alive is False` requirement preserves it. An unprobeable PID (`None`) + also fails closed. +* **Affirmative identity proof is required (review 657 B2).** An earlier + revision required only that the pre-existing registry columns were non-null — + which a legacy `legacy-pid-…` row satisfies trivially, so a row that proved + nothing about *which* process it described was retireable. Retirement now + needs both halves of a positive proof: + * **Attribution** — a launcher-minted `inst-…` `client_instance_id`, so the + row is known to belong to one specific application launch rather than + having been inferred from pid proximity. + * **Fencing** — `host_id`, `boot_id`, and `process_start_time`, which turn a + bare pid into a statement about one process: which machine it ran on, which + boot of that machine, and which incarnation of that pid number. + + **Consequence, stated plainly:** registrations written before these columns + existed, and any row on a legacy instance identity, are preserved + *permanently*. They are retired only after their worker re-registers under a + trusted identity — never on weaker evidence. That the alternative would leave + legacy rows outstanding indefinitely is not a reason to relax the proof. +* **A live pid is an absolute block.** Even across a boot boundary, where the + number provably cannot belong to the registered process, an occupied pid + preserves the row rather than being argued away by the fencing proof. + +Multiple processes belonging to one legitimate worker cohort are not treated as +multiple independent workers: the fleet model from #948/#978 is preserved +unchanged, and sharing a role or profile is never a duplicate. + +## Tools + +### `gitea_plan_stale_worker_retirement` + +Read-only. Controller and reconciler only. + +| Parameter | Meaning | +| --- | --- | +| `remote` | `dadeschools` or `prgs` | +| `host`, `org`, `repo` | Optional overrides (audit context) | +| `canonical_repository` | Expected repository binding; defaults to the process root | + +Returns `registry_fingerprint`, `candidate_fingerprint`, +`candidate_worker_identities`, per-worker `candidates` and `preserved` entries +(each with `reason_code`, `detail`, and structured `evidence`), +`preserved_reason_counts`, `assessed_count`, `candidate_count`, +`preserved_count`, and `protected_active_workflow_owners`. +`mutation_performed` is always `false` and `read_only` is always `true`. + +Planning is deterministic: the same authoritative registry contents produce the +same plan and the same tokens regardless of when they are observed. + +### `gitea_apply_stale_worker_retirement` + +Mutating. Controller and reconciler only. + +| Parameter | Meaning | +| --- | --- | +| `registry_fingerprint` | The exact stable token the plan returned | +| `candidate_fingerprint` | The exact candidate-set token the plan returned | +| `worker_identities` | The exact candidate identities (list, or JSON / comma-separated string) | +| `remote`, `host`, `org`, `repo` | As above | +| `canonical_repository` | Must match the value the plan used | + +Before touching the registry, apply fails closed on: profile permission, role +kind, master parity (`mutation_safe`), stable-runtime mode, capability +resolution refreshed immediately before mutation, worker-registry availability, +workflow-lease enumeration failure, and daemon-cohort uniqueness +(`classify_cohort`). + +A matching token is necessary but never sufficient. Inside one +`BEGIN IMMEDIATE` transaction (`WorkerRegistry.retire_stale_workers`) the +server: + +1. re-reads the authoritative rows; +2. recomputes `registry_fingerprint` from *those* rows and compares — a mismatch + returns `registry_revision_moved` with `retired_count: 0` and no write; +3. recomputes the eligibility plan from *those* rows — re-reading the active + workflow leases rather than reusing the set captured before the transaction + opened, so a lease acquired after planning still preserves its worker — and + compares `candidate_fingerprint`. A mismatch returns `candidate_set_moved` + with `retired_count: 0` and no write; a lease-enumeration failure raises and + rolls the transaction back; +4. revalidates every requested identity against that fresh plan; +5. retires each survivor with a guarded `UPDATE` that additionally asserts + `status`, `last_heartbeat_at`, `generation_id`, `session_id`, + `fencing_epoch`, and `pid` are unchanged. A guard that matches no row + preserves the worker with `row_changed_since_plan`. + +There is no window between a safety check and its matching write, so a worker +that comes back to life, changes ownership, or is retired concurrently cannot be +removed on the strength of a stale observation. Any exception — including a +commit failure — rolls the whole transaction back and returns +`transaction_failed` with `success: false`, `retired_count: 0`, and +`mutation_performed: false`; a partial write can never be reported as success. + +### Outcomes + +| `outcome` | Meaning | `mutation_performed` | +| --- | --- | --- | +| `planned` | Dry-run result | `false` | +| `applied` | Transaction ran; see `retired` / `preserved` | `true` only if something was retired | +| `registry_revision_moved` | Registry changed between plan and apply | `false` | +| `candidate_set_moved` | Eligibility verdict changed between plan and apply | `false` | +| `already_retired` | Every requested row is already retired (idempotent replay) | `false` | +| `nothing_requested` | Empty target list | `false` | +| `transaction_failed` | Rolled back; nothing retired | `false` | + +## What retirement does to the fleet snapshot + +A retired row keeps its history: `status` moves to `retired` and `retired_at`, +`retired_by`, `retirement_reason` are recorded. Nothing is deleted. Because +`retired` is not `active`, the #978 snapshot counts the row as **historical**, +not stale, so `stale_worker_count` falls and historical rows never make the live +fleet unsafe by themselves. + +**Retirement does not repair untrusted live identity.** Live workers registered +under legacy `pid-`/`proc-` instance identities are preserved untouched and +keep their `legacy_incomplete_identity` blockers. Retiring every stale row can +therefore legitimately produce: + +* `stale_worker_count: 0` +* residual live `legacy_incomplete_identity` blockers +* `live_fleet_safe: false` + +That is a truthful result, and the `post_apply` block reports the remaining +blockers rather than claiming the fleet became safe. Trusted +`client_instance_id` propagation through launchers is a separate enrollment +problem. + +## Permissions + +Plan and apply are authorized differently, and deliberately so (review 657 B1). + +| | Plan | Apply | +| --- | --- | --- | +| Capability | `gitea.read` | `gitea.worker_registry.retire` | +| Nature | observational; opens no transaction, writes nothing | mutation | +| Roles | `controller`, `reconciler` | `controller`, `reconciler` | + +An earlier revision authorized apply with `gitea.read` alone, reasoning that +the mutation lands in the local control-plane registry rather than in Gitea. +That the write is local makes it **no less a mutation**: sharing an +observational permission class with plan meant any profile that could *look* +could also *destroy*. Apply now requires its own capability. + +* **Denied:** author, reviewer, merger, and every ordinary read-only profile — + they lack the capability, so they fail closed on the permission itself rather + than on the role check alone. The role restriction remains as defence in + depth: a profile mistakenly granted the capability still cannot reach apply + from an author, reviewer, or merger role. +* The capability is checked at entry **and** re-resolved immediately before the + registry mutation, so a profile change mid-call cannot be outrun. +* **No new Gitea write permission is introduced.** + `gitea.worker_registry.retire` authorizes exactly one local control-plane + transition (`worker_registrations.status -> retired`) and grants no branch, + issue, PR, review, merge, or restart authority. No author permission is + broadened. +* The fleet snapshot remains observational: nothing here turns it into a gate on + ordinary author work. + +### Operator step + +No profile holds `gitea.worker_registry.retire` by default, so apply is inert +until an operator adds it to the `allowed_operations` of the controller or +reconciler profile in `profiles.json`. Removing it again immediately and +completely revokes apply, while leaving plan and every other capability +untouched. That grant is a configuration change and is outside the scope of the +code that implements this feature. + +## External-state fencing + +`BEGIN IMMEDIATE` locks the worker registry and nothing else, so two inputs the +decision depends on sit outside the transaction's isolation domain: the +workflow-lease table in a separate control-plane database, and OS process +liveness. Re-reading them once during revalidation is not sufficient — the +per-target loop runs afterwards, so a lease acquired (or a pid revived) after +revalidation but before a given row's `UPDATE` would go unnoticed, and the +registry-column guard cannot catch it because no registry column changed. + +Two mechanisms close that window, both applied per target immediately before +its own write: + +* **`external_fence_fn`** — a version token over active leases + (`external_state_fingerprint`), captured inside the transaction *before* the + authoritative read and re-compared before every guarded `UPDATE`. Any movement + raises, rolling back the whole transaction: once the world has changed, every + remaining per-row decision was computed against a world that no longer exists. + An unreadable lease store raises rather than returning a token, because + "unreadable" must not silently compare equal to "unchanged". +* **`liveness_fn`** — a re-probe of process liveness and fencing identity that + must affirmatively re-establish that this exact process is gone. It compares + `process_start_time`, so a pid number reused since the plan is refused rather + than accepted. + +A caller that supplies no `liveness_fn` retires nothing +(`liveness_reprobe_unavailable`) rather than proceeding unfenced. The guarded +`UPDATE` additionally asserts `host_id`, `boot_id`, and `process_start_time` are +unchanged, and all three participate in the CAS token, so fencing movement +alone is enough to abort. + +## Non-goals + +* Killing or restarting processes. +* Editing session files or configuration. +* Direct database cleanup outside the sanctioned transaction. +* Retiring live workers. +* Backfilling trusted identity for legacy workers. +* Rewriting worker ownership. +* Cleaning unrelated workflow-lease or issue-claim registries. diff --git a/gitea_mcp_server.py b/gitea_mcp_server.py index e365da7..279cc88 100644 --- a/gitea_mcp_server.py +++ b/gitea_mcp_server.py @@ -19623,6 +19623,701 @@ def gitea_snapshot_instance_fleet( return snapshot +# --- #980 CAS-protected stale worker retirement --------------------------- + + +#: #980 review 657 B1. Apply is a mutation and must not be authorized by the +#: observational ``gitea.read`` class. Declared once here and consumed by both +#: the entry gate and the immediately-pre-mutation re-check so the two can +#: never drift apart; ``task_capability_map`` maps the apply tasks to the same +#: literal. +RETIREMENT_MUTATION_PERMISSION = "gitea.worker_registry.retire" + + +def _retirement_role_block(task: str) -> dict | None: + """Refuse the retirement surface to any non-controller/reconciler role. + + Role is a *secondary* control. The primary authority for apply is the + dedicated ``gitea.worker_registry.retire`` capability (#980 review 657 B1); + this check additionally pins the surface to the two roles that own fleet + reconciliation, so a profile mistakenly granted the permission still cannot + reach it from an author, reviewer, or merger role. Plan remains + ``gitea.read`` and is observational. + """ + profile = get_profile() + role = _profile_role_kind(profile) + if role in {"controller", "reconciler"}: + return None + return { + "success": False, + "allowed": False, + "mutation_performed": False, + "retired_count": 0, + "denied_role": role, + "required_roles": ["controller", "reconciler"], + "requested_task": task, + "reasons": [ + f"{task} is restricted to controller and reconciler roles; active " + f"role_kind is {role!r}. Author, reviewer, and merger profiles keep " + "gitea.read for diagnosis elsewhere but never receive worker-" + "registry retirement, and no unrelated mutation permission is " + "granted." + ], + "exact_next_action": ( + "Re-run from a prgs-controller or prgs-reconciler namespace." + ), + } + + +def _retirement_runtime_block() -> list[str]: + """Fail-closed runtime reasons that must stop a retirement apply (#980). + + ``gitea.read`` deliberately bypasses the #420 parity gate and the #615 + stable-runtime gate, because a stale server may still be *inspected*. An + apply is a mutation, so both gates are re-asserted explicitly here rather + than inherited. + """ + reasons: list[str] = [] + try: + parity = _current_master_parity() + except Exception as exc: + return [ + f"master parity could not be assessed (fail closed): {_redact(str(exc))}" + ] + if not parity.get("mutation_safe"): + reasons.append( + "runtime parity is not mutation-safe: " + f"{parity.get('summary') or 'stale runtime'}" + ) + try: + reasons.extend( + stable_control_runtime.runtime_block_reasons( + _current_runtime_mode_report() + ) + ) + except Exception as exc: + reasons.append( + f"runtime mode could not be assessed (fail closed): {_redact(str(exc))}" + ) + return reasons + + +def _retirement_protected_owners() -> tuple[dict | None, dict]: + """Workflow owners that must never be retired as stale workers. + + A registration whose process still owns an active control-plane lease is a + live workflow participant needing its own reconciliation, not a stale + orphan. Failure to enumerate leases is ambiguity, so it aborts rather than + proceeding with an empty protection set. + """ + db, errs = _control_plane_db_or_error() + if db is None: + return ( + { + "success": False, + "mutation_performed": False, + "retired_count": 0, + "reasons": [ + "active workflow leases could not be enumerated, so " + "protected owners are unknown (fail closed)", + *errs, + ], + }, + {}, + ) + try: + leases = db.list_leases(statuses=["active"], limit=1000) + except Exception as exc: + return ( + { + "success": False, + "mutation_performed": False, + "retired_count": 0, + "reasons": [ + "active workflow leases could not be enumerated, so " + "protected owners are unknown (fail closed): " + f"{_redact(str(exc))}" + ], + }, + {}, + ) + session_ids: set[str] = set() + pids: set[int] = set() + for lease in leases: + owner = lease.get("session_id") or lease.get("owner_session_id") + if owner: + session_ids.add(str(owner)) + for key in ("owner_pid", "session_pid"): + value = lease.get(key) + if value is None: + continue + try: + pids.add(int(value)) + except (TypeError, ValueError): + continue + return None, { + "session_ids": sorted(session_ids), + "pids": sorted(pids), + "active_lease_count": len(leases), + } + + +def _retirement_plan(records, canonical_repository: str, protected: dict) -> dict: + """The single planning path shared by dry run and in-transaction revalidation.""" + import mcp_fleet_retirement + import mcp_process_fencing + + return mcp_fleet_retirement.plan_stale_worker_retirement( + records, + pid_alive_probe=issue_lock_store.is_process_alive, + canonical_repository=canonical_repository, + protected_session_ids=(protected or {}).get("session_ids"), + protected_pids=(protected or {}).get("pids"), + # #980 review 657 B2: the fencing evidence that makes a recorded pid + # interpretable. Probed live so a plan produced on one host can never + # authorize a retirement carried out on another. + current_host_id=mcp_process_fencing.current_host_id(), + current_boot_id=mcp_process_fencing.current_boot_id(), + start_time_probe=mcp_process_fencing.process_start_time, + ) + + +def _retirement_external_fence() -> str: + """Version token over lease + liveness state consumed by a retirement (#980 B3). + + Recomputed inside the retirement transaction immediately before every + guarded write. Because neither the control-plane lease database nor the OS + process table is covered by the registry's ``BEGIN IMMEDIATE``, this token + is the only thing that can detect either of them moving mid-transaction. + + A failure to read leases raises rather than returning a token: an + unreadable external world is indistinguishable from a changed one, so it + must abort the transaction rather than silently compare equal. + """ + import mcp_fleet_retirement + + db, errs = _control_plane_db_or_error() + if db is None: + raise RuntimeError( + "active workflow leases could not be read while fencing the " + f"retirement transaction: {errs}" + ) + leases = db.list_leases(statuses=["active"], limit=1000) + return mcp_fleet_retirement.external_state_fingerprint(leases) + + +def _retirement_liveness_reprobe(row: dict) -> dict: + """Immediate pre-write re-establishment that one specific process is gone. + + #980 review 657 B3: revalidation happens once per transaction, but the + per-target loop runs afterwards, so this re-probe closes the remaining + window between "this row was judged safe" and "this row is written". It + re-reads OS state rather than trusting the plan, and it compares the + fencing triple so a pid number reused since the plan cannot pass. + """ + import mcp_fleet_retirement + import mcp_process_fencing + + pid = row.get("pid") + try: + alive = issue_lock_store.is_process_alive(pid) + except Exception as exc: + return { + "safe": False, + "reason_code": mcp_fleet_retirement.REASON_PID_UNKNOWN, + "detail": ( + f"pid {pid!r} could not be re-probed immediately before the " + f"write: {_redact(str(exc))}" + ), + } + if alive: + live_start = mcp_process_fencing.process_start_time(pid) + recorded_start = (row.get("process_start_time") or "").strip() + reused = bool(live_start) and live_start != recorded_start + return { + "safe": False, + "reason_code": ( + mcp_fleet_retirement.REASON_PID_REUSED + if reused + else mcp_fleet_retirement.REASON_PID_ALIVE + ), + "detail": ( + f"pid {pid!r} is alive at write time" + + ( + f" but is a different incarnation ({live_start!r} != " + f"{recorded_start!r})" + if reused + else "" + ) + ), + "evidence": { + "pid": pid, + "pid_alive": True, + "live_process_start_time": live_start, + "recorded_process_start_time": recorded_start, + }, + } + + host_now = mcp_process_fencing.current_host_id() + boot_now = mcp_process_fencing.current_boot_id() + if not host_now or (row.get("host_id") or "").strip() != host_now: + return { + "safe": False, + "reason_code": mcp_fleet_retirement.REASON_HOST_UNPROVEN, + "detail": ( + "host identity no longer agrees with the registration at write " + f"time (row {row.get('host_id')!r}, now {host_now!r})" + ), + } + if not boot_now: + return { + "safe": False, + "reason_code": mcp_fleet_retirement.REASON_BOOT_UNKNOWN, + "detail": "boot identity became unobtainable before the write", + } + return { + "safe": True, + "evidence": { + "pid": pid, + "pid_alive": False, + "host_id": host_now, + "boot_id": boot_now, + }, + } + + +def _retirement_revalidation_plan(records, canonical_repository: str) -> dict: + """Revalidation planner used *inside* the retirement transaction. + + Deliberately re-reads the active workflow leases rather than reusing the + set captured before the transaction opened: a lease acquired after planning + must still preserve its worker. An enumeration failure raises, which rolls + the transaction back and retires nothing. + """ + block, protected = _retirement_protected_owners() + if block: + raise RuntimeError( + "active workflow leases could not be re-read inside the retirement " + "transaction; refusing to retire anything" + ) + return _retirement_plan(records, canonical_repository, protected) + + +@mcp.tool() +def gitea_plan_stale_worker_retirement( + remote: str = "dadeschools", + host: str | None = None, + org: str | None = None, + repo: str | None = None, + canonical_repository: str | None = None, +) -> dict: + """Read-only: authoritative retirement plan for stale worker rows (#980). + + Controller and reconciler only. Reads the worker registry, classifies every + registration with the same assessor the #978 fleet snapshot uses, and + returns the exact set of registrations that are conclusively stale orphans + together with a **stable** ``registry_fingerprint`` and an exact + ``candidate_fingerprint``. + + The fingerprint is derived only from canonical retirement-relevant registry + content — never from ``snapshot_at``, wall-clock, request, or report time — + so two plans over an unchanged registry agree and the apply compare-and-swap + can actually pass. Row order never affects it. + + Retires nothing. Live workers, workers whose PID cannot be probed, workers + with unparsable heartbeats, workers sharing identity evidence with a live or + unprobeable worker, foreign or unbound repositories, and workers that still + own an active workflow lease are all preserved with a reason code. + + Args: + remote: Known instance — 'dadeschools' or 'prgs'. + host: Optional host override. + org: Optional org override (audit context only). + repo: Optional repo override (audit context only). + canonical_repository: Expected repository binding for + foreign-repository classification (defaults to the process root). + """ + read_block = _profile_operation_gate("gitea.read") + if read_block: + return { + "success": False, + "read_only": True, + "mutation_performed": False, + "reasons": read_block, + "permission_report": _permission_block_report("gitea.read"), + } + + role_block = _retirement_role_block("gitea_plan_stale_worker_retirement") + if role_block: + role_block["read_only"] = True + return role_block + + registry = _worker_registry() + if registry is None: + return { + "success": False, + "read_only": True, + "mutation_performed": False, + "reasons": [ + "worker registry is unavailable; cannot produce an authoritative " + "retirement plan (fail closed)" + ], + "exact_next_action": ( + "Ensure GITEA_WORKER_REGISTRY_DB is writable and re-run after " + "workers have registered." + ), + } + + protected_block, protected = _retirement_protected_owners() + if protected_block: + protected_block["read_only"] = True + return protected_block + + try: + records = registry.list_workers(status=None) + except Exception as exc: + return { + "success": False, + "read_only": True, + "mutation_performed": False, + "reasons": [ + f"failed to read worker registry: {type(exc).__name__}: " + f"{_redact(str(exc))}" + ], + } + + canon = canonical_repository or PROJECT_ROOT + plan = _retirement_plan(records, canon, protected) + profile = get_profile() + plan["role_kind"] = _profile_role_kind(profile) + plan["profile"] = profile.get("profile_name") + plan["remote"] = _effective_remote(remote) + plan["repository"] = {"org": org, "repo": repo, "canonical_repository": canon} + plan["protected_active_workflow_owners"] = protected + plan["apply_tool"] = "gitea_apply_stale_worker_retirement" + plan["permission_scope"] = { + "read_only": True, + "granted_operations": ["gitea.read"], + "denied_unrelated_mutations": True, + "note": ( + "Planning is strictly observational. It does not authorize branch, " + "issue, PR, review, merge, or restart mutations, and it retires " + "nothing." + ), + } + plan["exact_next_action"] = ( + "Pass registry_fingerprint, candidate_fingerprint, and the exact " + "candidate_worker_identities to gitea_apply_stale_worker_retirement." + if plan.get("candidate_count") + else "No registration is conclusively stale; nothing to apply." + ) + return plan + + +@mcp.tool() +def gitea_apply_stale_worker_retirement( + registry_fingerprint: str, + candidate_fingerprint: str, + worker_identities: list | str, + remote: str = "dadeschools", + host: str | None = None, + org: str | None = None, + repo: str | None = None, + canonical_repository: str | None = None, +) -> dict: + """Retire conclusively stale worker registrations under CAS (#980). + + Controller and reconciler only. Requires the exact ``registry_fingerprint`` + and ``candidate_fingerprint`` returned by + ``gitea_plan_stale_worker_retirement`` plus the exact candidate identity + list. A matching token is necessary but never sufficient: inside a single + ``BEGIN IMMEDIATE`` transaction the registry is re-read, the fingerprint is + recomputed from those rows, the eligibility plan is recomputed from those + rows, and every target is independently revalidated immediately before its + own guarded ``UPDATE``. Any drift retires zero workers and reports + ``registry_revision_moved`` or ``candidate_set_moved``. + + Retiring stale rows does not repair untrusted live identity. Live workers + registered under legacy ``pid-``/``proc-`` instance identities remain + untouched and their ``legacy_incomplete_identity`` blockers remain + outstanding, so the result never claims the fleet became safe. + + Args: + registry_fingerprint: Stable token from the plan (CAS expectation). + candidate_fingerprint: Exact candidate-set token from the plan. + worker_identities: The exact candidate identities the plan returned + (list, or a JSON / comma-separated string). + remote: Known instance — 'dadeschools' or 'prgs'. + host: Optional host override. + org: Optional org override (audit context only). + repo: Optional repo override (audit context only). + canonical_repository: Expected repository binding (defaults to the + process root); must match the value the plan used. + """ + import json as _json + + import mcp_fleet_retirement + + # #980 review 657 B1: the dedicated mutation capability, not gitea.read. + # A profile holding only the observational permission cannot get past this + # line, and the role check below is defence in depth rather than the sole + # authority. + mutation_block = _profile_operation_gate(RETIREMENT_MUTATION_PERMISSION) + if mutation_block: + return { + "success": False, + "mutation_performed": False, + "retired_count": 0, + "requested_task": "apply_stale_worker_retirement", + "required_operation_permission": RETIREMENT_MUTATION_PERMISSION, + "reasons": mutation_block, + "permission_report": _permission_block_report( + RETIREMENT_MUTATION_PERMISSION + ), + } + + role_block = _retirement_role_block("gitea_apply_stale_worker_retirement") + if role_block: + return role_block + + runtime_reasons = _retirement_runtime_block() + if runtime_reasons: + return { + "success": False, + "mutation_performed": False, + "retired_count": 0, + "blocker_kind": "runtime_not_mutation_safe", + "reasons": runtime_reasons, + "exact_next_action": ( + "Restore runtime parity on the stable control checkout, then " + "re-plan and re-apply." + ), + } + + # Fresh identity + capability resolution immediately before mutation. + profile = get_profile() + role_kind = _profile_role_kind(profile) + required_permission = task_capability_map.required_permission( + "apply_stale_worker_retirement" + ) + required_role = task_capability_map.required_role( + "apply_stale_worker_retirement" + ) + permission_ok, permission_reason = gitea_config.check_operation( + required_permission, + profile.get("allowed_operations") or [], + profile.get("forbidden_operations") or [], + ) + if not permission_ok: + return { + "success": False, + "mutation_performed": False, + "retired_count": 0, + "requested_task": "apply_stale_worker_retirement", + "required_operation_permission": required_permission, + "required_role_kind": required_role, + "reasons": [ + "capability resolution immediately before apply refused this " + f"session: {permission_reason}" + ], + } + try: + resolved_host = host or REMOTES[_effective_remote(remote)]["host"] + authenticated_username = _authenticated_username(resolved_host) + except Exception: + authenticated_username = None + + registry = _worker_registry() + if registry is None: + return { + "success": False, + "mutation_performed": False, + "retired_count": 0, + "reasons": [ + "worker registry is unavailable; refusing to retire anything " + "(fail closed)" + ], + } + + raw = worker_identities + if isinstance(raw, str): + text = raw.strip() + try: + raw = _json.loads(text) + except Exception: + raw = [part.strip() for part in text.split(",") if part.strip()] + if isinstance(raw, str): + raw = [raw] + if not isinstance(raw, list): + return { + "success": False, + "mutation_performed": False, + "retired_count": 0, + "reasons": ["worker_identities must be a list of worker identities"], + } + targets = [str(item).strip() for item in raw if str(item).strip()] + + protected_block, protected = _retirement_protected_owners() + if protected_block: + return protected_block + + # #948 daemon-cohort uniqueness: a contested generation or a reused worker + # identity means ownership is ambiguous fleet-wide, so retire nothing. + try: + cohort = mcp_worker_identity.classify_cohort( + registry.list_workers(status=mcp_worker_identity.STATUS_ACTIVE), + pid_alive_probe=issue_lock_store.is_process_alive, + ) + except Exception as exc: + return { + "success": False, + "mutation_performed": False, + "retired_count": 0, + "reasons": [ + "daemon-cohort uniqueness could not be assessed (fail closed): " + f"{_redact(str(exc))}" + ], + } + if cohort.get("blocked"): + return { + "success": False, + "mutation_performed": False, + "retired_count": 0, + "blocker_kind": cohort.get("blocker_kind"), + "reasons": [ + "daemon-cohort uniqueness failed; worker ownership is contested", + *(cohort.get("reasons") or []), + ], + "cohort": { + "blocked_worker_identities": cohort.get("blocked_worker_identities"), + "duplicate_identities": cohort.get("duplicate_identities"), + "contested_generations": cohort.get("contested_generations"), + }, + } + + canon = canonical_repository or PROJECT_ROOT + acting = "/".join( + part + for part in (authenticated_username, profile.get("profile_name")) + if part + ) or "unknown" + + result = registry.retire_stale_workers( + expected_registry_fingerprint=registry_fingerprint, + expected_candidate_fingerprint=candidate_fingerprint, + worker_identities=targets, + fingerprint_fn=mcp_fleet_retirement.registry_fingerprint, + plan_fn=lambda rows: _retirement_revalidation_plan(rows, canon), + retired_by=acting, + retirement_reason=mcp_fleet_retirement.REASON_ELIGIBLE, + external_fence_fn=_retirement_external_fence, + liveness_fn=_retirement_liveness_reprobe, + ) + + result["role_kind"] = role_kind + result["profile"] = profile.get("profile_name") + result["remote"] = _effective_remote(remote) + result["repository"] = {"org": org, "repo": repo, "canonical_repository": canon} + result["protected_active_workflow_owners"] = protected + result["permission_scope"] = { + "granted_operations": [RETIREMENT_MUTATION_PERMISSION], + "mutation_capability": RETIREMENT_MUTATION_PERMISSION, + "plan_capability": "gitea.read", + "control_plane_mutation": "worker_registrations.status -> retired", + "denied_unrelated_mutations": True, + "note": ( + "This capability retires local worker-registry rows only. It grants " + "no branch, issue, PR, review, merge, or restart authority, and it " + "never kills or restarts a process." + ), + } + result["post_apply"] = _retirement_post_apply( + registry, result, canonical_repository=canon + ) + + try: + gitea_audit.write_event( + gitea_audit.build_event( + action="gitea_apply_stale_worker_retirement", + result=( + gitea_audit.SUCCEEDED + if result.get("mutation_performed") + else gitea_audit.BLOCKED + if not result.get("success") + else gitea_audit.ALLOWED + ), + remote=_effective_remote(remote), + repository=canon, + profile_name=profile.get("profile_name"), + audit_label=profile.get("audit_label"), + authenticated_username=authenticated_username, + task_role=role_kind, + operation="worker_registry.retire_stale_workers", + reason=result.get("outcome"), + request_metadata=mcp_fleet_retirement.summarize_plan(result), + ) + ) + except Exception: + pass + + return result + + +def _retirement_post_apply( + registry, result: dict, *, canonical_repository: str +) -> dict: + """Fresh fleet verification after a retirement attempt (#980 requirement 6).""" + import mcp_fleet_retirement + import mcp_fleet_snapshot as _fleet_after + + try: + after_records = registry.list_workers(status=None) + after = _fleet_after.snapshot_instance_fleet( + after_records, + pid_alive_probe=issue_lock_store.is_process_alive, + canonical_repository=canonical_repository, + ) + except Exception as exc: + return { + "available": False, + "reasons": [ + f"post-apply verification could not be produced: {_redact(str(exc))}" + ], + } + retired_ids = { + str(r.get("worker_identity")) for r in result.get("retired") or [] + } + return { + "available": True, + "registry_fingerprint": mcp_fleet_retirement.registry_fingerprint( + after_records + ), + "live_worker_count": after.get("live_worker_count"), + "stale_worker_count": after.get("stale_worker_count"), + "historical_worker_count": after.get("historical_worker_count"), + "retired_still_counted_live": sorted( + str(w.get("worker_identity")) + for w in after.get("live_workers") or [] + if str(w.get("worker_identity")) in retired_ids + ), + "retired_still_counted_stale": sorted( + str(w.get("worker_identity")) + for w in after.get("stale_workers") or [] + if str(w.get("worker_identity")) in retired_ids + ), + "live_fleet_safe": after.get("live_fleet_safe"), + "remaining_blockers": [ + {"classification": f.get("classification"), "detail": f.get("detail")} + for f in after.get("active_blockers") or [] + ], + "note": ( + "Stale retirement does not repair untrusted live identity; residual " + "legacy_incomplete_identity blockers keep live_fleet_safe false and " + "that is a truthful result." + ), + } + + @mcp.tool() 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..78b7097 --- /dev/null +++ b/mcp_fleet_retirement.py @@ -0,0 +1,776 @@ +"""CAS-protected retirement planning for stale worker registrations (#980). + +#978 (merged PR #979) made the fleet observable: every registered namespace +worker, its instance attribution, its heartbeat freshness, and a structured +classification. It deliberately stopped there — the snapshot is read-only and +the control plane still had no sanctioned way to retire registry rows whose +owning process is conclusively gone. + +This module is the *decision layer* for that retirement. It is pure: callers +supply registry rows, a clock, and a PID probe; nothing here opens SQLite, +scans process tables, or mutates state. The transactional apply lives in +:meth:`mcp_worker_identity.WorkerRegistry.retire_stale_workers`, which calls +back into these same pure functions so plan and apply can never disagree about +what "the registry looks like" or "which rows are eligible". + +Why a separate token +-------------------- + +``mcp_fleet_snapshot._consistency_token`` seeds its digest with ``snapshot_at`` +at second precision, so ``registry_revision`` changes on every call even when +no registry row changed. A compare-and-swap gated on it can never pass — a +dry-run/apply cycle spanning more than one second aborts unconditionally. That +token is still useful as an observation stamp, so it is left exactly as it is; +#980 gets its own :func:`registry_fingerprint`, derived *only* from canonical +retirement-relevant row content: + +* identical registry contents observed at any two times produce the same token, +* row order never affects the token (serialized rows are sorted), +* any create/delete/identity/liveness/ownership/registration-state change to a + retirement-relevant field changes the token. + +Fail-closed posture +------------------- + +A worker is retired only when the control plane *conclusively* establishes it +is a stale orphan. Missing evidence is never read as permission: an unprobeable +PID, an unparsable heartbeat, a row that shares identity evidence with a live +or unprobeable worker, a foreign or absent repository binding, or a worker that +still owns an active workflow lease all preserve the row. + +Trusted launcher identity (``inst-…`` provenance) is deliberately *not* part of +the conjunction. #980 places trusted ``client_instance_id`` propagation out of +scope and lists "backfilling trusted identity for legacy workers" as a non-goal; +requiring it here would preserve every legacy row forever and make the feature +inert. What *is* required is that the registry fields the conjunction reads are +actually present — see :data:`REQUIRED_IDENTITY_FIELDS`. +""" + +from __future__ import annotations + +import hashlib +from datetime import datetime +from typing import Any, Callable, Iterable, Mapping, Sequence + +import mcp_fleet_snapshot as fleet +import mcp_worker_identity as mwi + +# --- Outcomes ------------------------------------------------------------- + +OUTCOME_PLANNED = "planned" +OUTCOME_APPLIED = "applied" +OUTCOME_REGISTRY_MOVED = "registry_revision_moved" +OUTCOME_CANDIDATES_MOVED = "candidate_set_moved" +OUTCOME_ALREADY_RETIRED = "already_retired" +OUTCOME_NOTHING_REQUESTED = "nothing_requested" + +# --- Reason codes --------------------------------------------------------- + +#: The only reason code that authorizes retirement. +REASON_ELIGIBLE = "eligible_stale_orphan" + +REASON_ALREADY_TERMINAL = "already_terminal_registration" +REASON_AMBIGUOUS_OWNERSHIP = "ambiguous_ownership_state" +REASON_CONFLICTING_IDENTITY = "conflicting_identity_evidence" +REASON_FOREIGN_REPOSITORY = "repository_binding_ambiguous" +REASON_HEARTBEAT_FRESH = "heartbeat_not_expired" +REASON_INCOMPLETE_IDENTITY = "incomplete_registry_identity" +REASON_NOT_IN_PLAN = "not_in_current_plan" +REASON_PID_ALIVE = "pid_alive" +REASON_PID_UNKNOWN = "pid_liveness_unknown" +REASON_PROTECTED_OWNER = "protected_active_workflow_owner" +REASON_ROW_CHANGED = "row_changed_since_plan" +REASON_ROW_MISSING = "registration_missing" +REASON_UNPARSABLE_HEARTBEAT = "unparsable_heartbeat" +REASON_WORKER_LIVE = "worker_live" + +# --- #980 review 657 B2: affirmative identity/liveness proof --------------- + +#: The registration's instance identity is not launcher-minted (``inst-…``), +#: so nothing proves which application launch this row belongs to. +REASON_UNTRUSTED_PROVENANCE = "untrusted_identity_provenance" +#: The row does not record which host its pid belongs to, or records a +#: different host than the one probing. A local pid probe cannot speak for a +#: process on another machine. +REASON_HOST_UNPROVEN = "host_binding_unproven" +#: Boot identity is missing on the row or unobtainable here, so a recorded pid +#: cannot be compared against a live pid at all. +REASON_BOOT_UNKNOWN = "boot_identity_unknown" +#: The pid is alive but belongs to a different process incarnation than the one +#: registered — reported distinctly from a plain live worker. +REASON_PID_REUSED = "pid_reuse_detected" +#: Two active registrations claim one client instance within one namespace. +REASON_INSTANCE_CONFLICT = "client_instance_conflict" +#: The immediate pre-write re-probe could not re-establish death (#980 B3). +REASON_LIVENESS_REPROBE = "liveness_reprobe_refused" + +#: Instance-identity prefix minted by the trusted launcher. Kept in sync with +#: ``mcp_fleet_snapshot._TRUSTED_INSTANCE_PREFIX`` through +#: :func:`mcp_fleet_snapshot.assess_instance_identity`, which stays the single +#: authority on what "trusted" means — this module never re-implements it. +TRUSTED_INSTANCE_PREFIX = "inst-" + +#: Registry columns that must carry a usable value before the eligibility +#: conjunction can even be evaluated. Absence is ambiguity, not permission. +#: +#: #980 review 657 B2 added the fencing triple. Before it, "complete identity" +#: meant only that the pre-existing columns were non-null, which a legacy +#: ``legacy-pid-…`` row satisfies trivially — so a row that proved nothing about +#: *which* process it described was retireable. The triple is what makes a +#: recorded pid interpretable: which machine it ran on, which boot of that +#: machine, and which incarnation of that pid number. A registration written +#: before these columns existed carries NULL and is therefore preserved +#: permanently, which is the intended fail-closed outcome. +REQUIRED_IDENTITY_FIELDS: tuple[str, ...] = ( + "worker_identity", + "client_instance_id", + "session_id", + "generation_id", + "status", + "started_at", + "last_heartbeat_at", + "heartbeat_ttl_seconds", + "pid", + "host_id", + "boot_id", + "process_start_time", +) + +#: Canonical retirement-relevant content. Ordering here is fixed and part of +#: the token contract; adding a field changes every fingerprint, so a change +#: here is a deliberate contract revision. +#: +#: Deliberately excluded: ``token_fingerprint`` (credential-adjacent, never a +#: retirement input), the four ``*_revision`` columns (revision drift is an +#: independent restart concern and is not part of the eligibility conjunction), +#: and the ``retired_*`` bookkeeping columns this feature adds. +FINGERPRINT_FIELDS: tuple[str, ...] = ( + "worker_identity", + "client_name", + "client_instance_id", + "session_id", + "generation_id", + "role", + "profile", + "namespace", + "remote", + "repository_binding", + "pid", + "process_identity", + "transport", + "started_at", + "last_heartbeat_at", + "heartbeat_ttl_seconds", + "fencing_epoch", + "status", + "fleet_run_id", + "authenticated_account", + "instance_id_provenance", + # #980 review 657 B3: fencing evidence is a retirement input, so moving it + # must move the CAS token. Without these, a row whose host, boot, or + # process incarnation changed would hash identically to the row the plan + # approved. + "host_id", + "boot_id", + "process_start_time", +) + +_FINGERPRINT_VERSION = "registryfp-v1" +_CANDIDATE_VERSION = "candidatefp-v1" +_UNIT = "\x1f" +_RECORD = "\x1e" + + +def _canon(value: Any) -> str: + """Stable text for one field value, independent of Python/SQLite typing. + + ``900`` and ``900.0`` are the same TTL and must hash the same; a value that + round-trips through SQLite as REAL must not produce a different token than + the same value supplied by a caller as ``int``. + """ + if value is None: + return "" + if isinstance(value, bool): + return "true" if value else "false" + if isinstance(value, float): + if value != value or value in (float("inf"), float("-inf")): + return repr(value) + if value.is_integer(): + return str(int(value)) + return repr(value) + if isinstance(value, int): + return str(value) + return str(value) + + +def _serialize_row(row: Mapping[str, Any]) -> str: + return _UNIT.join(f"{name}={_canon(row.get(name))}" for name in FINGERPRINT_FIELDS) + + +def _digest(version: str, serialized: Sequence[str], prefix: str) -> str: + ordered = sorted(serialized) + material = _RECORD.join([version, str(len(ordered)), *ordered]) + return f"{prefix}-{hashlib.sha256(material.encode('utf-8')).hexdigest()[:32]}" + + +def registry_fingerprint(rows: Iterable[Mapping[str, Any]]) -> str: + """Content-derived compare-and-swap token for the worker registry (#980). + + Derived exclusively from :data:`FINGERPRINT_FIELDS` across every row. It + contains no ``snapshot_at``, wall-clock, request, or report-generation + time, so two observations of an unchanged registry always agree, and the + serialized rows are sorted so iteration order cannot perturb the digest. + """ + return _digest( + _FINGERPRINT_VERSION, + [_serialize_row(row) for row in rows], + "registryfp", + ) + + +def candidate_fingerprint(candidate_rows: Iterable[Mapping[str, Any]]) -> str: + """Exact-candidate-set token over the selected rows' canonical content. + + A matching :func:`registry_fingerprint` already implies these rows are + unchanged; this second token additionally pins *which* rows the operator + approved, so an apply can never widen or narrow the approved set. + """ + return _digest( + _CANDIDATE_VERSION, + [_serialize_row(row) for row in candidate_rows], + "candidatefp", + ) + + +def _probe_pid( + pid: Any, pid_alive_probe: Callable[[int | None], bool | None] | None +) -> bool | None: + if pid_alive_probe is None or pid is None: + return None + try: + result = pid_alive_probe(pid) + except Exception: + return None + return None if result is None else bool(result) + + +def _reuse_detected( + row: Mapping[str, Any], + live_start_time: str | None, + current_host_id: str | None, + current_boot_id: str | None, +) -> bool: + """Is the pid occupied by a *different* incarnation than the one recorded? + + Only meaningful when the recorded pid is comparable to the live one — same + machine, same boot. Across hosts or boots the number is unrelated by + construction and reuse is not the interesting question. + """ + recorded_host = (row.get("host_id") or "").strip() + recorded_boot = (row.get("boot_id") or "").strip() + if not current_host_id or recorded_host != current_host_id: + return False + if not current_boot_id or recorded_boot != current_boot_id: + return False + recorded_start = (row.get("process_start_time") or "").strip() + return bool(live_start_time) and live_start_time != recorded_start + + +def _instance_key(row: Mapping[str, Any]) -> tuple[str, str] | None: + """The (instance, namespace) pair #978 requires to be unique among live rows.""" + instance = _canon(row.get("client_instance_id")) + namespace = _canon(row.get("namespace")) + if not instance or not namespace: + return None + return (instance, namespace) + + +def _probe_start_time( + pid: Any, start_time_probe: Callable[[Any], str | None] | None +) -> str | None: + if start_time_probe is None or pid is None: + return None + try: + return start_time_probe(pid) + except Exception: + return None + + +def _evidence( + row: Mapping[str, Any], + snapshot_row: Mapping[str, Any], + pid_alive: bool | None, +) -> dict[str, Any]: + liveness = snapshot_row.get("liveness") or {} + return { + "worker_identity": row.get("worker_identity"), + "client_type": snapshot_row.get("client_type"), + "client_instance_id": row.get("client_instance_id"), + "fleet_run_id": row.get("fleet_run_id"), + "namespace": row.get("namespace"), + "profile": row.get("profile"), + "declared_role": row.get("role"), + "session_id": row.get("session_id"), + "generation_id": row.get("generation_id"), + "fencing_epoch": row.get("fencing_epoch"), + "process_identity": snapshot_row.get("process_identity"), + "pid": row.get("pid"), + "pid_alive": pid_alive, + "host_id": row.get("host_id"), + "boot_id": row.get("boot_id"), + "process_start_time": row.get("process_start_time"), + "repository_binding": row.get("repository_binding"), + "foreign_repository": bool(snapshot_row.get("foreign_repository")), + "status": row.get("status"), + "started_at": row.get("started_at"), + "last_heartbeat_at": row.get("last_heartbeat_at"), + "heartbeat_ttl_seconds": row.get("heartbeat_ttl_seconds"), + "heartbeat_age_seconds": liveness.get("heartbeat_age_seconds"), + "heartbeat_fresh": liveness.get("heartbeat_fresh"), + "live": bool(snapshot_row.get("live")), + "ownership_state": snapshot_row.get("ownership_state"), + "instance_id_provenance": snapshot_row.get("instance_id_provenance"), + "instance_identity_trusted": bool( + snapshot_row.get("instance_identity_trusted") + ), + } + + +def _conflict_keys( + row: Mapping[str, Any], snapshot_row: Mapping[str, Any] +) -> list[tuple[str, str]]: + keys: list[tuple[str, str]] = [] + for name, value in ( + ("session_id", row.get("session_id")), + ("generation_id", row.get("generation_id")), + ("process_identity", snapshot_row.get("process_identity")), + ("pid", row.get("pid")), + ): + text = _canon(value) + if text: + keys.append((name, text)) + return keys + + +def external_state_fingerprint( + leases: Iterable[Mapping[str, Any]], + *, + liveness: Iterable[tuple[Any, Any]] = (), +) -> str: + """Version token over the external state a retirement decision consumed. + + #980 review 657 B3: ``BEGIN IMMEDIATE`` on the worker registry does not + cover the control-plane lease table or the OS process table, so those + inputs can move while the transaction is open. This token lets the + transaction detect that movement: it is captured before the authoritative + read and re-compared immediately before every guarded write, and any + difference aborts rather than retiring against evidence that has changed. + + Only ownership-relevant lease fields participate, so unrelated churn (a + heartbeat timestamp advancing on an unrelated lease) does not cause + spurious aborts, while an acquire, release, or owner change always does. + """ + lease_units: list[str] = [] + for lease in leases: + lease_units.append( + _UNIT.join( + f"{name}={_canon(lease.get(name))}" + for name in ( + "lease_id", + "role", + "target", + "status", + "session_id", + "owner_session_id", + "owner_pid", + "session_pid", + "generation", + ) + ) + ) + for pid, alive in liveness: + lease_units.append(f"liveness{_UNIT}pid={_canon(pid)}{_UNIT}alive={_canon(alive)}") + return _digest("externalfp-v1", lease_units, "externalfp") + + +def assess_retirement_identity_proof( + row: Mapping[str, Any], + snapshot_row: Mapping[str, Any], + *, + current_host_id: str | None, + current_boot_id: str | None, + live_start_time: str | None, + pid_alive: bool | None, +) -> tuple[str, str] | None: + """Affirmative proof that this row names one specific, now-dead process. + + Returns ``None`` when the proof holds, or ``(reason_code, detail)`` naming + the first thing that could not be established. #980 review 657 B2: absence + of evidence is never read as staleness, so every branch here refuses on + *missing* information exactly as firmly as on contradictory information. + + The proof has two independent halves and needs both: + + * **Attribution** — a launcher-minted ``inst-…`` instance identity, so the + row is known to belong to one specific application launch rather than + having been inferred from pid proximity. + * **Fencing** — the row's host matches the host doing the probing, boot + identity is known on both sides, and the recorded process incarnation + agrees with whatever currently occupies that pid number. + + Requiring trusted attribution means pre-#978 ``legacy-pid-…`` rows are + preserved permanently. That is deliberate. The reviewer specifically + rejected the argument that legacy rows "would remain forever" as grounds + for a weaker proof, and #980 lists backfilling trusted identity for legacy + workers as a non-goal — so those rows are retired only after their worker + re-registers under a trusted identity, never on weaker evidence. + """ + if not snapshot_row.get("instance_identity_trusted"): + return ( + REASON_UNTRUSTED_PROVENANCE, + "client_instance_id " + f"{row.get('client_instance_id')!r} is not launcher-minted " + f"({snapshot_row.get('instance_id_provenance')!r}); nothing proves " + "which application launch this registration belongs to", + ) + + recorded_host = (row.get("host_id") or "").strip() + if not current_host_id: + return ( + REASON_HOST_UNPROVEN, + "this process cannot establish its own host identity, so a local " + "pid probe cannot be attributed to any machine", + ) + if recorded_host != current_host_id: + return ( + REASON_HOST_UNPROVEN, + f"registration is bound to host {recorded_host!r} but retirement is " + f"running on {current_host_id!r}; a local pid probe says nothing " + "about a process on another machine", + ) + + recorded_boot = (row.get("boot_id") or "").strip() + if not current_boot_id: + return ( + REASON_BOOT_UNKNOWN, + "the current boot identity could not be determined, so a recorded " + "pid cannot be compared against a live pid", + ) + + recorded_start = (row.get("process_start_time") or "").strip() + + if recorded_boot != current_boot_id: + # A different boot is the strongest possible death evidence: every pid + # from a previous boot is gone, and pid numbers restart, so whatever + # occupies this number now is unrelated by construction. + return None + + # Same boot: the pid number is directly comparable, so the recorded + # incarnation must still agree with whatever holds that number. + if pid_alive and live_start_time and live_start_time != recorded_start: + return ( + REASON_PID_REUSED, + f"pid {row.get('pid')!r} is alive but started at " + f"{live_start_time!r}, not the registered {recorded_start!r}; the " + "number was reused by an unrelated process and this registration's " + "own liveness is therefore unproven", + ) + if pid_alive: + return ( + REASON_PID_ALIVE, + f"recorded pid {row.get('pid')!r} is still running on this host and " + "boot", + ) + return None + + +def _missing_identity_fields(row: Mapping[str, Any]) -> list[str]: + missing: list[str] = [] + for name in REQUIRED_IDENTITY_FIELDS: + value = row.get(name) + if value is None or (isinstance(value, str) and not value.strip()): + missing.append(name) + return missing + + +def plan_stale_worker_retirement( + rows: Iterable[Mapping[str, Any]], + *, + now: datetime | None = None, + pid_alive_probe: Callable[[int | None], bool | None] | None = None, + canonical_repository: str | None = None, + protected_worker_identities: Iterable[str] | None = None, + protected_session_ids: Iterable[str] | None = None, + protected_pids: Iterable[Any] | None = None, + current_host_id: str | None = None, + current_boot_id: str | None = None, + start_time_probe: Callable[[Any], str | None] | None = None, +) -> dict[str, Any]: + """Decide, without mutating anything, which registrations may be retired. + + Every row lands in exactly one of ``candidates`` (eligible) or + ``preserved`` (with the reason code that stopped it), so the output + explains the whole registry rather than only the interesting part. + """ + all_rows = [dict(row) for row in rows] + protected_ids = {str(w) for w in (protected_worker_identities or []) if w} + protected_sessions = {str(s) for s in (protected_session_ids or []) if s} + protected_pid_set = {_canon(p) for p in (protected_pids or []) if p is not None} + + snapshots: dict[int, dict[str, Any]] = {} + pid_alive_by_index: dict[int, bool | None] = {} + start_time_by_index: dict[int, str | None] = {} + for index, row in enumerate(all_rows): + pid_alive = _probe_pid(row.get("pid"), pid_alive_probe) + pid_alive_by_index[index] = pid_alive + start_time_by_index[index] = _probe_start_time( + row.get("pid"), start_time_probe + ) + snapshots[index] = fleet.build_worker_snapshot_row( + row, + now=now, + pid_alive_probe=(lambda _pid, _value=pid_alive: _value), + canonical_repository=canonical_repository, + ) + + # Identity evidence owned by a worker that is live, or whose liveness could + # not be established, is ambiguous: anything sharing it is preserved. + ambiguous_keys: set[tuple[str, str]] = set() + identity_counts: dict[str, int] = {} + for index, row in enumerate(all_rows): + identity = _canon(row.get("worker_identity")) + if identity: + identity_counts[identity] = identity_counts.get(identity, 0) + 1 + snapshot_row = snapshots[index] + liveness = snapshot_row.get("liveness") or {} + unresolved = ( + pid_alive_by_index[index] is None + or liveness.get("heartbeat_fresh") is None + ) + if snapshot_row.get("live") or ( + str(row.get("status") or "") == mwi.STATUS_ACTIVE and unresolved + ): + ambiguous_keys.update(_conflict_keys(row, snapshot_row)) + + # Two active registrations claiming one (client_instance_id, namespace) + # violate the #978 uniqueness invariant — but only when one of them might + # still be running. Several *dead* rows accumulating on one slot across + # restarts is ordinary history and every one of them is safely retirable; + # a slot shared with a live or unprobeable worker is genuinely ambiguous, + # because which row that process belongs to cannot be settled from the + # registry alone. + # + # The key is deliberately the (instance, namespace) pair, not the instance + # alone: one legitimate cohort is exactly one instance spread across + # distinct namespaces, so keying on the instance would make every cohort + # look self-conflicting and preserve the whole fleet forever. + instance_members: dict[tuple[str, str], list[int]] = {} + for index, row in enumerate(all_rows): + if str(row.get("status") or "") != mwi.STATUS_ACTIVE: + continue + key = _instance_key(row) + if key is None: + continue + instance_members.setdefault(key, []).append(index) + instance_conflicts: dict[tuple[str, str], bool] = {} + for key, members in instance_members.items(): + if len(members) < 2: + continue + contested = any( + snapshots[i].get("live") or pid_alive_by_index[i] is None + for i in members + ) + if contested: + instance_conflicts[key] = True + + candidates: list[dict[str, Any]] = [] + candidate_rows: list[Mapping[str, Any]] = [] + preserved: list[dict[str, Any]] = [] + + for index, row in enumerate(all_rows): + snapshot_row = snapshots[index] + pid_alive = pid_alive_by_index[index] + liveness = snapshot_row.get("liveness") or {} + evidence = _evidence(row, snapshot_row, pid_alive) + blocked: tuple[str, str] | None = None + + identity = _canon(row.get("worker_identity")) + missing = _missing_identity_fields(row) + shared = sorted( + f"{name}={value}" + for name, value in _conflict_keys(row, snapshot_row) + if (name, value) in ambiguous_keys + ) + binding = (row.get("repository_binding") or "").strip() + protected_hits: list[str] = [] + if identity and identity in protected_ids: + protected_hits.append(f"worker_identity={identity}") + if _canon(row.get("session_id")) in protected_sessions: + protected_hits.append(f"session_id={_canon(row.get('session_id'))}") + if _canon(row.get("pid")) in protected_pid_set: + protected_hits.append(f"pid={_canon(row.get('pid'))}") + + if identity and identity_counts.get(identity, 0) > 1: + blocked = ( + REASON_CONFLICTING_IDENTITY, + f"worker identity {identity!r} appears on more than one registry row", + ) + elif str(row.get("status") or "") != mwi.STATUS_ACTIVE: + blocked = ( + REASON_ALREADY_TERMINAL, + f"registration status is {row.get('status')!r}; nothing to retire", + ) + elif missing: + blocked = ( + REASON_INCOMPLETE_IDENTITY, + "registry row is missing field(s) the retirement conjunction " + f"reads: {missing}", + ) + elif mwi._parse_ts(row.get("last_heartbeat_at")) is None: + blocked = ( + REASON_UNPARSABLE_HEARTBEAT, + "last_heartbeat_at is not a parsable UTC stamp; liveness is unknown", + ) + elif snapshot_row.get("live"): + blocked = (REASON_WORKER_LIVE, "worker is live and must not be retired") + elif pid_alive is None: + blocked = ( + REASON_PID_UNKNOWN, + f"pid {row.get('pid')!r} could not be probed; liveness is unproven", + ) + elif pid_alive and _reuse_detected( + row, start_time_by_index[index], current_host_id, current_boot_id + ): + # Reported before the generic live-pid branch so the operator sees + # *why* the number is occupied: an unrelated process inherited it, + # which means this registration's own liveness is unproven rather + # than positively established. + blocked = ( + REASON_PID_REUSED, + f"pid {row.get('pid')!r} is alive but started at " + f"{start_time_by_index[index]!r}, not the registered " + f"{row.get('process_start_time')!r}; the number was reused", + ) + elif pid_alive: + blocked = ( + REASON_PID_ALIVE, + f"recorded pid {row.get('pid')!r} is still running", + ) + elif liveness.get("heartbeat_fresh") is not False: + blocked = ( + REASON_HEARTBEAT_FRESH, + "heartbeat has not expired under the canonical TTL policy", + ) + elif snapshot_row.get("ownership_state") != "stale": + blocked = ( + REASON_AMBIGUOUS_OWNERSHIP, + "ownership_state is " + f"{snapshot_row.get('ownership_state')!r}, not 'stale'", + ) + elif not binding or snapshot_row.get("foreign_repository"): + blocked = ( + REASON_FOREIGN_REPOSITORY, + "repository binding is absent or does not match the canonical " + "repository; retirement scope is ambiguous", + ) + elif shared: + blocked = ( + REASON_CONFLICTING_IDENTITY, + "identity evidence is shared with a live or unprobeable worker: " + f"{shared}", + ) + elif protected_hits: + blocked = ( + REASON_PROTECTED_OWNER, + "worker still owns active workflow state requiring separate " + f"reconciliation: {protected_hits}", + ) + elif instance_conflicts.get(_instance_key(row)): + blocked = ( + REASON_INSTANCE_CONFLICT, + "another active registration claims client_instance_id " + f"{row.get('client_instance_id')!r} in namespace " + f"{row.get('namespace')!r}; instance ownership is ambiguous", + ) + else: + # Affirmative identity + fencing proof runs last: everything above + # establishes the row is *inert*, and this establishes it is + # unambiguously *this* worker (#980 review 657 B2). + blocked = assess_retirement_identity_proof( + row, + snapshot_row, + current_host_id=current_host_id, + current_boot_id=current_boot_id, + live_start_time=start_time_by_index[index], + pid_alive=pid_alive, + ) + + if blocked is not None: + preserved.append( + { + "worker_identity": row.get("worker_identity"), + "reason_code": blocked[0], + "detail": blocked[1], + "evidence": evidence, + } + ) + continue + + candidates.append( + { + "worker_identity": row.get("worker_identity"), + "reason_code": REASON_ELIGIBLE, + "detail": ( + "dead pid, expired heartbeat, stale ownership, unambiguous " + "identity, canonical repository binding, no active workflow " + "ownership" + ), + "evidence": evidence, + } + ) + candidate_rows.append(row) + + counts: dict[str, int] = {} + for entry in preserved: + counts[entry["reason_code"]] = counts.get(entry["reason_code"], 0) + 1 + + return { + "success": True, + "read_only": True, + "mutation_performed": False, + "outcome": OUTCOME_PLANNED, + "registry_fingerprint": registry_fingerprint(all_rows), + "candidate_fingerprint": candidate_fingerprint(candidate_rows), + "assessed_count": len(all_rows), + "candidate_count": len(candidates), + "preserved_count": len(preserved), + "candidates": candidates, + "candidate_worker_identities": [c["worker_identity"] for c in candidates], + "preserved": preserved, + "preserved_reason_counts": counts, + "protected_inputs": { + "worker_identities": sorted(protected_ids), + "session_ids": sorted(protected_sessions), + "pids": sorted(protected_pid_set), + }, + "canonical_repository": canonical_repository, + "fencing_context": { + "current_host_id": current_host_id, + "current_boot_id": current_boot_id, + "start_time_probe_available": start_time_probe is not None, + }, + } + + +def summarize_plan(plan: Mapping[str, Any]) -> dict[str, Any]: + """Compact, log-safe view of a plan or apply result.""" + return { + "outcome": plan.get("outcome"), + "registry_fingerprint": plan.get("registry_fingerprint"), + "candidate_fingerprint": plan.get("candidate_fingerprint"), + "assessed_count": plan.get("assessed_count"), + "candidate_count": plan.get("candidate_count"), + "retired_count": plan.get("retired_count"), + "preserved_count": plan.get("preserved_count"), + "mutation_performed": plan.get("mutation_performed"), + } diff --git a/mcp_process_fencing.py b/mcp_process_fencing.py new file mode 100644 index 0000000..3f89881 --- /dev/null +++ b/mcp_process_fencing.py @@ -0,0 +1,205 @@ +"""Host, boot, and process-start fencing evidence for retirement safety (#980). + +Review 657 B2/B3 established that a local ``os.kill(pid, 0)`` probe is not, on +its own, evidence that a *particular registered worker* is gone: + +* **Host ambiguity.** A registry row written on host A records pid 1234. Probing + pid 1234 on host B answers a question nobody asked. "Not running here" is not + "not running". +* **PID reuse.** Pid 1234 may be alive and belong to an unrelated process that + the kernel handed the number to after the original exited. The naive probe + reads that as "the worker is live" (safe, over-preserving) — but the converse + matters too: evidence captured about pid 1234 at time T must not be honoured + at time T+n if the process behind that number changed in between. +* **Boot boundaries.** Every pid from a previous boot is conclusively gone, and + pid numbers restart, so a recorded pid is only comparable to a live pid when + both belong to the same boot. + +This module supplies the three pieces of evidence that turn a bare pid into a +statement about one specific process: + +``host_id`` + Which machine the pid belongs to. +``boot_id`` + Which boot of that machine the pid belongs to. Pids are only comparable + within a single boot. +``process_start_time`` + Which *incarnation* of that pid number. Two processes on the same host and + boot sharing a pid number cannot share a start time, so comparing start + times defeats reuse. + +Every probe here is read-only, never raises, and returns ``None`` when the +evidence cannot be established. ``None`` means *unknown*, and the retirement +conjunction is required to treat unknown as "preserve", never as "safe". +""" + +from __future__ import annotations + +import os +import platform +import subprocess + +# --- Host identity -------------------------------------------------------- + + +def current_host_id() -> str | None: + """Stable identifier for the machine this process runs on. + + Deliberately the kernel node name rather than anything network-derived: it + does not change when an interface goes down or a VPN reassigns an address, + and retirement must not become unsafe because DNS moved. + """ + try: + node = (platform.node() or "").strip() + except Exception: + return None + return node or None + + +# --- Boot identity -------------------------------------------------------- + + +def _linux_boot_id() -> str | None: + try: + with open("/proc/sys/kernel/random/boot_id", encoding="utf-8") as handle: + value = handle.read().strip() + except Exception: + return None + return value or None + + +def _darwin_boot_id() -> str | None: + """macOS boot identity, derived from ``kern.boottime``. + + ``sysctl`` prints e.g. ``{ sec = 1785400000, usec = 123456 } Wed Jul 30 ...``. + Only the integer seconds are kept: the trailing human-readable date is + locale-dependent and would make the identifier unstable across environments + for the very same boot. + """ + try: + completed = subprocess.run( + ["/usr/sbin/sysctl", "-n", "kern.boottime"], + capture_output=True, + text=True, + timeout=5, + check=False, + ) + except Exception: + return None + if completed.returncode != 0: + return None + text = (completed.stdout or "").strip() + marker = "sec = " + start = text.find(marker) + if start < 0: + return None + tail = text[start + len(marker) :] + digits = "" + for char in tail: + if char.isdigit(): + digits += char + else: + break + return f"boot-{digits}" if digits else None + + +def current_boot_id() -> str | None: + """Identifier for the current boot, or ``None`` when it cannot be proven.""" + system = "" + try: + system = (platform.system() or "").strip().lower() + except Exception: + system = "" + if system == "linux": + return _linux_boot_id() + if system == "darwin": + return _darwin_boot_id() + # An unrecognised platform yields no boot evidence rather than a guess. + return None + + +# --- Process start time --------------------------------------------------- + + +def _linux_process_start_time(pid: int) -> str | None: + """Field 22 of ``/proc//stat`` — start time in clock ticks since boot. + + The executable name in field 2 is parenthesised and may itself contain + spaces and parentheses, so the fields are located from the *last* ``)`` + rather than by splitting the whole line. + """ + try: + with open(f"/proc/{pid}/stat", encoding="utf-8") as handle: + raw = handle.read() + except Exception: + return None + close = raw.rfind(")") + if close < 0: + return None + fields = raw[close + 1 :].split() + # After the ')' the next field is state (field 3), so field 22 is index 19. + if len(fields) < 20: + return None + value = fields[19].strip() + return f"ticks-{value}" if value else None + + +def _darwin_process_start_time(pid: int) -> str | None: + """macOS process start, from ``ps -o lstart=``. + + ``lstart`` is the absolute wall-clock start of that pid's current + incarnation. Two processes reusing one pid number report different values, + which is exactly the discrimination reuse detection needs. + """ + try: + completed = subprocess.run( + ["/bin/ps", "-o", "lstart=", "-p", str(int(pid))], + capture_output=True, + text=True, + timeout=5, + check=False, + ) + except Exception: + return None + if completed.returncode != 0: + return None + value = " ".join((completed.stdout or "").split()) + return f"lstart-{value}" if value else None + + +def process_start_time(pid: int | None) -> str | None: + """Start-time token for *pid*'s current incarnation, or ``None``. + + ``None`` is returned both when the pid does not exist and when the platform + cannot answer. Callers must not read either case as evidence of death — a + dead pid is established by the liveness probe, and this value only ever + *withdraws* a retirement that pid-level evidence would otherwise allow. + """ + if pid is None: + return None + try: + numeric = int(pid) + except (TypeError, ValueError): + return None + if numeric <= 0: + return None + system = "" + try: + system = (platform.system() or "").strip().lower() + except Exception: + system = "" + if system == "linux": + return _linux_process_start_time(numeric) + if system == "darwin": + return _darwin_process_start_time(numeric) + return None + + +def current_process_fencing(pid: int | None = None) -> dict[str, str | None]: + """The full fencing triple for *pid* (defaults to this process).""" + target = os.getpid() if pid is None else pid + return { + "host_id": current_host_id(), + "boot_id": current_boot_id(), + "process_start_time": process_start_time(target), + } diff --git a/mcp_worker_identity.py b/mcp_worker_identity.py index f9363e9..60d4193 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,13 +322,53 @@ _SCHEMA_OPTIONAL_COLUMNS: tuple[tuple[str, str], ...] = ( ("parity_revision", "TEXT"), ("live_revision", "TEXT"), ("instance_id_provenance", "TEXT"), + # #980 review 657 B2/B3: fencing evidence that turns a bare pid into a + # statement about one specific process. A registration written before these + # columns existed carries NULL and can never satisfy the retirement identity + # proof, which is the intended fail-closed outcome — absence of evidence is + # not evidence of staleness. + ("host_id", "TEXT"), + ("boot_id", "TEXT"), + ("process_start_time", "TEXT"), + # #980 retirement bookkeeping. Deliberately outside the CAS fingerprint + # field set: they record *that* a retirement happened, and are written only + # by the retirement transaction itself. + ("retired_at", "TEXT"), + ("retired_by", "TEXT"), + ("retirement_reason", "TEXT"), ) +def _default_process_fencing(pid: int | None) -> dict[str, str | None]: + """Fencing triple for *pid*, degrading to unknowns rather than raising. + + Imported lazily so this storage module keeps no import-time dependency on + the probe layer, and so a platform where the probes are unavailable still + registers workers — it simply records no fencing evidence, and those rows + are then permanently ineligible for retirement. + """ + try: + import mcp_process_fencing + + return mcp_process_fencing.current_process_fencing(pid) + except Exception: + return {"host_id": None, "boot_id": None, "process_start_time": None} + + class WorkerRegistryError(RuntimeError): """Raised for registry misuse that is a programming error, not a refusal.""" +class _ExternalStateMoved(RuntimeError): + """External safety evidence changed inside the retirement transaction. + + Raised so the surrounding ``with self._tx()`` rolls back: once lease state + or process liveness has moved, every remaining per-row decision was computed + against a world that no longer exists, so the whole attempt is abandoned + rather than partially applied. + """ + + def _utc_now() -> datetime: return datetime.now(timezone.utc) @@ -791,6 +838,9 @@ class WorkerRegistry: parity_revision: str | None = None, live_revision: str | None = None, instance_id_provenance: str | None = None, + host_id: str | None = None, + boot_id: str | None = None, + process_start_time: str | None = None, ) -> dict[str, Any]: """Atomically register one worker identity. @@ -823,6 +873,20 @@ class WorkerRegistry: proc_id = process_identity or ( f"pid-{int(pid)}" if pid is not None else None ) + # #980 B2/B3: capture the fencing triple for the pid being registered. + # A caller may supply it (tests, or a launcher that already probed); + # otherwise it is probed here, at the only moment the process is known + # to be the one that owns this registration. Any probe that cannot + # answer stores NULL, which permanently withholds retirement eligibility + # from the row rather than granting it on absent evidence. + fencing = _default_process_fencing(pid) + host_id = host_id if host_id is not None else fencing["host_id"] + boot_id = boot_id if boot_id is not None else fencing["boot_id"] + process_start_time = ( + process_start_time + if process_start_time is not None + else fencing["process_start_time"] + ) with self._tx() as conn: existing = conn.execute( "SELECT * FROM worker_registrations WHERE worker_identity = ?", @@ -905,8 +969,9 @@ class WorkerRegistry: fencing_epoch, status, fleet_run_id, authenticated_account, process_identity, startup_revision, loaded_revision, parity_revision, - live_revision, instance_id_provenance - ) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?) + live_revision, instance_id_provenance, + host_id, boot_id, process_start_time + ) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?) """, ( worker_identity, @@ -935,6 +1000,9 @@ class WorkerRegistry: (parity_revision or "").strip() or None, (live_revision or "").strip() or None, (instance_id_provenance or "").strip() or None, + (host_id or "").strip() or None, + (boot_id or "").strip() or None, + (process_start_time or "").strip() or None, ), ) row = conn.execute( @@ -1217,6 +1285,363 @@ 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, + external_fence_fn: Callable[[], str] | None = None, + liveness_fn: Callable[[dict[str, Any]], dict[str, Any]] | None = None, + ) -> dict[str, Any]: + """Compare-and-swap retirement of conclusively stale registrations (#980). + + The whole decision happens inside one ``BEGIN IMMEDIATE`` transaction: + the authoritative rows are re-read, the stable registry fingerprint is + recomputed from *those* rows, the eligibility plan is recomputed from + *those* rows, and only then are the approved targets retired — each + with a per-row guarded ``UPDATE`` that also asserts the row's identity, + liveness, and ownership columns are byte-identical to what the + revalidation just read. There is no window in which a safety check and + its matching write are separated by another statement, so a worker that + comes back to life, changes ownership, or is retired concurrently + cannot be deleted on the strength of a stale observation. + + ``fingerprint_fn`` and ``plan_fn`` are injected rather than imported so + the storage layer never depends on the decision layer; production wires + in :func:`mcp_fleet_retirement.registry_fingerprint` and + :func:`mcp_fleet_retirement.plan_stale_worker_retirement`, which is + exactly what the plan surface used. + + Any exception — including a failure to commit — rolls the transaction + back and is reported as ``transaction_failed`` with zero retirements + and ``mutation_performed`` false, so a partial write can never be + reported as success. + + **External state (#980 review 657 B3).** ``BEGIN IMMEDIATE`` locks this + registry and nothing else, so two critical inputs live outside the + transaction's isolation domain: the workflow-lease table in a separate + control-plane database, and OS process liveness. Re-reading them once + during revalidation is not enough — the per-target loop takes time, so a + lease acquired (or a pid revived) after ``plan_fn`` returned but before + *this* row's ``UPDATE`` would go unnoticed, and the registry-column guard + cannot catch it because no registry column changed. + + Two mechanisms close that gap, both applied per target and immediately + before its own write: + + ``external_fence_fn`` + A version token over all external state the decision consumed. It is + captured inside the transaction before revalidation and re-read + before every guarded ``UPDATE``; any movement aborts the whole + transaction rather than retiring against evidence that has changed. + ``liveness_fn`` + A per-row re-probe of process liveness and fencing identity (host, + boot, start time). It runs immediately before the row's write and + must affirmatively re-establish that this exact process is gone. + + Both default to ``None`` only so the storage layer stays independent of + the decision and control-plane layers; production always supplies them, + and a caller that omits ``liveness_fn`` gets no retirement at all rather + than an unfenced one. + """ + requested = [str(w) for w in (worker_identities or []) if str(w).strip()] + stamp = _ts(now or _utc_now()) + result: dict[str, Any] | None = None + + def _base(outcome: str) -> dict[str, Any]: + return { + "success": True, + "outcome": outcome, + "mutation_performed": False, + "retired": [], + "retired_count": 0, + "preserved": [], + "preserved_count": 0, + "requested_count": len(requested), + "expected_registry_fingerprint": expected_registry_fingerprint, + "expected_candidate_fingerprint": expected_candidate_fingerprint, + "acting_identity": retired_by, + "reasons": [], + } + + if not requested: + outcome = _base("nothing_requested") + outcome["reasons"] = [ + "no worker identities were supplied; nothing to retire" + ] + return outcome + + try: + with self._tx() as conn: + # Captured *after* BEGIN IMMEDIATE and *before* the + # authoritative read, so every later comparison is against the + # external state this decision was actually built on. + fence_at_plan = external_fence_fn() if external_fence_fn else None + rows = [ + self._row_to_record(r) + for r in conn.execute( + "SELECT * FROM worker_registrations" + ).fetchall() + ] + by_identity = { + str(row.get("worker_identity")): row for row in rows + } + current_registry_fingerprint = fingerprint_fn(rows) + + if current_registry_fingerprint != expected_registry_fingerprint: + already = [ + wid + for wid in requested + if str( + (by_identity.get(wid) or {}).get("status") or "" + ) + == STATUS_RETIRED + ] + idempotent = len(already) == len(requested) + result = _base( + "already_retired" if idempotent else "registry_revision_moved" + ) + result["idempotent"] = idempotent + result["current_registry_fingerprint"] = ( + current_registry_fingerprint + ) + result["reasons"] = [ + "the worker registry changed between plan and apply; " + "retiring zero workers" + if not idempotent + else "every requested registration is already retired; " + "safe no-op" + ] + result["preserved"] = [ + { + "worker_identity": wid, + "reason_code": "registry_revision_moved", + "detail": ( + "aborted before any retirement: registry " + "fingerprint moved" + ), + } + for wid in requested + ] + result["preserved_count"] = len(requested) + return result + + fresh_plan = plan_fn(rows) + current_candidate_fingerprint = fresh_plan.get( + "candidate_fingerprint" + ) + if current_candidate_fingerprint != expected_candidate_fingerprint: + result = _base("candidate_set_moved") + result["current_registry_fingerprint"] = ( + current_registry_fingerprint + ) + result["current_candidate_fingerprint"] = ( + current_candidate_fingerprint + ) + result["reasons"] = [ + "the retirement candidate set changed between plan and " + "apply; retiring zero workers" + ] + result["preserved"] = [ + { + "worker_identity": wid, + "reason_code": "candidate_set_moved", + "detail": ( + "aborted before any retirement: candidate " + "fingerprint moved" + ), + } + for wid in requested + ] + result["preserved_count"] = len(requested) + return result + + eligible = { + str(c.get("worker_identity")): c + for c in fresh_plan.get("candidates") or [] + } + preserved_index = { + str(p.get("worker_identity")): p + for p in fresh_plan.get("preserved") or [] + } + + retired: list[dict[str, Any]] = [] + preserved: list[dict[str, Any]] = [] + for wid in requested: + row = by_identity.get(wid) + if row is None: + preserved.append( + { + "worker_identity": wid, + "reason_code": "registration_missing", + "detail": "no registration row with this identity", + } + ) + continue + if wid not in eligible: + blocked = preserved_index.get(wid) or {} + preserved.append( + { + "worker_identity": wid, + "reason_code": blocked.get("reason_code") + or "not_in_current_plan", + "detail": blocked.get("detail") + or ( + "revalidation immediately before retirement " + "no longer finds this worker eligible" + ), + "evidence": blocked.get("evidence"), + } + ) + continue + + # --- External-state fence, immediately before this write --- + # + # Lease state and OS liveness are outside this transaction, + # so they are re-checked here rather than trusted from + # revalidation. Movement aborts the whole transaction: a + # changed world invalidates every remaining decision, not + # only this row's. + if external_fence_fn is not None: + fence_now = external_fence_fn() + if fence_now != fence_at_plan: + raise _ExternalStateMoved( + "external safety state (workflow leases or " + "process liveness) changed inside the retirement " + "transaction; rolling back and retiring nothing" + ) + + if liveness_fn is None: + preserved.append( + { + "worker_identity": wid, + "reason_code": "liveness_reprobe_unavailable", + "detail": ( + "no immediate pre-write liveness re-probe was " + "supplied; refusing to retire on revalidation " + "evidence alone (fail closed)" + ), + } + ) + continue + + verdict = liveness_fn(dict(row)) + if not verdict.get("safe"): + preserved.append( + { + "worker_identity": wid, + "reason_code": verdict.get("reason_code") + or "liveness_reprobe_refused", + "detail": verdict.get("detail") + or ( + "immediate pre-write re-probe could not " + "re-establish that this process is gone" + ), + "evidence": verdict.get("evidence"), + } + ) + continue + + cursor = conn.execute( + "UPDATE worker_registrations " + "SET status = ?, retired_at = ?, retired_by = ?, " + " retirement_reason = ? " + "WHERE worker_identity = ? " + " AND status = ? " + " AND last_heartbeat_at = ? " + " AND generation_id = ? " + " AND session_id = ? " + " AND fencing_epoch = ? " + " AND IFNULL(pid, -1) = IFNULL(?, -1) " + " AND IFNULL(host_id, '') = IFNULL(?, '') " + " AND IFNULL(boot_id, '') = IFNULL(?, '') " + " AND IFNULL(process_start_time, '') = IFNULL(?, '')", + ( + STATUS_RETIRED, + stamp, + retired_by, + retirement_reason + or (eligible[wid].get("reason_code") or ""), + wid, + STATUS_ACTIVE, + row.get("last_heartbeat_at"), + row.get("generation_id"), + row.get("session_id"), + row.get("fencing_epoch"), + row.get("pid"), + row.get("host_id"), + row.get("boot_id"), + row.get("process_start_time"), + ), + ) + if cursor.rowcount == 1: + retired.append( + { + "worker_identity": wid, + "reason_code": eligible[wid].get("reason_code"), + "retired_at": stamp, + "retired_by": retired_by, + "evidence": eligible[wid].get("evidence"), + } + ) + else: + preserved.append( + { + "worker_identity": wid, + "reason_code": "row_changed_since_plan", + "detail": ( + "guarded update matched no row; the " + "registration changed inside the retirement " + "transaction" + ), + } + ) + + result = _base("applied") + result["mutation_performed"] = bool(retired) + result["retired"] = retired + result["retired_count"] = len(retired) + result["preserved"] = preserved + result["preserved_count"] = len(preserved) + result["current_registry_fingerprint"] = ( + current_registry_fingerprint + ) + result["current_candidate_fingerprint"] = ( + current_candidate_fingerprint + ) + result["retired_at"] = stamp if retired else None + result["external_fence"] = fence_at_plan + except Exception as exc: # rolled back by _tx; report, never half-claim + moved = isinstance(exc, _ExternalStateMoved) + failure = _base("external_state_moved" if moved else "transaction_failed") + failure["success"] = False + failure["reasons"] = [ + str(exc) + if moved + else "retirement transaction failed and was rolled back; zero " + f"registrations were retired: {type(exc).__name__}: {exc}" + ] + failure["preserved"] = [ + { + "worker_identity": wid, + "reason_code": ( + "external_state_moved" if moved else "transaction_failed" + ), + "detail": "transaction rolled back before any commit", + } + for wid in requested + ] + failure["preserved_count"] = len(requested) + return failure + + return result + def _public_record(record: dict[str, Any]) -> dict[str, Any]: """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..58ed59d 100644 --- a/task_capability_map.py +++ b/task_capability_map.py @@ -175,6 +175,44 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = { "permission": "gitea.read", "role": "controller", }, + # #980: CAS-protected retirement of conclusively stale worker + # registrations. + # + # Planning is observational and stays on ``gitea.read``: it opens no + # transaction, writes nothing, and returns only what a fleet snapshot + # already exposes to the same roles. + # + # Applying is a mutation and review 657 B1 established that ``gitea.read`` + # cannot authorize it. The mutation landing in the local control-plane + # registry rather than in Gitea makes it *no less* a mutation, and sharing + # an observational permission class with plan meant any profile that could + # look could also destroy. It now requires its own permission, + # ``gitea.worker_registry.retire``, which no profile holds by default — so + # author, reviewer, merger, and ordinary read-only profiles fail closed on + # the permission itself rather than relying on the role check alone. The + # role restriction (controller/reconciler), runtime parity, cohort + # uniqueness, and the exact registry + candidate fingerprints all remain, + # and are now defence in depth behind the capability rather than a + # substitute for it. + # + # Granting the permission is a deliberate operator act in profiles.json; + # removing it from a profile immediately and completely revokes apply. + "plan_stale_worker_retirement": { + "permission": "gitea.read", + "role": "controller", + }, + "gitea_plan_stale_worker_retirement": { + "permission": "gitea.read", + "role": "controller", + }, + "apply_stale_worker_retirement": { + "permission": "gitea.worker_registry.retire", + "role": "controller", + }, + "gitea_apply_stale_worker_retirement": { + "permission": "gitea.worker_registry.retire", + "role": "controller", + }, # #644: Phase 2 Web Console recovery tasks. "clear_stale_binding": { "permission": "gitea.read", 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..a477bde --- /dev/null +++ b/tests/test_issue_980_stale_worker_retirement.py @@ -0,0 +1,1673 @@ +"""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 hashlib +import os +import sqlite3 +import tempfile +import unittest +from datetime import datetime, timedelta, timezone +from typing import Any +from unittest import mock + +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" + +#: Synthetic fencing identity (#980 review 657 B2/B3). Retirement is only ever +#: safe relative to a host and a boot, so the fixtures name both explicitly +#: rather than inheriting whatever machine the suite happens to run on. +HOST = "synthetic-host-a" +OTHER_HOST = "synthetic-host-b" +BOOT = "boot-1000" +OTHER_BOOT = "boot-2000" + + +def _start(pid: int | None) -> str | None: + """Deterministic start-time token for one pid incarnation.""" + return None if pid is None else f"start-{pid}" + +_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", + "host_id", + "boot_id", + "process_start_time", +) + + +def _registry() -> mwi.WorkerRegistry: + handle, path = tempfile.mkstemp(suffix=".sqlite3") + os.close(handle) + os.unlink(path) + return mwi.WorkerRegistry(path) + + +def _instance_for(identity: str) -> str: + """A launcher-minted instance id unique to one synthetic worker. + + Independent workers are independent *application launches*, so they get + distinct instance ids by default. Sharing one id across several rows is the + #978 cohort/uniqueness situation and tests that mean it say so explicitly + by passing ``instance=``. + """ + suffix = hashlib.sha256(identity.encode("utf-8")).hexdigest()[:12] + return f"inst-codex-20260730T070000Z-{suffix}" + + +def _row( + *, + identity: str, + pid: int | None, + instance: str | None = None, + 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 or _instance_for(identity), + "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, + # Complete fencing evidence by default, so a test that wants an + # incomplete row has to say so explicitly rather than getting one by + # omission (#980 review 657 B2). + "host_id": HOST, + "boot_id": BOOT, + "process_start_time": _start(pid), + } + 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): + kwargs.setdefault("current_host_id", HOST) + kwargs.setdefault("current_boot_id", BOOT) + kwargs.setdefault("start_time_probe", _start) + return retire.plan_stale_worker_retirement( + rows, + now=now, + pid_alive_probe=probe, + canonical_repository=REPO, + **kwargs, + ) + + +def _safe_liveness(row: dict[str, Any]) -> dict[str, Any]: + """Pre-write re-probe that agrees the process is gone.""" + return {"safe": True, "evidence": {"pid": row.get("pid"), "pid_alive": False}} + + +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( + # An independent worker, so a distinct application launch and + # therefore a distinct instance: sharing one instance *and* + # namespace with a live worker is the #978 uniqueness + # violation, which is covered separately. + identity="w-live", + pid=103, + session="s3", + generation="g3", + instance="inst-codex-20260730T070000Z-ffffffffffff", + 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, + liveness_fn=_safe_liveness, + external_fence_fn=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, + liveness_fn=liveness_fn, + external_fence_fn=external_fence_fn, + ) + + 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", + instance="inst-codex-20260730T070000Z-ffffffffffff", + 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, + liveness_fn=_safe_liveness, + ) + 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, + liveness_fn=_safe_liveness, + ) + 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_role(task), "controller") + + def test_plan_is_observational_but_apply_needs_its_own_capability(self): + """#980 review 657 B1: apply must not share plan's read permission.""" + import task_capability_map as tcm + + for task in ( + "plan_stale_worker_retirement", + "gitea_plan_stale_worker_retirement", + ): + with self.subTest(task=task): + self.assertEqual(tcm.required_permission(task), "gitea.read") + + for task in ( + "apply_stale_worker_retirement", + "gitea_apply_stale_worker_retirement", + ): + with self.subTest(task=task): + self.assertEqual( + tcm.required_permission(task), + "gitea.worker_registry.retire", + "apply is a mutation and cannot be authorized by gitea.read", + ) + self.assertNotEqual(tcm.required_permission(task), "gitea.read") + + 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) + + +RETIRE_PERMISSION = "gitea.worker_registry.retire" + + +def _profile(name: str, role: str, operations: list[str]) -> dict[str, Any]: + return { + "profile_name": name, + "role": role, + "role_kind": role, + "allowed_operations": list(operations), + "forbidden_operations": [], + } + + +class DedicatedMutationCapabilityTests(unittest.TestCase): + """#980 review 657 B1: apply requires its own mutation capability. + + These exercise the *real* permission gate. Only the two environment gates + that ``_profile_operation_gate`` also consults — master parity and runtime + mode — are neutralised, because they are unrelated to the capability being + proven and would otherwise make the result depend on the machine the suite + runs on. + """ + + def _apply_as(self, profile: dict[str, Any], **patches): + import gitea_mcp_server as server + + with mock.patch.object(server, "_master_parity_block", return_value=[]), \ + mock.patch.object(server, "_runtime_mode_block", return_value=[]), \ + mock.patch.object(server, "get_profile", return_value=profile), \ + mock.patch.object( + server, "_retirement_runtime_block", return_value=[] + ): + with mock.patch.multiple(server, **patches) if patches else _null(): + return server.gitea_apply_stale_worker_retirement( + registry_fingerprint="registryfp-synthetic", + candidate_fingerprint="candidatefp-synthetic", + worker_identities=["w-1"], + remote="prgs", + ) + + def _assert_denied_without_mutation(self, result): + self.assertFalse(result.get("success")) + self.assertEqual(result.get("retired_count"), 0) + self.assertFalse(result.get("mutation_performed", True)) + + def test_read_only_capability_is_insufficient(self): + result = self._apply_as( + _profile("prgs-controller", "controller", ["gitea.read"]) + ) + self._assert_denied_without_mutation(result) + self.assertEqual( + result.get("required_operation_permission"), RETIRE_PERMISSION + ) + + def test_author_reviewer_and_merger_are_denied(self): + cases = { + "prgs-author": ("author", ["gitea.read", "gitea.pr.create"]), + "prgs-reviewer": ("reviewer", ["gitea.read", "gitea.pr.review"]), + "prgs-merger": ("merger", ["gitea.read", "gitea.pr.merge"]), + } + for name, (role, operations) in cases.items(): + with self.subTest(profile=name): + result = self._apply_as(_profile(name, role, operations)) + self._assert_denied_without_mutation(result) + + def test_author_is_denied_even_if_granted_the_capability(self): + """Role remains defence in depth behind the capability.""" + result = self._apply_as( + _profile("prgs-author", "author", ["gitea.read", RETIRE_PERMISSION]) + ) + self._assert_denied_without_mutation(result) + self.assertEqual(result.get("denied_role"), "author") + + def test_removing_the_capability_fails_closed(self): + granted = _profile( + "prgs-controller", "controller", ["gitea.read", RETIRE_PERMISSION] + ) + revoked = _profile("prgs-controller", "controller", ["gitea.read"]) + self._assert_denied_without_mutation(self._apply_as(revoked)) + # The same profile *with* the capability gets a different refusal — + # proving the denial above came from the capability, not from something + # incidental that would deny either way. + with_capability = self._apply_as(granted) + self.assertNotEqual( + with_capability.get("required_operation_permission"), + RETIRE_PERMISSION, + "a granted profile must not be refused on the mutation capability", + ) + + def test_granted_controller_reaches_the_isolated_apply_path(self): + import gitea_mcp_server as server + + registry = _registry() + _register(registry, _row(identity="w-1", pid=101)) + plan = _plan(registry.list_workers(status=None)) + self.assertEqual(plan["candidate_worker_identities"], ["w-1"]) + + with mock.patch.object(server, "_master_parity_block", return_value=[]), \ + mock.patch.object(server, "_runtime_mode_block", return_value=[]), \ + mock.patch.object( + server, + "get_profile", + return_value=_profile( + "prgs-controller", + "controller", + ["gitea.read", RETIRE_PERMISSION], + ), + ), \ + mock.patch.object( + server, "_retirement_runtime_block", return_value=[] + ), \ + mock.patch.object( + server, "_worker_registry", return_value=registry + ), \ + mock.patch.object( + server, + "_retirement_protected_owners", + return_value=(None, {"session_ids": [], "pids": []}), + ), \ + mock.patch.object( + server, "_retirement_external_fence", return_value="fence-1" + ), \ + mock.patch.object( + server, "_retirement_liveness_reprobe", _safe_liveness + ): + result = server.gitea_apply_stale_worker_retirement( + registry_fingerprint=plan["registry_fingerprint"], + candidate_fingerprint=plan["candidate_fingerprint"], + worker_identities=["w-1"], + remote="prgs", + canonical_repository=REPO, + ) + + # It reached the CAS/registry layer rather than a permission refusal. + self.assertNotEqual( + result.get("required_operation_permission"), RETIRE_PERMISSION + ) + self.assertIn("outcome", result) + self.assertEqual( + result.get("permission_scope", {}).get("mutation_capability"), + RETIRE_PERMISSION, + ) + + def test_plan_stays_on_the_observational_capability(self): + import gitea_mcp_server as server + + registry = _registry() + _register(registry, _row(identity="w-1", pid=101)) + with mock.patch.object(server, "_master_parity_block", return_value=[]), \ + mock.patch.object(server, "_runtime_mode_block", return_value=[]), \ + mock.patch.object( + server, + "get_profile", + return_value=_profile( + "prgs-controller", "controller", ["gitea.read"] + ), + ), \ + mock.patch.object( + server, "_worker_registry", return_value=registry + ), \ + mock.patch.object( + server, + "_retirement_protected_owners", + return_value=(None, {"session_ids": [], "pids": []}), + ): + result = server.gitea_plan_stale_worker_retirement( + remote="prgs", canonical_repository=REPO + ) + self.assertTrue(result.get("read_only")) + self.assertFalse(result.get("mutation_performed", True)) + + +class _null: + def __enter__(self): + return None + + def __exit__(self, *exc): + return False + + +class AffirmativeIdentityProofTests(unittest.TestCase): + """#980 review 657 B2: retirement needs proof, not merely absent evidence. + + Every case here is a registration whose *registry columns are all present*. + Before the correction that alone made a row eligible, which is exactly the + defect: non-null columns say nothing about which process a row describes. + """ + + def _only(self, plan, key="preserved"): + self.assertEqual(len(plan[key]), 1, plan[key]) + return plan[key][0] + + def test_complete_trusted_identity_that_is_safely_stale_is_eligible(self): + plan = _plan([_row(identity="w-1", pid=101)]) + self.assertEqual(plan["candidate_worker_identities"], ["w-1"]) + self.assertEqual( + self._only(plan, "candidates")["reason_code"], retire.REASON_ELIGIBLE + ) + + def test_missing_instance_identity_is_preserved(self): + plan = _plan([_row(identity="w-1", pid=101, client_instance_id=None)]) + self.assertEqual(plan["candidate_worker_identities"], []) + self.assertEqual( + self._only(plan)["reason_code"], retire.REASON_INCOMPLETE_IDENTITY + ) + + def test_legacy_row_without_trusted_provenance_is_preserved(self): + """The headline B2 case: complete columns, untrustworthy identity.""" + plan = _plan( + [ + _row( + identity="w-1", + pid=101, + client_instance_id="legacy-pid-80287", + instance_id_provenance=fleet.INSTANCE_ID_PROVENANCE_LEGACY, + ) + ] + ) + self.assertEqual(plan["candidate_worker_identities"], []) + self.assertEqual( + self._only(plan)["reason_code"], retire.REASON_UNTRUSTED_PROVENANCE + ) + + def test_missing_fencing_columns_are_preserved(self): + for field in ("host_id", "boot_id", "process_start_time"): + with self.subTest(missing=field): + plan = _plan([_row(identity="w-1", pid=101, **{field: None})]) + self.assertEqual(plan["candidate_worker_identities"], []) + self.assertEqual( + self._only(plan)["reason_code"], + retire.REASON_INCOMPLETE_IDENTITY, + ) + + def test_pid_reuse_is_detected_and_preserved(self): + """Same pid, same host and boot, but a different process incarnation.""" + plan = _plan( + [_row(identity="w-1", pid=101, process_start_time="start-ORIGINAL")], + probe=_alive, + ) + self.assertEqual(plan["candidate_worker_identities"], []) + self.assertEqual(self._only(plan)["reason_code"], retire.REASON_PID_REUSED) + + def test_same_pid_on_a_different_host_is_preserved(self): + plan = _plan([_row(identity="w-1", pid=101, host_id=OTHER_HOST)]) + self.assertEqual(plan["candidate_worker_identities"], []) + self.assertEqual(self._only(plan)["reason_code"], retire.REASON_HOST_UNPROVEN) + + def test_host_reuse_without_a_provable_local_host_is_preserved(self): + plan = _plan([_row(identity="w-1", pid=101)], current_host_id=None) + self.assertEqual(plan["candidate_worker_identities"], []) + self.assertEqual(self._only(plan)["reason_code"], retire.REASON_HOST_UNPROVEN) + + def test_unknown_current_boot_identity_is_preserved(self): + plan = _plan([_row(identity="w-1", pid=101)], current_boot_id=None) + self.assertEqual(plan["candidate_worker_identities"], []) + self.assertEqual(self._only(plan)["reason_code"], retire.REASON_BOOT_UNKNOWN) + + def test_a_prior_boot_satisfies_the_fencing_proof_when_the_pid_is_dead(self): + """Pids do not survive a reboot, so the recorded process is gone.""" + plan = _plan([_row(identity="w-1", pid=101, boot_id=OTHER_BOOT)]) + self.assertEqual(plan["candidate_worker_identities"], ["w-1"]) + + def test_a_live_pid_always_blocks_even_across_a_boot_boundary(self): + """Deliberately conservative: an occupied pid number is never retired. + + A number occupied after a reboot cannot belong to the registered + process, so retiring would arguably be safe — but "the pid is alive" + stays an absolute block rather than something the fencing proof can + argue away. + """ + plan = _plan( + [_row(identity="w-1", pid=101, boot_id=OTHER_BOOT)], probe=_alive + ) + self.assertEqual(plan["candidate_worker_identities"], []) + self.assertEqual(self._only(plan)["reason_code"], retire.REASON_PID_ALIVE) + + def test_conflicting_session_ownership_is_preserved(self): + rows = [ + _row(identity="w-dead", pid=101, session="shared-session"), + _row( + identity="w-live", + pid=102, + session="shared-session", + heartbeat=NOW - timedelta(seconds=5), + ), + ] + plan = _plan(rows, probe=_selective({102})) + self.assertEqual(plan["candidate_worker_identities"], []) + reasons = {p["worker_identity"]: p["reason_code"] for p in plan["preserved"]} + self.assertEqual(reasons["w-dead"], retire.REASON_CONFLICTING_IDENTITY) + + def test_conflicting_client_instance_identity_is_preserved(self): + shared = _instance_for("shared-launch") + rows = [ + _row( + identity="w-dead", + pid=101, + session="s1", + generation="g1", + instance=shared, + namespace="author", + ), + _row( + identity="w-live", + pid=102, + session="s2", + generation="g2", + instance=shared, + namespace="author", + heartbeat=NOW - timedelta(seconds=5), + ), + ] + plan = _plan(rows, probe=_selective({102})) + self.assertEqual(plan["candidate_worker_identities"], []) + reasons = {p["worker_identity"]: p["reason_code"] for p in plan["preserved"]} + self.assertEqual(reasons["w-dead"], retire.REASON_INSTANCE_CONFLICT) + + def test_conflicting_generation_and_fencing_evidence_is_preserved(self): + rows = [ + _row(identity="w-dead", pid=101, session="s1", generation="shared-gen"), + _row( + identity="w-live", + pid=102, + session="s2", + generation="shared-gen", + heartbeat=NOW - timedelta(seconds=5), + ), + ] + plan = _plan(rows, probe=_selective({102})) + self.assertEqual(plan["candidate_worker_identities"], []) + reasons = {p["worker_identity"]: p["reason_code"] for p in plan["preserved"]} + self.assertEqual(reasons["w-dead"], retire.REASON_CONFLICTING_IDENTITY) + + def test_live_lease_ownership_is_preserved(self): + plan = _plan( + [_row(identity="w-1", pid=101, session="s-leased")], + protected_session_ids=["s-leased"], + ) + self.assertEqual(plan["candidate_worker_identities"], []) + self.assertEqual( + self._only(plan)["reason_code"], retire.REASON_PROTECTED_OWNER + ) + + def test_unknown_lease_state_aborts_rather_than_assuming_none(self): + """An unreadable lease store must not degrade to an empty protection set.""" + import gitea_mcp_server as server + + with mock.patch.object( + server, "_control_plane_db_or_error", return_value=(None, ["db down"]) + ): + with self.assertRaises(Exception): + server._retirement_external_fence() + + def test_active_worker_is_preserved(self): + plan = _plan( + [_row(identity="w-1", pid=101, heartbeat=NOW - timedelta(seconds=5))], + probe=_alive, + ) + self.assertEqual(plan["candidate_worker_identities"], []) + self.assertEqual(self._only(plan)["reason_code"], retire.REASON_WORKER_LIVE) + + def test_multiple_processes_in_one_cohort_are_distinguished_from_conflicts(self): + """One launch, five namespaces: a cohort, not five conflicting claims.""" + shared = _instance_for("one-cohort") + rows = [ + _row( + identity=f"w-{ns}", + pid=200 + index, + session=f"s-{ns}", + generation="gen-cohort", + namespace=ns, + instance=shared, + ) + for index, ns in enumerate( + ("author", "reviewer", "merger", "controller", "reconciler") + ) + ] + plan = _plan(rows) + self.assertEqual(len(plan["candidate_worker_identities"]), 5) + + def test_multiple_independent_workers_are_each_assessed_alone(self): + rows = [ + _row(identity="w-1", pid=101, session="s1", generation="g1"), + _row(identity="w-2", pid=102, session="s2", generation="g2"), + ] + plan = _plan(rows) + self.assertEqual( + sorted(plan["candidate_worker_identities"]), ["w-1", "w-2"] + ) + + def test_mixed_eligible_and_uncertain_candidates_split_correctly(self): + rows = [ + _row(identity="w-ok", pid=101, session="s1", generation="g1"), + _row( + identity="w-legacy", + pid=102, + session="s2", + generation="g2", + client_instance_id="legacy-pid-777", + instance_id_provenance=fleet.INSTANCE_ID_PROVENANCE_LEGACY, + ), + _row( + identity="w-nohost", + pid=103, + session="s3", + generation="g3", + host_id=None, + ), + _row( + identity="w-otherhost", + pid=104, + session="s4", + generation="g4", + host_id=OTHER_HOST, + ), + ] + plan = _plan(rows) + self.assertEqual(plan["candidate_worker_identities"], ["w-ok"]) + reasons = {p["worker_identity"]: p["reason_code"] for p in plan["preserved"]} + self.assertEqual(reasons["w-legacy"], retire.REASON_UNTRUSTED_PROVENANCE) + self.assertEqual(reasons["w-nohost"], retire.REASON_INCOMPLETE_IDENTITY) + self.assertEqual(reasons["w-otherhost"], retire.REASON_HOST_UNPROVEN) + + def test_every_uncertain_case_reports_a_structured_reason(self): + rows = [ + _row(identity="w-legacy", pid=101, client_instance_id="legacy-pid-1", + instance_id_provenance=fleet.INSTANCE_ID_PROVENANCE_LEGACY), + _row(identity="w-nohost", pid=102, session="s2", host_id=None), + _row(identity="w-otherhost", pid=103, session="s3", host_id=OTHER_HOST), + ] + plan = _plan(rows) + for entry in plan["preserved"]: + with self.subTest(worker=entry["worker_identity"]): + self.assertTrue(entry["reason_code"]) + self.assertTrue(entry["detail"]) + self.assertIn("evidence", entry) + + +class ExternalStateFencingTests(unittest.TestCase): + """#980 review 657 B3: evidence outside the registry TX cannot go stale. + + ``BEGIN IMMEDIATE`` locks the worker registry and nothing else. The lease + table lives in a different database and process liveness lives in the + kernel, so both can move while the transaction is open. These prove the + fence catches that movement and that nothing is ever committed against + evidence that changed. + """ + + def setUp(self) -> None: + self.registry = _registry() + + def _rows(self): + return self.registry.list_workers(status=None) + + def _apply(self, plan, **kwargs): + kwargs.setdefault("liveness_fn", _safe_liveness) + kwargs.setdefault("plan_fn", lambda rows: _plan(rows, probe=_dead)) + return 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, + retired_by="synthetic-user/prgs-reconciler", + now=NOW, + **kwargs, + ) + + def _assert_nothing_committed(self, result): + self.assertEqual(result["retired_count"], 0) + self.assertFalse(result["mutation_performed"]) + self.assertTrue( + all(r["status"] == mwi.STATUS_ACTIVE for r in self._rows()), + "no row may be left retired after an aborted attempt", + ) + + # --- lease movement --------------------------------------------------- + + def test_lease_acquired_between_plan_and_apply_preserves_its_worker(self): + _register(self.registry, _row(identity="w-1", pid=101, session="s-1")) + plan = _plan(self._rows()) + self.assertEqual(plan["candidate_worker_identities"], ["w-1"]) + # A lease appears before apply: revalidation re-reads leases and the + # candidate set no longer matches what was approved. + result = self._apply( + plan, + plan_fn=lambda rows: _plan( + rows, probe=_dead, protected_session_ids=["s-1"] + ), + ) + self.assertEqual(result["outcome"], "candidate_set_moved") + self._assert_nothing_committed(result) + + def test_lease_changing_during_apply_validation_aborts(self): + """The fence moves *after* revalidation, mid per-target loop.""" + _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()) + self.assertEqual(len(plan["candidate_worker_identities"]), 2) + + fences = iter(["fence-1", "fence-1", "fence-CHANGED", "fence-CHANGED"]) + + result = self._apply(plan, external_fence_fn=lambda: next(fences)) + self.assertEqual(result["outcome"], "external_state_moved") + self.assertFalse(result["success"]) + self._assert_nothing_committed(result) + + def test_unchanged_external_fence_allows_the_retirement(self): + _register(self.registry, _row(identity="w-1", pid=101)) + plan = _plan(self._rows()) + result = self._apply(plan, external_fence_fn=lambda: "stable-fence") + self.assertEqual(result["outcome"], "applied") + self.assertEqual(result["retired_count"], 1) + self.assertTrue(result["mutation_performed"]) + + def test_unreadable_external_state_aborts_rather_than_comparing_equal(self): + _register(self.registry, _row(identity="w-1", pid=101)) + plan = _plan(self._rows()) + + def exploding_fence(): + raise RuntimeError("lease store unavailable") + + result = self._apply(plan, external_fence_fn=exploding_fence) + self.assertFalse(result["success"]) + self._assert_nothing_committed(result) + + # --- liveness movement ------------------------------------------------ + + def test_worker_becoming_active_before_the_write_is_preserved(self): + _register(self.registry, _row(identity="w-1", pid=101)) + plan = _plan(self._rows()) + + def revived(row): + return { + "safe": False, + "reason_code": retire.REASON_PID_ALIVE, + "detail": "pid came back to life before the write", + } + + result = self._apply(plan, liveness_fn=revived) + self.assertEqual(result["outcome"], "applied") + self._assert_nothing_committed(result) + self.assertEqual( + result["preserved"][0]["reason_code"], retire.REASON_PID_ALIVE + ) + + def test_pid_reuse_before_the_write_is_preserved(self): + _register(self.registry, _row(identity="w-1", pid=101)) + plan = _plan(self._rows()) + + def reused(row): + return { + "safe": False, + "reason_code": retire.REASON_PID_REUSED, + "detail": "pid number reused by an unrelated process", + } + + result = self._apply(plan, liveness_fn=reused) + self._assert_nothing_committed(result) + self.assertEqual( + result["preserved"][0]["reason_code"], retire.REASON_PID_REUSED + ) + + def test_host_or_boot_movement_before_the_write_is_preserved(self): + for reason in (retire.REASON_HOST_UNPROVEN, retire.REASON_BOOT_UNKNOWN): + with self.subTest(reason=reason): + registry = _registry() + _register(registry, _row(identity="w-1", pid=101)) + plan = _plan(registry.list_workers(status=None)) + 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=_dead), + retired_by="synthetic-user/prgs-reconciler", + now=NOW, + liveness_fn=lambda row: { + "safe": False, + "reason_code": reason, + "detail": "fencing identity moved before the write", + }, + ) + self.assertEqual(result["retired_count"], 0) + self.assertFalse(result["mutation_performed"]) + + def test_omitting_the_reprobe_retires_nothing(self): + """A caller that supplies no pre-write probe gets no retirement.""" + _register(self.registry, _row(identity="w-1", pid=101)) + plan = _plan(self._rows()) + result = self._apply(plan, liveness_fn=None) + self._assert_nothing_committed(result) + self.assertEqual( + result["preserved"][0]["reason_code"], "liveness_reprobe_unavailable" + ) + + # --- registry movement ------------------------------------------------ + + def test_heartbeat_movement_retires_zero(self): + _register(self.registry, _row(identity="w-1", pid=101)) + plan = _plan(self._rows()) + with self.registry._tx() as conn: + conn.execute( + "UPDATE worker_registrations SET last_heartbeat_at = ? " + "WHERE worker_identity = ?", + (mwi._ts(NOW), "w-1"), + ) + result = self._apply(plan) + self.assertEqual(result["outcome"], "registry_revision_moved") + self._assert_nothing_committed(result) + + def test_generation_and_fencing_movement_retires_zero(self): + for column, value in ( + ("generation_id", "gen-MOVED"), + ("fencing_epoch", 99), + ("host_id", OTHER_HOST), + ("boot_id", OTHER_BOOT), + ("process_start_time", "start-MOVED"), + ): + with self.subTest(column=column): + registry = _registry() + _register(registry, _row(identity="w-1", pid=101)) + plan = _plan(registry.list_workers(status=None)) + with registry._tx() as conn: + conn.execute( + f"UPDATE worker_registrations SET {column} = ? " + "WHERE worker_identity = ?", + (value, "w-1"), + ) + 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=_dead), + retired_by="synthetic-user/prgs-reconciler", + now=NOW, + liveness_fn=_safe_liveness, + ) + self.assertEqual(result["outcome"], "registry_revision_moved") + self.assertEqual(result["retired_count"], 0) + self.assertFalse(result["mutation_performed"]) + + def test_fencing_columns_are_part_of_the_cas_token(self): + base = [_row(identity="w-1", pid=101)] + for column, value in ( + ("host_id", OTHER_HOST), + ("boot_id", OTHER_BOOT), + ("process_start_time", "start-OTHER"), + ): + with self.subTest(column=column): + moved = [_row(identity="w-1", pid=101, **{column: value})] + self.assertNotEqual( + retire.registry_fingerprint(base), + retire.registry_fingerprint(moved), + ) + + # --- concurrency and partial-failure ---------------------------------- + + def test_concurrent_applies_cannot_both_report_success(self): + _register(self.registry, _row(identity="w-1", pid=101)) + plan = _plan(self._rows()) + first = self._apply(plan) + second = self._apply(plan) + self.assertEqual(first["outcome"], "applied") + self.assertEqual(first["retired_count"], 1) + self.assertTrue(first["mutation_performed"]) + # The second sees a moved registry: idempotent no-op, never a second + # claimed mutation. + self.assertIn(second["outcome"], {"already_retired", "registry_revision_moved"}) + self.assertEqual(second["retired_count"], 0) + self.assertFalse(second["mutation_performed"]) + + def test_later_candidate_failure_rolls_back_the_earlier_one(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()) + self.assertEqual(len(plan["candidate_worker_identities"]), 2) + seen = {"n": 0} + + def fail_on_second(row): + seen["n"] += 1 + if seen["n"] > 1: + raise sqlite3.OperationalError("database is locked") + return {"safe": True} + + result = self._apply(plan, liveness_fn=fail_on_second) + self.assertFalse(result["success"]) + self.assertEqual(result["outcome"], "transaction_failed") + self._assert_nothing_committed(result) + + def test_database_failure_reports_no_mutation(self): + _register(self.registry, _row(identity="w-1", pid=101)) + plan = _plan(self._rows()) + + def exploding_plan(rows): + raise sqlite3.OperationalError("database is locked") + + result = self._apply(plan, plan_fn=exploding_plan) + self.assertFalse(result["success"]) + self.assertEqual(result["outcome"], "transaction_failed") + self._assert_nothing_committed(result) + + def test_reported_counts_match_committed_state(self): + _register(self.registry, _row(identity="w-ok", pid=101, session="s1")) + _register(self.registry, _row(identity="w-hold", pid=102, session="s2")) + plan = _plan(self._rows()) + + def selective(row): + if row.get("worker_identity") == "w-hold": + return { + "safe": False, + "reason_code": retire.REASON_PID_ALIVE, + "detail": "came back before the write", + } + return {"safe": True} + + result = self._apply(plan, liveness_fn=selective) + self.assertEqual(result["retired_count"], 1) + self.assertTrue(result["mutation_performed"]) + committed = {r["worker_identity"]: r["status"] for r in self._rows()} + self.assertEqual(committed["w-ok"], mwi.STATUS_RETIRED) + self.assertEqual(committed["w-hold"], mwi.STATUS_ACTIVE) + self.assertEqual( + len(result["retired"]), result["retired_count"], "counts must agree" + ) + + def test_external_fence_is_captured_before_the_authoritative_read(self): + """Ordering proof: BEGIN IMMEDIATE, then fence, then read.""" + _register(self.registry, _row(identity="w-1", pid=101)) + plan = _plan(self._rows()) + order: list[str] = [] + + def fence(): + order.append("fence") + return "stable" + + def plan_fn(rows): + order.append("revalidate") + return _plan(rows, probe=_dead) + + result = self._apply(plan, external_fence_fn=fence, plan_fn=plan_fn) + self.assertEqual(result["outcome"], "applied") + self.assertEqual(order[0], "fence") + self.assertIn("revalidate", order) + self.assertGreater( + order.count("fence"), 1, "the fence must be re-read before writes" + ) + + +class ExternalStateFingerprintTests(unittest.TestCase): + """The lease/liveness version token itself.""" + + def test_identical_lease_state_produces_one_token(self): + leases = [{"lease_id": "l1", "status": "active", "session_id": "s1"}] + self.assertEqual( + retire.external_state_fingerprint(leases), + retire.external_state_fingerprint(list(leases)), + ) + + def test_acquiring_a_lease_moves_the_token(self): + before = retire.external_state_fingerprint([]) + after = retire.external_state_fingerprint( + [{"lease_id": "l1", "status": "active", "session_id": "s1"}] + ) + self.assertNotEqual(before, after) + + def test_owner_change_moves_the_token(self): + first = retire.external_state_fingerprint( + [{"lease_id": "l1", "status": "active", "session_id": "s1"}] + ) + second = retire.external_state_fingerprint( + [{"lease_id": "l1", "status": "active", "session_id": "s2"}] + ) + self.assertNotEqual(first, second) + + def test_liveness_movement_moves_the_token(self): + first = retire.external_state_fingerprint([], liveness=[(101, False)]) + second = retire.external_state_fingerprint([], liveness=[(101, True)]) + self.assertNotEqual(first, second) + + +if __name__ == "__main__": # pragma: no cover + unittest.main()