diff --git a/docs/stale-worker-retirement.md b/docs/stale-worker-retirement.md index 6c7b90d..043034d 100644 --- a/docs/stale-worker-retirement.md +++ b/docs/stale-worker-retirement.md @@ -59,22 +59,42 @@ or contradictory evidence preserves the row. | `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`. -Two properties are worth stating explicitly: +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. -* **Trusted launcher identity is not required.** #980 places trusted - `client_instance_id` propagation out of scope and lists backfilling trusted - identity for legacy workers as a non-goal. Requiring `inst-…` provenance here - would preserve every legacy row forever and make the feature inert. What is - required is that the registry *fields* the conjunction reads are present. +* **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 @@ -182,16 +202,75 @@ problem. ## Permissions -* **Allowed:** `controller`, `reconciler`. -* **Denied:** author, reviewer, merger — they keep `gitea.read` for diagnosis - elsewhere and are refused this surface by role. -* The Gitea operation gate stays `gitea.read` because the mutation lands in the - **local control-plane worker registry**, not in Gitea — the same model the - #601 lease lifecycle uses. **No new Gitea write permission is introduced for - any profile**, and no author permission is broadened. +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. diff --git a/gitea_mcp_server.py b/gitea_mcp_server.py index 8da17e7..279cc88 100644 --- a/gitea_mcp_server.py +++ b/gitea_mcp_server.py @@ -19626,13 +19626,23 @@ def gitea_snapshot_instance_fleet( # --- #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. - Mirrors ``gitea_snapshot_instance_fleet``: ``gitea.read`` is the operation - gate (the mutation lands in the local worker registry, not in Gitea), and - the role restriction is what actually keeps author, reviewer, and merger - profiles out. No unrelated permission is granted to anyone. + 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) @@ -19755,6 +19765,7 @@ def _retirement_protected_owners() -> tuple[dict | None, dict]: 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, @@ -19762,9 +19773,119 @@ def _retirement_plan(records, canonical_repository: str, protected: dict) -> dic 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. @@ -19937,14 +20058,22 @@ def gitea_apply_stale_worker_retirement( import mcp_fleet_retirement - read_block = _profile_operation_gate("gitea.read") - if read_block: + # #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, - "reasons": read_block, - "permission_report": _permission_block_report("gitea.read"), + "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") @@ -20081,6 +20210,8 @@ def gitea_apply_stale_worker_retirement( 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 @@ -20089,7 +20220,9 @@ def gitea_apply_stale_worker_retirement( result["repository"] = {"org": org, "repo": repo, "canonical_repository": canon} result["protected_active_workflow_owners"] = protected result["permission_scope"] = { - "granted_operations": ["gitea.read"], + "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": ( diff --git a/mcp_fleet_retirement.py b/mcp_fleet_retirement.py index 133f2fc..78b7097 100644 --- a/mcp_fleet_retirement.py +++ b/mcp_fleet_retirement.py @@ -84,8 +84,43 @@ 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", @@ -96,6 +131,9 @@ REQUIRED_IDENTITY_FIELDS: tuple[str, ...] = ( "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 @@ -128,6 +166,13 @@ FINGERPRINT_FIELDS: tuple[str, ...] = ( "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" @@ -209,6 +254,48 @@ def _probe_pid( 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], @@ -229,6 +316,9 @@ def _evidence( "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"), @@ -262,6 +352,138 @@ def _conflict_keys( 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: @@ -280,6 +502,9 @@ def plan_stale_worker_retirement( 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. @@ -294,9 +519,13 @@ def plan_stale_worker_retirement( 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, @@ -323,6 +552,37 @@ def plan_stale_worker_retirement( ): 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]] = [] @@ -378,6 +638,19 @@ def plan_stale_worker_retirement( 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, @@ -412,6 +685,25 @@ def plan_stale_worker_retirement( "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( @@ -462,6 +754,11 @@ def plan_stale_worker_retirement( "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, + }, } 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 7c6f83e..60d4193 100644 --- a/mcp_worker_identity.py +++ b/mcp_worker_identity.py @@ -322,6 +322,14 @@ _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. @@ -331,10 +339,36 @@ _SCHEMA_OPTIONAL_COLUMNS: tuple[tuple[str, str], ...] = ( ) +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) @@ -804,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. @@ -836,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 = ?", @@ -918,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, @@ -948,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( @@ -1241,6 +1296,8 @@ class WorkerRegistry: 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). @@ -1265,6 +1322,33 @@ class WorkerRegistry: 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()) @@ -1295,6 +1379,10 @@ class WorkerRegistry: 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( @@ -1413,6 +1501,53 @@ class WorkerRegistry: ) 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 = ?, " @@ -1423,7 +1558,10 @@ class WorkerRegistry: " AND generation_id = ? " " AND session_id = ? " " AND fencing_epoch = ? " - " AND IFNULL(pid, -1) = IFNULL(?, -1)", + " 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, @@ -1437,6 +1575,9 @@ class WorkerRegistry: 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: @@ -1475,17 +1616,23 @@ class WorkerRegistry: 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 - failure = _base("transaction_failed") + moved = isinstance(exc, _ExternalStateMoved) + failure = _base("external_state_moved" if moved else "transaction_failed") failure["success"] = False failure["reasons"] = [ - "retirement transaction failed and was rolled back; zero " + 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": "transaction_failed", + "reason_code": ( + "external_state_moved" if moved else "transaction_failed" + ), "detail": "transaction rolled back before any commit", } for wid in requested diff --git a/task_capability_map.py b/task_capability_map.py index 10723b5..58ed59d 100644 --- a/task_capability_map.py +++ b/task_capability_map.py @@ -176,15 +176,27 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = { "role": "controller", }, # #980: CAS-protected retirement of conclusively stale worker - # registrations. The mutation lands in the local control-plane worker - # registry, not in Gitea, so — exactly like the #601 lease lifecycle — the - # Gitea operation gate stays ``gitea.read`` and no new Gitea write - # permission is introduced for any profile. The real authority is enforced - # in the tools themselves: role_kind must be controller or reconciler, the - # runtime must be parity-clean and cohort-unique, and apply additionally - # requires the exact stable registry + candidate fingerprints returned by - # the plan. Author, reviewer, and merger profiles keep gitea.read for - # diagnosis elsewhere and are refused this surface. + # 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", @@ -194,11 +206,11 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = { "role": "controller", }, "apply_stale_worker_retirement": { - "permission": "gitea.read", + "permission": "gitea.worker_registry.retire", "role": "controller", }, "gitea_apply_stale_worker_retirement": { - "permission": "gitea.read", + "permission": "gitea.worker_registry.retire", "role": "controller", }, # #644: Phase 2 Web Console recovery tasks. diff --git a/tests/test_issue_980_stale_worker_retirement.py b/tests/test_issue_980_stale_worker_retirement.py index b35be53..a477bde 100644 --- a/tests/test_issue_980_stale_worker_retirement.py +++ b/tests/test_issue_980_stale_worker_retirement.py @@ -16,11 +16,14 @@ 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 @@ -31,6 +34,19 @@ 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", @@ -54,6 +70,9 @@ _ROW_COLUMNS = ( "fleet_run_id", "authenticated_account", "instance_id_provenance", + "host_id", + "boot_id", + "process_start_time", ) @@ -64,11 +83,23 @@ def _registry() -> mwi.WorkerRegistry: 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 = "inst-codex-20260730T070000Z-0123456789ab", + instance: str | None = None, session: str = "sess-a", generation: str = "gen-a", namespace: str = "author", @@ -82,7 +113,7 @@ def _row( record = { "worker_identity": identity, "client_name": "codex", - "client_instance_id": instance, + "client_instance_id": instance or _instance_for(identity), "session_id": session, "generation_id": generation, "role": namespace, @@ -102,6 +133,12 @@ def _row( "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 @@ -123,6 +160,9 @@ def _selective(alive_pids: set[int]): 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, @@ -132,6 +172,11 @@ def _plan(rows, *, probe=_dead, now=NOW, **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] @@ -355,10 +400,15 @@ class PlanEligibilityTests(unittest.TestCase): _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"), @@ -415,7 +465,16 @@ class ApplyCasTests(unittest.TestCase): def _rows(self): return self.registry.list_workers(status=None) - def _apply(self, plan, *, probe=None, identities=None, **kwargs): + 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"], @@ -429,6 +488,8 @@ class ApplyCasTests(unittest.TestCase): 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): @@ -520,6 +581,7 @@ class ApplyCasTests(unittest.TestCase): pid=102, session="s2", generation="g2", + instance="inst-codex-20260730T070000Z-ffffffffffff", heartbeat=NOW - timedelta(seconds=5), ), ) @@ -586,6 +648,7 @@ class ApplyCasTests(unittest.TestCase): 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"]) @@ -652,6 +715,7 @@ class PostRetirementFleetCompatibilityTests(unittest.TestCase): 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) @@ -706,9 +770,31 @@ class CapabilityExposureTests(unittest.TestCase): "gitea_apply_stale_worker_retirement", ): with self.subTest(task=task): - self.assertEqual(tcm.required_permission(task), "gitea.read") self.assertEqual(tcm.required_role(task), "controller") + def test_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 @@ -802,5 +888,786 @@ class SurfaceRegistrationTests(unittest.TestCase): 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()