Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
c0c6d14b73 |
@@ -189,6 +189,11 @@ mutation. Diagnostic reads remain available where `gitea.read` allows.
|
||||
* `mcp_fleet_snapshot` — pure snapshot + classification (#978).
|
||||
* `gitea_snapshot_instance_fleet` — sanctioned MCP tool (#978).
|
||||
* `gitea_get_runtime_context` — single-process view (not fleet-wide).
|
||||
* [stale-worker-retirement.md](stale-worker-retirement.md) — CAS-protected
|
||||
retirement of conclusively stale registrations (#980). Note that the
|
||||
`registry_revision` this snapshot returns is **time-seeded** and unsuitable
|
||||
for compare-and-swap; retirement derives its own stable
|
||||
`registry_fingerprint` from registry content alone.
|
||||
|
||||
## Non-goals
|
||||
|
||||
|
||||
@@ -51,6 +51,7 @@ that gates each call, not which tools exist.
|
||||
- `gitea_adopt_merger_pr_lease`
|
||||
- `gitea_adopt_workflow_lease`
|
||||
- `gitea_allocate_next_work`
|
||||
- `gitea_apply_stale_worker_retirement`
|
||||
- `gitea_assess_already_landed_reconciliation`
|
||||
- `gitea_assess_conflict_fix_classification`
|
||||
- `gitea_assess_conflict_fix_push`
|
||||
@@ -124,6 +125,7 @@ that gates each call, not which tools exist.
|
||||
- `gitea_observability_link_issue`
|
||||
- `gitea_observability_list_projects`
|
||||
- `gitea_observability_reconcile_incident`
|
||||
- `gitea_plan_stale_worker_retirement`
|
||||
- `gitea_post_heartbeat`
|
||||
- `gitea_publish_unpublished_issue_branch`
|
||||
- `gitea_quarantine_contaminated_review`
|
||||
|
||||
@@ -0,0 +1,203 @@
|
||||
# Stale worker retirement (#980)
|
||||
|
||||
The #978 instance-fleet snapshot made registry accuracy observable but
|
||||
deliberately read-only: a registry full of rows whose owning processes are long
|
||||
gone stays full. This document describes the sanctioned way to retire those
|
||||
rows — a dry-run-first, compare-and-swap-protected workflow available only to
|
||||
controller and reconciler namespaces.
|
||||
|
||||
Related: [instance-fleet-identity.md](instance-fleet-identity.md) (#978),
|
||||
[post-restart-reconcile.md](post-restart-reconcile.md) (#662).
|
||||
|
||||
## Why a dedicated registry token
|
||||
|
||||
`mcp_fleet_snapshot._consistency_token` seeds its digest with `snapshot_at`,
|
||||
formatted at second precision. Its output — surfaced as `registry_revision` and
|
||||
`consistency_token` on the snapshot — therefore changes on **every call**, even
|
||||
when no registry row changed. Any compare-and-swap gated on it can never pass:
|
||||
a dry-run/apply cycle spanning more than one second aborts unconditionally.
|
||||
|
||||
That token remains useful as an observation stamp and is unchanged. #980 adds a
|
||||
separate, *stable* token instead:
|
||||
|
||||
| Token | Module | Derived from | Stable across time? |
|
||||
| --- | --- | --- | --- |
|
||||
| `registry_revision` / `consistency_token` | `mcp_fleet_snapshot` | `snapshot_at` + a subset of row fields | **No** — moves every second |
|
||||
| `registry_fingerprint` | `mcp_fleet_retirement` | canonical retirement-relevant row content only | **Yes** |
|
||||
| `candidate_fingerprint` | `mcp_fleet_retirement` | canonical content of the selected candidate rows | **Yes** |
|
||||
|
||||
`registry_fingerprint` guarantees:
|
||||
|
||||
* identical canonical registry contents always produce the same token, whenever
|
||||
they are observed;
|
||||
* row iteration order never affects the token (serialized rows are sorted);
|
||||
* any retirement-relevant change moves it — row creation or deletion, identity
|
||||
change, heartbeat or TTL change, ownership change, registration-state change,
|
||||
PID change, or repository-binding change.
|
||||
|
||||
The exact field set is `mcp_fleet_retirement.FINGERPRINT_FIELDS`. Deliberately
|
||||
excluded: `token_fingerprint` (credential-adjacent, never a retirement input),
|
||||
the four `*_revision` columns (revision drift is an independent restart concern
|
||||
and is not part of the eligibility conjunction), and the `retired_*` bookkeeping
|
||||
columns this feature adds. Numeric values are canonicalised, so a TTL that
|
||||
round-trips through SQLite as `900.0` hashes identically to `900`.
|
||||
|
||||
## Eligibility — the conjunction
|
||||
|
||||
A registration is retired only when **every** one of these holds. Any missing
|
||||
or contradictory evidence preserves the row.
|
||||
|
||||
| Requirement | Preserve reason code when it fails |
|
||||
| --- | --- |
|
||||
| `status` is `active` | `already_terminal_registration` |
|
||||
| Every field the conjunction reads is present (`REQUIRED_IDENTITY_FIELDS`) | `incomplete_registry_identity` |
|
||||
| `last_heartbeat_at` parses as a UTC stamp | `unparsable_heartbeat` |
|
||||
| Worker is not live | `worker_live` |
|
||||
| PID probe returns a definite answer | `pid_liveness_unknown` |
|
||||
| PID probe says the process is gone | `pid_alive` |
|
||||
| Heartbeat has expired under the canonical TTL | `heartbeat_not_expired` |
|
||||
| `ownership_state` is exactly `stale` | `ambiguous_ownership_state` |
|
||||
| Repository binding present and canonical | `repository_binding_ambiguous` |
|
||||
| No identity evidence shared with a live or unprobeable worker | `conflicting_identity_evidence` |
|
||||
| Not an active workflow-lease owner | `protected_active_workflow_owner` |
|
||||
|
||||
Eligible rows carry `eligible_stale_orphan`.
|
||||
|
||||
Two properties are worth stating explicitly:
|
||||
|
||||
* **`pid_alive` can only withdraw liveness, never grant it**
|
||||
(`WorkerRegistry.is_live`, #948 AC7). A heartbeat-lapsed but still-running
|
||||
process therefore classifies as `stale` in the snapshot, yet #980's added
|
||||
`pid_alive is False` requirement preserves it. An unprobeable PID (`None`)
|
||||
also fails closed.
|
||||
* **Trusted launcher identity is not required.** #980 places trusted
|
||||
`client_instance_id` propagation out of scope and lists backfilling trusted
|
||||
identity for legacy workers as a non-goal. Requiring `inst-…` provenance here
|
||||
would preserve every legacy row forever and make the feature inert. What is
|
||||
required is that the registry *fields* the conjunction reads are present.
|
||||
|
||||
Multiple processes belonging to one legitimate worker cohort are not treated as
|
||||
multiple independent workers: the fleet model from #948/#978 is preserved
|
||||
unchanged, and sharing a role or profile is never a duplicate.
|
||||
|
||||
## Tools
|
||||
|
||||
### `gitea_plan_stale_worker_retirement`
|
||||
|
||||
Read-only. Controller and reconciler only.
|
||||
|
||||
| Parameter | Meaning |
|
||||
| --- | --- |
|
||||
| `remote` | `dadeschools` or `prgs` |
|
||||
| `host`, `org`, `repo` | Optional overrides (audit context) |
|
||||
| `canonical_repository` | Expected repository binding; defaults to the process root |
|
||||
|
||||
Returns `registry_fingerprint`, `candidate_fingerprint`,
|
||||
`candidate_worker_identities`, per-worker `candidates` and `preserved` entries
|
||||
(each with `reason_code`, `detail`, and structured `evidence`),
|
||||
`preserved_reason_counts`, `assessed_count`, `candidate_count`,
|
||||
`preserved_count`, and `protected_active_workflow_owners`.
|
||||
`mutation_performed` is always `false` and `read_only` is always `true`.
|
||||
|
||||
Planning is deterministic: the same authoritative registry contents produce the
|
||||
same plan and the same tokens regardless of when they are observed.
|
||||
|
||||
### `gitea_apply_stale_worker_retirement`
|
||||
|
||||
Mutating. Controller and reconciler only.
|
||||
|
||||
| Parameter | Meaning |
|
||||
| --- | --- |
|
||||
| `registry_fingerprint` | The exact stable token the plan returned |
|
||||
| `candidate_fingerprint` | The exact candidate-set token the plan returned |
|
||||
| `worker_identities` | The exact candidate identities (list, or JSON / comma-separated string) |
|
||||
| `remote`, `host`, `org`, `repo` | As above |
|
||||
| `canonical_repository` | Must match the value the plan used |
|
||||
|
||||
Before touching the registry, apply fails closed on: profile permission, role
|
||||
kind, master parity (`mutation_safe`), stable-runtime mode, capability
|
||||
resolution refreshed immediately before mutation, worker-registry availability,
|
||||
workflow-lease enumeration failure, and daemon-cohort uniqueness
|
||||
(`classify_cohort`).
|
||||
|
||||
A matching token is necessary but never sufficient. Inside one
|
||||
`BEGIN IMMEDIATE` transaction (`WorkerRegistry.retire_stale_workers`) the
|
||||
server:
|
||||
|
||||
1. re-reads the authoritative rows;
|
||||
2. recomputes `registry_fingerprint` from *those* rows and compares — a mismatch
|
||||
returns `registry_revision_moved` with `retired_count: 0` and no write;
|
||||
3. recomputes the eligibility plan from *those* rows — re-reading the active
|
||||
workflow leases rather than reusing the set captured before the transaction
|
||||
opened, so a lease acquired after planning still preserves its worker — and
|
||||
compares `candidate_fingerprint`. A mismatch returns `candidate_set_moved`
|
||||
with `retired_count: 0` and no write; a lease-enumeration failure raises and
|
||||
rolls the transaction back;
|
||||
4. revalidates every requested identity against that fresh plan;
|
||||
5. retires each survivor with a guarded `UPDATE` that additionally asserts
|
||||
`status`, `last_heartbeat_at`, `generation_id`, `session_id`,
|
||||
`fencing_epoch`, and `pid` are unchanged. A guard that matches no row
|
||||
preserves the worker with `row_changed_since_plan`.
|
||||
|
||||
There is no window between a safety check and its matching write, so a worker
|
||||
that comes back to life, changes ownership, or is retired concurrently cannot be
|
||||
removed on the strength of a stale observation. Any exception — including a
|
||||
commit failure — rolls the whole transaction back and returns
|
||||
`transaction_failed` with `success: false`, `retired_count: 0`, and
|
||||
`mutation_performed: false`; a partial write can never be reported as success.
|
||||
|
||||
### Outcomes
|
||||
|
||||
| `outcome` | Meaning | `mutation_performed` |
|
||||
| --- | --- | --- |
|
||||
| `planned` | Dry-run result | `false` |
|
||||
| `applied` | Transaction ran; see `retired` / `preserved` | `true` only if something was retired |
|
||||
| `registry_revision_moved` | Registry changed between plan and apply | `false` |
|
||||
| `candidate_set_moved` | Eligibility verdict changed between plan and apply | `false` |
|
||||
| `already_retired` | Every requested row is already retired (idempotent replay) | `false` |
|
||||
| `nothing_requested` | Empty target list | `false` |
|
||||
| `transaction_failed` | Rolled back; nothing retired | `false` |
|
||||
|
||||
## What retirement does to the fleet snapshot
|
||||
|
||||
A retired row keeps its history: `status` moves to `retired` and `retired_at`,
|
||||
`retired_by`, `retirement_reason` are recorded. Nothing is deleted. Because
|
||||
`retired` is not `active`, the #978 snapshot counts the row as **historical**,
|
||||
not stale, so `stale_worker_count` falls and historical rows never make the live
|
||||
fleet unsafe by themselves.
|
||||
|
||||
**Retirement does not repair untrusted live identity.** Live workers registered
|
||||
under legacy `pid-`/`proc-` instance identities are preserved untouched and
|
||||
keep their `legacy_incomplete_identity` blockers. Retiring every stale row can
|
||||
therefore legitimately produce:
|
||||
|
||||
* `stale_worker_count: 0`
|
||||
* residual live `legacy_incomplete_identity` blockers
|
||||
* `live_fleet_safe: false`
|
||||
|
||||
That is a truthful result, and the `post_apply` block reports the remaining
|
||||
blockers rather than claiming the fleet became safe. Trusted
|
||||
`client_instance_id` propagation through launchers is a separate enrollment
|
||||
problem.
|
||||
|
||||
## Permissions
|
||||
|
||||
* **Allowed:** `controller`, `reconciler`.
|
||||
* **Denied:** author, reviewer, merger — they keep `gitea.read` for diagnosis
|
||||
elsewhere and are refused this surface by role.
|
||||
* The Gitea operation gate stays `gitea.read` because the mutation lands in the
|
||||
**local control-plane worker registry**, not in Gitea — the same model the
|
||||
#601 lease lifecycle uses. **No new Gitea write permission is introduced for
|
||||
any profile**, and no author permission is broadened.
|
||||
* The fleet snapshot remains observational: nothing here turns it into a gate on
|
||||
ordinary author work.
|
||||
|
||||
## Non-goals
|
||||
|
||||
* Killing or restarting processes.
|
||||
* Editing session files or configuration.
|
||||
* Direct database cleanup outside the sanctioned transaction.
|
||||
* Retiring live workers.
|
||||
* Backfilling trusted identity for legacy workers.
|
||||
* Rewriting worker ownership.
|
||||
* Cleaning unrelated workflow-lease or issue-claim registries.
|
||||
@@ -19623,6 +19623,568 @@ def gitea_snapshot_instance_fleet(
|
||||
return snapshot
|
||||
|
||||
|
||||
# --- #980 CAS-protected stale worker retirement ---------------------------
|
||||
|
||||
|
||||
def _retirement_role_block(task: str) -> dict | None:
|
||||
"""Refuse the retirement surface to any non-controller/reconciler role.
|
||||
|
||||
Mirrors ``gitea_snapshot_instance_fleet``: ``gitea.read`` is the operation
|
||||
gate (the mutation lands in the local worker registry, not in Gitea), and
|
||||
the role restriction is what actually keeps author, reviewer, and merger
|
||||
profiles out. No unrelated permission is granted to anyone.
|
||||
"""
|
||||
profile = get_profile()
|
||||
role = _profile_role_kind(profile)
|
||||
if role in {"controller", "reconciler"}:
|
||||
return None
|
||||
return {
|
||||
"success": False,
|
||||
"allowed": False,
|
||||
"mutation_performed": False,
|
||||
"retired_count": 0,
|
||||
"denied_role": role,
|
||||
"required_roles": ["controller", "reconciler"],
|
||||
"requested_task": task,
|
||||
"reasons": [
|
||||
f"{task} is restricted to controller and reconciler roles; active "
|
||||
f"role_kind is {role!r}. Author, reviewer, and merger profiles keep "
|
||||
"gitea.read for diagnosis elsewhere but never receive worker-"
|
||||
"registry retirement, and no unrelated mutation permission is "
|
||||
"granted."
|
||||
],
|
||||
"exact_next_action": (
|
||||
"Re-run from a prgs-controller or prgs-reconciler namespace."
|
||||
),
|
||||
}
|
||||
|
||||
|
||||
def _retirement_runtime_block() -> list[str]:
|
||||
"""Fail-closed runtime reasons that must stop a retirement apply (#980).
|
||||
|
||||
``gitea.read`` deliberately bypasses the #420 parity gate and the #615
|
||||
stable-runtime gate, because a stale server may still be *inspected*. An
|
||||
apply is a mutation, so both gates are re-asserted explicitly here rather
|
||||
than inherited.
|
||||
"""
|
||||
reasons: list[str] = []
|
||||
try:
|
||||
parity = _current_master_parity()
|
||||
except Exception as exc:
|
||||
return [
|
||||
f"master parity could not be assessed (fail closed): {_redact(str(exc))}"
|
||||
]
|
||||
if not parity.get("mutation_safe"):
|
||||
reasons.append(
|
||||
"runtime parity is not mutation-safe: "
|
||||
f"{parity.get('summary') or 'stale runtime'}"
|
||||
)
|
||||
try:
|
||||
reasons.extend(
|
||||
stable_control_runtime.runtime_block_reasons(
|
||||
_current_runtime_mode_report()
|
||||
)
|
||||
)
|
||||
except Exception as exc:
|
||||
reasons.append(
|
||||
f"runtime mode could not be assessed (fail closed): {_redact(str(exc))}"
|
||||
)
|
||||
return reasons
|
||||
|
||||
|
||||
def _retirement_protected_owners() -> tuple[dict | None, dict]:
|
||||
"""Workflow owners that must never be retired as stale workers.
|
||||
|
||||
A registration whose process still owns an active control-plane lease is a
|
||||
live workflow participant needing its own reconciliation, not a stale
|
||||
orphan. Failure to enumerate leases is ambiguity, so it aborts rather than
|
||||
proceeding with an empty protection set.
|
||||
"""
|
||||
db, errs = _control_plane_db_or_error()
|
||||
if db is None:
|
||||
return (
|
||||
{
|
||||
"success": False,
|
||||
"mutation_performed": False,
|
||||
"retired_count": 0,
|
||||
"reasons": [
|
||||
"active workflow leases could not be enumerated, so "
|
||||
"protected owners are unknown (fail closed)",
|
||||
*errs,
|
||||
],
|
||||
},
|
||||
{},
|
||||
)
|
||||
try:
|
||||
leases = db.list_leases(statuses=["active"], limit=1000)
|
||||
except Exception as exc:
|
||||
return (
|
||||
{
|
||||
"success": False,
|
||||
"mutation_performed": False,
|
||||
"retired_count": 0,
|
||||
"reasons": [
|
||||
"active workflow leases could not be enumerated, so "
|
||||
"protected owners are unknown (fail closed): "
|
||||
f"{_redact(str(exc))}"
|
||||
],
|
||||
},
|
||||
{},
|
||||
)
|
||||
session_ids: set[str] = set()
|
||||
pids: set[int] = set()
|
||||
for lease in leases:
|
||||
owner = lease.get("session_id") or lease.get("owner_session_id")
|
||||
if owner:
|
||||
session_ids.add(str(owner))
|
||||
for key in ("owner_pid", "session_pid"):
|
||||
value = lease.get(key)
|
||||
if value is None:
|
||||
continue
|
||||
try:
|
||||
pids.add(int(value))
|
||||
except (TypeError, ValueError):
|
||||
continue
|
||||
return None, {
|
||||
"session_ids": sorted(session_ids),
|
||||
"pids": sorted(pids),
|
||||
"active_lease_count": len(leases),
|
||||
}
|
||||
|
||||
|
||||
def _retirement_plan(records, canonical_repository: str, protected: dict) -> dict:
|
||||
"""The single planning path shared by dry run and in-transaction revalidation."""
|
||||
import mcp_fleet_retirement
|
||||
|
||||
return mcp_fleet_retirement.plan_stale_worker_retirement(
|
||||
records,
|
||||
pid_alive_probe=issue_lock_store.is_process_alive,
|
||||
canonical_repository=canonical_repository,
|
||||
protected_session_ids=(protected or {}).get("session_ids"),
|
||||
protected_pids=(protected or {}).get("pids"),
|
||||
)
|
||||
|
||||
|
||||
def _retirement_revalidation_plan(records, canonical_repository: str) -> dict:
|
||||
"""Revalidation planner used *inside* the retirement transaction.
|
||||
|
||||
Deliberately re-reads the active workflow leases rather than reusing the
|
||||
set captured before the transaction opened: a lease acquired after planning
|
||||
must still preserve its worker. An enumeration failure raises, which rolls
|
||||
the transaction back and retires nothing.
|
||||
"""
|
||||
block, protected = _retirement_protected_owners()
|
||||
if block:
|
||||
raise RuntimeError(
|
||||
"active workflow leases could not be re-read inside the retirement "
|
||||
"transaction; refusing to retire anything"
|
||||
)
|
||||
return _retirement_plan(records, canonical_repository, protected)
|
||||
|
||||
|
||||
@mcp.tool()
|
||||
def gitea_plan_stale_worker_retirement(
|
||||
remote: str = "dadeschools",
|
||||
host: str | None = None,
|
||||
org: str | None = None,
|
||||
repo: str | None = None,
|
||||
canonical_repository: str | None = None,
|
||||
) -> dict:
|
||||
"""Read-only: authoritative retirement plan for stale worker rows (#980).
|
||||
|
||||
Controller and reconciler only. Reads the worker registry, classifies every
|
||||
registration with the same assessor the #978 fleet snapshot uses, and
|
||||
returns the exact set of registrations that are conclusively stale orphans
|
||||
together with a **stable** ``registry_fingerprint`` and an exact
|
||||
``candidate_fingerprint``.
|
||||
|
||||
The fingerprint is derived only from canonical retirement-relevant registry
|
||||
content — never from ``snapshot_at``, wall-clock, request, or report time —
|
||||
so two plans over an unchanged registry agree and the apply compare-and-swap
|
||||
can actually pass. Row order never affects it.
|
||||
|
||||
Retires nothing. Live workers, workers whose PID cannot be probed, workers
|
||||
with unparsable heartbeats, workers sharing identity evidence with a live or
|
||||
unprobeable worker, foreign or unbound repositories, and workers that still
|
||||
own an active workflow lease are all preserved with a reason code.
|
||||
|
||||
Args:
|
||||
remote: Known instance — 'dadeschools' or 'prgs'.
|
||||
host: Optional host override.
|
||||
org: Optional org override (audit context only).
|
||||
repo: Optional repo override (audit context only).
|
||||
canonical_repository: Expected repository binding for
|
||||
foreign-repository classification (defaults to the process root).
|
||||
"""
|
||||
read_block = _profile_operation_gate("gitea.read")
|
||||
if read_block:
|
||||
return {
|
||||
"success": False,
|
||||
"read_only": True,
|
||||
"mutation_performed": False,
|
||||
"reasons": read_block,
|
||||
"permission_report": _permission_block_report("gitea.read"),
|
||||
}
|
||||
|
||||
role_block = _retirement_role_block("gitea_plan_stale_worker_retirement")
|
||||
if role_block:
|
||||
role_block["read_only"] = True
|
||||
return role_block
|
||||
|
||||
registry = _worker_registry()
|
||||
if registry is None:
|
||||
return {
|
||||
"success": False,
|
||||
"read_only": True,
|
||||
"mutation_performed": False,
|
||||
"reasons": [
|
||||
"worker registry is unavailable; cannot produce an authoritative "
|
||||
"retirement plan (fail closed)"
|
||||
],
|
||||
"exact_next_action": (
|
||||
"Ensure GITEA_WORKER_REGISTRY_DB is writable and re-run after "
|
||||
"workers have registered."
|
||||
),
|
||||
}
|
||||
|
||||
protected_block, protected = _retirement_protected_owners()
|
||||
if protected_block:
|
||||
protected_block["read_only"] = True
|
||||
return protected_block
|
||||
|
||||
try:
|
||||
records = registry.list_workers(status=None)
|
||||
except Exception as exc:
|
||||
return {
|
||||
"success": False,
|
||||
"read_only": True,
|
||||
"mutation_performed": False,
|
||||
"reasons": [
|
||||
f"failed to read worker registry: {type(exc).__name__}: "
|
||||
f"{_redact(str(exc))}"
|
||||
],
|
||||
}
|
||||
|
||||
canon = canonical_repository or PROJECT_ROOT
|
||||
plan = _retirement_plan(records, canon, protected)
|
||||
profile = get_profile()
|
||||
plan["role_kind"] = _profile_role_kind(profile)
|
||||
plan["profile"] = profile.get("profile_name")
|
||||
plan["remote"] = _effective_remote(remote)
|
||||
plan["repository"] = {"org": org, "repo": repo, "canonical_repository": canon}
|
||||
plan["protected_active_workflow_owners"] = protected
|
||||
plan["apply_tool"] = "gitea_apply_stale_worker_retirement"
|
||||
plan["permission_scope"] = {
|
||||
"read_only": True,
|
||||
"granted_operations": ["gitea.read"],
|
||||
"denied_unrelated_mutations": True,
|
||||
"note": (
|
||||
"Planning is strictly observational. It does not authorize branch, "
|
||||
"issue, PR, review, merge, or restart mutations, and it retires "
|
||||
"nothing."
|
||||
),
|
||||
}
|
||||
plan["exact_next_action"] = (
|
||||
"Pass registry_fingerprint, candidate_fingerprint, and the exact "
|
||||
"candidate_worker_identities to gitea_apply_stale_worker_retirement."
|
||||
if plan.get("candidate_count")
|
||||
else "No registration is conclusively stale; nothing to apply."
|
||||
)
|
||||
return plan
|
||||
|
||||
|
||||
@mcp.tool()
|
||||
def gitea_apply_stale_worker_retirement(
|
||||
registry_fingerprint: str,
|
||||
candidate_fingerprint: str,
|
||||
worker_identities: list | str,
|
||||
remote: str = "dadeschools",
|
||||
host: str | None = None,
|
||||
org: str | None = None,
|
||||
repo: str | None = None,
|
||||
canonical_repository: str | None = None,
|
||||
) -> dict:
|
||||
"""Retire conclusively stale worker registrations under CAS (#980).
|
||||
|
||||
Controller and reconciler only. Requires the exact ``registry_fingerprint``
|
||||
and ``candidate_fingerprint`` returned by
|
||||
``gitea_plan_stale_worker_retirement`` plus the exact candidate identity
|
||||
list. A matching token is necessary but never sufficient: inside a single
|
||||
``BEGIN IMMEDIATE`` transaction the registry is re-read, the fingerprint is
|
||||
recomputed from those rows, the eligibility plan is recomputed from those
|
||||
rows, and every target is independently revalidated immediately before its
|
||||
own guarded ``UPDATE``. Any drift retires zero workers and reports
|
||||
``registry_revision_moved`` or ``candidate_set_moved``.
|
||||
|
||||
Retiring stale rows does not repair untrusted live identity. Live workers
|
||||
registered under legacy ``pid-``/``proc-`` instance identities remain
|
||||
untouched and their ``legacy_incomplete_identity`` blockers remain
|
||||
outstanding, so the result never claims the fleet became safe.
|
||||
|
||||
Args:
|
||||
registry_fingerprint: Stable token from the plan (CAS expectation).
|
||||
candidate_fingerprint: Exact candidate-set token from the plan.
|
||||
worker_identities: The exact candidate identities the plan returned
|
||||
(list, or a JSON / comma-separated string).
|
||||
remote: Known instance — 'dadeschools' or 'prgs'.
|
||||
host: Optional host override.
|
||||
org: Optional org override (audit context only).
|
||||
repo: Optional repo override (audit context only).
|
||||
canonical_repository: Expected repository binding (defaults to the
|
||||
process root); must match the value the plan used.
|
||||
"""
|
||||
import json as _json
|
||||
|
||||
import mcp_fleet_retirement
|
||||
|
||||
read_block = _profile_operation_gate("gitea.read")
|
||||
if read_block:
|
||||
return {
|
||||
"success": False,
|
||||
"mutation_performed": False,
|
||||
"retired_count": 0,
|
||||
"reasons": read_block,
|
||||
"permission_report": _permission_block_report("gitea.read"),
|
||||
}
|
||||
|
||||
role_block = _retirement_role_block("gitea_apply_stale_worker_retirement")
|
||||
if role_block:
|
||||
return role_block
|
||||
|
||||
runtime_reasons = _retirement_runtime_block()
|
||||
if runtime_reasons:
|
||||
return {
|
||||
"success": False,
|
||||
"mutation_performed": False,
|
||||
"retired_count": 0,
|
||||
"blocker_kind": "runtime_not_mutation_safe",
|
||||
"reasons": runtime_reasons,
|
||||
"exact_next_action": (
|
||||
"Restore runtime parity on the stable control checkout, then "
|
||||
"re-plan and re-apply."
|
||||
),
|
||||
}
|
||||
|
||||
# Fresh identity + capability resolution immediately before mutation.
|
||||
profile = get_profile()
|
||||
role_kind = _profile_role_kind(profile)
|
||||
required_permission = task_capability_map.required_permission(
|
||||
"apply_stale_worker_retirement"
|
||||
)
|
||||
required_role = task_capability_map.required_role(
|
||||
"apply_stale_worker_retirement"
|
||||
)
|
||||
permission_ok, permission_reason = gitea_config.check_operation(
|
||||
required_permission,
|
||||
profile.get("allowed_operations") or [],
|
||||
profile.get("forbidden_operations") or [],
|
||||
)
|
||||
if not permission_ok:
|
||||
return {
|
||||
"success": False,
|
||||
"mutation_performed": False,
|
||||
"retired_count": 0,
|
||||
"requested_task": "apply_stale_worker_retirement",
|
||||
"required_operation_permission": required_permission,
|
||||
"required_role_kind": required_role,
|
||||
"reasons": [
|
||||
"capability resolution immediately before apply refused this "
|
||||
f"session: {permission_reason}"
|
||||
],
|
||||
}
|
||||
try:
|
||||
resolved_host = host or REMOTES[_effective_remote(remote)]["host"]
|
||||
authenticated_username = _authenticated_username(resolved_host)
|
||||
except Exception:
|
||||
authenticated_username = None
|
||||
|
||||
registry = _worker_registry()
|
||||
if registry is None:
|
||||
return {
|
||||
"success": False,
|
||||
"mutation_performed": False,
|
||||
"retired_count": 0,
|
||||
"reasons": [
|
||||
"worker registry is unavailable; refusing to retire anything "
|
||||
"(fail closed)"
|
||||
],
|
||||
}
|
||||
|
||||
raw = worker_identities
|
||||
if isinstance(raw, str):
|
||||
text = raw.strip()
|
||||
try:
|
||||
raw = _json.loads(text)
|
||||
except Exception:
|
||||
raw = [part.strip() for part in text.split(",") if part.strip()]
|
||||
if isinstance(raw, str):
|
||||
raw = [raw]
|
||||
if not isinstance(raw, list):
|
||||
return {
|
||||
"success": False,
|
||||
"mutation_performed": False,
|
||||
"retired_count": 0,
|
||||
"reasons": ["worker_identities must be a list of worker identities"],
|
||||
}
|
||||
targets = [str(item).strip() for item in raw if str(item).strip()]
|
||||
|
||||
protected_block, protected = _retirement_protected_owners()
|
||||
if protected_block:
|
||||
return protected_block
|
||||
|
||||
# #948 daemon-cohort uniqueness: a contested generation or a reused worker
|
||||
# identity means ownership is ambiguous fleet-wide, so retire nothing.
|
||||
try:
|
||||
cohort = mcp_worker_identity.classify_cohort(
|
||||
registry.list_workers(status=mcp_worker_identity.STATUS_ACTIVE),
|
||||
pid_alive_probe=issue_lock_store.is_process_alive,
|
||||
)
|
||||
except Exception as exc:
|
||||
return {
|
||||
"success": False,
|
||||
"mutation_performed": False,
|
||||
"retired_count": 0,
|
||||
"reasons": [
|
||||
"daemon-cohort uniqueness could not be assessed (fail closed): "
|
||||
f"{_redact(str(exc))}"
|
||||
],
|
||||
}
|
||||
if cohort.get("blocked"):
|
||||
return {
|
||||
"success": False,
|
||||
"mutation_performed": False,
|
||||
"retired_count": 0,
|
||||
"blocker_kind": cohort.get("blocker_kind"),
|
||||
"reasons": [
|
||||
"daemon-cohort uniqueness failed; worker ownership is contested",
|
||||
*(cohort.get("reasons") or []),
|
||||
],
|
||||
"cohort": {
|
||||
"blocked_worker_identities": cohort.get("blocked_worker_identities"),
|
||||
"duplicate_identities": cohort.get("duplicate_identities"),
|
||||
"contested_generations": cohort.get("contested_generations"),
|
||||
},
|
||||
}
|
||||
|
||||
canon = canonical_repository or PROJECT_ROOT
|
||||
acting = "/".join(
|
||||
part
|
||||
for part in (authenticated_username, profile.get("profile_name"))
|
||||
if part
|
||||
) or "unknown"
|
||||
|
||||
result = registry.retire_stale_workers(
|
||||
expected_registry_fingerprint=registry_fingerprint,
|
||||
expected_candidate_fingerprint=candidate_fingerprint,
|
||||
worker_identities=targets,
|
||||
fingerprint_fn=mcp_fleet_retirement.registry_fingerprint,
|
||||
plan_fn=lambda rows: _retirement_revalidation_plan(rows, canon),
|
||||
retired_by=acting,
|
||||
retirement_reason=mcp_fleet_retirement.REASON_ELIGIBLE,
|
||||
)
|
||||
|
||||
result["role_kind"] = role_kind
|
||||
result["profile"] = profile.get("profile_name")
|
||||
result["remote"] = _effective_remote(remote)
|
||||
result["repository"] = {"org": org, "repo": repo, "canonical_repository": canon}
|
||||
result["protected_active_workflow_owners"] = protected
|
||||
result["permission_scope"] = {
|
||||
"granted_operations": ["gitea.read"],
|
||||
"control_plane_mutation": "worker_registrations.status -> retired",
|
||||
"denied_unrelated_mutations": True,
|
||||
"note": (
|
||||
"This capability retires local worker-registry rows only. It grants "
|
||||
"no branch, issue, PR, review, merge, or restart authority, and it "
|
||||
"never kills or restarts a process."
|
||||
),
|
||||
}
|
||||
result["post_apply"] = _retirement_post_apply(
|
||||
registry, result, canonical_repository=canon
|
||||
)
|
||||
|
||||
try:
|
||||
gitea_audit.write_event(
|
||||
gitea_audit.build_event(
|
||||
action="gitea_apply_stale_worker_retirement",
|
||||
result=(
|
||||
gitea_audit.SUCCEEDED
|
||||
if result.get("mutation_performed")
|
||||
else gitea_audit.BLOCKED
|
||||
if not result.get("success")
|
||||
else gitea_audit.ALLOWED
|
||||
),
|
||||
remote=_effective_remote(remote),
|
||||
repository=canon,
|
||||
profile_name=profile.get("profile_name"),
|
||||
audit_label=profile.get("audit_label"),
|
||||
authenticated_username=authenticated_username,
|
||||
task_role=role_kind,
|
||||
operation="worker_registry.retire_stale_workers",
|
||||
reason=result.get("outcome"),
|
||||
request_metadata=mcp_fleet_retirement.summarize_plan(result),
|
||||
)
|
||||
)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
return result
|
||||
|
||||
|
||||
def _retirement_post_apply(
|
||||
registry, result: dict, *, canonical_repository: str
|
||||
) -> dict:
|
||||
"""Fresh fleet verification after a retirement attempt (#980 requirement 6)."""
|
||||
import mcp_fleet_retirement
|
||||
import mcp_fleet_snapshot as _fleet_after
|
||||
|
||||
try:
|
||||
after_records = registry.list_workers(status=None)
|
||||
after = _fleet_after.snapshot_instance_fleet(
|
||||
after_records,
|
||||
pid_alive_probe=issue_lock_store.is_process_alive,
|
||||
canonical_repository=canonical_repository,
|
||||
)
|
||||
except Exception as exc:
|
||||
return {
|
||||
"available": False,
|
||||
"reasons": [
|
||||
f"post-apply verification could not be produced: {_redact(str(exc))}"
|
||||
],
|
||||
}
|
||||
retired_ids = {
|
||||
str(r.get("worker_identity")) for r in result.get("retired") or []
|
||||
}
|
||||
return {
|
||||
"available": True,
|
||||
"registry_fingerprint": mcp_fleet_retirement.registry_fingerprint(
|
||||
after_records
|
||||
),
|
||||
"live_worker_count": after.get("live_worker_count"),
|
||||
"stale_worker_count": after.get("stale_worker_count"),
|
||||
"historical_worker_count": after.get("historical_worker_count"),
|
||||
"retired_still_counted_live": sorted(
|
||||
str(w.get("worker_identity"))
|
||||
for w in after.get("live_workers") or []
|
||||
if str(w.get("worker_identity")) in retired_ids
|
||||
),
|
||||
"retired_still_counted_stale": sorted(
|
||||
str(w.get("worker_identity"))
|
||||
for w in after.get("stale_workers") or []
|
||||
if str(w.get("worker_identity")) in retired_ids
|
||||
),
|
||||
"live_fleet_safe": after.get("live_fleet_safe"),
|
||||
"remaining_blockers": [
|
||||
{"classification": f.get("classification"), "detail": f.get("detail")}
|
||||
for f in after.get("active_blockers") or []
|
||||
],
|
||||
"note": (
|
||||
"Stale retirement does not repair untrusted live identity; residual "
|
||||
"legacy_incomplete_identity blockers keep live_fleet_safe false and "
|
||||
"that is a truthful result."
|
||||
),
|
||||
}
|
||||
|
||||
|
||||
@mcp.tool()
|
||||
def gitea_get_runtime_context(
|
||||
remote: str = "dadeschools",
|
||||
|
||||
@@ -0,0 +1,479 @@
|
||||
"""CAS-protected retirement planning for stale worker registrations (#980).
|
||||
|
||||
#978 (merged PR #979) made the fleet observable: every registered namespace
|
||||
worker, its instance attribution, its heartbeat freshness, and a structured
|
||||
classification. It deliberately stopped there — the snapshot is read-only and
|
||||
the control plane still had no sanctioned way to retire registry rows whose
|
||||
owning process is conclusively gone.
|
||||
|
||||
This module is the *decision layer* for that retirement. It is pure: callers
|
||||
supply registry rows, a clock, and a PID probe; nothing here opens SQLite,
|
||||
scans process tables, or mutates state. The transactional apply lives in
|
||||
:meth:`mcp_worker_identity.WorkerRegistry.retire_stale_workers`, which calls
|
||||
back into these same pure functions so plan and apply can never disagree about
|
||||
what "the registry looks like" or "which rows are eligible".
|
||||
|
||||
Why a separate token
|
||||
--------------------
|
||||
|
||||
``mcp_fleet_snapshot._consistency_token`` seeds its digest with ``snapshot_at``
|
||||
at second precision, so ``registry_revision`` changes on every call even when
|
||||
no registry row changed. A compare-and-swap gated on it can never pass — a
|
||||
dry-run/apply cycle spanning more than one second aborts unconditionally. That
|
||||
token is still useful as an observation stamp, so it is left exactly as it is;
|
||||
#980 gets its own :func:`registry_fingerprint`, derived *only* from canonical
|
||||
retirement-relevant row content:
|
||||
|
||||
* identical registry contents observed at any two times produce the same token,
|
||||
* row order never affects the token (serialized rows are sorted),
|
||||
* any create/delete/identity/liveness/ownership/registration-state change to a
|
||||
retirement-relevant field changes the token.
|
||||
|
||||
Fail-closed posture
|
||||
-------------------
|
||||
|
||||
A worker is retired only when the control plane *conclusively* establishes it
|
||||
is a stale orphan. Missing evidence is never read as permission: an unprobeable
|
||||
PID, an unparsable heartbeat, a row that shares identity evidence with a live
|
||||
or unprobeable worker, a foreign or absent repository binding, or a worker that
|
||||
still owns an active workflow lease all preserve the row.
|
||||
|
||||
Trusted launcher identity (``inst-…`` provenance) is deliberately *not* part of
|
||||
the conjunction. #980 places trusted ``client_instance_id`` propagation out of
|
||||
scope and lists "backfilling trusted identity for legacy workers" as a non-goal;
|
||||
requiring it here would preserve every legacy row forever and make the feature
|
||||
inert. What *is* required is that the registry fields the conjunction reads are
|
||||
actually present — see :data:`REQUIRED_IDENTITY_FIELDS`.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
from datetime import datetime
|
||||
from typing import Any, Callable, Iterable, Mapping, Sequence
|
||||
|
||||
import mcp_fleet_snapshot as fleet
|
||||
import mcp_worker_identity as mwi
|
||||
|
||||
# --- Outcomes -------------------------------------------------------------
|
||||
|
||||
OUTCOME_PLANNED = "planned"
|
||||
OUTCOME_APPLIED = "applied"
|
||||
OUTCOME_REGISTRY_MOVED = "registry_revision_moved"
|
||||
OUTCOME_CANDIDATES_MOVED = "candidate_set_moved"
|
||||
OUTCOME_ALREADY_RETIRED = "already_retired"
|
||||
OUTCOME_NOTHING_REQUESTED = "nothing_requested"
|
||||
|
||||
# --- Reason codes ---------------------------------------------------------
|
||||
|
||||
#: The only reason code that authorizes retirement.
|
||||
REASON_ELIGIBLE = "eligible_stale_orphan"
|
||||
|
||||
REASON_ALREADY_TERMINAL = "already_terminal_registration"
|
||||
REASON_AMBIGUOUS_OWNERSHIP = "ambiguous_ownership_state"
|
||||
REASON_CONFLICTING_IDENTITY = "conflicting_identity_evidence"
|
||||
REASON_FOREIGN_REPOSITORY = "repository_binding_ambiguous"
|
||||
REASON_HEARTBEAT_FRESH = "heartbeat_not_expired"
|
||||
REASON_INCOMPLETE_IDENTITY = "incomplete_registry_identity"
|
||||
REASON_NOT_IN_PLAN = "not_in_current_plan"
|
||||
REASON_PID_ALIVE = "pid_alive"
|
||||
REASON_PID_UNKNOWN = "pid_liveness_unknown"
|
||||
REASON_PROTECTED_OWNER = "protected_active_workflow_owner"
|
||||
REASON_ROW_CHANGED = "row_changed_since_plan"
|
||||
REASON_ROW_MISSING = "registration_missing"
|
||||
REASON_UNPARSABLE_HEARTBEAT = "unparsable_heartbeat"
|
||||
REASON_WORKER_LIVE = "worker_live"
|
||||
|
||||
#: Registry columns that must carry a usable value before the eligibility
|
||||
#: conjunction can even be evaluated. Absence is ambiguity, not permission.
|
||||
REQUIRED_IDENTITY_FIELDS: tuple[str, ...] = (
|
||||
"worker_identity",
|
||||
"client_instance_id",
|
||||
"session_id",
|
||||
"generation_id",
|
||||
"status",
|
||||
"started_at",
|
||||
"last_heartbeat_at",
|
||||
"heartbeat_ttl_seconds",
|
||||
"pid",
|
||||
)
|
||||
|
||||
#: Canonical retirement-relevant content. Ordering here is fixed and part of
|
||||
#: the token contract; adding a field changes every fingerprint, so a change
|
||||
#: here is a deliberate contract revision.
|
||||
#:
|
||||
#: Deliberately excluded: ``token_fingerprint`` (credential-adjacent, never a
|
||||
#: retirement input), the four ``*_revision`` columns (revision drift is an
|
||||
#: independent restart concern and is not part of the eligibility conjunction),
|
||||
#: and the ``retired_*`` bookkeeping columns this feature adds.
|
||||
FINGERPRINT_FIELDS: tuple[str, ...] = (
|
||||
"worker_identity",
|
||||
"client_name",
|
||||
"client_instance_id",
|
||||
"session_id",
|
||||
"generation_id",
|
||||
"role",
|
||||
"profile",
|
||||
"namespace",
|
||||
"remote",
|
||||
"repository_binding",
|
||||
"pid",
|
||||
"process_identity",
|
||||
"transport",
|
||||
"started_at",
|
||||
"last_heartbeat_at",
|
||||
"heartbeat_ttl_seconds",
|
||||
"fencing_epoch",
|
||||
"status",
|
||||
"fleet_run_id",
|
||||
"authenticated_account",
|
||||
"instance_id_provenance",
|
||||
)
|
||||
|
||||
_FINGERPRINT_VERSION = "registryfp-v1"
|
||||
_CANDIDATE_VERSION = "candidatefp-v1"
|
||||
_UNIT = "\x1f"
|
||||
_RECORD = "\x1e"
|
||||
|
||||
|
||||
def _canon(value: Any) -> str:
|
||||
"""Stable text for one field value, independent of Python/SQLite typing.
|
||||
|
||||
``900`` and ``900.0`` are the same TTL and must hash the same; a value that
|
||||
round-trips through SQLite as REAL must not produce a different token than
|
||||
the same value supplied by a caller as ``int``.
|
||||
"""
|
||||
if value is None:
|
||||
return ""
|
||||
if isinstance(value, bool):
|
||||
return "true" if value else "false"
|
||||
if isinstance(value, float):
|
||||
if value != value or value in (float("inf"), float("-inf")):
|
||||
return repr(value)
|
||||
if value.is_integer():
|
||||
return str(int(value))
|
||||
return repr(value)
|
||||
if isinstance(value, int):
|
||||
return str(value)
|
||||
return str(value)
|
||||
|
||||
|
||||
def _serialize_row(row: Mapping[str, Any]) -> str:
|
||||
return _UNIT.join(f"{name}={_canon(row.get(name))}" for name in FINGERPRINT_FIELDS)
|
||||
|
||||
|
||||
def _digest(version: str, serialized: Sequence[str], prefix: str) -> str:
|
||||
ordered = sorted(serialized)
|
||||
material = _RECORD.join([version, str(len(ordered)), *ordered])
|
||||
return f"{prefix}-{hashlib.sha256(material.encode('utf-8')).hexdigest()[:32]}"
|
||||
|
||||
|
||||
def registry_fingerprint(rows: Iterable[Mapping[str, Any]]) -> str:
|
||||
"""Content-derived compare-and-swap token for the worker registry (#980).
|
||||
|
||||
Derived exclusively from :data:`FINGERPRINT_FIELDS` across every row. It
|
||||
contains no ``snapshot_at``, wall-clock, request, or report-generation
|
||||
time, so two observations of an unchanged registry always agree, and the
|
||||
serialized rows are sorted so iteration order cannot perturb the digest.
|
||||
"""
|
||||
return _digest(
|
||||
_FINGERPRINT_VERSION,
|
||||
[_serialize_row(row) for row in rows],
|
||||
"registryfp",
|
||||
)
|
||||
|
||||
|
||||
def candidate_fingerprint(candidate_rows: Iterable[Mapping[str, Any]]) -> str:
|
||||
"""Exact-candidate-set token over the selected rows' canonical content.
|
||||
|
||||
A matching :func:`registry_fingerprint` already implies these rows are
|
||||
unchanged; this second token additionally pins *which* rows the operator
|
||||
approved, so an apply can never widen or narrow the approved set.
|
||||
"""
|
||||
return _digest(
|
||||
_CANDIDATE_VERSION,
|
||||
[_serialize_row(row) for row in candidate_rows],
|
||||
"candidatefp",
|
||||
)
|
||||
|
||||
|
||||
def _probe_pid(
|
||||
pid: Any, pid_alive_probe: Callable[[int | None], bool | None] | None
|
||||
) -> bool | None:
|
||||
if pid_alive_probe is None or pid is None:
|
||||
return None
|
||||
try:
|
||||
result = pid_alive_probe(pid)
|
||||
except Exception:
|
||||
return None
|
||||
return None if result is None else bool(result)
|
||||
|
||||
|
||||
def _evidence(
|
||||
row: Mapping[str, Any],
|
||||
snapshot_row: Mapping[str, Any],
|
||||
pid_alive: bool | None,
|
||||
) -> dict[str, Any]:
|
||||
liveness = snapshot_row.get("liveness") or {}
|
||||
return {
|
||||
"worker_identity": row.get("worker_identity"),
|
||||
"client_type": snapshot_row.get("client_type"),
|
||||
"client_instance_id": row.get("client_instance_id"),
|
||||
"fleet_run_id": row.get("fleet_run_id"),
|
||||
"namespace": row.get("namespace"),
|
||||
"profile": row.get("profile"),
|
||||
"declared_role": row.get("role"),
|
||||
"session_id": row.get("session_id"),
|
||||
"generation_id": row.get("generation_id"),
|
||||
"fencing_epoch": row.get("fencing_epoch"),
|
||||
"process_identity": snapshot_row.get("process_identity"),
|
||||
"pid": row.get("pid"),
|
||||
"pid_alive": pid_alive,
|
||||
"repository_binding": row.get("repository_binding"),
|
||||
"foreign_repository": bool(snapshot_row.get("foreign_repository")),
|
||||
"status": row.get("status"),
|
||||
"started_at": row.get("started_at"),
|
||||
"last_heartbeat_at": row.get("last_heartbeat_at"),
|
||||
"heartbeat_ttl_seconds": row.get("heartbeat_ttl_seconds"),
|
||||
"heartbeat_age_seconds": liveness.get("heartbeat_age_seconds"),
|
||||
"heartbeat_fresh": liveness.get("heartbeat_fresh"),
|
||||
"live": bool(snapshot_row.get("live")),
|
||||
"ownership_state": snapshot_row.get("ownership_state"),
|
||||
"instance_id_provenance": snapshot_row.get("instance_id_provenance"),
|
||||
"instance_identity_trusted": bool(
|
||||
snapshot_row.get("instance_identity_trusted")
|
||||
),
|
||||
}
|
||||
|
||||
|
||||
def _conflict_keys(
|
||||
row: Mapping[str, Any], snapshot_row: Mapping[str, Any]
|
||||
) -> list[tuple[str, str]]:
|
||||
keys: list[tuple[str, str]] = []
|
||||
for name, value in (
|
||||
("session_id", row.get("session_id")),
|
||||
("generation_id", row.get("generation_id")),
|
||||
("process_identity", snapshot_row.get("process_identity")),
|
||||
("pid", row.get("pid")),
|
||||
):
|
||||
text = _canon(value)
|
||||
if text:
|
||||
keys.append((name, text))
|
||||
return keys
|
||||
|
||||
|
||||
def _missing_identity_fields(row: Mapping[str, Any]) -> list[str]:
|
||||
missing: list[str] = []
|
||||
for name in REQUIRED_IDENTITY_FIELDS:
|
||||
value = row.get(name)
|
||||
if value is None or (isinstance(value, str) and not value.strip()):
|
||||
missing.append(name)
|
||||
return missing
|
||||
|
||||
|
||||
def plan_stale_worker_retirement(
|
||||
rows: Iterable[Mapping[str, Any]],
|
||||
*,
|
||||
now: datetime | None = None,
|
||||
pid_alive_probe: Callable[[int | None], bool | None] | None = None,
|
||||
canonical_repository: str | None = None,
|
||||
protected_worker_identities: Iterable[str] | None = None,
|
||||
protected_session_ids: Iterable[str] | None = None,
|
||||
protected_pids: Iterable[Any] | None = None,
|
||||
) -> dict[str, Any]:
|
||||
"""Decide, without mutating anything, which registrations may be retired.
|
||||
|
||||
Every row lands in exactly one of ``candidates`` (eligible) or
|
||||
``preserved`` (with the reason code that stopped it), so the output
|
||||
explains the whole registry rather than only the interesting part.
|
||||
"""
|
||||
all_rows = [dict(row) for row in rows]
|
||||
protected_ids = {str(w) for w in (protected_worker_identities or []) if w}
|
||||
protected_sessions = {str(s) for s in (protected_session_ids or []) if s}
|
||||
protected_pid_set = {_canon(p) for p in (protected_pids or []) if p is not None}
|
||||
|
||||
snapshots: dict[int, dict[str, Any]] = {}
|
||||
pid_alive_by_index: dict[int, bool | None] = {}
|
||||
for index, row in enumerate(all_rows):
|
||||
pid_alive = _probe_pid(row.get("pid"), pid_alive_probe)
|
||||
pid_alive_by_index[index] = pid_alive
|
||||
snapshots[index] = fleet.build_worker_snapshot_row(
|
||||
row,
|
||||
now=now,
|
||||
pid_alive_probe=(lambda _pid, _value=pid_alive: _value),
|
||||
canonical_repository=canonical_repository,
|
||||
)
|
||||
|
||||
# Identity evidence owned by a worker that is live, or whose liveness could
|
||||
# not be established, is ambiguous: anything sharing it is preserved.
|
||||
ambiguous_keys: set[tuple[str, str]] = set()
|
||||
identity_counts: dict[str, int] = {}
|
||||
for index, row in enumerate(all_rows):
|
||||
identity = _canon(row.get("worker_identity"))
|
||||
if identity:
|
||||
identity_counts[identity] = identity_counts.get(identity, 0) + 1
|
||||
snapshot_row = snapshots[index]
|
||||
liveness = snapshot_row.get("liveness") or {}
|
||||
unresolved = (
|
||||
pid_alive_by_index[index] is None
|
||||
or liveness.get("heartbeat_fresh") is None
|
||||
)
|
||||
if snapshot_row.get("live") or (
|
||||
str(row.get("status") or "") == mwi.STATUS_ACTIVE and unresolved
|
||||
):
|
||||
ambiguous_keys.update(_conflict_keys(row, snapshot_row))
|
||||
|
||||
candidates: list[dict[str, Any]] = []
|
||||
candidate_rows: list[Mapping[str, Any]] = []
|
||||
preserved: list[dict[str, Any]] = []
|
||||
|
||||
for index, row in enumerate(all_rows):
|
||||
snapshot_row = snapshots[index]
|
||||
pid_alive = pid_alive_by_index[index]
|
||||
liveness = snapshot_row.get("liveness") or {}
|
||||
evidence = _evidence(row, snapshot_row, pid_alive)
|
||||
blocked: tuple[str, str] | None = None
|
||||
|
||||
identity = _canon(row.get("worker_identity"))
|
||||
missing = _missing_identity_fields(row)
|
||||
shared = sorted(
|
||||
f"{name}={value}"
|
||||
for name, value in _conflict_keys(row, snapshot_row)
|
||||
if (name, value) in ambiguous_keys
|
||||
)
|
||||
binding = (row.get("repository_binding") or "").strip()
|
||||
protected_hits: list[str] = []
|
||||
if identity and identity in protected_ids:
|
||||
protected_hits.append(f"worker_identity={identity}")
|
||||
if _canon(row.get("session_id")) in protected_sessions:
|
||||
protected_hits.append(f"session_id={_canon(row.get('session_id'))}")
|
||||
if _canon(row.get("pid")) in protected_pid_set:
|
||||
protected_hits.append(f"pid={_canon(row.get('pid'))}")
|
||||
|
||||
if identity and identity_counts.get(identity, 0) > 1:
|
||||
blocked = (
|
||||
REASON_CONFLICTING_IDENTITY,
|
||||
f"worker identity {identity!r} appears on more than one registry row",
|
||||
)
|
||||
elif str(row.get("status") or "") != mwi.STATUS_ACTIVE:
|
||||
blocked = (
|
||||
REASON_ALREADY_TERMINAL,
|
||||
f"registration status is {row.get('status')!r}; nothing to retire",
|
||||
)
|
||||
elif missing:
|
||||
blocked = (
|
||||
REASON_INCOMPLETE_IDENTITY,
|
||||
"registry row is missing field(s) the retirement conjunction "
|
||||
f"reads: {missing}",
|
||||
)
|
||||
elif mwi._parse_ts(row.get("last_heartbeat_at")) is None:
|
||||
blocked = (
|
||||
REASON_UNPARSABLE_HEARTBEAT,
|
||||
"last_heartbeat_at is not a parsable UTC stamp; liveness is unknown",
|
||||
)
|
||||
elif snapshot_row.get("live"):
|
||||
blocked = (REASON_WORKER_LIVE, "worker is live and must not be retired")
|
||||
elif pid_alive is None:
|
||||
blocked = (
|
||||
REASON_PID_UNKNOWN,
|
||||
f"pid {row.get('pid')!r} could not be probed; liveness is unproven",
|
||||
)
|
||||
elif pid_alive:
|
||||
blocked = (
|
||||
REASON_PID_ALIVE,
|
||||
f"recorded pid {row.get('pid')!r} is still running",
|
||||
)
|
||||
elif liveness.get("heartbeat_fresh") is not False:
|
||||
blocked = (
|
||||
REASON_HEARTBEAT_FRESH,
|
||||
"heartbeat has not expired under the canonical TTL policy",
|
||||
)
|
||||
elif snapshot_row.get("ownership_state") != "stale":
|
||||
blocked = (
|
||||
REASON_AMBIGUOUS_OWNERSHIP,
|
||||
"ownership_state is "
|
||||
f"{snapshot_row.get('ownership_state')!r}, not 'stale'",
|
||||
)
|
||||
elif not binding or snapshot_row.get("foreign_repository"):
|
||||
blocked = (
|
||||
REASON_FOREIGN_REPOSITORY,
|
||||
"repository binding is absent or does not match the canonical "
|
||||
"repository; retirement scope is ambiguous",
|
||||
)
|
||||
elif shared:
|
||||
blocked = (
|
||||
REASON_CONFLICTING_IDENTITY,
|
||||
"identity evidence is shared with a live or unprobeable worker: "
|
||||
f"{shared}",
|
||||
)
|
||||
elif protected_hits:
|
||||
blocked = (
|
||||
REASON_PROTECTED_OWNER,
|
||||
"worker still owns active workflow state requiring separate "
|
||||
f"reconciliation: {protected_hits}",
|
||||
)
|
||||
|
||||
if blocked is not None:
|
||||
preserved.append(
|
||||
{
|
||||
"worker_identity": row.get("worker_identity"),
|
||||
"reason_code": blocked[0],
|
||||
"detail": blocked[1],
|
||||
"evidence": evidence,
|
||||
}
|
||||
)
|
||||
continue
|
||||
|
||||
candidates.append(
|
||||
{
|
||||
"worker_identity": row.get("worker_identity"),
|
||||
"reason_code": REASON_ELIGIBLE,
|
||||
"detail": (
|
||||
"dead pid, expired heartbeat, stale ownership, unambiguous "
|
||||
"identity, canonical repository binding, no active workflow "
|
||||
"ownership"
|
||||
),
|
||||
"evidence": evidence,
|
||||
}
|
||||
)
|
||||
candidate_rows.append(row)
|
||||
|
||||
counts: dict[str, int] = {}
|
||||
for entry in preserved:
|
||||
counts[entry["reason_code"]] = counts.get(entry["reason_code"], 0) + 1
|
||||
|
||||
return {
|
||||
"success": True,
|
||||
"read_only": True,
|
||||
"mutation_performed": False,
|
||||
"outcome": OUTCOME_PLANNED,
|
||||
"registry_fingerprint": registry_fingerprint(all_rows),
|
||||
"candidate_fingerprint": candidate_fingerprint(candidate_rows),
|
||||
"assessed_count": len(all_rows),
|
||||
"candidate_count": len(candidates),
|
||||
"preserved_count": len(preserved),
|
||||
"candidates": candidates,
|
||||
"candidate_worker_identities": [c["worker_identity"] for c in candidates],
|
||||
"preserved": preserved,
|
||||
"preserved_reason_counts": counts,
|
||||
"protected_inputs": {
|
||||
"worker_identities": sorted(protected_ids),
|
||||
"session_ids": sorted(protected_sessions),
|
||||
"pids": sorted(protected_pid_set),
|
||||
},
|
||||
"canonical_repository": canonical_repository,
|
||||
}
|
||||
|
||||
|
||||
def summarize_plan(plan: Mapping[str, Any]) -> dict[str, Any]:
|
||||
"""Compact, log-safe view of a plan or apply result."""
|
||||
return {
|
||||
"outcome": plan.get("outcome"),
|
||||
"registry_fingerprint": plan.get("registry_fingerprint"),
|
||||
"candidate_fingerprint": plan.get("candidate_fingerprint"),
|
||||
"assessed_count": plan.get("assessed_count"),
|
||||
"candidate_count": plan.get("candidate_count"),
|
||||
"retired_count": plan.get("retired_count"),
|
||||
"preserved_count": plan.get("preserved_count"),
|
||||
"mutation_performed": plan.get("mutation_performed"),
|
||||
}
|
||||
+279
-1
@@ -39,7 +39,7 @@ import threading
|
||||
import time
|
||||
from contextlib import contextmanager
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any, Iterator
|
||||
from typing import Any, Callable, Iterator, Sequence
|
||||
|
||||
|
||||
# --- Provenance verdicts -------------------------------------------------
|
||||
@@ -224,6 +224,13 @@ def _heartbeat_expectation_drift(
|
||||
STATUS_ACTIVE = "active"
|
||||
STATUS_SUPERSEDED = "superseded"
|
||||
STATUS_RELEASED = "released"
|
||||
#: #980 terminal state for a registration whose owning process is conclusively
|
||||
#: gone. Distinct from ``released`` (the worker said goodbye) and
|
||||
#: ``superseded`` (a newer generation took over): ``retired`` records that the
|
||||
#: *control plane* concluded the row was a stale orphan and retired it under a
|
||||
#: compare-and-swap. Like every non-active status it is not live, so a retired
|
||||
#: row counts as historical rather than stale in the #978 fleet snapshot.
|
||||
STATUS_RETIRED = "retired"
|
||||
|
||||
_TRUE_VALUES = frozenset({"1", "true", "yes", "client_managed"})
|
||||
_FALSE_VALUES = frozenset({"0", "false", "no", "manual", "manual_launch"})
|
||||
@@ -315,6 +322,12 @@ _SCHEMA_OPTIONAL_COLUMNS: tuple[tuple[str, str], ...] = (
|
||||
("parity_revision", "TEXT"),
|
||||
("live_revision", "TEXT"),
|
||||
("instance_id_provenance", "TEXT"),
|
||||
# #980 retirement bookkeeping. Deliberately outside the CAS fingerprint
|
||||
# field set: they record *that* a retirement happened, and are written only
|
||||
# by the retirement transaction itself.
|
||||
("retired_at", "TEXT"),
|
||||
("retired_by", "TEXT"),
|
||||
("retirement_reason", "TEXT"),
|
||||
)
|
||||
|
||||
|
||||
@@ -1217,6 +1230,271 @@ class WorkerRegistry:
|
||||
"reasons": [],
|
||||
}
|
||||
|
||||
def retire_stale_workers(
|
||||
self,
|
||||
*,
|
||||
expected_registry_fingerprint: str,
|
||||
expected_candidate_fingerprint: str,
|
||||
worker_identities: Sequence[str],
|
||||
fingerprint_fn: Callable[[list[dict[str, Any]]], str],
|
||||
plan_fn: Callable[[list[dict[str, Any]]], dict[str, Any]],
|
||||
retired_by: str | None = None,
|
||||
retirement_reason: str = "",
|
||||
now: datetime | None = None,
|
||||
) -> dict[str, Any]:
|
||||
"""Compare-and-swap retirement of conclusively stale registrations (#980).
|
||||
|
||||
The whole decision happens inside one ``BEGIN IMMEDIATE`` transaction:
|
||||
the authoritative rows are re-read, the stable registry fingerprint is
|
||||
recomputed from *those* rows, the eligibility plan is recomputed from
|
||||
*those* rows, and only then are the approved targets retired — each
|
||||
with a per-row guarded ``UPDATE`` that also asserts the row's identity,
|
||||
liveness, and ownership columns are byte-identical to what the
|
||||
revalidation just read. There is no window in which a safety check and
|
||||
its matching write are separated by another statement, so a worker that
|
||||
comes back to life, changes ownership, or is retired concurrently
|
||||
cannot be deleted on the strength of a stale observation.
|
||||
|
||||
``fingerprint_fn`` and ``plan_fn`` are injected rather than imported so
|
||||
the storage layer never depends on the decision layer; production wires
|
||||
in :func:`mcp_fleet_retirement.registry_fingerprint` and
|
||||
:func:`mcp_fleet_retirement.plan_stale_worker_retirement`, which is
|
||||
exactly what the plan surface used.
|
||||
|
||||
Any exception — including a failure to commit — rolls the transaction
|
||||
back and is reported as ``transaction_failed`` with zero retirements
|
||||
and ``mutation_performed`` false, so a partial write can never be
|
||||
reported as success.
|
||||
"""
|
||||
requested = [str(w) for w in (worker_identities or []) if str(w).strip()]
|
||||
stamp = _ts(now or _utc_now())
|
||||
result: dict[str, Any] | None = None
|
||||
|
||||
def _base(outcome: str) -> dict[str, Any]:
|
||||
return {
|
||||
"success": True,
|
||||
"outcome": outcome,
|
||||
"mutation_performed": False,
|
||||
"retired": [],
|
||||
"retired_count": 0,
|
||||
"preserved": [],
|
||||
"preserved_count": 0,
|
||||
"requested_count": len(requested),
|
||||
"expected_registry_fingerprint": expected_registry_fingerprint,
|
||||
"expected_candidate_fingerprint": expected_candidate_fingerprint,
|
||||
"acting_identity": retired_by,
|
||||
"reasons": [],
|
||||
}
|
||||
|
||||
if not requested:
|
||||
outcome = _base("nothing_requested")
|
||||
outcome["reasons"] = [
|
||||
"no worker identities were supplied; nothing to retire"
|
||||
]
|
||||
return outcome
|
||||
|
||||
try:
|
||||
with self._tx() as conn:
|
||||
rows = [
|
||||
self._row_to_record(r)
|
||||
for r in conn.execute(
|
||||
"SELECT * FROM worker_registrations"
|
||||
).fetchall()
|
||||
]
|
||||
by_identity = {
|
||||
str(row.get("worker_identity")): row for row in rows
|
||||
}
|
||||
current_registry_fingerprint = fingerprint_fn(rows)
|
||||
|
||||
if current_registry_fingerprint != expected_registry_fingerprint:
|
||||
already = [
|
||||
wid
|
||||
for wid in requested
|
||||
if str(
|
||||
(by_identity.get(wid) or {}).get("status") or ""
|
||||
)
|
||||
== STATUS_RETIRED
|
||||
]
|
||||
idempotent = len(already) == len(requested)
|
||||
result = _base(
|
||||
"already_retired" if idempotent else "registry_revision_moved"
|
||||
)
|
||||
result["idempotent"] = idempotent
|
||||
result["current_registry_fingerprint"] = (
|
||||
current_registry_fingerprint
|
||||
)
|
||||
result["reasons"] = [
|
||||
"the worker registry changed between plan and apply; "
|
||||
"retiring zero workers"
|
||||
if not idempotent
|
||||
else "every requested registration is already retired; "
|
||||
"safe no-op"
|
||||
]
|
||||
result["preserved"] = [
|
||||
{
|
||||
"worker_identity": wid,
|
||||
"reason_code": "registry_revision_moved",
|
||||
"detail": (
|
||||
"aborted before any retirement: registry "
|
||||
"fingerprint moved"
|
||||
),
|
||||
}
|
||||
for wid in requested
|
||||
]
|
||||
result["preserved_count"] = len(requested)
|
||||
return result
|
||||
|
||||
fresh_plan = plan_fn(rows)
|
||||
current_candidate_fingerprint = fresh_plan.get(
|
||||
"candidate_fingerprint"
|
||||
)
|
||||
if current_candidate_fingerprint != expected_candidate_fingerprint:
|
||||
result = _base("candidate_set_moved")
|
||||
result["current_registry_fingerprint"] = (
|
||||
current_registry_fingerprint
|
||||
)
|
||||
result["current_candidate_fingerprint"] = (
|
||||
current_candidate_fingerprint
|
||||
)
|
||||
result["reasons"] = [
|
||||
"the retirement candidate set changed between plan and "
|
||||
"apply; retiring zero workers"
|
||||
]
|
||||
result["preserved"] = [
|
||||
{
|
||||
"worker_identity": wid,
|
||||
"reason_code": "candidate_set_moved",
|
||||
"detail": (
|
||||
"aborted before any retirement: candidate "
|
||||
"fingerprint moved"
|
||||
),
|
||||
}
|
||||
for wid in requested
|
||||
]
|
||||
result["preserved_count"] = len(requested)
|
||||
return result
|
||||
|
||||
eligible = {
|
||||
str(c.get("worker_identity")): c
|
||||
for c in fresh_plan.get("candidates") or []
|
||||
}
|
||||
preserved_index = {
|
||||
str(p.get("worker_identity")): p
|
||||
for p in fresh_plan.get("preserved") or []
|
||||
}
|
||||
|
||||
retired: list[dict[str, Any]] = []
|
||||
preserved: list[dict[str, Any]] = []
|
||||
for wid in requested:
|
||||
row = by_identity.get(wid)
|
||||
if row is None:
|
||||
preserved.append(
|
||||
{
|
||||
"worker_identity": wid,
|
||||
"reason_code": "registration_missing",
|
||||
"detail": "no registration row with this identity",
|
||||
}
|
||||
)
|
||||
continue
|
||||
if wid not in eligible:
|
||||
blocked = preserved_index.get(wid) or {}
|
||||
preserved.append(
|
||||
{
|
||||
"worker_identity": wid,
|
||||
"reason_code": blocked.get("reason_code")
|
||||
or "not_in_current_plan",
|
||||
"detail": blocked.get("detail")
|
||||
or (
|
||||
"revalidation immediately before retirement "
|
||||
"no longer finds this worker eligible"
|
||||
),
|
||||
"evidence": blocked.get("evidence"),
|
||||
}
|
||||
)
|
||||
continue
|
||||
|
||||
cursor = conn.execute(
|
||||
"UPDATE worker_registrations "
|
||||
"SET status = ?, retired_at = ?, retired_by = ?, "
|
||||
" retirement_reason = ? "
|
||||
"WHERE worker_identity = ? "
|
||||
" AND status = ? "
|
||||
" AND last_heartbeat_at = ? "
|
||||
" AND generation_id = ? "
|
||||
" AND session_id = ? "
|
||||
" AND fencing_epoch = ? "
|
||||
" AND IFNULL(pid, -1) = IFNULL(?, -1)",
|
||||
(
|
||||
STATUS_RETIRED,
|
||||
stamp,
|
||||
retired_by,
|
||||
retirement_reason
|
||||
or (eligible[wid].get("reason_code") or ""),
|
||||
wid,
|
||||
STATUS_ACTIVE,
|
||||
row.get("last_heartbeat_at"),
|
||||
row.get("generation_id"),
|
||||
row.get("session_id"),
|
||||
row.get("fencing_epoch"),
|
||||
row.get("pid"),
|
||||
),
|
||||
)
|
||||
if cursor.rowcount == 1:
|
||||
retired.append(
|
||||
{
|
||||
"worker_identity": wid,
|
||||
"reason_code": eligible[wid].get("reason_code"),
|
||||
"retired_at": stamp,
|
||||
"retired_by": retired_by,
|
||||
"evidence": eligible[wid].get("evidence"),
|
||||
}
|
||||
)
|
||||
else:
|
||||
preserved.append(
|
||||
{
|
||||
"worker_identity": wid,
|
||||
"reason_code": "row_changed_since_plan",
|
||||
"detail": (
|
||||
"guarded update matched no row; the "
|
||||
"registration changed inside the retirement "
|
||||
"transaction"
|
||||
),
|
||||
}
|
||||
)
|
||||
|
||||
result = _base("applied")
|
||||
result["mutation_performed"] = bool(retired)
|
||||
result["retired"] = retired
|
||||
result["retired_count"] = len(retired)
|
||||
result["preserved"] = preserved
|
||||
result["preserved_count"] = len(preserved)
|
||||
result["current_registry_fingerprint"] = (
|
||||
current_registry_fingerprint
|
||||
)
|
||||
result["current_candidate_fingerprint"] = (
|
||||
current_candidate_fingerprint
|
||||
)
|
||||
result["retired_at"] = stamp if retired else None
|
||||
except Exception as exc: # rolled back by _tx; report, never half-claim
|
||||
failure = _base("transaction_failed")
|
||||
failure["success"] = False
|
||||
failure["reasons"] = [
|
||||
"retirement transaction failed and was rolled back; zero "
|
||||
f"registrations were retired: {type(exc).__name__}: {exc}"
|
||||
]
|
||||
failure["preserved"] = [
|
||||
{
|
||||
"worker_identity": wid,
|
||||
"reason_code": "transaction_failed",
|
||||
"detail": "transaction rolled back before any commit",
|
||||
}
|
||||
for wid in requested
|
||||
]
|
||||
failure["preserved_count"] = len(requested)
|
||||
return failure
|
||||
|
||||
return result
|
||||
|
||||
|
||||
def _public_record(record: dict[str, Any]) -> dict[str, Any]:
|
||||
"""Registry row minus anything that should not travel to an LLM surface."""
|
||||
|
||||
@@ -175,6 +175,32 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
|
||||
"permission": "gitea.read",
|
||||
"role": "controller",
|
||||
},
|
||||
# #980: CAS-protected retirement of conclusively stale worker
|
||||
# registrations. The mutation lands in the local control-plane worker
|
||||
# registry, not in Gitea, so — exactly like the #601 lease lifecycle — the
|
||||
# Gitea operation gate stays ``gitea.read`` and no new Gitea write
|
||||
# permission is introduced for any profile. The real authority is enforced
|
||||
# in the tools themselves: role_kind must be controller or reconciler, the
|
||||
# runtime must be parity-clean and cohort-unique, and apply additionally
|
||||
# requires the exact stable registry + candidate fingerprints returned by
|
||||
# the plan. Author, reviewer, and merger profiles keep gitea.read for
|
||||
# diagnosis elsewhere and are refused this surface.
|
||||
"plan_stale_worker_retirement": {
|
||||
"permission": "gitea.read",
|
||||
"role": "controller",
|
||||
},
|
||||
"gitea_plan_stale_worker_retirement": {
|
||||
"permission": "gitea.read",
|
||||
"role": "controller",
|
||||
},
|
||||
"apply_stale_worker_retirement": {
|
||||
"permission": "gitea.read",
|
||||
"role": "controller",
|
||||
},
|
||||
"gitea_apply_stale_worker_retirement": {
|
||||
"permission": "gitea.read",
|
||||
"role": "controller",
|
||||
},
|
||||
# #644: Phase 2 Web Console recovery tasks.
|
||||
"clear_stale_binding": {
|
||||
"permission": "gitea.read",
|
||||
|
||||
@@ -0,0 +1,806 @@
|
||||
"""CAS-protected stale worker retirement (#980).
|
||||
|
||||
Covers the acceptance criteria: the registry CAS token is stable across time
|
||||
and row order but moves on any retirement-relevant change, plan mutates
|
||||
nothing, apply fails closed on drift, every target is revalidated inside the
|
||||
retirement transaction, and live / ambiguous / incomplete rows are preserved.
|
||||
|
||||
The regression test that matters most is
|
||||
``test_snapshot_at_would_have_moved_the_token``: it demonstrates the exact
|
||||
defect the R3-C review found — ``mcp_fleet_snapshot._consistency_token``
|
||||
changes when only the observation time changed — and proves the new token does
|
||||
not.
|
||||
|
||||
All identifiers are synthetic.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
import tempfile
|
||||
import unittest
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from typing import Any
|
||||
|
||||
import mcp_fleet_retirement as retire
|
||||
import mcp_fleet_snapshot as fleet
|
||||
import mcp_worker_identity as mwi
|
||||
|
||||
|
||||
NOW = datetime(2026, 7, 30, 7, 0, 0, tzinfo=timezone.utc)
|
||||
TTL = 900.0
|
||||
REPO = "/synthetic/repo/Gitea-Tools"
|
||||
|
||||
_ROW_COLUMNS = (
|
||||
"worker_identity",
|
||||
"client_name",
|
||||
"client_instance_id",
|
||||
"session_id",
|
||||
"generation_id",
|
||||
"role",
|
||||
"profile",
|
||||
"namespace",
|
||||
"remote",
|
||||
"repository_binding",
|
||||
"pid",
|
||||
"process_identity",
|
||||
"transport",
|
||||
"token_fingerprint",
|
||||
"started_at",
|
||||
"last_heartbeat_at",
|
||||
"heartbeat_ttl_seconds",
|
||||
"fencing_epoch",
|
||||
"status",
|
||||
"fleet_run_id",
|
||||
"authenticated_account",
|
||||
"instance_id_provenance",
|
||||
)
|
||||
|
||||
|
||||
def _registry() -> mwi.WorkerRegistry:
|
||||
handle, path = tempfile.mkstemp(suffix=".sqlite3")
|
||||
os.close(handle)
|
||||
os.unlink(path)
|
||||
return mwi.WorkerRegistry(path)
|
||||
|
||||
|
||||
def _row(
|
||||
*,
|
||||
identity: str,
|
||||
pid: int | None,
|
||||
instance: str = "inst-codex-20260730T070000Z-0123456789ab",
|
||||
session: str = "sess-a",
|
||||
generation: str = "gen-a",
|
||||
namespace: str = "author",
|
||||
heartbeat: datetime | None = None,
|
||||
status: str = mwi.STATUS_ACTIVE,
|
||||
repository_binding: str | None = REPO,
|
||||
ttl: float = TTL,
|
||||
**overrides: Any,
|
||||
) -> dict[str, Any]:
|
||||
"""One synthetic ``worker_registrations`` row."""
|
||||
record = {
|
||||
"worker_identity": identity,
|
||||
"client_name": "codex",
|
||||
"client_instance_id": instance,
|
||||
"session_id": session,
|
||||
"generation_id": generation,
|
||||
"role": namespace,
|
||||
"profile": f"prgs-{namespace}",
|
||||
"namespace": namespace,
|
||||
"remote": "prgs",
|
||||
"repository_binding": repository_binding,
|
||||
"pid": pid,
|
||||
"process_identity": f"pid-{pid}" if pid is not None else None,
|
||||
"transport": "stdio",
|
||||
"token_fingerprint": None,
|
||||
"started_at": "2026-07-30T06:00:00Z",
|
||||
"last_heartbeat_at": mwi._ts(heartbeat or (NOW - timedelta(seconds=7200))),
|
||||
"heartbeat_ttl_seconds": ttl,
|
||||
"fencing_epoch": 1,
|
||||
"status": status,
|
||||
"fleet_run_id": "run-canary",
|
||||
"authenticated_account": "synthetic-user",
|
||||
"instance_id_provenance": fleet.INSTANCE_ID_PROVENANCE_TRUSTED,
|
||||
}
|
||||
record.update(overrides)
|
||||
return record
|
||||
|
||||
|
||||
def _dead(_pid: int | None) -> bool:
|
||||
return False
|
||||
|
||||
|
||||
def _alive(_pid: int | None) -> bool:
|
||||
return True
|
||||
|
||||
|
||||
def _selective(alive_pids: set[int]):
|
||||
def probe(pid: int | None) -> bool:
|
||||
return pid in alive_pids
|
||||
|
||||
return probe
|
||||
|
||||
|
||||
def _plan(rows, *, probe=_dead, now=NOW, **kwargs):
|
||||
return retire.plan_stale_worker_retirement(
|
||||
rows,
|
||||
now=now,
|
||||
pid_alive_probe=probe,
|
||||
canonical_repository=REPO,
|
||||
**kwargs,
|
||||
)
|
||||
|
||||
|
||||
def _register(registry: mwi.WorkerRegistry, row: dict[str, Any]) -> None:
|
||||
"""Insert a synthetic row directly, bypassing register()'s live-now stamps."""
|
||||
columns = [name for name in _ROW_COLUMNS if name in row]
|
||||
placeholders = ", ".join("?" for _ in columns)
|
||||
with registry._tx() as conn:
|
||||
conn.execute(
|
||||
f"INSERT INTO worker_registrations ({', '.join(columns)}) "
|
||||
f"VALUES ({placeholders})",
|
||||
[row[name] for name in columns],
|
||||
)
|
||||
|
||||
|
||||
class RegistryFingerprintStabilityTests(unittest.TestCase):
|
||||
"""AC: identical contents observed at different times produce one token."""
|
||||
|
||||
def test_same_contents_different_observation_times_same_token(self):
|
||||
rows = [_row(identity="w-1", pid=101), _row(identity="w-2", pid=102)]
|
||||
self.assertEqual(
|
||||
retire.registry_fingerprint(rows), retire.registry_fingerprint(rows)
|
||||
)
|
||||
# Recompute after the clock has moved a full hour: the token is derived
|
||||
# from content only, so it cannot notice.
|
||||
plan_early = _plan(rows, now=NOW)
|
||||
plan_late = _plan(rows, now=NOW + timedelta(hours=1))
|
||||
self.assertEqual(
|
||||
plan_early["registry_fingerprint"], plan_late["registry_fingerprint"]
|
||||
)
|
||||
|
||||
def test_snapshot_at_would_have_moved_the_token(self):
|
||||
"""Regression: the old time-seeded derivation moved, the new one does not."""
|
||||
rows = [_row(identity="w-1", pid=101)]
|
||||
old_early = fleet._consistency_token(rows, "2026-07-30T07:00:00Z")
|
||||
old_late = fleet._consistency_token(rows, "2026-07-30T07:00:01Z")
|
||||
self.assertNotEqual(
|
||||
old_early,
|
||||
old_late,
|
||||
"the #980 defect: one second of observation drift changed the token",
|
||||
)
|
||||
self.assertEqual(
|
||||
retire.registry_fingerprint(rows), retire.registry_fingerprint(rows)
|
||||
)
|
||||
|
||||
def test_row_order_does_not_change_the_token(self):
|
||||
rows = [
|
||||
_row(identity="w-1", pid=101),
|
||||
_row(identity="w-2", pid=102),
|
||||
_row(identity="w-3", pid=103),
|
||||
]
|
||||
self.assertEqual(
|
||||
retire.registry_fingerprint(rows),
|
||||
retire.registry_fingerprint(list(reversed(rows))),
|
||||
)
|
||||
|
||||
def test_numeric_typing_does_not_change_the_token(self):
|
||||
as_float = [_row(identity="w-1", pid=101, heartbeat_ttl_seconds=900.0)]
|
||||
as_int = [_row(identity="w-1", pid=101, heartbeat_ttl_seconds=900)]
|
||||
self.assertEqual(
|
||||
retire.registry_fingerprint(as_float),
|
||||
retire.registry_fingerprint(as_int),
|
||||
)
|
||||
|
||||
def test_retirement_relevant_changes_move_the_token(self):
|
||||
base = [_row(identity="w-1", pid=101)]
|
||||
baseline = retire.registry_fingerprint(base)
|
||||
mutations = {
|
||||
"row added": base + [_row(identity="w-2", pid=102)],
|
||||
"row removed": [],
|
||||
"status changed": [_row(identity="w-1", pid=101, status="released")],
|
||||
"identity changed": [
|
||||
_row(
|
||||
identity="w-1",
|
||||
pid=101,
|
||||
instance="inst-codex-20260730T070000Z-ffffffffffff",
|
||||
)
|
||||
],
|
||||
"ownership changed": [_row(identity="w-1", pid=101, session="sess-other")],
|
||||
"generation changed": [
|
||||
_row(identity="w-1", pid=101, generation="gen-other")
|
||||
],
|
||||
"pid changed": [_row(identity="w-1", pid=999)],
|
||||
"repository binding changed": [
|
||||
_row(identity="w-1", pid=101, repository_binding="/elsewhere")
|
||||
],
|
||||
}
|
||||
for label, rows in mutations.items():
|
||||
with self.subTest(change=label):
|
||||
self.assertNotEqual(baseline, retire.registry_fingerprint(rows))
|
||||
|
||||
def test_heartbeat_change_moves_the_token(self):
|
||||
base = [_row(identity="w-1", pid=101)]
|
||||
moved = [_row(identity="w-1", pid=101, heartbeat=NOW - timedelta(seconds=30))]
|
||||
self.assertNotEqual(
|
||||
retire.registry_fingerprint(base), retire.registry_fingerprint(moved)
|
||||
)
|
||||
|
||||
def test_ttl_change_moves_the_token(self):
|
||||
base = [_row(identity="w-1", pid=101)]
|
||||
moved = [_row(identity="w-1", pid=101, ttl=60.0)]
|
||||
self.assertNotEqual(
|
||||
retire.registry_fingerprint(base), retire.registry_fingerprint(moved)
|
||||
)
|
||||
|
||||
|
||||
class PlanEligibilityTests(unittest.TestCase):
|
||||
"""AC: only conclusively stale orphans are selected; everything else stays."""
|
||||
|
||||
def test_plan_selects_dead_stale_orphan(self):
|
||||
plan = _plan([_row(identity="w-1", pid=101)])
|
||||
self.assertEqual(plan["candidate_count"], 1)
|
||||
self.assertEqual(plan["candidate_worker_identities"], ["w-1"])
|
||||
self.assertEqual(plan["candidates"][0]["reason_code"], retire.REASON_ELIGIBLE)
|
||||
self.assertFalse(plan["mutation_performed"])
|
||||
self.assertTrue(plan["read_only"])
|
||||
|
||||
def test_plan_performs_no_mutation(self):
|
||||
registry = _registry()
|
||||
_register(registry, _row(identity="w-1", pid=101))
|
||||
before = registry.list_workers(status=None)
|
||||
plan = _plan(before)
|
||||
after = registry.list_workers(status=None)
|
||||
self.assertEqual(plan["candidate_count"], 1)
|
||||
self.assertEqual(before, after)
|
||||
self.assertEqual([r["status"] for r in after], [mwi.STATUS_ACTIVE])
|
||||
self.assertFalse(plan["mutation_performed"])
|
||||
|
||||
def test_live_worker_is_preserved(self):
|
||||
row = _row(identity="w-live", pid=101, heartbeat=NOW - timedelta(seconds=10))
|
||||
plan = _plan([row], probe=_alive)
|
||||
self.assertEqual(plan["candidate_count"], 0)
|
||||
self.assertEqual(
|
||||
plan["preserved"][0]["reason_code"], retire.REASON_WORKER_LIVE
|
||||
)
|
||||
|
||||
def test_fresh_heartbeat_with_dead_pid_still_fails_closed(self):
|
||||
"""A dead pid withdraws liveness; the unexpired heartbeat still preserves."""
|
||||
row = _row(identity="w-fresh", pid=101, heartbeat=NOW - timedelta(seconds=10))
|
||||
plan = _plan([row], probe=_dead)
|
||||
self.assertEqual(plan["candidate_count"], 0)
|
||||
self.assertEqual(
|
||||
plan["preserved"][0]["reason_code"], retire.REASON_HEARTBEAT_FRESH
|
||||
)
|
||||
|
||||
def test_unprobeable_pid_is_preserved(self):
|
||||
plan = _plan([_row(identity="w-1", pid=101)], probe=lambda _pid: None)
|
||||
self.assertEqual(plan["candidate_count"], 0)
|
||||
self.assertEqual(
|
||||
plan["preserved"][0]["reason_code"], retire.REASON_PID_UNKNOWN
|
||||
)
|
||||
|
||||
def test_incomplete_legacy_identity_is_preserved(self):
|
||||
"""AC: incomplete legacy identities remain fail-closed."""
|
||||
rows = [
|
||||
_row(identity="w-nopid", pid=None, instance="legacy-pid-27833"),
|
||||
_row(identity="w-nosession", pid=102, session=""),
|
||||
]
|
||||
plan = _plan(rows)
|
||||
self.assertEqual(plan["candidate_count"], 0)
|
||||
self.assertEqual(
|
||||
{p["reason_code"] for p in plan["preserved"]},
|
||||
{retire.REASON_INCOMPLETE_IDENTITY},
|
||||
)
|
||||
detail = next(
|
||||
p["detail"] for p in plan["preserved"] if p["worker_identity"] == "w-nopid"
|
||||
)
|
||||
self.assertIn("pid", detail)
|
||||
|
||||
def test_unparsable_heartbeat_is_preserved(self):
|
||||
rows = [_row(identity="w-1", pid=101, last_heartbeat_at="not-a-stamp")]
|
||||
plan = _plan(rows)
|
||||
self.assertEqual(plan["candidate_count"], 0)
|
||||
self.assertEqual(
|
||||
plan["preserved"][0]["reason_code"], retire.REASON_UNPARSABLE_HEARTBEAT
|
||||
)
|
||||
|
||||
def test_identity_shared_with_live_worker_is_preserved(self):
|
||||
"""Two rows, one live: the dead one shares session evidence, so it stays."""
|
||||
rows = [
|
||||
_row(
|
||||
identity="w-live",
|
||||
pid=101,
|
||||
session="sess-shared",
|
||||
heartbeat=NOW - timedelta(seconds=5),
|
||||
),
|
||||
_row(identity="w-dead", pid=102, session="sess-shared"),
|
||||
]
|
||||
plan = _plan(rows, probe=_selective({101}))
|
||||
self.assertEqual(plan["candidate_count"], 0)
|
||||
codes = {p["reason_code"] for p in plan["preserved"]}
|
||||
self.assertIn(retire.REASON_CONFLICTING_IDENTITY, codes)
|
||||
|
||||
def test_foreign_or_missing_repository_binding_is_preserved(self):
|
||||
rows = [
|
||||
_row(identity="w-foreign", pid=101, repository_binding="/other/repo"),
|
||||
_row(identity="w-unbound", pid=102, repository_binding=None),
|
||||
]
|
||||
plan = _plan(rows)
|
||||
self.assertEqual(plan["candidate_count"], 0)
|
||||
self.assertEqual(
|
||||
{p["reason_code"] for p in plan["preserved"]},
|
||||
{retire.REASON_FOREIGN_REPOSITORY},
|
||||
)
|
||||
|
||||
def test_active_workflow_owner_is_preserved(self):
|
||||
rows = [_row(identity="w-1", pid=101, session="sess-leased")]
|
||||
plan = _plan(rows, protected_session_ids=["sess-leased"])
|
||||
self.assertEqual(plan["candidate_count"], 0)
|
||||
self.assertEqual(
|
||||
plan["preserved"][0]["reason_code"], retire.REASON_PROTECTED_OWNER
|
||||
)
|
||||
|
||||
def test_terminal_rows_are_not_retired_again(self):
|
||||
rows = [_row(identity="w-1", pid=101, status=mwi.STATUS_RELEASED)]
|
||||
plan = _plan(rows)
|
||||
self.assertEqual(plan["candidate_count"], 0)
|
||||
self.assertEqual(
|
||||
plan["preserved"][0]["reason_code"], retire.REASON_ALREADY_TERMINAL
|
||||
)
|
||||
|
||||
def test_mixed_fleet_produces_correct_per_worker_outcomes(self):
|
||||
rows = [
|
||||
_row(identity="w-stale-1", pid=101, session="s1", generation="g1"),
|
||||
_row(identity="w-stale-2", pid=102, session="s2", generation="g2"),
|
||||
_row(
|
||||
identity="w-live",
|
||||
pid=103,
|
||||
session="s3",
|
||||
generation="g3",
|
||||
heartbeat=NOW - timedelta(seconds=5),
|
||||
),
|
||||
_row(identity="w-nopid", pid=None, session="s4", generation="g4"),
|
||||
_row(
|
||||
identity="w-foreign",
|
||||
pid=105,
|
||||
session="s5",
|
||||
generation="g5",
|
||||
repository_binding="/other",
|
||||
),
|
||||
_row(
|
||||
identity="w-terminal",
|
||||
pid=106,
|
||||
session="s6",
|
||||
generation="g6",
|
||||
status=mwi.STATUS_SUPERSEDED,
|
||||
),
|
||||
]
|
||||
plan = _plan(rows, probe=_selective({103}))
|
||||
self.assertEqual(
|
||||
sorted(plan["candidate_worker_identities"]), ["w-stale-1", "w-stale-2"]
|
||||
)
|
||||
by_identity = {
|
||||
p["worker_identity"]: p["reason_code"] for p in plan["preserved"]
|
||||
}
|
||||
self.assertEqual(by_identity["w-live"], retire.REASON_WORKER_LIVE)
|
||||
self.assertEqual(by_identity["w-nopid"], retire.REASON_INCOMPLETE_IDENTITY)
|
||||
self.assertEqual(by_identity["w-foreign"], retire.REASON_FOREIGN_REPOSITORY)
|
||||
self.assertEqual(by_identity["w-terminal"], retire.REASON_ALREADY_TERMINAL)
|
||||
self.assertEqual(plan["assessed_count"], 6)
|
||||
self.assertEqual(plan["preserved_count"], 4)
|
||||
|
||||
def test_plan_is_deterministic_for_identical_contents(self):
|
||||
rows = [_row(identity="w-1", pid=101), _row(identity="w-2", pid=102)]
|
||||
first = _plan(rows)
|
||||
second = _plan(list(reversed(rows)), now=NOW + timedelta(minutes=5))
|
||||
self.assertEqual(first["registry_fingerprint"], second["registry_fingerprint"])
|
||||
self.assertEqual(
|
||||
first["candidate_fingerprint"], second["candidate_fingerprint"]
|
||||
)
|
||||
self.assertEqual(
|
||||
sorted(first["candidate_worker_identities"]),
|
||||
sorted(second["candidate_worker_identities"]),
|
||||
)
|
||||
|
||||
|
||||
class ApplyCasTests(unittest.TestCase):
|
||||
"""AC: apply is compare-and-swap protected and revalidates every target."""
|
||||
|
||||
def setUp(self) -> None:
|
||||
self.registry = _registry()
|
||||
self.probe = _dead
|
||||
|
||||
def _rows(self):
|
||||
return self.registry.list_workers(status=None)
|
||||
|
||||
def _apply(self, plan, *, probe=None, identities=None, **kwargs):
|
||||
chosen = probe or self.probe
|
||||
return self.registry.retire_stale_workers(
|
||||
expected_registry_fingerprint=plan["registry_fingerprint"],
|
||||
expected_candidate_fingerprint=plan["candidate_fingerprint"],
|
||||
worker_identities=(
|
||||
identities
|
||||
if identities is not None
|
||||
else plan["candidate_worker_identities"]
|
||||
),
|
||||
fingerprint_fn=retire.registry_fingerprint,
|
||||
plan_fn=lambda rows: _plan(rows, probe=chosen, **kwargs),
|
||||
retired_by="synthetic-user/prgs-reconciler",
|
||||
now=NOW,
|
||||
)
|
||||
|
||||
def test_matching_token_retires_the_planned_set(self):
|
||||
_register(self.registry, _row(identity="w-1", pid=101, session="s1"))
|
||||
_register(self.registry, _row(identity="w-2", pid=102, session="s2"))
|
||||
plan = _plan(self._rows(), probe=self.probe)
|
||||
result = self._apply(plan)
|
||||
self.assertEqual(result["outcome"], "applied")
|
||||
self.assertTrue(result["mutation_performed"])
|
||||
self.assertEqual(result["retired_count"], 2)
|
||||
statuses = {r["worker_identity"]: r["status"] for r in self._rows()}
|
||||
self.assertEqual(
|
||||
statuses, {"w-1": mwi.STATUS_RETIRED, "w-2": mwi.STATUS_RETIRED}
|
||||
)
|
||||
retired_row = self.registry.get("w-1")
|
||||
self.assertEqual(retired_row["retired_at"], mwi._ts(NOW))
|
||||
self.assertEqual(retired_row["retired_by"], "synthetic-user/prgs-reconciler")
|
||||
|
||||
def test_moved_registry_token_retires_zero_workers(self):
|
||||
_register(self.registry, _row(identity="w-1", pid=101))
|
||||
plan = _plan(self._rows(), probe=self.probe)
|
||||
# An unrelated registration lands between plan and apply.
|
||||
_register(self.registry, _row(identity="w-2", pid=102, session="s2"))
|
||||
result = self._apply(plan)
|
||||
self.assertEqual(result["outcome"], retire.OUTCOME_REGISTRY_MOVED)
|
||||
self.assertEqual(result["retired_count"], 0)
|
||||
self.assertFalse(result["mutation_performed"])
|
||||
self.assertTrue(all(r["status"] == mwi.STATUS_ACTIVE for r in self._rows()))
|
||||
|
||||
def test_moved_candidate_set_retires_zero_workers(self):
|
||||
"""Registry unchanged, but the eligibility verdict is no longer the same."""
|
||||
_register(self.registry, _row(identity="w-1", pid=101))
|
||||
plan = _plan(self._rows(), probe=self.probe)
|
||||
# The row did not change; the process came back (pid probes alive), so
|
||||
# revalidation inside the transaction finds no candidates at all.
|
||||
result = self._apply(plan, probe=_alive)
|
||||
self.assertEqual(result["outcome"], retire.OUTCOME_CANDIDATES_MOVED)
|
||||
self.assertEqual(result["retired_count"], 0)
|
||||
self.assertFalse(result["mutation_performed"])
|
||||
self.assertEqual(self.registry.get("w-1")["status"], mwi.STATUS_ACTIVE)
|
||||
|
||||
def test_worker_that_becomes_live_between_plan_and_apply_is_preserved(self):
|
||||
_register(self.registry, _row(identity="w-dead", pid=101, session="s1"))
|
||||
_register(self.registry, _row(identity="w-back", pid=102, session="s2"))
|
||||
plan = _plan(self._rows(), probe=self.probe)
|
||||
self.assertEqual(len(plan["candidate_worker_identities"]), 2)
|
||||
# w-back's process is alive by the time apply runs, so the candidate
|
||||
# fingerprint moves and nothing at all is retired.
|
||||
result = self._apply(plan, probe=_selective({102}))
|
||||
self.assertEqual(result["outcome"], retire.OUTCOME_CANDIDATES_MOVED)
|
||||
self.assertEqual(result["retired_count"], 0)
|
||||
self.assertEqual(self.registry.get("w-back")["status"], mwi.STATUS_ACTIVE)
|
||||
self.assertEqual(self.registry.get("w-dead")["status"], mwi.STATUS_ACTIVE)
|
||||
|
||||
def test_worker_that_becomes_ambiguous_between_plan_and_apply_is_preserved(self):
|
||||
_register(self.registry, _row(identity="w-1", pid=101, session="s1"))
|
||||
plan = _plan(self._rows(), probe=self.probe)
|
||||
result = self._apply(plan, probe=lambda _pid: None)
|
||||
self.assertEqual(result["outcome"], retire.OUTCOME_CANDIDATES_MOVED)
|
||||
self.assertEqual(result["retired_count"], 0)
|
||||
self.assertEqual(self.registry.get("w-1")["status"], mwi.STATUS_ACTIVE)
|
||||
|
||||
def test_apply_revalidates_and_will_not_narrow_the_approved_set(self):
|
||||
"""A protection appearing after plan moves the CAS, so nothing is retired."""
|
||||
_register(self.registry, _row(identity="w-1", pid=101, session="s1"))
|
||||
_register(
|
||||
self.registry, _row(identity="w-protected", pid=102, session="sess-leased")
|
||||
)
|
||||
unprotected_plan = _plan(self._rows(), probe=self.probe)
|
||||
self.assertEqual(unprotected_plan["candidate_count"], 2)
|
||||
result = self._apply(
|
||||
unprotected_plan, protected_session_ids=["sess-leased"]
|
||||
)
|
||||
self.assertEqual(result["outcome"], retire.OUTCOME_CANDIDATES_MOVED)
|
||||
self.assertEqual(result["retired_count"], 0)
|
||||
self.assertEqual(
|
||||
self.registry.get("w-protected")["status"], mwi.STATUS_ACTIVE
|
||||
)
|
||||
|
||||
def test_target_outside_the_current_plan_is_preserved(self):
|
||||
_register(
|
||||
self.registry,
|
||||
_row(identity="w-1", pid=101, session="s1", generation="g1"),
|
||||
)
|
||||
_register(
|
||||
self.registry,
|
||||
_row(
|
||||
identity="w-live",
|
||||
pid=102,
|
||||
session="s2",
|
||||
generation="g2",
|
||||
heartbeat=NOW - timedelta(seconds=5),
|
||||
),
|
||||
)
|
||||
probe = _selective({102})
|
||||
plan = _plan(self._rows(), probe=probe)
|
||||
self.assertEqual(plan["candidate_worker_identities"], ["w-1"])
|
||||
# A caller that appends a live identity to the approved list gets it
|
||||
# preserved with the live reason code; only the real candidate retires.
|
||||
result = self._apply(plan, probe=probe, identities=["w-1", "w-live"])
|
||||
self.assertEqual(result["outcome"], "applied")
|
||||
self.assertEqual(result["retired_count"], 1)
|
||||
self.assertEqual(result["retired"][0]["worker_identity"], "w-1")
|
||||
preserved = {
|
||||
p["worker_identity"]: p["reason_code"] for p in result["preserved"]
|
||||
}
|
||||
self.assertEqual(preserved["w-live"], retire.REASON_WORKER_LIVE)
|
||||
self.assertEqual(self.registry.get("w-live")["status"], mwi.STATUS_ACTIVE)
|
||||
|
||||
def test_missing_registration_is_reported_not_invented(self):
|
||||
_register(self.registry, _row(identity="w-1", pid=101))
|
||||
plan = _plan(self._rows(), probe=self.probe)
|
||||
result = self._apply(plan, identities=["w-1", "w-ghost"])
|
||||
self.assertEqual(result["retired_count"], 1)
|
||||
preserved = {
|
||||
p["worker_identity"]: p["reason_code"] for p in result["preserved"]
|
||||
}
|
||||
self.assertEqual(preserved["w-ghost"], "registration_missing")
|
||||
|
||||
def test_reapplying_a_completed_plan_is_a_safe_no_op(self):
|
||||
_register(self.registry, _row(identity="w-1", pid=101))
|
||||
plan = _plan(self._rows(), probe=self.probe)
|
||||
first = self._apply(plan)
|
||||
self.assertEqual(first["retired_count"], 1)
|
||||
second = self._apply(plan)
|
||||
self.assertEqual(second["outcome"], retire.OUTCOME_ALREADY_RETIRED)
|
||||
self.assertTrue(second["idempotent"])
|
||||
self.assertEqual(second["retired_count"], 0)
|
||||
self.assertFalse(second["mutation_performed"])
|
||||
self.assertEqual(self.registry.get("w-1")["status"], mwi.STATUS_RETIRED)
|
||||
|
||||
def test_empty_target_list_reports_nothing_requested(self):
|
||||
_register(self.registry, _row(identity="w-1", pid=101))
|
||||
plan = _plan(self._rows(), probe=self.probe)
|
||||
result = self._apply(plan, identities=[])
|
||||
self.assertEqual(result["outcome"], retire.OUTCOME_NOTHING_REQUESTED)
|
||||
self.assertFalse(result["mutation_performed"])
|
||||
self.assertEqual(self.registry.get("w-1")["status"], mwi.STATUS_ACTIVE)
|
||||
|
||||
def test_transaction_failure_cannot_report_partial_success(self):
|
||||
_register(self.registry, _row(identity="w-1", pid=101, session="s1"))
|
||||
_register(self.registry, _row(identity="w-2", pid=102, session="s2"))
|
||||
plan = _plan(self._rows(), probe=self.probe)
|
||||
calls = {"n": 0}
|
||||
|
||||
def exploding_plan(rows):
|
||||
calls["n"] += 1
|
||||
raise RuntimeError("synthetic revalidation failure")
|
||||
|
||||
result = self.registry.retire_stale_workers(
|
||||
expected_registry_fingerprint=plan["registry_fingerprint"],
|
||||
expected_candidate_fingerprint=plan["candidate_fingerprint"],
|
||||
worker_identities=plan["candidate_worker_identities"],
|
||||
fingerprint_fn=retire.registry_fingerprint,
|
||||
plan_fn=exploding_plan,
|
||||
retired_by="synthetic-user/prgs-reconciler",
|
||||
now=NOW,
|
||||
)
|
||||
self.assertEqual(calls["n"], 1)
|
||||
self.assertFalse(result["success"])
|
||||
self.assertEqual(result["outcome"], "transaction_failed")
|
||||
self.assertEqual(result["retired_count"], 0)
|
||||
self.assertFalse(result["mutation_performed"])
|
||||
self.assertTrue(
|
||||
all(r["status"] == mwi.STATUS_ACTIVE for r in self._rows()),
|
||||
"a failed transaction must roll back every retirement",
|
||||
)
|
||||
|
||||
def test_structured_result_fields_are_accurate(self):
|
||||
_register(self.registry, _row(identity="w-1", pid=101))
|
||||
plan = _plan(self._rows(), probe=self.probe)
|
||||
blocked = self._apply(plan, probe=_alive)
|
||||
self.assertFalse(blocked["mutation_performed"])
|
||||
self.assertEqual(blocked["retired"], [])
|
||||
self.assertEqual(blocked["requested_count"], 1)
|
||||
self.assertEqual(
|
||||
blocked["expected_registry_fingerprint"], plan["registry_fingerprint"]
|
||||
)
|
||||
applied = self._apply(plan)
|
||||
self.assertTrue(applied["mutation_performed"])
|
||||
self.assertEqual(applied["acting_identity"], "synthetic-user/prgs-reconciler")
|
||||
self.assertEqual(applied["retired"][0]["reason_code"], retire.REASON_ELIGIBLE)
|
||||
self.assertIn("evidence", applied["retired"][0])
|
||||
|
||||
|
||||
class PostRetirementFleetCompatibilityTests(unittest.TestCase):
|
||||
"""AC: retired rows leave the stale count and never fake fleet safety."""
|
||||
|
||||
def test_retired_rows_are_historical_not_stale(self):
|
||||
registry = _registry()
|
||||
_register(
|
||||
registry, _row(identity="w-1", pid=101, session="s1", generation="g1")
|
||||
)
|
||||
_register(
|
||||
registry,
|
||||
_row(
|
||||
identity="w-live",
|
||||
pid=102,
|
||||
session="s2",
|
||||
generation="g2",
|
||||
instance="legacy-pid-27833",
|
||||
instance_id_provenance=fleet.INSTANCE_ID_PROVENANCE_LEGACY,
|
||||
heartbeat=NOW - timedelta(seconds=5),
|
||||
),
|
||||
)
|
||||
probe = _selective({102})
|
||||
before = fleet.snapshot_instance_fleet(
|
||||
registry.list_workers(status=None),
|
||||
now=NOW,
|
||||
pid_alive_probe=probe,
|
||||
canonical_repository=REPO,
|
||||
)
|
||||
self.assertEqual(before["stale_worker_count"], 1)
|
||||
|
||||
plan = _plan(registry.list_workers(status=None), probe=probe)
|
||||
result = registry.retire_stale_workers(
|
||||
expected_registry_fingerprint=plan["registry_fingerprint"],
|
||||
expected_candidate_fingerprint=plan["candidate_fingerprint"],
|
||||
worker_identities=plan["candidate_worker_identities"],
|
||||
fingerprint_fn=retire.registry_fingerprint,
|
||||
plan_fn=lambda rows: _plan(rows, probe=probe),
|
||||
retired_by="synthetic-user/prgs-reconciler",
|
||||
now=NOW,
|
||||
)
|
||||
self.assertEqual(result["retired_count"], 1)
|
||||
|
||||
after = fleet.snapshot_instance_fleet(
|
||||
registry.list_workers(status=None),
|
||||
now=NOW,
|
||||
pid_alive_probe=probe,
|
||||
canonical_repository=REPO,
|
||||
)
|
||||
self.assertEqual(after["stale_worker_count"], 0)
|
||||
self.assertEqual(after["live_worker_count"], 1)
|
||||
self.assertEqual(after["historical_worker_count"], 1)
|
||||
# The surviving live worker still carries a legacy instance identity, so
|
||||
# the fleet must not be declared safe.
|
||||
self.assertFalse(after["live_fleet_safe"])
|
||||
self.assertIn(
|
||||
fleet.CLASS_LEGACY_INCOMPLETE,
|
||||
{f["classification"] for f in after["active_blockers"]},
|
||||
)
|
||||
|
||||
def test_existing_snapshot_shape_is_unchanged(self):
|
||||
rows = [_row(identity="w-1", pid=101)]
|
||||
snapshot = fleet.snapshot_instance_fleet(
|
||||
rows, now=NOW, pid_alive_probe=_dead, canonical_repository=REPO
|
||||
)
|
||||
for key in (
|
||||
"consistency_token",
|
||||
"registry_revision",
|
||||
"snapshot_at",
|
||||
"live_workers",
|
||||
"stale_workers",
|
||||
"historical_workers",
|
||||
"findings",
|
||||
"live_fleet_safe",
|
||||
):
|
||||
self.assertIn(key, snapshot)
|
||||
self.assertTrue(snapshot["registry_revision"].startswith("fleetrev-"))
|
||||
self.assertTrue(snapshot["read_only"])
|
||||
self.assertFalse(snapshot.get("mutation_performed", False))
|
||||
|
||||
|
||||
class CapabilityExposureTests(unittest.TestCase):
|
||||
"""AC: only controller/reconciler reach the surface; author policy unchanged."""
|
||||
|
||||
def test_capability_map_declares_controller_role(self):
|
||||
import task_capability_map as tcm
|
||||
|
||||
for task in (
|
||||
"plan_stale_worker_retirement",
|
||||
"gitea_plan_stale_worker_retirement",
|
||||
"apply_stale_worker_retirement",
|
||||
"gitea_apply_stale_worker_retirement",
|
||||
):
|
||||
with self.subTest(task=task):
|
||||
self.assertEqual(tcm.required_permission(task), "gitea.read")
|
||||
self.assertEqual(tcm.required_role(task), "controller")
|
||||
|
||||
def test_ordinary_author_capability_resolution_is_unaffected(self):
|
||||
import task_capability_map as tcm
|
||||
|
||||
self.assertEqual(tcm.required_permission("work_issue"), "gitea.pr.create")
|
||||
self.assertEqual(tcm.required_role("work_issue"), "author")
|
||||
self.assertEqual(tcm.required_permission("create_pr"), "gitea.pr.create")
|
||||
self.assertEqual(tcm.required_role("create_pr"), "author")
|
||||
self.assertEqual(tcm.required_permission("commit_files"), "gitea.repo.commit")
|
||||
self.assertEqual(tcm.required_permission("push_branch"), "gitea.branch.push")
|
||||
self.assertEqual(
|
||||
tcm.required_permission("snapshot_instance_fleet"), "gitea.read"
|
||||
)
|
||||
|
||||
def test_no_new_gitea_write_permission_is_introduced(self):
|
||||
import task_capability_map as tcm
|
||||
|
||||
for task in ("plan_stale_worker_retirement", "apply_stale_worker_retirement"):
|
||||
with self.subTest(task=task):
|
||||
self.assertNotIn(
|
||||
tcm.required_permission(task),
|
||||
{
|
||||
"gitea.pr.approve",
|
||||
"gitea.pr.merge",
|
||||
"gitea.pr.review",
|
||||
"gitea.branch.push",
|
||||
"gitea.repo.commit",
|
||||
},
|
||||
)
|
||||
|
||||
|
||||
class SurfaceRegistrationTests(unittest.TestCase):
|
||||
"""AC: the dry-run and apply modes are exposed as sanctioned native tools."""
|
||||
|
||||
def test_tools_are_registered_and_role_restricted(self):
|
||||
import inspect
|
||||
|
||||
import gitea_mcp_server as server
|
||||
|
||||
for name in (
|
||||
"gitea_plan_stale_worker_retirement",
|
||||
"gitea_apply_stale_worker_retirement",
|
||||
):
|
||||
with self.subTest(tool=name):
|
||||
tool = getattr(server, name)
|
||||
self.assertTrue(callable(tool))
|
||||
source = inspect.getsource(tool)
|
||||
# Both tools route their role check through the shared guard.
|
||||
self.assertIn(f'_retirement_role_block("{name}")', source)
|
||||
|
||||
guard = inspect.getsource(server._retirement_role_block)
|
||||
self.assertIn('{"controller", "reconciler"}', guard)
|
||||
self.assertIn("required_roles", guard)
|
||||
|
||||
apply_source = inspect.getsource(server.gitea_apply_stale_worker_retirement)
|
||||
self.assertIn("registry_fingerprint", apply_source)
|
||||
self.assertIn("candidate_fingerprint", apply_source)
|
||||
self.assertIn("_retirement_runtime_block", apply_source)
|
||||
self.assertIn("classify_cohort", apply_source)
|
||||
# Revalidation inside the transaction re-reads lease protection rather
|
||||
# than reusing the pre-transaction snapshot.
|
||||
self.assertIn("_retirement_revalidation_plan", apply_source)
|
||||
revalidation = inspect.getsource(server._retirement_revalidation_plan)
|
||||
self.assertIn("_retirement_protected_owners", revalidation)
|
||||
self.assertIn("raise RuntimeError", revalidation)
|
||||
|
||||
def test_plan_tool_declares_itself_read_only(self):
|
||||
import inspect
|
||||
|
||||
import gitea_mcp_server as server
|
||||
|
||||
source = inspect.getsource(server.gitea_plan_stale_worker_retirement)
|
||||
self.assertIn("read_only", source)
|
||||
self.assertNotIn("retire_stale_workers", source)
|
||||
|
||||
def test_docs_exist(self):
|
||||
root = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
|
||||
path = os.path.join(root, "docs", "stale-worker-retirement.md")
|
||||
self.assertTrue(os.path.isfile(path))
|
||||
with open(path, encoding="utf-8") as handle:
|
||||
text = handle.read()
|
||||
for token in (
|
||||
"registry_fingerprint",
|
||||
"candidate_fingerprint",
|
||||
"gitea_plan_stale_worker_retirement",
|
||||
"gitea_apply_stale_worker_retirement",
|
||||
"eligible_stale_orphan",
|
||||
"registry_revision_moved",
|
||||
"does not repair untrusted live identity",
|
||||
):
|
||||
with self.subTest(token=token):
|
||||
self.assertIn(token, text)
|
||||
|
||||
|
||||
if __name__ == "__main__": # pragma: no cover
|
||||
unittest.main()
|
||||
Reference in New Issue
Block a user