feat(fleet): add CAS-protected stale worker retirement capability (Closes #980) #982
@@ -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.
|
||||
|
||||
+142
-9
@@ -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": (
|
||||
|
||||
@@ -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,
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -0,0 +1,205 @@
|
||||
"""Host, boot, and process-start fencing evidence for retirement safety (#980).
|
||||
|
||||
Review 657 B2/B3 established that a local ``os.kill(pid, 0)`` probe is not, on
|
||||
its own, evidence that a *particular registered worker* is gone:
|
||||
|
||||
* **Host ambiguity.** A registry row written on host A records pid 1234. Probing
|
||||
pid 1234 on host B answers a question nobody asked. "Not running here" is not
|
||||
"not running".
|
||||
* **PID reuse.** Pid 1234 may be alive and belong to an unrelated process that
|
||||
the kernel handed the number to after the original exited. The naive probe
|
||||
reads that as "the worker is live" (safe, over-preserving) — but the converse
|
||||
matters too: evidence captured about pid 1234 at time T must not be honoured
|
||||
at time T+n if the process behind that number changed in between.
|
||||
* **Boot boundaries.** Every pid from a previous boot is conclusively gone, and
|
||||
pid numbers restart, so a recorded pid is only comparable to a live pid when
|
||||
both belong to the same boot.
|
||||
|
||||
This module supplies the three pieces of evidence that turn a bare pid into a
|
||||
statement about one specific process:
|
||||
|
||||
``host_id``
|
||||
Which machine the pid belongs to.
|
||||
``boot_id``
|
||||
Which boot of that machine the pid belongs to. Pids are only comparable
|
||||
within a single boot.
|
||||
``process_start_time``
|
||||
Which *incarnation* of that pid number. Two processes on the same host and
|
||||
boot sharing a pid number cannot share a start time, so comparing start
|
||||
times defeats reuse.
|
||||
|
||||
Every probe here is read-only, never raises, and returns ``None`` when the
|
||||
evidence cannot be established. ``None`` means *unknown*, and the retirement
|
||||
conjunction is required to treat unknown as "preserve", never as "safe".
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
import platform
|
||||
import subprocess
|
||||
|
||||
# --- Host identity --------------------------------------------------------
|
||||
|
||||
|
||||
def current_host_id() -> str | None:
|
||||
"""Stable identifier for the machine this process runs on.
|
||||
|
||||
Deliberately the kernel node name rather than anything network-derived: it
|
||||
does not change when an interface goes down or a VPN reassigns an address,
|
||||
and retirement must not become unsafe because DNS moved.
|
||||
"""
|
||||
try:
|
||||
node = (platform.node() or "").strip()
|
||||
except Exception:
|
||||
return None
|
||||
return node or None
|
||||
|
||||
|
||||
# --- Boot identity --------------------------------------------------------
|
||||
|
||||
|
||||
def _linux_boot_id() -> str | None:
|
||||
try:
|
||||
with open("/proc/sys/kernel/random/boot_id", encoding="utf-8") as handle:
|
||||
value = handle.read().strip()
|
||||
except Exception:
|
||||
return None
|
||||
return value or None
|
||||
|
||||
|
||||
def _darwin_boot_id() -> str | None:
|
||||
"""macOS boot identity, derived from ``kern.boottime``.
|
||||
|
||||
``sysctl`` prints e.g. ``{ sec = 1785400000, usec = 123456 } Wed Jul 30 ...``.
|
||||
Only the integer seconds are kept: the trailing human-readable date is
|
||||
locale-dependent and would make the identifier unstable across environments
|
||||
for the very same boot.
|
||||
"""
|
||||
try:
|
||||
completed = subprocess.run(
|
||||
["/usr/sbin/sysctl", "-n", "kern.boottime"],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
timeout=5,
|
||||
check=False,
|
||||
)
|
||||
except Exception:
|
||||
return None
|
||||
if completed.returncode != 0:
|
||||
return None
|
||||
text = (completed.stdout or "").strip()
|
||||
marker = "sec = "
|
||||
start = text.find(marker)
|
||||
if start < 0:
|
||||
return None
|
||||
tail = text[start + len(marker) :]
|
||||
digits = ""
|
||||
for char in tail:
|
||||
if char.isdigit():
|
||||
digits += char
|
||||
else:
|
||||
break
|
||||
return f"boot-{digits}" if digits else None
|
||||
|
||||
|
||||
def current_boot_id() -> str | None:
|
||||
"""Identifier for the current boot, or ``None`` when it cannot be proven."""
|
||||
system = ""
|
||||
try:
|
||||
system = (platform.system() or "").strip().lower()
|
||||
except Exception:
|
||||
system = ""
|
||||
if system == "linux":
|
||||
return _linux_boot_id()
|
||||
if system == "darwin":
|
||||
return _darwin_boot_id()
|
||||
# An unrecognised platform yields no boot evidence rather than a guess.
|
||||
return None
|
||||
|
||||
|
||||
# --- Process start time ---------------------------------------------------
|
||||
|
||||
|
||||
def _linux_process_start_time(pid: int) -> str | None:
|
||||
"""Field 22 of ``/proc/<pid>/stat`` — start time in clock ticks since boot.
|
||||
|
||||
The executable name in field 2 is parenthesised and may itself contain
|
||||
spaces and parentheses, so the fields are located from the *last* ``)``
|
||||
rather than by splitting the whole line.
|
||||
"""
|
||||
try:
|
||||
with open(f"/proc/{pid}/stat", encoding="utf-8") as handle:
|
||||
raw = handle.read()
|
||||
except Exception:
|
||||
return None
|
||||
close = raw.rfind(")")
|
||||
if close < 0:
|
||||
return None
|
||||
fields = raw[close + 1 :].split()
|
||||
# After the ')' the next field is state (field 3), so field 22 is index 19.
|
||||
if len(fields) < 20:
|
||||
return None
|
||||
value = fields[19].strip()
|
||||
return f"ticks-{value}" if value else None
|
||||
|
||||
|
||||
def _darwin_process_start_time(pid: int) -> str | None:
|
||||
"""macOS process start, from ``ps -o lstart=``.
|
||||
|
||||
``lstart`` is the absolute wall-clock start of that pid's current
|
||||
incarnation. Two processes reusing one pid number report different values,
|
||||
which is exactly the discrimination reuse detection needs.
|
||||
"""
|
||||
try:
|
||||
completed = subprocess.run(
|
||||
["/bin/ps", "-o", "lstart=", "-p", str(int(pid))],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
timeout=5,
|
||||
check=False,
|
||||
)
|
||||
except Exception:
|
||||
return None
|
||||
if completed.returncode != 0:
|
||||
return None
|
||||
value = " ".join((completed.stdout or "").split())
|
||||
return f"lstart-{value}" if value else None
|
||||
|
||||
|
||||
def process_start_time(pid: int | None) -> str | None:
|
||||
"""Start-time token for *pid*'s current incarnation, or ``None``.
|
||||
|
||||
``None`` is returned both when the pid does not exist and when the platform
|
||||
cannot answer. Callers must not read either case as evidence of death — a
|
||||
dead pid is established by the liveness probe, and this value only ever
|
||||
*withdraws* a retirement that pid-level evidence would otherwise allow.
|
||||
"""
|
||||
if pid is None:
|
||||
return None
|
||||
try:
|
||||
numeric = int(pid)
|
||||
except (TypeError, ValueError):
|
||||
return None
|
||||
if numeric <= 0:
|
||||
return None
|
||||
system = ""
|
||||
try:
|
||||
system = (platform.system() or "").strip().lower()
|
||||
except Exception:
|
||||
system = ""
|
||||
if system == "linux":
|
||||
return _linux_process_start_time(numeric)
|
||||
if system == "darwin":
|
||||
return _darwin_process_start_time(numeric)
|
||||
return None
|
||||
|
||||
|
||||
def current_process_fencing(pid: int | None = None) -> dict[str, str | None]:
|
||||
"""The full fencing triple for *pid* (defaults to this process)."""
|
||||
target = os.getpid() if pid is None else pid
|
||||
return {
|
||||
"host_id": current_host_id(),
|
||||
"boot_id": current_boot_id(),
|
||||
"process_start_time": process_start_time(target),
|
||||
}
|
||||
+153
-6
@@ -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
|
||||
|
||||
+23
-11
@@ -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.
|
||||
|
||||
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user