Compare commits

Author SHA1 Message Date
jcwalker3 6da68fffb8 Merge branch 'master' into feat/issue-643-request-preview-initiate 2026-07-25 02:44:58 -05:00
sysadmin 8598537a35 Merge pull request 'feat: enforce MCP restart class permissions' (#886) from feat/issue-663-restart-classes into master 2026-07-25 02:34:24 -05:00
sysadminandClaude Opus 5 53ce1b1a5e fix(webui): compensate stray allocations and keep preview side-effect free
Addresses both blockers from the PR #902 review (review 589) for issue #643.

B1 — an allocator-created assignment could be orphaned and reported as no
mutation.

apply_request re-previews, then re-runs the allocator with apply=True. The CAS
fingerprint hashes only {kind, number} plus exclusions, so a competing lease
taken on the requested unit inside the window leaves the fingerprint identical:
the pin passes, the selection loop skips the now-claimed unit and commits an
assignment on the *next* one, and _selection_matches then fails on egress. The
old code returned mutation_performed False with that lease still committed and
owned by a synthetic session nothing heartbeats. The window contains a second
full load_queue_snapshot(), so it is seconds wide, and foreign sessions acting
on this repo concurrently are an observed condition.

The egress mismatch now releases the assignment the allocator created before
refusing. When the release succeeds the refusal reports mutation_performed
False and a compensation record; when it fails the response carries
mutation_performed True, an explicit orphaned_assignment, and a
gitea_release_workflow_lease reclaim action, because claiming nothing changed
while a lease is live is the defect rather than a report of it.

The same state was reachable through _run_allocator's bare except Exception:
allocate_next_work only catches InvalidWorkKindError, LeaseRequiredError and
ControlPlaneError, so anything raised after assign_and_lease committed arrived
as "no result" with a durable lease. A None result on the apply path now sweeps
and releases whatever this flow's session owns. That sweep is only possible
because of the B2 fix below — the session id is now stable across the flow, so
the lease is findable.

B2 — the "read-only" preview wrote to the control-plane DB.

allocate_next_work called db.upsert_session and db.expire_stale_leases
unconditionally, before the apply branch was consulted, and default_allocator
minted a fresh webui-request-<hex> per call. Every preview therefore appended a
never-reused session row and mutated global lease state while the payload said
dry_run True / mutation_performed False, driven by an operator refreshing a
form. One apply wrote two rows and bound the lease to the second, which is why
B1's orphan had no reclaimable owner.

allocate_next_work gains a keyword-only side_effect_free flag, default False so
every existing caller is byte-for-byte unchanged. Under the flag both writes are
suppressed and expired leases are instead filtered out of the claim map in
memory, which reaches the same selection the sweep would have produced without
persisting anything; a claim whose expiry cannot be parsed is kept, since an
unreadable expiry is not evidence that work is free. side_effect_free with
apply=True fails closed rather than silently reserving. request_service mints
one session id per request flow and threads it through both the dry-run and the
apply, and the dry-run now routes through the side-effect-free path.

Coverage.

test_allocator_drift_on_apply_is_not_read_as_an_assignment asserted the defect —
it built the orphan state and then required mutation_performed to be False, which
a leak satisfies. It now requires the compensating release. Added: release-failure
surfacing a reclaim action, an assignment with no lease id, the post-commit
exception route, and session-id identity across the flow. The two areas the review
named as having zero coverage now have it: default_allocator past its two
fail-closed early returns (side_effect_free routing, session-id pass-through and
minting, scope and fingerprint propagation) and default_claims_source (scoped
read, and an unreadable substrate denying rather than reading as "nothing
claimed"). allocator_service gains side-effect-free tests against a real
temp-file DB plus the expiry-filter unit tests.

Every new guard was mutation-tested by reverting it one at a time: B1
compensation removed → 3 failures; post-commit sweep removed → 1; session id
re-minted per call → 2; side_effect_free ignored → 3; in-memory expiry filter
removed → 1.

Verification: WEBUI_TEST_OFFLINE=1 python -m pytest tests/ -q from this
branches/ worktree gives 23 failed, 5263 passed, 6 skipped, 899 subtests. The
sorted FAILED set is identical to the reviewer's clean-master baseline at
2f4dec83 (23 failed, 5190 passed) — no new, changed, or disappeared failure —
and 21 tests were added over the reviewed head's 5242. Zero conflict markers;
py_compile passes; git diff --check clean.

Refs #643, PR #902

Co-Authored-By: Claude Opus 5 (1M context) <[email protected]>
Claude-Session: https://claude.ai/code/session_01V6xFqovhbArPv61j9KCGkL
2026-07-25 03:32:24 -04:00
sysadminandClaude Opus 5 220361ad94 fix(restart): require both authorizations for apply, correct coordinator doc
Addresses the two blockers raised in the PR #886 review (comment 16559) for
issue #663.

B1 — apply_authorized ignored restart-class authorization.

The #663 restart-class matrix and the #661 drain-proof hard gate are
independent authorizations that first coexisted when PR #882 landed on
master and this branch merged it. The union preserved both, but the apply
decision consulted only the drain gate:

    payload["apply_authorized"] = gate.allow

so a clean drain proof — or an authorized break-glass, which needs no proof
at all — reported apply_authorized: True for a class the least-privilege
matrix had just denied, in the same payload carrying allow_restart: False
and "role 'author' may not request full_mcp_restart". One environment
variable therefore collapsed the whole nine-class matrix for the apply
decision, including host_restart.

The apply decision is now the conjunction of both authorizations, and
apply_gate carries drain_gate_allow and restart_class_authorized so a denial
is attributable to the authorization that produced it. Break-glass keeps its
purpose — bypassing the drain proof — and never bypasses the class matrix.
No existing fail-closed behaviour is weakened: allow_restart, drain-proof
verification, fingerprint binding, and requester authorization are untouched.

B2 — docs/mcp-restart-coordinator.md described pre-#661 behaviour.

The document still called the drain proof "a separate child" and omitted
drain_proof_json and request_break_glass from the published signature, so a
safety document asserted there was no gate where a gate now exists. It now
documents both parameters, states that the gate executes inside this tool,
and records dry-run versus apply behaviour, authorization ordering, the
break-glass scope, and fail-closed conditions as implemented.

Regression coverage.

tests/test_issue_886_apply_authorization_conjunction.py exercises the MCP
tool itself, which previously had no test at all — that absence is why the
defect shipped. It pins both conjunction directions, proves a clean proof
cannot override a role, approval, unknown-class, or missing-target denial,
proves break-glass does not collapse the matrix for any worker role or
restricted class, and proves the existing scoped and unscoped paths and the
#661 denials still hold. Against the pre-fix tree 24 of these fail; against
this commit all 19 pass with 45 subtests.

tests/test_mcp_restart_governance_docs.py now binds the published signature
to inspect.signature() of the real tool and forbids the stale pre-#661
phrasing, so the drift that produced B2 cannot return unnoticed.

Verification: targeted restart/drain/governance/webui suites 194 passed,
113 subtests. Full suite 23 failed, 5230 passed, 6 skipped, 912 subtests —
the failure set is identical to the reviewed baseline at 9bc021e
(23 failed, 5201 passed), with +29 passing from the added tests and no new
or changed failure. Zero conflict markers; py_compile passes; the #882
union remains intact in both directions.

Refs #663, PR #886

Co-Authored-By: Claude Opus 5 (1M context) <[email protected]>
Claude-Session: https://claude.ai/code/session_01V6xFqovhbArPv61j9KCGkL
2026-07-25 02:36:41 -04:00
sysadminandClaude Opus 4.8 433f66add8 feat(webui): request preview, authorization, and workflow initiation (Closes #643)
Operators had to paste a role prompt into a terminal to start work, and
nothing enforced that the allocator had been consulted first, so two sessions
could reach for the same issue and each believe it was theirs. This adds a
request surface: a desired role, an issue or PR, and a stated intent, answered
by an authorization decision and - on confirmation - an exclusive assignment
from the allocator.

Preview (POST /api/v1/requests/preview, and the /requests form) runs five
checks and reports authorize/deny with a reason for each: console
authorization, capability resolution for the desired role, lease availability,
whether the allocator would independently select this work unit, and head
pinning for PR work. It is read-only - it calls the allocator with apply=false
and writes only an audit line. An unauthorized principal never reaches the
allocator or the control-plane DB, so a denial cannot enumerate the queue.

Initiation (POST /api/v1/requests/apply) never assigns the requested item
directly. It runs a dry-run first and proceeds only when the allocator would
independently pick that exact work unit, carrying the dry-run's
candidate_set_fingerprint as a CAS pin; otherwise it returns wait or blocked
and mutates nothing. An active claim on the work unit rejects a duplicate
assign before one is attempted. A returned assignment carries a handoff block
naming the required profile, namespace, and the actions that stay forbidden.

Authorization reuses the #633 model rather than adding a second one. The new
initiate_workflow action is operator-class because its outcome is a claim, not
a Gitea verdict: requesting reviewer or merger work reserves that work but
grants no right to approve or merge. Execution is gated by a new per-action
execution_env_flag (WEBUI_REQUESTS_EXECUTION), deliberately in place of raising
ACTIVE_PHASE - a phase bump would enable execution for every phase-2 action at
once, including ones whose execution path is not implemented. Actions that
declare no flag are unchanged and still report execution_enabled false.

Every preview and apply emits a console audit record correlated to the
resulting assignment by correlation.request_id.

Fail-closed throughout: an unreadable control-plane DB, an incomplete queue
inventory (#758), an allocator that raises, an unpinned PR head, a moved PR
head, and an unconfirmed apply all deny without mutating.

Files:
- webui/request_service.py (new) - request model, preview, initiation
- webui/request_views.py (new) - form and preview rendering, escaped
- tests/test_webui_request_initiation.py (new) - 52 tests
- webui/console_authz.py - initiate_workflow action, execution_wired()
- webui/app.py - /requests, /api/v1/requests/preview, /api/v1/requests/apply
- webui/nav.py - Requests nav entry
- webui/traffic_loader.py - public candidates_from_queue_snapshot alias
- docs/webui-requests.md (new), docs/webui-authz-audit.md

Validation: full suite on this branch 5242 passed, 6 skipped, 899 subtests, 23
failed. Clean master baseline at 2f4dec83 in an equivalent branches/ worktree:
5190 passed, 6 skipped, 867 subtests, the same 23 tests failed. The branch adds
52 passing tests and introduces no new full-suite failure signature.

Closes #643

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-25 01:47:38 -04:00
sysadminandClaude Opus 4.8 9bc021e9c0 Merge master into feat/issue-663-restart-classes (resolve #886 conflict)
Brings PR #886 up to date with master @ 2f4dec8323
(8 commits behind), resolving the single conflicted file.

Conflict: gitea_mcp_server.py, both hunks inside gitea_request_mcp_restart.
Both sides were purely additive to the same tool, so both are kept in full:

- Branch side (#663, restart classes): parameters restart_class,
  target_session_id, target_role, target_connector; payload keys
  controller_approval_authorized, requester_role, requester_permissions.
- Master side (#661 via PR #882, drain-proof hard gate): parameters
  drain_proof_json, request_break_glass; the explanatory comment describing
  the apply-path hard gate and break-glass authorization.

No behaviour from either side was dropped, reordered, or reimplemented. Every
parameter from both sides is already consumed by the auto-merged function body
(restart_class and the three target_* arguments flow into the coordinator call;
drain_proof_json and request_break_glass drive the dry_run=False hard gate), so
the union is the only resolution that keeps the merged function coherent.

Validation on the merged tree:

  python -m pytest tests/test_drain_proof.py tests/test_restart_classes.py \
    tests/test_restart_coordinator.py tests/test_mcp_restart_paths.py \
    tests/test_mcp_restart_governance_docs.py tests/test_webui_sanctioned_restart.py \
    tests/test_issue_662_post_restart_reconcile.py -q
  # 165 passed, 68 subtests passed

py_compile on gitea_mcp_server.py passes and no conflict markers remain.

Closes #663

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-25 01:16:09 -04:00
sysadmin 2f4dec8323 Merge pull request 'feat(restart): pre-restart drain proof and hard gate (Closes #661)' (#882) from feat/issue-661-drain-proof-hard-gate into master 2026-07-24 23:40:02 -05:00
sysadmin 3a9d634c17 fix(drain-proof): bind acknowledgement coverage to session identity (#661)
Review 582 (REQUEST_CHANGES at 95178349) found a residual fail-open of the
same class the PR set out to close. Acknowledgement coverage was decided by
comparing a count against a count:

    covers_live_sessions = live_count_known and acked_count >= sessions_live_other

Nothing bound an acknowledgement to the identity of a session that actually
owed one, so acknowledgements supplied for the requesting session and for a
session that does not exist satisfied the obligations of two live sessions
that never answered - minting a clean, correctly signed proof and an allow
verdict from the restart gate.

Coverage is now derived from authoritative impact-report evidence:

- New `_required_ack_sessions()` derives the required session ids from the
  report itself, via `ack_state` keys and/or `affected_sessions` filtered on
  `live and not is_requester`. The requester is excluded only on explicit
  `is_requester` evidence, never inferred.
- When both views are present they must name the same set, and the result is
  reconciled against `counts.sessions_live_other`. Missing, malformed,
  duplicated, contradictory, or unreconcilable identity evidence fails closed
  and outranks every permitting path, including the timeout policy.
- Coverage requires every required id to carry an explicit acknowledgement
  token keyed by that id. Acknowledgements for the requester, for unknown
  ids, or for fabricated ids never increase coverage.
- Caller-supplied acknowledgement cardinality is no longer proof of anything.

Failure propagates unchanged through `acks_or_timeout` -> `proof.clean` ->
`failed_checks` -> `gate_apply_restart` verdict `deny` / `allow=False`.

The earlier missing-acknowledgement remediation is preserved in full: absent,
None, non-mapping, empty, partial, stale, and unparseable acks still fail
closed, `ack_timeout_policy_applied` stays strict `value is True`, and the
legitimate zero-live-sessions and explicit-timeout paths still pass.

Reviewer's reproduction, before and after this commit:

    sessions_live_other = 2
    report ack_state    = {'other-0': 'pending', 'other-1': 'pending'}
    supplied acks       = {'req': 'ack', 'totally-bogus-session': 'ack'}

    before: acks_or_timeout = True  | proof.clean = True  | gate allow
    after:  acks_or_timeout = False | proof.clean = False | gate deny

Tests: 22 new cases in `AcknowledgementIdentityBindingTests` covering the
reviewer's exact exploit, wrong-ids-with-sufficient-count, partial identity
match, requester-only acks, fabricated ids, unproven per-session states,
missing/malformed/contradictory identity evidence, count mismatch, and the
preserved success paths.

Verification:
- `pytest tests/test_drain_proof.py` -> 61 passed, 56 subtests
  (baseline at 95178349: 39 passed, 26 subtests)
- Restart surface (6 modules) -> 154 passed, 68 subtests, exit 0
  (baseline at 95178349: 132 passed, 38 subtests)
- Full `pytest tests/` -> 5177 passed vs baseline 5155 passed; the 23
  failures are identical in both runs and pre-exist at 95178349.

Scope: drain_proof.py, tests/test_drain_proof.py.

Co-Authored-By: Claude Opus 5 (1M context) <[email protected]>
Claude-Session: https://claude.ai/code/session_01VRUZAf3Fr5n3kqhhiayN6C
(cherry picked from commit 4193b63f415b066ee292386c2c89bc3d2651a0cc)
2026-07-25 00:04:17 -04:00
jcwalker3 930dc24632 Merge branch 'master' into feat/issue-663-restart-classes 2026-07-24 22:34:33 -05:00
jcwalker3 2068bae341 Merge branch 'master' into feat/issue-661-drain-proof-hard-gate 2026-07-24 22:34:27 -05:00
sysadmin 7af40fb5ff Merge pull request 'fix(allocator): exclude vision/roadmap/umbrella coordination containers (Closes #854)' (#883) from fix/issue-854-semantic-container-exclusion into master 2026-07-24 22:27:58 -05:00
sysadmin 9517834913 Merge commit '578c44b685a7ff5b01006c5e398bfac9863e0d8d' into feat/issue-661-drain-proof-hard-gate 2026-07-24 22:37:49 -04:00
sysadminandClaude Opus 5 824c42f7e3 fix(drain): fail closed on missing or unproven acknowledgement evidence (#661)
The acks_or_timeout check treated an absent `acks` key as proof that no
session needed to acknowledge: `drain_state.get("acks") or {}` collapsed
absent, None, and empty into the same value, and the resulting empty mapping
satisfied `no_sessions_to_ack`. The impact report's counts.sessions_live_other
was never consulted, so absence of evidence was read as evidence of absence.

Reproduced at head 1cbbde0089: with
sessions_live_other = 3 and the acknowledgement key absent, acks_or_timeout
passed with detail "no other live sessions required to acknowledge", the proof
minted clean, and gate_apply_restart returned verdict allow — a restart
authorized against three live sessions with zero acknowledgement evidence, and
the resulting artifact carried a valid signature.

Whether acknowledgement is required is now derived from the impact report,
never from the shape of the drain state:

- _live_session_count() reads counts.sessions_live_other and returns None for a
  missing, malformed, negative, or bool value, so an unreadable report fails
  closed instead of reading as "nobody was live".
- Absent, None, non-mapping, empty, partially-covering, and unparseable or
  stale acknowledgement data all fail closed while live sessions require
  acknowledgement.
- _is_acknowledged() no longer coerces with str(); only an explicit
  "ack"/"acked"/"acknowledged" string counts, so None, timestamps, and
  "pending"/"stale" markers are never read as an acknowledgement.
- Present-but-unacknowledged entries fail closed even when the report claims
  zero live sessions: that contradiction is not safe to resolve in favour of
  the restart.
- ack_timeout_policy_applied stays strict (`value is True`), so an absent, null,
  or non-boolean value cannot open the gate on its own.

The genuine no-other-live-sessions case still passes, now justified by the
report proving sessions_live_other == 0 rather than by the absence of data.

Adds AcknowledgementFailClosedTests: 14 cases / 26 subtests covering missing,
null, empty, malformed, stale, partial-coverage, and unproven-count inputs,
the valid-acknowledgement and zero-live-session paths, timeout-policy
strictness, and that a failed check blocks proof.clean, verification, and the
restart gate.

Restart-surface suite: 132 passed, 38 subtests (branch baseline 118 passed,
12 subtests; +14 new tests, no regressions).

Co-Authored-By: Claude Opus 5 (1M context) <[email protected]>
Claude-Session: https://claude.ai/code/session_01VEaP3TohHLFWkp3Z2mmuZw
2026-07-24 22:36:25 -04:00
jcwalker3 578c44b685 Merge branch 'master' into feat/issue-661-drain-proof-hard-gate 2026-07-24 21:28:04 -05:00
jcwalker3 41622c5985 Merge branch 'master' into feat/issue-663-restart-classes 2026-07-24 21:27:15 -05:00
jcwalker3 9f686253eb Merge branch 'master' into feat/issue-661-drain-proof-hard-gate 2026-07-24 21:06:49 -05:00
jcwalker3 301c78de20 Merge branch 'master' into feat/issue-663-restart-classes 2026-07-24 21:06:21 -05:00
sysadmin 714190e02a feat: enforce MCP restart class permissions (#663) 2026-07-24 18:10:22 -04:00
sysadmin 1cbbde0089 feat(restart): pre-restart drain proof and hard gate (#661)
Add `drain_proof.py`: a machine-verifiable DrainProof artifact plus a
fail-closed verifier and the hard gate the sanctioned restart-apply path
must consult, so a restart can never proceed on a stale or false "ready"
claim (#655 umbrella, child of #658 coordinator / #659 drain / #660
checkpoints).

- DrainProof: HMAC-SHA256 keyed proof-id over a canonical body using a
  per-process secret -> non-forgeable within the process; a proof minted in
  a prior daemon process will not verify after restart. Short TTL (120s).
- build_drain_proof(): mints the proof from the #658 impact report + the
  drain-mode outcomes. Checklist: no in-flight mutations, assignments
  stopped, checkpoints complete, handoffs ok, leases handled, acks-or-
  timeout. Every check fails closed on missing/ambiguous evidence; the
  no-in-flight-mutations and leases-handled checks are derived from the
  authoritative impact report, not self-reported.
- verify_drain_proof(): fail-closed — rejects missing, malformed, expired,
  signature-mismatched (forged/tampered/prior-process), unclean, or
  stale-fingerprint proofs; recomputes cleanliness from the checks rather
  than trusting the flag.
- gate_apply_restart(): allow only on a valid clean proof; deny -> durable
  incident descriptor; break-glass is the only bypass and is never silent.
- Checkpoint completeness is a supplied input, not a hard dependency on the
  (still-unmerged #660) checkpoint schema.

Wire the gate into gitea_request_mcp_restart: dry_run=False now enforces the
hard gate (drain_proof_json required; break-glass via request_break_glass +
GITEA_BREAKGLASS_RESTART_AUTHORIZATION env). The tool still performs no
actual restart — execution remains a further child.

Tests: tests/test_drain_proof.py — 25 cases covering AC#1-4 (apply without
proof denied, successful drain verifiable, open unsafe mutation fails,
pass/fail/expired), forgery/tamper/wrong-secret/stale-fingerprint rejection,
break-glass bypass, and secret hygiene. 25/25 pass (coordinator suite
unaffected: 40/40 together).

Links #652 #653 #655 #658 #659 #660.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
Claude-Session: https://claude.ai/code/session_01E7Fv9Bp2XWgvaWa4M1kdR7
(cherry picked from commit e7bcc952bb3e820fda95acbecefeaebfa5f8fcff)
2026-07-24 17:17:13 -04:00
22 changed files with 6540 additions and 75 deletions
+74
View File
@@ -23,6 +23,7 @@ import json
import os
import uuid
from dataclasses import dataclass, field
from datetime import datetime, timezone
from typing import Any, Mapping, Sequence
from control_plane_db import (
@@ -738,6 +739,46 @@ def normalize_exclude_issue_numbers(
return sorted(out)
def _claim_expires_at(claim: Any) -> datetime | None:
"""Parse a claim's ``expires_at``, or ``None`` when it is absent/malformed."""
if not isinstance(claim, Mapping):
return None
text = str(claim.get("expires_at") or "").strip()
if not text:
return None
if text.endswith("Z"):
text = text[:-1] + "+00:00"
try:
parsed = datetime.fromisoformat(text)
except ValueError:
return None
if parsed.tzinfo is None:
parsed = parsed.replace(tzinfo=timezone.utc)
return parsed.astimezone(timezone.utc)
def _drop_expired_claims(
claims: Mapping[tuple[str, int], dict[str, Any]],
*,
now: datetime | None = None,
) -> dict[tuple[str, int], dict[str, Any]]:
"""Claims minus those whose lease has already expired (#643).
The read-only mirror of ``expire_stale_leases``: the sweep marks such rows
``expired`` so they stop being returned as claims, and this reaches the same
view without writing. A claim with no parseable ``expires_at`` is **kept** —
an unreadable expiry is not evidence that work is free.
"""
moment = now or datetime.now(timezone.utc)
kept: dict[tuple[str, int], dict[str, Any]] = {}
for key, claim in (claims or {}).items():
expires_at = _claim_expires_at(claim)
if expires_at is not None and expires_at <= moment:
continue
kept[key] = claim
return kept
def candidate_set_fingerprint(
candidates: Sequence[WorkCandidate],
*,
@@ -826,12 +867,22 @@ def allocate_next_work(
exclude_issue_numbers: Sequence[int] | None = None,
expected_candidate_set_fingerprint: str | None = None,
allocation_mode: str | None = None,
side_effect_free: bool = False,
) -> dict[str, Any]:
"""Select and optionally reserve the next work unit via control-plane DB.
*apply=False* (default): dry-run selection only — no lease/assignment.
*apply=True*: atomic ``assign_and_lease`` for the selected candidate.
*side_effect_free* (#643): a dry run that writes **nothing** to the
control-plane DB. A plain ``apply=False`` still registered a session row and
swept stale leases globally, so a caller advertising a read-only preview was
mutating on every call. Under this flag both writes are suppressed and stale
leases are instead filtered out of the claim map in memory, which yields the
same selection the sweep would have produced without persisting anything.
Incompatible with *apply* — the combination fails closed rather than
silently reserving.
*allocation_mode* (#840): ``cross_role`` (default for controller) inspects
the complete queue and returns one authoritative selection naming the
required downstream role/profile/action. ``role_scoped`` keeps prior
@@ -885,7 +936,24 @@ def allocate_next_work(
"allocation_mode": (allocation_mode or "").strip() or None,
}
# A side-effect-free run may never reserve: reserving is a write, and the
# flag is the caller's assertion that this call writes nothing (#643).
if side_effect_free and apply:
return {
"success": False,
"outcome": OUTCOME_NO_SAFE,
"apply": True,
"reasons": [
"side_effect_free is incompatible with apply=True; an "
"assignment is a write (fail closed, #643)"
],
"skipped": [],
"assignment": None,
"substrate": "control_plane_db",
}
session_id = (session_id or "").strip() or f"alloc-{uuid.uuid4().hex[:12]}"
if not side_effect_free:
try:
db.upsert_session(
session_id=session_id,
@@ -953,6 +1021,12 @@ def allocate_next_work(
"assignment": None,
"substrate": "control_plane_db",
}
if side_effect_free:
# ``list_active_claims`` filters on status alone, so without the
# global sweep an already-expired lease would still read as a live
# claim and the preview would report work as taken that is free.
# Drop those in memory: same view the sweep produces, no write.
claims = _drop_expired_claims(claims)
try:
exclude_nums = normalize_exclude_issue_numbers(exclude_issue_numbers)
+53
View File
@@ -0,0 +1,53 @@
# MCP restart classes and blast-radius permissions (#663)
This is the machine-enforced class matrix used by
`restart_coordinator.RESTART_CLASS_POLICIES`. It implements the narrower-first
recovery ladder from #655 and the authorization policy from #656, using the
path inventory from #657 and the impact coordinator from #658. Product and
delivery lineage: vision #652 and roadmap #653.
Unknown class names are denied. The coordinator requires both the class
permission and an eligible request role. Approval gates are additional: a
caller cannot turn a request permission into execution authority.
| Restart class | Required permission | Expected blast radius | Drain requirement | Approval requirement | Audit requirement | Recovery behavior |
|---|---|---|---|---|---|---|
| `client_reconnect` | `mcp.reconnect.client` | none | none | self service | class, actor, client namespace, reason, outcome | Reconnect only the caller's client transport. No daemon or peer work changes. |
| `session_reconnect` | `mcp.reconnect.session` | low | requesting-session safe point | self service | class, actor, session, reason, outcome | Rebind identity, capability, and workspace state for one session. |
| `worker_restart` | `mcp.restart.worker.request` | low | target worker | controller approval + automated gates | class, actor, worker, approval, scoped drain, outcome | Restart one worker after its own leases and mutations drain. |
| `role_runtime_restart` | `mcp.restart.role_runtime.request` | medium | target role runtime | controller approval + automated gates | class, actor, role namespace, approval, scoped drain, outcome | Restart and re-probe one role runtime; unrelated roles remain available. |
| `connector_restart` | `mcp.restart.connector.request` | medium | target connector | controller approval + automated gates | class, actor, connector, approval, scoped drain, outcome | Restart one connector while unrelated runtimes remain available. |
| `configuration_reload` | `mcp.reload.configuration.request` | low | mutation quiesce | controller approval + automated gates | class, actor, configuration revision, approval, outcome | Gracefully reload configuration without replacing the daemon. |
| `rolling_mcp_restart` | `mcp.restart.rolling.request` | medium | one instance at a time | controller approval + automated gates | class, actor, instance order, approval, per-instance drains, outcome | Drain, restart, verify, and restore each instance before advancing. |
| `full_mcp_restart` | `mcp.restart.full.request` | high | all sessions and mutations | controller approval + automated gates | class, actor, full impact, approval, full drain proof, outcome | Replace the complete MCP runtime only after a verified full drain. |
| `host_restart` | `mcp.restart.host.request` | high | all host work | controller approval + infrastructure operator | class, actor, host/change or incident id, approval, full drain proof, outcome | Hand off to infrastructure ownership and reconcile every runtime afterward. |
## Drain boundary
Only `full_mcp_restart` and `host_restart` set `full_drain_required=true`.
Reconnects and configuration reloads do not disrupt peer sessions. Worker,
role-runtime, and connector restarts evaluate only their explicitly named
target. Rolling restart drains one instance at a time. Missing required target
scope denies the request rather than silently widening it to a full restart.
## Permission and approval boundary
Author, reviewer, merger, and reconciler roles may self-request reconnects and
request scoped worker/role/connector/reload recovery. They cannot request
rolling, full, or host restart classes. Controller/operator/admin roles may
request the broader classes, while execution remains operator/admin-owned.
Controller approval is independently required for every class above a session
reconnect. Host restart additionally requires infrastructure-operator proof.
The MCP request tool derives class permissions from its authenticated runtime
role. It does not accept caller-supplied permissions. Controller and operator
authorization are read from the already-running daemon environment, never
from a request argument.
## Audit and failure behavior
Every impact audit and every console restart/reload audit includes a
`restart_class` field. The impact audit also includes the exact
`required_permission`. Unknown classes, missing permissions, ineligible roles,
missing approval, missing scoped targets, and incomplete inventory all deny
fail closed. Manual process kills remain forbidden and contaminating (#630).
+67 -10
View File
@@ -6,11 +6,23 @@ console (#642 / #652) can see the blast radius *before* concurrent LLM work is
disrupted. Uncoordinated restarts destroy in-flight author/reviewer/merger work
and give operators no way to see what they are about to break.
This lands the coordinator + impact DTO + a dry-run MCP tool. It is the single
This lands the coordinator + impact DTO + the MCP tool. It is the single
sanctioned entry point for restart evaluation post-#657 (which inventoried the
restart/reload/kill paths). The **mutative apply** path — actually performing a
restart — is a later child gated by a drain proof and is explicitly out of
scope here.
restart/reload/kill paths).
The **drain-proof hard gate now executes inside this tool** (#661, via PR #882):
an apply request (`dry_run=False`) is evaluated against a drain proof here and
denied when that proof is missing, expired, unclean, tampered with, or stale.
It is no longer a separate child operation. What remains a later child is only
the **execution** step — actually stopping and restoring a process. This tool
still never restarts anything: `apply_supported` is always `false` and
`restart_performed` is always `false`.
The coordinator now routes every request through the restart-class policy
matrix defined for #663. See
[`mcp-restart-classes.md`](./mcp-restart-classes.md) for permissions, expected
blast radius, scoped drain and approval requirements, audit fields, and
recovery behavior for all nine classes.
## Components
@@ -19,7 +31,8 @@ scope here.
| `restart_coordinator.evaluate_restart_impact` | `restart_coordinator.py` | Pure classification: inventory → impact report DTO. No I/O, no restart. |
| `RestartImpactReport` / `SessionImpact` / `LeaseImpact` | `restart_coordinator.py` | Console-facing DTO (`.as_dict()` is JSON-serializable). |
| `ControlPlaneDB.list_sessions` | `control_plane_db.py` | Read-only session inventory (the process-level unit a restart kills). |
| `gitea_request_mcp_restart` | `gitea_mcp_server.py` | MCP tool: gathers inventory from the #613 DB, calls the coordinator, returns the report. Dry-run only. |
| `gitea_request_mcp_restart` | `gitea_mcp_server.py` | MCP tool: gathers inventory from the #613 DB, calls the coordinator, returns the report, and on `dry_run=False` runs the #661 drain-proof hard gate. Never restarts a process. |
| `drain_proof.gate_apply_restart` | `drain_proof.py` | The #661 hard gate: verifies a drain proof against the current impact fingerprint, or records an authorized break-glass bypass. |
## Dimensions evaluated
@@ -77,17 +90,61 @@ authorization is present.
```text
gitea_request_mcp_restart(remote, host, org, repo,
dry_run=True, request_override=False,
session_id=None, limit=200)
session_id=None, limit=200,
restart_class="full_mcp_restart",
target_session_id=None, target_role=None,
target_connector=None,
drain_proof_json=None,
request_break_glass=False)
```
Read-only, dry-run, and it **never restarts anything**. `apply_supported` is
always `false`; passing `dry_run=False` performs no restart and reports that
apply is gated by a drain proof (a separate child).
It **never restarts anything**: `apply_supported` is always `false` and
`restart_performed` is always `false`.
### Dry-run versus apply
| Call | Behavior |
|------|----------|
| `dry_run=True` (default) | Read-only impact preview. No drain proof is required or consulted. |
| `dry_run=False` | The #661 drain-proof hard gate runs **in this tool**. The outcome is reported under `apply_gate` / `apply_authorized`; a denial also returns a durable `incident` descriptor. Still no restart. |
### Authorization ordering
An apply requires **both** authorizations, and they are independent:
1. **Restart-class authorization** (#663) — the requester's role and permissions
must allow the requested class, the class's approval requirement must be
satisfied, and any target-scoped class must name its target. Failing any of
these makes `allow_restart` `false`.
2. **Drain-proof gate** (#661) — a valid, unexpired, clean proof bound to the
current impact fingerprint, or an authorized break-glass.
`apply_authorized` is the conjunction: `gate.allow and allow_restart`. A clean
drain proof therefore cannot override a class or requester-role denial, and a
denied class never reports an authorized apply. `apply_gate` carries
`drain_gate_allow` and `restart_class_authorized` so a denial is attributable to
the authorization that produced it.
### Break-glass
Break-glass bypasses the **drain proof only** — never the restart-class matrix.
It is honoured solely when `request_break_glass` is set *and* the environment
carries `GITEA_BREAKGLASS_RESTART_AUTHORIZATION`; like operator override, the
tool argument expresses caller intent and cannot be self-asserted by a worker
session. `break_glass_requested` and `break_glass_authorized` are both reported,
so a bypass is never silent.
### Fail closed on apply
A missing, malformed, expired, unclean, tampered, or fingerprint-stale drain
proof denies the apply and returns an `incident` descriptor. An unknown restart
class denies before any of this. Ambiguity always denies.
## Audit
Every evaluation carries an `audit_record` (event, coordinator version, verdict,
allow decision, blast radius, counts, timestamp) so restart decisions are
restart class, required permission, allow decision, blast radius, counts,
timestamp) so restart decisions are
auditable. No secrets flow through the coordinator — session ids, pids, and
profiles are operational metadata only.
+41 -9
View File
@@ -94,6 +94,7 @@ already define, and a regression test asserts each mapping matches.
| `record_analytics_usage` | operator | gated_write | `runtime.record_analytics_usage` | Yes | No | No | 2 |
| `system.reload_namespace` | controller | privileged | `runtime.reload_namespace` | Yes | No | No | 2 |
| `system.restart_namespace` | admin | destructive | `runtime.restart_namespace` | Yes | **Yes** | **Yes** | 2 |
| `initiate_workflow` | operator | gated_write | `gitea.read` | Yes | No | No | 2 |
**Dual control** means the acting principal may not be the sole authority: a
second distinct principal must confirm. **Break-glass** means the action is
@@ -112,6 +113,12 @@ by the console — both hand off to a host supervisor, and neither exposes a raw
process kill. See
[`sanctioned-restart-controls.md`](sanctioned-restart-controls.md) (#642).
`initiate_workflow` (#643) is operator-class because its outcome is a *claim*,
not a Gitea verdict. Requesting reviewer or merger work reserves that work
through the allocator; it does not grant the right to approve or merge, which
stays with the MCP role profile and its own capability gates. See
[`webui-requests.md`](webui-requests.md).
### Authorization decision
`authorize(action_id, principal, for_execution=False)` returns a decision
@@ -126,9 +133,24 @@ record and **denies by default**. The deny reasons are closed and enumerated:
| `phase_not_active` | Execution requested for an action whose phase is not open. |
| `allowed_preview_only` | Authorized — preview only, execution still disabled. |
There is no implicit allow branch. Even the allow result reports
`execution_enabled: false` while the console is in Phase 1, so no caller can
read an allow as permission to mutate.
There is no implicit allow branch.
`execution_enabled` on the decision reports whether the action has a live
execution path at all, and is computed by `execution_wired(action)`. There are
exactly two ways to be wired:
1. the action's `phase` is at or below `ACTIVE_PHASE`; or
2. the action declares an `execution_env_flag` **and** that variable is set.
Every action that declares no flag therefore reports `execution_enabled: false`
while the console is in Phase 1, so no caller can read an allow as permission
to mutate. The per-action flag exists because raising `ACTIVE_PHASE` would
enable execution for every action of that phase at once, including ones whose
execution path is not implemented. One implemented action goes live on its own
flag instead of dragging its unimplemented phase-mates with it.
`initiate_workflow` is the only action that currently declares a flag
(`WEBUI_REQUESTS_EXECUTION`), and it stays denied until an operator sets it.
## Secret redaction
@@ -235,13 +257,22 @@ second one. The integration points are already wired and observable:
instead of adding a parallel check.
- **`GET /api/console/security-model`** publishes the RBAC matrix, redaction
policy, and audit policy as JSON for operators and tests.
- **`POST /api/v1/requests/preview` and `.../apply`** (#643) are the first
actions to use this model for a real execution path. Preview always returns a
decision and an audited `previewed` record; apply requires `confirm=true`,
emits `succeeded` or `denied`, and reserves work only through the allocator.
See [`webui-requests.md`](webui-requests.md).
To open Phase 2, a child issue must: raise `ACTIVE_PHASE`, implement the
confirmation and dual-control flow the matrix already declares, emit a
`succeeded` or `failed` record alongside the `gitea_audit` mutation record, and
keep `viewer` unable to reach any of it. Turning on execution without the
confirmation flow contradicts a declared requirement and is a review failure,
not a shortcut.
A Phase 2 action must: use `execution_wired` rather than a private enable flag,
implement the confirmation and dual-control flow the matrix already declares,
emit a `succeeded` or `failed` record alongside the `gitea_audit` mutation
record, and keep `viewer` unable to reach any of it. Turning on execution
without the confirmation flow contradicts a declared requirement and is a
review failure, not a shortcut.
Raising `ACTIVE_PHASE` remains the way to open a whole phase at once, and is
deliberately *not* what #643 did: an action-scoped opt-in cannot enable an
action whose execution path nobody wrote.
## Local-dev mode
@@ -294,6 +325,7 @@ Until Phase 2 wires it, probe protection rests on network placement alone, as
| `WEBUI_ROLE_MAP` | unset | JSON subject → role map |
| `WEBUI_REQUIRE_PROBE_AUTH` | unset | Require auth for non-public probes |
| `WEBUI_CONSOLE_AUDIT_LOG` | unset | Append-only audit sink path |
| `WEBUI_REQUESTS_EXECUTION` | unset | Opt in to `initiate_workflow` execution (#643) |
All are read server-side only. None is ever rendered into a page or returned by
an API.
+160
View File
@@ -0,0 +1,160 @@
# Web console requests: intent preview and workflow initiation (#643)
**Phase 2. Preview is always live and always read-only. Initiation is wired but
denied until an operator opts in.**
Before this surface, starting role work meant pasting a prompt into a terminal
and trusting the operator to have checked the allocator first. Nothing enforced
that check, so two sessions could reach for the same issue and each believe it
was theirs. This page replaces the paste with a *request*: a desired role, an
issue or PR, and a stated intent, answered by an authorization decision and —
on confirmation — an exclusive assignment from the allocator.
| Concern | Module |
|---------|--------|
| Request model, preview, initiation | `webui/request_service.py` |
| Form and preview rendering | `webui/request_views.py` |
| Authorization | `webui/console_authz.py` (`initiate_workflow`) |
| Audit | `webui/console_audit.py` |
| Ownership substrate | `allocator_service.py` + `control_plane_db.py` |
## Surfaces
| Path | Method | Purpose |
|------|--------|---------|
| `/requests` | GET | Request form |
| `/requests` | POST | Render an intent preview. **Never assigns.** |
| `/api/v1/requests/preview` | POST | Intent preview as JSON |
| `/api/v1/requests/apply` | POST | Initiate — confirmed, audited, allocator-owned |
The HTML form has no initiate button on purpose. Initiating requires a
confirmed POST to `/api/v1/requests/apply`, so a stray form submission cannot
reserve work as a side effect.
## The request
```json
{
"desired_role": "author",
"work_kind": "issue",
"work_number": 643,
"intent_summary": "implement request preview and initiation",
"remote": "prgs",
"org": "Scaled-Tech-Consulting",
"repo": "Gitea-Tools",
"expected_head_sha": null
}
```
`desired_role` is one of `author`, `reviewer`, `merger`, `reconciler`,
`controller`. `work_kind` is `issue` or `pr`. `remote`/`org`/`repo` default to
the first project in the registry when omitted; when neither the request nor
the registry resolves them, the request is rejected rather than pointed at some
other repository. `intent_summary` is required — it is what the audit record
states as the reason — and is truncated to 500 characters.
Parsing rejects rather than corrects. An unknown role, an unknown work kind, a
non-positive number, or a missing intent each return `400` with a `reason_code`
and the offending `field`.
## Preview
Five checks, each with its own verdict, reason code, and detail:
| Check | Passes when |
|-------|-------------|
| `authorization` | The console principal holds `operator` or above |
| `capability` | The desired role maps to a declared profile and MCP namespace |
| `lease_availability` | No active claim holds the work unit |
| `next_safe_action` | The allocator would independently select this exact work unit |
| `head_pin` | PR work resolves to a head SHA, and a supplied SHA still matches |
A preview also returns the role's `allowed_actions` and `prohibited_actions`
(from `allocator_service.ROLE_ACTIONS`), the `required_profile` and
`required_namespace` the work must run under, and a `correlation_id` that ties
the preview to its audit record and to any assignment that follows.
Preview is read-only in the strict sense: it calls the allocator with
`apply=false` and writes nothing but an audit line. An unauthorized principal
never reaches the allocator or the control-plane DB at all, so a denial cannot
be used to enumerate the queue.
## Initiation
`POST /api/v1/requests/apply` refuses in this order, and every refusal returns
before any assignment is attempted:
| Condition | Outcome | Status |
|-----------|---------|--------|
| Unparseable request | `invalid_request` | 400 |
| Not authorized, or execution not wired | `denied` | 403 |
| `confirm` not set | `denied` / `confirmation_required` | 409 |
| Work unit already claimed | `blocked` / `duplicate_assignment` | 409 |
| Allocator would select other work | `wait` / `not_next_safe_work` | 409 |
| Allocator declines on apply | `blocked` or `wait` | 409 |
| Evidence unavailable | `wait` / `evidence_unavailable` | 503 |
| Assigned | `assigned_work` | 201 |
A success returns the assignment plus a `handoff` block naming the profile, the
namespace, and the actions that stay forbidden — enough for the operator to
continue in the right MCP namespace without guessing.
### Why apply runs the allocator twice
The allocator is the only source of exclusive ownership (#600 / #613), and it
selects work; it does not take orders. So `apply` runs a dry-run first and
proceeds only when the allocator would independently pick the requested work
unit. If it would not, the request reports `wait` and mutates nothing.
A request is therefore a *confirmation* of the allocator's decision, never an
override of it. The apply call carries the dry-run's
`candidate_set_fingerprint` as a CAS pin (#776), so a queue that changed
between the two calls fails closed rather than assigning against a stale view.
The result is checked again on the way out: an assignment naming a different
work unit is not read as success.
### Fail-closed defaults
- An unreadable control-plane DB denies. It is never treated as "nothing holds
this work unit".
- An incomplete queue inventory denies (#758). Ranking a partial candidate set
can select the wrong work.
- An allocator that raises denies.
- PR work with no resolvable head SHA denies; a supplied SHA that no longer
matches denies with `head_moved`.
## Enabling initiation
Execution is wired off. Set `WEBUI_REQUESTS_EXECUTION=1` to enable it for the
`initiate_workflow` action only — see
[`webui-authz-audit.md`](webui-authz-audit.md) for why this is an
action-scoped flag rather than a phase bump. With the variable unset, `apply`
returns `403` with `reason_code: unauthorized` no matter who asks.
Enabling execution does **not** enable approvals or merges. Those are phase 3
console actions and remain forbidden in every path here; the console reserves
work and hands off, and the MCP role profile enforces what that role may then
do.
## Audit
Every preview and every apply emits a console audit record (schema in
[`webui-authz-audit.md`](webui-authz-audit.md)):
| Event | `result` |
|-------|----------|
| Preview | `previewed` |
| Refusal at any stage | `denied` |
| Assignment created | `succeeded` |
`correlation.request_id` carries the request's `correlation_id`, and a
successful record's `metadata` carries `assignment_id` and `lease_id`, so an
assignment can be traced back to the intent that produced it. The operator's
`intent_summary` travels in `metadata` and passes through the standard
redaction pass before persistence like every other field.
## Non-goals
- No browser-initiated approve or merge, in this phase or any other.
- No bypass of allocator exclusive ownership; no self-selection of work.
- No auto-start from raw monitoring incidents (#612 stays downstream).
+1015
View File
File diff suppressed because it is too large Load Diff
+112 -8
View File
@@ -2068,6 +2068,7 @@ import lease_lifecycle # noqa: E402
import lease_policy # noqa: E402
import workflow_dashboard # noqa: E402 # #605 live queue/lease dashboard
import restart_coordinator # noqa: E402 # #658 MCP restart coordinator/impact
import drain_proof # noqa: E402 # #661 pre-restart drain proof and hard gate
import incident_bridge # noqa: E402
import sentry_observability # noqa: E402 (#606 optional Sentry observability)
import sentry_incident_bridge # noqa: E402 (#607 Sentry→Gitea incident bridge)
@@ -22342,19 +22343,40 @@ def gitea_request_mcp_restart(
request_override: bool = False,
session_id: str | None = None,
limit: int = 200,
restart_class: str = "full_mcp_restart",
target_session_id: str | None = None,
target_role: str | None = None,
target_connector: str | None = None,
drain_proof_json: str | None = None,
request_break_glass: bool = False,
) -> dict:
"""Evaluate a proposed MCP restart and return an impact preview (#658).
Central restart coordinator: gathers live control-plane state (sessions,
Central restart coordinator: resolves the requested restart class, gathers
live control-plane state (sessions,
leases/locks, in-flight issue/PR work, mutations, worktrees) and returns a
blast-radius impact report with a ``safe`` / ``unsafe`` / ``override``
verdict, so the console (#642/#652) and operators can see what a restart
would disrupt *before* any concurrent LLM work is destroyed.
This tool is **dry-run and never restarts anything.** The mutative apply
path is a separate child gated by a drain proof (non-goal here); calling
with ``dry_run=False`` still performs no restart and reports that apply is
not yet available.
This tool **never restarts a process.** In dry-run (the default) it returns
only the impact preview. With ``dry_run=False`` it enforces the #661 hard
gate: the apply request must present a valid, unexpired, clean drain proof
(``drain_proof_json``) or it is denied and a durable incident descriptor is
returned under ``incident``. Break-glass is the only bypass and is honoured
only when ``request_break_glass`` is set *and* the environment carries
``GITEA_BREAKGLASS_RESTART_AUTHORIZATION``. Even an authorized gate performs
no restart here; actual execution is a further child. The gate outcome is
reported under ``apply_gate``.
``apply_authorized`` requires **both** authorizations to pass: the #661 drain
gate *and* the #663 restart-class matrix (``allow_restart``). They are
independent the drain gate proves the blast radius was drained and knows
nothing about whether this requester may request this class so a class the
matrix denied never reports an authorized apply. Break-glass bypasses the
drain proof only; it never bypasses the class matrix. ``apply_gate`` carries
``drain_gate_allow`` and ``restart_class_authorized`` so a denial is
attributable to the authorization that produced it.
Operator override authority is read from the process environment
(``GITEA_OPERATOR_RESTART_OVERRIDE_AUTHORIZATION``), never self-asserted by
@@ -22439,6 +22461,9 @@ def gitea_request_mcp_restart(
profile = get_profile()
profile_name = (profile.get("profile_name") or "").strip() or "session"
requester_role = (
profile.get("role_kind") or profile.get("role") or ""
).strip().lower()
sid = (session_id or "").strip() or f"{profile_name}-{os.getpid()}"
# Override authority is read from the environment only — a worker session
@@ -22448,6 +22473,15 @@ def gitea_request_mcp_restart(
(os.environ.get("GITEA_OPERATOR_RESTART_OVERRIDE_AUTHORIZATION") or "").strip()
)
operator_override = bool(request_override and operator_authorized)
controller_approved = bool(
(
os.environ.get("GITEA_CONTROLLER_RESTART_APPROVAL_AUTHORIZATION")
or ""
).strip()
)
requester_permissions = restart_coordinator.permissions_for_role(
requester_role
)
inventory = {
"sessions": sessions,
@@ -22462,6 +22496,14 @@ def gitea_request_mcp_restart(
operator_override=operator_override,
requesting_session_id=sid,
dry_run=True, # coordinator is always analysis-only (#658)
restart_class=restart_class,
requester_role=requester_role,
requester_permissions=requester_permissions,
controller_approved=controller_approved,
operator_authorized=operator_authorized,
target_session_id=target_session_id,
target_role=target_role,
target_connector=target_connector,
)
payload = report.as_dict()
@@ -22473,12 +22515,74 @@ def gitea_request_mcp_restart(
payload["requesting_session_id"] = sid
payload["operator_override_requested"] = bool(request_override)
payload["operator_override_authorized"] = operator_authorized
payload["controller_approval_authorized"] = controller_approved
payload["requester_role"] = requester_role
payload["requester_permissions"] = list(requester_permissions)
# Actual restart execution remains a further child; this tool never restarts
# a process. What #661 adds is the *hard gate*: an apply request (dry_run
# False) must present a valid, unexpired, clean drain proof, or it is denied
# and a durable incident is raised. Break-glass is the only bypass and its
# authorization is read from the environment, never self-asserted.
payload["apply_supported"] = False
if not dry_run:
payload["reasons"] = list(payload.get("reasons") or []) + [
"apply requested but not supported: sanctioned restart apply is "
"gated by a drain proof (separate child); no restart performed (#658)"
proof_obj: dict | None = None
proof_parse_error: str | None = None
if drain_proof_json:
try:
parsed = json.loads(drain_proof_json)
proof_obj = parsed if isinstance(parsed, dict) else None
if proof_obj is None:
proof_parse_error = "drain_proof_json is not a JSON object"
except (ValueError, TypeError) as exc:
proof_parse_error = f"invalid drain_proof_json: {_redact(str(exc))}"
break_glass_authorized = bool(
(
os.environ.get("GITEA_BREAKGLASS_RESTART_AUTHORIZATION") or ""
).strip()
)
break_glass = bool(request_break_glass and break_glass_authorized)
expected_fp = drain_proof.impact_fingerprint(report.as_dict())
gate = drain_proof.gate_apply_restart(
proof=proof_obj,
break_glass=break_glass,
expected_impact_fingerprint=expected_fp,
requesting_session_id=sid,
)
gate_payload = gate.as_dict()
if proof_parse_error and not break_glass:
gate_payload["reasons"] = [proof_parse_error] + list(
gate_payload.get("reasons") or []
)
# The #663 restart-class matrix and the #661 drain gate are two
# independent authorizations, and an apply requires BOTH. ``gate.allow``
# proves only that the blast radius was drained — or that break-glass
# was authorized — and knows nothing about whether this requester may
# request this class at all. Conjoining them keeps a class the matrix
# denied from ever reporting an authorized apply, and keeps break-glass
# scoped to what it is for: bypassing the drain proof, never the
# least-privilege class matrix.
restart_class_authorized = bool(report.allow_restart)
gate_payload["drain_gate_allow"] = bool(gate.allow)
gate_payload["restart_class_authorized"] = restart_class_authorized
if not restart_class_authorized:
gate_payload["reasons"] = list(gate_payload.get("reasons") or []) + [
"restart class authorization denied; apply denied regardless of "
"drain proof or break-glass (fail closed, #663)",
*(report.authorization_reasons or []),
]
payload["apply_gate"] = gate_payload
payload["apply_authorized"] = bool(gate.allow and restart_class_authorized)
payload["break_glass_requested"] = bool(request_break_glass)
payload["break_glass_authorized"] = break_glass_authorized
# Even an authorized gate performs no restart here: execution is a later
# child. The gate proves the apply path *would* be permitted.
payload["reasons"] = list(payload.get("reasons") or []) + list(
gate_payload.get("reasons") or []
)
if not gate.allow and gate.incident is not None:
payload["incident"] = gate.incident
return payload
+371 -15
View File
@@ -28,11 +28,12 @@ from __future__ import annotations
from dataclasses import dataclass, field
from datetime import datetime, timezone
from enum import Enum
from typing import Any, Mapping, Sequence
import lease_lifecycle
COORDINATOR_VERSION = "1.0.0-issue-658"
COORDINATOR_VERSION = "1.1.0-issue-663"
# Restart verdicts. Exactly the three the acceptance criteria name.
VERDICT_SAFE = "safe"
@@ -54,6 +55,194 @@ LEASE_FRESHNESS_LIVE = "active"
DEFAULT_SESSION_HEARTBEAT_STALE_SECONDS = 900
class RestartClass(str, Enum):
"""The only restart/recovery classes accepted by the coordinator."""
CLIENT_RECONNECT = "client_reconnect"
SESSION_RECONNECT = "session_reconnect"
WORKER_RESTART = "worker_restart"
ROLE_RUNTIME_RESTART = "role_runtime_restart"
CONNECTOR_RESTART = "connector_restart"
CONFIGURATION_RELOAD = "configuration_reload"
ROLLING_MCP_RESTART = "rolling_mcp_restart"
FULL_MCP_RESTART = "full_mcp_restart"
HOST_RESTART = "host_restart"
@dataclass(frozen=True)
class RestartClassPolicy:
"""Least-privilege policy for one :class:`RestartClass`."""
restart_class: RestartClass
required_permission: str
expected_blast_radius: str
drain_requirement: str
full_drain_required: bool
approval_requirement: str
audit_requirement: str
recovery_behavior: str
request_roles: tuple[str, ...]
execution_roles: tuple[str, ...]
def as_dict(self) -> dict[str, Any]:
return {
"restart_class": self.restart_class.value,
"required_permission": self.required_permission,
"expected_blast_radius": self.expected_blast_radius,
"drain_requirement": self.drain_requirement,
"full_drain_required": self.full_drain_required,
"approval_requirement": self.approval_requirement,
"audit_requirement": self.audit_requirement,
"recovery_behavior": self.recovery_behavior,
"request_roles": list(self.request_roles),
"execution_roles": list(self.execution_roles),
}
WORKER_ROLES = ("author", "reviewer", "merger", "reconciler")
CONTROL_ROLES = ("controller", "operator", "admin")
ALL_REQUEST_ROLES = WORKER_ROLES + CONTROL_ROLES
RESTART_CLASS_POLICIES: dict[RestartClass, RestartClassPolicy] = {
RestartClass.CLIENT_RECONNECT: RestartClassPolicy(
RestartClass.CLIENT_RECONNECT,
"mcp.reconnect.client",
BLAST_NONE,
"none",
False,
"self_service",
"record class, actor, client namespace, reason, and outcome",
"Reconnect only the caller's client transport; no daemon or peer session changes.",
ALL_REQUEST_ROLES,
ALL_REQUEST_ROLES,
),
RestartClass.SESSION_RECONNECT: RestartClassPolicy(
RestartClass.SESSION_RECONNECT,
"mcp.reconnect.session",
BLAST_LOW,
"requesting_session_safe_point",
False,
"self_service",
"record class, actor, session id, reason, and outcome",
"Rebind identity, capability, and workspace state for one session.",
ALL_REQUEST_ROLES,
ALL_REQUEST_ROLES,
),
RestartClass.WORKER_RESTART: RestartClassPolicy(
RestartClass.WORKER_RESTART,
"mcp.restart.worker.request",
BLAST_LOW,
"target_worker",
False,
"controller_approval_and_automated_gates",
"record class, actor, target worker, approval, drain proof, and outcome",
"Restart one worker after its own lease and mutation scope is drained.",
ALL_REQUEST_ROLES,
("operator", "admin"),
),
RestartClass.ROLE_RUNTIME_RESTART: RestartClassPolicy(
RestartClass.ROLE_RUNTIME_RESTART,
"mcp.restart.role_runtime.request",
BLAST_MEDIUM,
"target_role_runtime",
False,
"controller_approval_and_automated_gates",
"record class, actor, role namespace, approval, drain proof, and outcome",
"Restart only the selected role runtime and then re-probe that namespace.",
ALL_REQUEST_ROLES,
("operator", "admin"),
),
RestartClass.CONNECTOR_RESTART: RestartClassPolicy(
RestartClass.CONNECTOR_RESTART,
"mcp.restart.connector.request",
BLAST_MEDIUM,
"target_connector",
False,
"controller_approval_and_automated_gates",
"record class, actor, connector id, approval, drain proof, and outcome",
"Restart one connector while unrelated role runtimes remain available.",
ALL_REQUEST_ROLES,
("operator", "admin"),
),
RestartClass.CONFIGURATION_RELOAD: RestartClassPolicy(
RestartClass.CONFIGURATION_RELOAD,
"mcp.reload.configuration.request",
BLAST_LOW,
"mutation_quiesce",
False,
"controller_approval_and_automated_gates",
"record class, actor, configuration revision, approval, and outcome",
"Gracefully reload configuration without replacing the daemon process.",
ALL_REQUEST_ROLES,
("operator", "admin"),
),
RestartClass.ROLLING_MCP_RESTART: RestartClassPolicy(
RestartClass.ROLLING_MCP_RESTART,
"mcp.restart.rolling.request",
BLAST_MEDIUM,
"one_instance_at_a_time",
False,
"controller_approval_and_automated_gates",
"record class, actor, instance order, approval, per-instance drains, and outcome",
"Drain, restart, verify, and restore one instance before advancing to the next.",
CONTROL_ROLES,
("operator", "admin"),
),
RestartClass.FULL_MCP_RESTART: RestartClassPolicy(
RestartClass.FULL_MCP_RESTART,
"mcp.restart.full.request",
BLAST_HIGH,
"all_sessions_and_mutations",
True,
"controller_approval_and_automated_gates",
"record class, actor, full impact report, approval, drain proof, and outcome",
"Stop and restore the complete MCP runtime only after a verified full drain.",
CONTROL_ROLES,
("operator", "admin"),
),
RestartClass.HOST_RESTART: RestartClassPolicy(
RestartClass.HOST_RESTART,
"mcp.restart.host.request",
BLAST_HIGH,
"all_host_work",
True,
"controller_approval_plus_infrastructure_operator",
"record class, actor, host, incident or change id, approval, drain proof, and outcome",
"Hand off to infrastructure ownership; reconcile every runtime after the host returns.",
("controller", "operator", "admin"),
("operator", "admin"),
),
}
def resolve_restart_class(value: RestartClass | str) -> RestartClass:
"""Resolve a restart class or fail closed for an unknown value."""
if isinstance(value, RestartClass):
return value
try:
return RestartClass(str(value).strip())
except ValueError as exc:
raise ValueError(f"unknown restart class {value!r}; deny (fail closed)") from exc
def restart_class_policy(value: RestartClass | str) -> RestartClassPolicy:
"""Return the canonical policy for *value*."""
return RESTART_CLASS_POLICIES[resolve_restart_class(value)]
def permissions_for_role(role: str | None) -> tuple[str, ...]:
"""Return request permissions granted to a workflow role by this policy."""
normalized = str(role or "").strip().lower()
return tuple(
policy.required_permission
for policy in RESTART_CLASS_POLICIES.values()
if normalized in policy.request_roles
)
def _utc_now() -> datetime:
return datetime.now(timezone.utc)
@@ -75,6 +264,7 @@ class SessionImpact:
heartbeat_stale: bool
is_requester: bool
live: bool
connector: str | None = None
def as_dict(self) -> dict[str, Any]:
return {
@@ -87,6 +277,7 @@ class SessionImpact:
"heartbeat_stale": self.heartbeat_stale,
"is_requester": self.is_requester,
"live": self.live,
"connector": self.connector,
}
@@ -105,6 +296,7 @@ class LeaseImpact:
disruptive: bool
is_mutation: bool
is_critical_section: bool
connector: str | None = None
def as_dict(self) -> dict[str, Any]:
return {
@@ -119,6 +311,7 @@ class LeaseImpact:
"disruptive": self.disruptive,
"is_mutation": self.is_mutation,
"is_critical_section": self.is_critical_section,
"connector": self.connector,
}
@@ -127,6 +320,13 @@ class RestartImpactReport:
"""Impact preview DTO returned to the console / operator (#642/#652)."""
coordinator_version: str
restart_class: str
restart_policy: dict[str, Any]
policy_enforced: bool
permission_authorized: bool
role_authorized: bool
approval_satisfied: bool
authorization_reasons: list[str]
evaluated_at: str
dry_run: bool
restart_performed: bool
@@ -153,6 +353,13 @@ class RestartImpactReport:
def as_dict(self) -> dict[str, Any]:
return {
"coordinator_version": self.coordinator_version,
"restart_class": self.restart_class,
"restart_policy": dict(self.restart_policy),
"policy_enforced": self.policy_enforced,
"permission_authorized": self.permission_authorized,
"role_authorized": self.role_authorized,
"approval_satisfied": self.approval_satisfied,
"authorization_reasons": list(self.authorization_reasons),
"evaluated_at": self.evaluated_at,
"dry_run": self.dry_run,
"restart_performed": self.restart_performed,
@@ -206,6 +413,7 @@ def _classify_session(
requesting_session_id and session_id == requesting_session_id
),
live=live,
connector=(str(row.get("connector") or "").strip() or None),
)
@@ -258,6 +466,7 @@ def _classify_lease(row: Mapping[str, Any]) -> LeaseImpact:
disruptive=disruptive,
is_mutation=is_mutation,
is_critical_section=disruptive,
connector=(str(row.get("connector") or "").strip() or None),
)
@@ -279,6 +488,14 @@ def evaluate_restart_impact(
requesting_session_id: str | None = None,
dry_run: bool = True,
session_heartbeat_stale_seconds: int = DEFAULT_SESSION_HEARTBEAT_STALE_SECONDS,
restart_class: RestartClass | str | None = None,
requester_role: str | None = None,
requester_permissions: Sequence[str] | None = None,
controller_approved: bool = False,
operator_authorized: bool = False,
target_session_id: str | None = None,
target_role: str | None = None,
target_connector: str | None = None,
) -> RestartImpactReport:
"""Evaluate a proposed MCP restart and return an impact preview.
@@ -301,6 +518,61 @@ def evaluate_restart_impact(
"""
moment = now or _utc_now()
reasons: list[str] = []
authorization_reasons: list[str] = []
# ``None`` preserves the pre-#663 impact-only API for callers that have not
# yet been migrated. All MCP requests pass an explicit class and therefore
# take the fail-closed policy path.
policy_enforced = restart_class is not None
try:
resolved_class = resolve_restart_class(
restart_class or RestartClass.FULL_MCP_RESTART
)
policy = RESTART_CLASS_POLICIES[resolved_class]
unknown_class = False
except ValueError as exc:
resolved_class = None
policy = None
unknown_class = True
authorization_reasons.append(str(exc))
normalized_role = str(requester_role or "").strip().lower()
granted = {str(p).strip() for p in (requester_permissions or ())}
if policy_enforced and policy is not None:
permission_authorized = policy.required_permission in granted
role_authorized = normalized_role in policy.request_roles
if not permission_authorized:
authorization_reasons.append(
f"missing required permission {policy.required_permission!r}"
)
if not role_authorized:
authorization_reasons.append(
f"role {normalized_role or 'unknown'!r} may not request "
f"{policy.restart_class.value}"
)
elif unknown_class:
permission_authorized = False
role_authorized = False
else:
permission_authorized = True
role_authorized = True
if policy_enforced and policy is not None:
approval = policy.approval_requirement
if approval == "self_service":
approval_satisfied = True
elif approval == "controller_approval_plus_infrastructure_operator":
approval_satisfied = bool(controller_approved and operator_authorized)
else:
approval_satisfied = bool(controller_approved)
if not approval_satisfied:
authorization_reasons.append(
f"approval requirement not satisfied: {approval}"
)
elif unknown_class:
approval_satisfied = False
else:
approval_satisfied = True
inventory_complete = bool(inventory.get("inventory_complete", False))
incomplete_reasons = [str(r) for r in (inventory.get("incomplete_reasons") or [])]
@@ -323,15 +595,67 @@ def evaluate_restart_impact(
]
lease_impacts = [_classify_lease(l) for l in leases_raw]
# Only *other* live sessions and live leases constitute blast radius: a
# restart that would kill only the requesting session with no other work in
# flight is safe.
other_live_sessions = [
s for s in session_impacts if s.live and not s.is_requester
# Route impact through the selected class. Narrow classes never inherit a
# full-runtime drain merely because unrelated work exists.
target_complete = True
if resolved_class in {
RestartClass.CLIENT_RECONNECT,
RestartClass.SESSION_RECONNECT,
RestartClass.CONFIGURATION_RELOAD,
}:
scoped_sessions: list[SessionImpact] = []
scoped_leases: list[LeaseImpact] = []
elif resolved_class == RestartClass.WORKER_RESTART:
selected_session = (target_session_id or "").strip()
target_complete = bool(selected_session)
scoped_sessions = [
s for s in session_impacts if s.session_id == selected_session
]
disruptive_leases = [l for l in lease_impacts if l.disruptive]
critical_sections = [l for l in lease_impacts if l.is_critical_section]
mutations = [l for l in lease_impacts if l.is_mutation]
scoped_leases = [
l for l in lease_impacts if l.session_id == selected_session
]
elif resolved_class == RestartClass.ROLE_RUNTIME_RESTART:
selected_role = (target_role or "").strip().lower()
target_complete = bool(selected_role)
scoped_sessions = [
s for s in session_impacts if str(s.role or "").lower() == selected_role
]
scoped_leases = [
l for l in lease_impacts if str(l.role or "").lower() == selected_role
]
elif resolved_class == RestartClass.CONNECTOR_RESTART:
selected_connector = (target_connector or "").strip()
target_complete = bool(selected_connector)
scoped_sessions = [
s for s in session_impacts if s.connector == selected_connector
]
scoped_leases = [
l for l in lease_impacts if l.connector == selected_connector
]
else:
scoped_sessions = list(session_impacts)
scoped_leases = list(lease_impacts)
if policy_enforced and not target_complete:
authorization_reasons.append(
f"target required for {resolved_class.value if resolved_class else 'unknown class'}"
)
other_live_sessions = [
s for s in scoped_sessions if s.live and not s.is_requester
]
disruptive_leases = [l for l in scoped_leases if l.disruptive]
critical_sections = [l for l in scoped_leases if l.is_critical_section]
mutations = [l for l in scoped_leases if l.is_mutation]
terminal_lock_in_scope = (
terminal_lock
if resolved_class
not in {
RestartClass.CLIENT_RECONNECT,
RestartClass.SESSION_RECONNECT,
}
else None
)
affected_issues = sorted(
{
@@ -348,9 +672,24 @@ def evaluate_restart_impact(
}
)
disruptive = bool(disruptive_leases or other_live_sessions or terminal_lock)
disruptive = bool(
disruptive_leases or other_live_sessions or terminal_lock_in_scope
)
if not inventory_complete:
authorization_ok = bool(
not unknown_class
and permission_authorized
and role_authorized
and approval_satisfied
and target_complete
)
if policy_enforced and not authorization_ok:
verdict = VERDICT_UNSAFE
allow_restart = False
reasons.append("restart class authorization denied (fail closed)")
reasons.extend(authorization_reasons)
elif not inventory_complete:
verdict = VERDICT_UNSAFE
allow_restart = False
reasons.append(
@@ -381,7 +720,7 @@ def evaluate_restart_impact(
f"{len(critical_sections)} critical section(s) in flight "
"(active lease with a live owner)"
)
if terminal_lock:
if terminal_lock_in_scope:
reasons.append("active terminal (merge) lock present")
override_would_allow = bool(inventory_complete and disruptive)
@@ -411,6 +750,12 @@ def evaluate_restart_impact(
audit_record = {
"event": "restart_impact_evaluated",
"coordinator_version": COORDINATOR_VERSION,
"restart_class": (
resolved_class.value if resolved_class else str(restart_class or "")
),
"required_permission": (
policy.required_permission if policy is not None else None
),
"evaluated_at": moment.isoformat(),
"dry_run": dry_run,
"operator_override": bool(operator_override),
@@ -424,6 +769,15 @@ def evaluate_restart_impact(
return RestartImpactReport(
coordinator_version=COORDINATOR_VERSION,
restart_class=(
resolved_class.value if resolved_class else str(restart_class or "")
),
restart_policy=policy.as_dict() if policy is not None else {},
policy_enforced=policy_enforced,
permission_authorized=permission_authorized,
role_authorized=role_authorized,
approval_satisfied=approval_satisfied,
authorization_reasons=authorization_reasons,
evaluated_at=moment.isoformat(),
dry_run=dry_run,
restart_performed=False,
@@ -440,9 +794,11 @@ def evaluate_restart_impact(
affected_issues=affected_issues,
affected_prs=affected_prs,
mutations=mutations,
terminal_lock=dict(terminal_lock)
if isinstance(terminal_lock, Mapping)
else terminal_lock,
terminal_lock=(
dict(terminal_lock_in_scope)
if isinstance(terminal_lock_in_scope, Mapping)
else terminal_lock_in_scope
),
ack_state=ack_state,
prior_recovery_attempts=prior_recovery_attempts,
counts=counts,
+158
View File
@@ -7,6 +7,7 @@ import tempfile
import threading
import unittest
from concurrent.futures import ThreadPoolExecutor, as_completed
from datetime import datetime, timezone
from allocator_service import (
OUTCOME_ASSIGNED,
@@ -15,6 +16,7 @@ from allocator_service import (
OUTCOME_PREVIEW,
OUTCOME_WAIT,
WorkCandidate,
_drop_expired_claims,
allocate_next_work,
candidate_from_dict,
classify_skip,
@@ -362,5 +364,161 @@ class AllocatorServiceTest(unittest.TestCase):
self.assertIn("unavailable", res["reasons"][0].lower())
class SideEffectFreeAllocationTest(unittest.TestCase):
"""``side_effect_free`` dry runs write nothing to the control plane (#643).
A plain ``apply=False`` still called ``upsert_session`` and
``expire_stale_leases`` before the apply branch was consulted, so a caller
advertising a read-only preview mutated on every call one unreferenced
session row per preview, plus a global lease sweep.
"""
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory()
self.db = ControlPlaneDB(os.path.join(self._tmp.name, "cp.sqlite3"))
def tearDown(self) -> None:
self._tmp.cleanup()
def _alloc(self, **kwargs):
defaults = dict(
db=self.db,
session_id="s-preview",
role="author",
remote="prgs",
org="org",
repo="repo",
candidates=[
WorkCandidate(kind="issue", number=643, labels=("status:ready",))
],
apply=False,
profile_name="prgs-author",
username="jcwalker3",
)
defaults.update(kwargs)
return allocate_next_work(**defaults)
def _session_ids(self) -> set[str]:
return {str(r.get("session_id")) for r in self.db.list_sessions()}
def test_side_effect_free_preview_writes_no_session_row(self):
before = self._session_ids()
result = self._alloc(side_effect_free=True)
self.assertEqual(result["outcome"], OUTCOME_PREVIEW)
self.assertEqual(self._session_ids(), before)
self.assertNotIn("s-preview", self._session_ids())
def test_plain_dry_run_still_registers_a_session(self):
# The default is unchanged for every existing caller.
self._alloc()
self.assertIn("s-preview", self._session_ids())
def test_repeated_previews_do_not_accumulate_rows(self):
for index in range(5):
self._alloc(side_effect_free=True, session_id=f"s-{index}")
self.assertEqual(self._session_ids(), set())
def test_side_effect_free_does_not_sweep_stale_leases(self):
self.db.upsert_session(session_id="owner", role="author", pid=1)
assigned = self.db.assign_and_lease(
session_id="owner",
role="author",
remote="prgs",
org="org",
repo="repo",
kind="issue",
number=999,
lease_ttl_seconds=-60, # already expired
)
self.assertEqual(assigned.outcome, "assigned")
self._alloc(side_effect_free=True)
# The expired row is still 'active' in the DB: nothing swept it.
statuses = {
r["lease_id"]: r["status"]
for r in self.db.list_leases(
remote="prgs", org="org", repo="repo",
statuses=("active", "expired"),
)
}
self.assertEqual(statuses.get(assigned.lease_id), "active")
def test_expired_claims_are_filtered_in_memory_so_work_stays_selectable(self):
"""The read-only mirror of the sweep: expired claims must not block."""
self.db.upsert_session(session_id="owner", role="author", pid=1)
self.db.assign_and_lease(
session_id="owner",
role="author",
remote="prgs",
org="org",
repo="repo",
kind="issue",
number=643,
lease_ttl_seconds=-60, # expired: must not withhold #643
)
result = self._alloc(side_effect_free=True)
self.assertEqual(result["outcome"], OUTCOME_PREVIEW)
self.assertEqual(result["selected"]["number"], 643)
def test_a_live_claim_still_withholds_the_work(self):
self.db.upsert_session(session_id="owner", role="author", pid=1)
self.db.assign_and_lease(
session_id="owner",
role="author",
remote="prgs",
org="org",
repo="repo",
kind="issue",
number=643,
lease_ttl_seconds=3600,
)
result = self._alloc(side_effect_free=True)
self.assertNotEqual(result["outcome"], OUTCOME_ASSIGNED)
self.assertNotEqual((result.get("selected") or {}).get("number"), 643)
def test_side_effect_free_with_apply_fails_closed(self):
result = self._alloc(side_effect_free=True, apply=True)
self.assertFalse(result["success"])
self.assertEqual(result["outcome"], OUTCOME_NO_SAFE)
self.assertIsNone(result["assignment"])
self.assertIn("incompatible with apply", result["reasons"][0])
# And it reserved nothing.
self.assertEqual(
self.db.list_leases(remote="prgs", org="org", repo="repo"), []
)
class DropExpiredClaimsTest(unittest.TestCase):
"""The in-memory expiry filter behind side-effect-free previews (#643)."""
def test_unparseable_expiry_is_kept_rather_than_assumed_free(self):
claims = {
("issue", 1): {"lease_id": "l1", "expires_at": "not-a-date"},
("issue", 2): {"lease_id": "l2"},
("issue", 3): {"lease_id": "l3", "expires_at": None},
}
self.assertEqual(_drop_expired_claims(claims), claims)
def test_expired_dropped_and_future_kept(self):
now = datetime(2026, 7, 25, 12, 0, tzinfo=timezone.utc)
claims = {
("issue", 1): {"expires_at": "2026-07-25T11:59:59+00:00"},
("issue", 2): {"expires_at": "2026-07-25T12:00:01+00:00"},
("issue", 3): {"expires_at": "2026-07-25T12:00:00+00:00"}, # boundary
}
kept = _drop_expired_claims(claims, now=now)
self.assertEqual(set(kept), {("issue", 2)})
def test_naive_and_zulu_timestamps_are_treated_as_utc(self):
now = datetime(2026, 7, 25, 12, 0, tzinfo=timezone.utc)
claims = {
("issue", 1): {"expires_at": "2026-07-25T11:00:00"}, # naive, past
("issue", 2): {"expires_at": "2026-07-25T13:00:00Z"}, # zulu, future
}
kept = _drop_expired_claims(claims, now=now)
self.assertEqual(set(kept), {("issue", 2)})
if __name__ == "__main__":
unittest.main()
+907
View File
@@ -0,0 +1,907 @@
"""Tests for the pre-restart drain proof and hard gate (#661).
Covers the acceptance criteria:
1. Restart apply without a proof fails closed.
2. A successful drain produces a verifiable proof.
3. An open unsafe mutation makes the proof fail (multi-session fixture).
4. Pass / fail / expired verification paths.
Plus the security posture: forged/tampered proofs are rejected, break-glass is
the only bypass and is never silent, a stale blast-radius fingerprint rejects a
proof, and no per-process secret ever leaks into a serialized artifact.
"""
from __future__ import annotations
import os
import unittest
from datetime import datetime, timedelta, timezone
import drain_proof as dp
import restart_coordinator as rc
NOW = datetime(2026, 7, 24, 6, 0, 0, tzinfo=timezone.utc)
SECRET = b"unit-test-drain-proof-secret-0123456789abcdef"
def _live_pid() -> int:
return os.getpid()
def _clean_drain_state() -> dict:
"""Every drain action succeeded, no sessions outstanding."""
return {
"assignments_stopped": True,
"checkpoints_complete": True,
"handoffs_verified": True,
"leases_handled": True,
"acks": {}, # no other live sessions to acknowledge
"ack_timeout_policy_applied": False,
}
def _safe_report() -> dict:
"""Impact report with no other live work: a restart here is safe."""
report = rc.evaluate_restart_impact(
{"sessions": [], "leases": [], "inventory_complete": True},
now=NOW,
requesting_session_id="prgs-controller-1-req",
)
return report.as_dict()
def _unsafe_mutation_report() -> dict:
"""Multi-session report: a second session holds a live author mutation."""
sessions = [
{
"session_id": "prgs-controller-1-req",
"role": "controller",
"profile": "prgs-controller",
"pid": _live_pid(),
"status": "active",
"last_heartbeat_at": NOW.isoformat(),
},
{
"session_id": "prgs-author-99",
"role": "author",
"profile": "prgs-author",
"pid": _live_pid(),
"status": "active",
"last_heartbeat_at": NOW.isoformat(),
},
]
leases = [
{
"lease_id": "lease-mut",
"session_id": "prgs-author-99",
"role": "author",
"phase": "implementing",
"work_kind": "issue",
"work_number": 661,
"worktree_path": "branches/issue-661",
"freshness": {"freshness": "active"},
}
]
report = rc.evaluate_restart_impact(
{"sessions": sessions, "leases": leases, "inventory_complete": True},
now=NOW,
requesting_session_id="prgs-controller-1-req",
)
return report.as_dict()
class BuildDrainProofTests(unittest.TestCase):
def test_clean_drain_produces_verifiable_clean_proof(self):
"""AC#2: a successful drain produces a verifiable proof."""
proof = dp.build_drain_proof(
impact_report=_safe_report(),
drain_state=_clean_drain_state(),
requesting_session_id="prgs-controller-1-req",
now=NOW,
secret=SECRET,
)
self.assertTrue(proof.clean)
self.assertEqual(proof.failed_checks, [])
self.assertEqual(
{c.name for c in proof.checks}, set(dp.REQUIRED_CHECKS)
)
result = dp.verify_drain_proof(
proof.as_dict(), now=NOW, secret=SECRET
)
self.assertTrue(result.valid, result.reasons)
self.assertFalse(result.expired)
self.assertFalse(result.tampered)
def test_open_mutation_makes_proof_unclean(self):
"""AC#3: an unsafe mutation still in flight fails the proof."""
proof = dp.build_drain_proof(
impact_report=_unsafe_mutation_report(),
drain_state=_clean_drain_state(),
now=NOW,
secret=SECRET,
)
self.assertFalse(proof.clean)
self.assertIn(dp.CHECK_NO_INFLIGHT_MUTATIONS, proof.failed_checks)
# Leases-handled also fails: the report still shows a disruptive lease.
self.assertIn(dp.CHECK_LEASES_HANDLED, proof.failed_checks)
result = dp.verify_drain_proof(proof.as_dict(), now=NOW, secret=SECRET)
self.assertFalse(result.valid)
def test_incomplete_inventory_fails_no_mutations_check(self):
proof = dp.build_drain_proof(
impact_report={"inventory_complete": False},
drain_state=_clean_drain_state(),
now=NOW,
secret=SECRET,
)
self.assertFalse(proof.clean)
self.assertIn(dp.CHECK_NO_INFLIGHT_MUTATIONS, proof.failed_checks)
def test_missing_checkpoint_flag_fails_closed(self):
state = _clean_drain_state()
del state["checkpoints_complete"]
proof = dp.build_drain_proof(
impact_report=_safe_report(), drain_state=state, now=NOW, secret=SECRET
)
self.assertFalse(proof.clean)
self.assertIn(dp.CHECK_CHECKPOINTS_COMPLETE, proof.failed_checks)
def test_non_true_flags_fail_closed(self):
"""A truthy-but-not-True value (e.g. the string 'yes') must not pass."""
state = _clean_drain_state()
state["assignments_stopped"] = "yes"
proof = dp.build_drain_proof(
impact_report=_safe_report(), drain_state=state, now=NOW, secret=SECRET
)
self.assertIn(dp.CHECK_ASSIGNMENTS_STOPPED, proof.failed_checks)
def test_ack_timeout_policy_satisfies_ack_check(self):
state = _clean_drain_state()
state["acks"] = {"prgs-author-99": "pending"}
state["ack_timeout_policy_applied"] = True
proof = dp.build_drain_proof(
impact_report=_safe_report(), drain_state=state, now=NOW, secret=SECRET
)
names = {c.name: c.passed for c in proof.checks}
self.assertTrue(names[dp.CHECK_ACKS_OR_TIMEOUT])
def test_outstanding_acks_without_timeout_fail(self):
state = _clean_drain_state()
state["acks"] = {"prgs-author-99": "pending"}
state["ack_timeout_policy_applied"] = False
proof = dp.build_drain_proof(
impact_report=_safe_report(), drain_state=state, now=NOW, secret=SECRET
)
self.assertIn(dp.CHECK_ACKS_OR_TIMEOUT, proof.failed_checks)
def test_all_acked_satisfies_ack_check(self):
state = _clean_drain_state()
state["acks"] = {"prgs-author-99": "acked", "prgs-author-2": "acknowledged"}
proof = dp.build_drain_proof(
impact_report=_safe_report(), drain_state=state, now=NOW, secret=SECRET
)
names = {c.name: c.passed for c in proof.checks}
self.assertTrue(names[dp.CHECK_ACKS_OR_TIMEOUT])
class VerifyDrainProofTests(unittest.TestCase):
def _clean_proof_dict(self) -> dict:
return dp.build_drain_proof(
impact_report=_safe_report(),
drain_state=_clean_drain_state(),
now=NOW,
secret=SECRET,
).as_dict()
def test_missing_proof_is_invalid(self):
result = dp.verify_drain_proof(None, now=NOW, secret=SECRET)
self.assertFalse(result.valid)
self.assertIsNone(result.proof_id)
def test_expired_proof_is_invalid(self):
"""AC#4: an expired proof fails verification."""
proof = self._clean_proof_dict()
later = NOW + timedelta(seconds=dp.DEFAULT_PROOF_TTL_SECONDS + 1)
result = dp.verify_drain_proof(proof, now=later, secret=SECRET)
self.assertFalse(result.valid)
self.assertTrue(result.expired)
def test_proof_valid_just_before_expiry(self):
proof = self._clean_proof_dict()
almost = NOW + timedelta(seconds=dp.DEFAULT_PROOF_TTL_SECONDS - 1)
result = dp.verify_drain_proof(proof, now=almost, secret=SECRET)
self.assertTrue(result.valid, result.reasons)
def test_wrong_secret_rejected(self):
"""A proof minted in a prior process (different secret) will not verify."""
proof = self._clean_proof_dict()
result = dp.verify_drain_proof(proof, now=NOW, secret=b"other-secret")
self.assertFalse(result.valid)
self.assertTrue(result.tampered)
def test_flipping_clean_flag_is_detected(self):
"""Forging clean=True on an unclean proof breaks the signature."""
unclean = dp.build_drain_proof(
impact_report=_unsafe_mutation_report(),
drain_state=_clean_drain_state(),
now=NOW,
secret=SECRET,
).as_dict()
self.assertFalse(unclean["clean"])
unclean["clean"] = True # forge
result = dp.verify_drain_proof(unclean, now=NOW, secret=SECRET)
self.assertFalse(result.valid)
self.assertTrue(result.tampered)
def test_tampering_a_check_is_detected(self):
unclean = dp.build_drain_proof(
impact_report=_unsafe_mutation_report(),
drain_state=_clean_drain_state(),
now=NOW,
secret=SECRET,
).as_dict()
for c in unclean["checks"]:
if c["name"] == dp.CHECK_NO_INFLIGHT_MUTATIONS:
c["passed"] = True # forge the failing check to pass
result = dp.verify_drain_proof(unclean, now=NOW, secret=SECRET)
self.assertFalse(result.valid)
self.assertTrue(result.tampered)
def test_missing_required_check_rejected(self):
proof = self._clean_proof_dict()
proof["checks"] = [
c for c in proof["checks"] if c["name"] != dp.CHECK_HANDOFFS_OK
]
result = dp.verify_drain_proof(proof, now=NOW, secret=SECRET)
self.assertFalse(result.valid)
def test_stale_fingerprint_rejected(self):
proof = self._clean_proof_dict()
result = dp.verify_drain_proof(
proof,
now=NOW,
secret=SECRET,
expected_impact_fingerprint="deadbeef",
)
self.assertFalse(result.valid)
def test_matching_fingerprint_accepted(self):
report = _safe_report()
proof = dp.build_drain_proof(
impact_report=report,
drain_state=_clean_drain_state(),
now=NOW,
secret=SECRET,
).as_dict()
fp = dp.impact_fingerprint(report)
result = dp.verify_drain_proof(
proof, now=NOW, secret=SECRET, expected_impact_fingerprint=fp
)
self.assertTrue(result.valid, result.reasons)
class GateApplyRestartTests(unittest.TestCase):
def _clean_proof_dict(self) -> dict:
return dp.build_drain_proof(
impact_report=_safe_report(),
drain_state=_clean_drain_state(),
now=NOW,
secret=SECRET,
).as_dict()
def test_apply_without_proof_denied(self):
"""AC#1: restart apply without a proof fails closed + raises incident."""
decision = dp.gate_apply_restart(proof=None, now=NOW, secret=SECRET)
self.assertFalse(decision.allow)
self.assertEqual(decision.verdict, dp.GATE_DENY)
self.assertIsNotNone(decision.incident)
self.assertEqual(
decision.incident["kind"], "restart_drain_gate_denied"
)
def test_apply_with_valid_proof_allowed(self):
decision = dp.gate_apply_restart(
proof=self._clean_proof_dict(), now=NOW, secret=SECRET
)
self.assertTrue(decision.allow)
self.assertEqual(decision.verdict, dp.GATE_ALLOW)
self.assertIsNone(decision.incident)
def test_apply_with_expired_proof_denied_with_incident(self):
later = NOW + timedelta(seconds=dp.DEFAULT_PROOF_TTL_SECONDS + 5)
decision = dp.gate_apply_restart(
proof=self._clean_proof_dict(), now=later, secret=SECRET
)
self.assertFalse(decision.allow)
self.assertIsNotNone(decision.incident)
def test_apply_with_unclean_proof_denied(self):
"""AC#3 at the gate: an unsafe-mutation proof is denied."""
unclean = dp.build_drain_proof(
impact_report=_unsafe_mutation_report(),
drain_state=_clean_drain_state(),
now=NOW,
secret=SECRET,
).as_dict()
decision = dp.gate_apply_restart(proof=unclean, now=NOW, secret=SECRET)
self.assertFalse(decision.allow)
self.assertIsNotNone(decision.incident)
def test_break_glass_allows_without_proof_but_records_bypass(self):
decision = dp.gate_apply_restart(
proof=None, now=NOW, secret=SECRET, break_glass=True
)
self.assertTrue(decision.allow)
self.assertEqual(decision.verdict, dp.GATE_BREAK_GLASS)
self.assertTrue(decision.break_glass)
self.assertIsNone(decision.incident)
self.assertTrue(decision.audit_record["break_glass"])
def test_denied_gate_carries_stale_fingerprint_reason(self):
decision = dp.gate_apply_restart(
proof=self._clean_proof_dict(),
now=NOW,
secret=SECRET,
expected_impact_fingerprint="not-the-fingerprint",
)
self.assertFalse(decision.allow)
class SecretHygieneTests(unittest.TestCase):
def test_secret_never_serialized(self):
proof = dp.build_drain_proof(
impact_report=_safe_report(),
drain_state=_clean_drain_state(),
now=NOW,
secret=SECRET,
)
blob = dp._canonical(proof.as_dict())
self.assertNotIn(SECRET.decode(), blob)
# The signature is a hex digest, not the raw secret.
self.assertNotIn(SECRET.hex(), blob)
def test_incident_descriptor_has_no_secret(self):
decision = dp.gate_apply_restart(proof=None, now=NOW, secret=SECRET)
blob = dp._canonical(decision.incident)
self.assertNotIn(SECRET.decode(), blob)
def _drained_report_with_live_sessions(count: int) -> dict:
"""Report with ``count`` other live sessions but nothing in flight.
Every other checklist item passes against this report, so a failure
isolates the acknowledgement check rather than tripping on mutations.
"""
sessions = [
{
"session_id": "prgs-controller-1-req",
"role": "controller",
"profile": "prgs-controller",
"pid": _live_pid(),
"status": "active",
"last_heartbeat_at": NOW.isoformat(),
}
]
for index in range(count):
sessions.append(
{
"session_id": f"prgs-author-{index}",
"role": "author",
"profile": "prgs-author",
"pid": _live_pid(),
"status": "active",
"last_heartbeat_at": NOW.isoformat(),
}
)
report = rc.evaluate_restart_impact(
{"sessions": sessions, "leases": [], "inventory_complete": True},
now=NOW,
requesting_session_id="prgs-controller-1-req",
)
return report.as_dict()
class AcknowledgementFailClosedTests(unittest.TestCase):
"""Acknowledgement evidence must fail closed unless explicitly verified.
Regression cover for the reviewed fail-open on PR #882: an absent ``acks``
key collapsed to ``{}`` and was read as "no other live sessions required to
acknowledge", so a proof minted clean and the restart gate allowed while the
impact report still showed other live sessions.
"""
def _state(self, **overrides) -> dict:
state = _clean_drain_state()
state.pop("acks", None)
state["ack_timeout_policy_applied"] = False
state.update(overrides)
return state
def _acks_check(self, proof) -> dp.DrainCheck:
return next(c for c in proof.checks if c.name == dp.CHECK_ACKS_OR_TIMEOUT)
def _build(self, report: dict, state: dict):
return dp.build_drain_proof(
impact_report=report, drain_state=state, now=NOW, secret=SECRET
)
def assertAcksFailClosed(self, report: dict, state: dict) -> None:
proof = self._build(report, state)
self.assertFalse(self._acks_check(proof).passed)
self.assertIn(dp.CHECK_ACKS_OR_TIMEOUT, proof.failed_checks)
self.assertFalse(proof.clean)
# --- missing / null / empty / malformed ------------------------------
def test_missing_acks_key_with_live_sessions_fails_closed(self):
"""The exact reviewed defect: absent key, three other live sessions."""
report = _drained_report_with_live_sessions(3)
self.assertEqual(report["counts"]["sessions_live_other"], 3)
state = self._state()
self.assertNotIn("acks", state)
proof = self._build(report, state)
check = self._acks_check(proof)
self.assertFalse(check.passed)
self.assertNotIn("no other live sessions", check.detail)
self.assertIn("fail closed", check.detail)
self.assertFalse(proof.clean)
self.assertEqual(proof.failed_checks, [dp.CHECK_ACKS_OR_TIMEOUT])
def test_none_acks_with_live_sessions_fails_closed(self):
self.assertAcksFailClosed(
_drained_report_with_live_sessions(2), self._state(acks=None)
)
def test_empty_acks_with_live_sessions_fails_closed(self):
self.assertAcksFailClosed(
_drained_report_with_live_sessions(1), self._state(acks={})
)
def test_malformed_acks_fail_closed(self):
for malformed in ([], "ack", 7, ("ack",), True):
with self.subTest(malformed=malformed):
self.assertAcksFailClosed(
_drained_report_with_live_sessions(1),
self._state(acks=malformed),
)
# --- stale / unproven values -----------------------------------------
def test_stale_or_unproven_ack_values_fail_closed(self):
for value in ("pending", "stale", "unknown", "", None, True, 1, NOW):
with self.subTest(value=value):
self.assertAcksFailClosed(
_drained_report_with_live_sessions(1),
self._state(acks={"prgs-author-0": value}),
)
def test_partial_coverage_fails_closed(self):
"""Fewer acknowledgements than the report's live-session count."""
self.assertAcksFailClosed(
_drained_report_with_live_sessions(3),
self._state(acks={"prgs-author-0": "ack"}),
)
def test_one_unacked_entry_among_many_fails_closed(self):
self.assertAcksFailClosed(
_drained_report_with_live_sessions(2),
self._state(acks={"prgs-author-0": "ack", "prgs-author-1": "pending"}),
)
def test_unproven_live_session_count_fails_closed(self):
"""A missing or malformed count cannot prove nobody had to acknowledge."""
malformed_counts = (
None,
{},
{"sessions_live_other": None},
{"sessions_live_other": "3"},
{"sessions_live_other": -1},
{"sessions_live_other": True},
)
for counts in malformed_counts:
with self.subTest(counts=counts):
report = _drained_report_with_live_sessions(0)
if counts is None:
report.pop("counts", None)
else:
report["counts"] = counts
self.assertAcksFailClosed(report, self._state())
# --- valid evidence still passes -------------------------------------
def test_complete_valid_acks_pass(self):
report = _drained_report_with_live_sessions(2)
state = self._state(
acks={"prgs-author-0": "ack", "prgs-author-1": "acknowledged"}
)
proof = self._build(report, state)
self.assertTrue(self._acks_check(proof).passed)
self.assertTrue(proof.clean)
self.assertEqual(proof.failed_checks, [])
def test_no_other_live_sessions_still_passes(self):
"""Intended behavior retained: zero live sessions needs no acks."""
report = _drained_report_with_live_sessions(0)
self.assertEqual(report["counts"]["sessions_live_other"], 0)
proof = self._build(report, self._state())
check = self._acks_check(proof)
self.assertTrue(check.passed)
self.assertIn("sessions_live_other=0", check.detail)
self.assertTrue(proof.clean)
# --- timeout policy cannot become a second fail-open ------------------
def test_unproven_timeout_policy_cannot_open_the_gate(self):
for value in (None, "true", "yes", 1, "True", [], {}):
with self.subTest(value=value):
self.assertAcksFailClosed(
_drained_report_with_live_sessions(2),
self._state(ack_timeout_policy_applied=value),
)
def test_explicit_timeout_policy_permits(self):
proof = self._build(
_drained_report_with_live_sessions(2),
self._state(ack_timeout_policy_applied=True),
)
check = self._acks_check(proof)
self.assertTrue(check.passed)
self.assertIn("timeout policy", check.detail)
self.assertTrue(proof.clean)
# --- the gate itself must deny ---------------------------------------
def test_failed_ack_check_denies_the_restart_gate(self):
report = _drained_report_with_live_sessions(3)
proof = self._build(report, self._state())
self.assertFalse(proof.clean)
decision = dp.gate_apply_restart(
proof=proof.as_dict(),
now=NOW,
secret=SECRET,
expected_impact_fingerprint=dp.impact_fingerprint(report),
)
self.assertFalse(decision.allow)
self.assertEqual(decision.verdict, dp.GATE_DENY)
self.assertIsNotNone(decision.incident)
def test_unclean_ack_proof_fails_verification(self):
report = _drained_report_with_live_sessions(3)
proof = self._build(report, self._state())
result = dp.verify_drain_proof(
proof.as_dict(),
now=NOW,
secret=SECRET,
expected_impact_fingerprint=dp.impact_fingerprint(report),
)
self.assertFalse(result.valid)
self.assertFalse(result.clean)
def _identity_report(*, requester: str, others: tuple[str, ...]) -> dict:
"""Report with explicitly named requester and other live sessions.
Unlike :func:`_drained_report_with_live_sessions`, the session ids are
chosen by the caller so a test can supply acknowledgements for the *wrong*
identities while keeping the count correct.
"""
sessions = [
{
"session_id": requester,
"role": "controller",
"profile": "prgs-controller",
"pid": _live_pid(),
"status": "active",
"last_heartbeat_at": NOW.isoformat(),
}
]
for session_id in others:
sessions.append(
{
"session_id": session_id,
"role": "author",
"profile": "prgs-author",
"pid": _live_pid(),
"status": "active",
"last_heartbeat_at": NOW.isoformat(),
}
)
report = rc.evaluate_restart_impact(
{"sessions": sessions, "leases": [], "inventory_complete": True},
now=NOW,
requesting_session_id=requester,
)
return report.as_dict()
class AcknowledgementIdentityBindingTests(unittest.TestCase):
"""Acknowledgement coverage must be bound to session identity, not counted.
Regression cover for the second reviewed fail-open on PR #882 (review 582,
blocker B1): coverage compared ``acked_count`` against
``counts.sessions_live_other``, so acknowledgements supplied for the
requesting session and for ids that do not exist satisfied the obligations
of the live sessions that never answered. The required identities are
carried by the report itself ``ack_state`` keys and ``affected_sessions``
filtered on ``live and not is_requester`` and only an acknowledgement
keyed by one of those ids may count for it.
"""
def _state(self, **overrides) -> dict:
state = _clean_drain_state()
state.pop("acks", None)
state["ack_timeout_policy_applied"] = False
state.update(overrides)
return state
def _acks_check(self, proof) -> dp.DrainCheck:
return next(c for c in proof.checks if c.name == dp.CHECK_ACKS_OR_TIMEOUT)
def _build(self, report: dict, state: dict):
return dp.build_drain_proof(
impact_report=report, drain_state=state, now=NOW, secret=SECRET
)
def assertAcksFailClosed(self, report: dict, state: dict) -> dp.DrainCheck:
"""Failure must propagate through the check, the proof, and the gate."""
proof = self._build(report, state)
check = self._acks_check(proof)
self.assertFalse(check.passed)
self.assertFalse(proof.clean)
self.assertIn(dp.CHECK_ACKS_OR_TIMEOUT, proof.failed_checks)
decision = dp.gate_apply_restart(
proof=proof.as_dict(),
now=NOW,
secret=SECRET,
expected_impact_fingerprint=dp.impact_fingerprint(report),
)
self.assertEqual(decision.verdict, dp.GATE_DENY)
self.assertFalse(decision.allow)
return check
# --- the reviewer's exact reproduction --------------------------------
def test_requester_plus_unknown_id_cannot_satisfy_two_live_sessions(self):
"""Review 582 B1 verbatim: requester + a nonexistent session.
``sessions_live_other=2`` with ``ack_state`` naming ``other-0`` and
``other-1``; the drain state supplies an acknowledgement from the
requesting session itself and from a session that does not exist. The
count matches, the identities do not.
"""
report = _identity_report(requester="req", others=("other-0", "other-1"))
self.assertEqual(report["counts"]["sessions_live_other"], 2)
self.assertEqual(
report["ack_state"], {"other-0": "pending", "other-1": "pending"}
)
state = self._state(acks={"req": "ack", "totally-bogus-session": "ack"})
check = self.assertAcksFailClosed(report, state)
self.assertIn("other-0", check.detail)
self.assertIn("other-1", check.detail)
self.assertIn("fail closed", check.detail)
# --- wrong / unknown / requester identities ---------------------------
def test_sufficient_count_of_wrong_ids_fails_closed(self):
"""Right cardinality, wrong identities: two acks, neither required."""
report = _identity_report(requester="req", others=("other-0", "other-1"))
state = self._state(acks={"ghost-a": "ack", "ghost-b": "ack"})
check = self.assertAcksFailClosed(report, state)
self.assertIn("do not count", check.detail)
def test_more_acks_than_required_still_fails_on_wrong_ids(self):
"""Coverage cannot be bought with volume: five acks, none required."""
report = _identity_report(requester="req", others=("other-0", "other-1"))
state = self._state(acks={f"ghost-{i}": "acknowledged" for i in range(5)})
self.assertAcksFailClosed(report, state)
def test_partial_identity_match_fails_closed(self):
"""One required id acknowledged, the rest padded with unknown ids."""
report = _identity_report(
requester="req", others=("other-0", "other-1", "other-2")
)
state = self._state(
acks={"other-0": "ack", "ghost-1": "ack", "ghost-2": "ack"}
)
check = self.assertAcksFailClosed(report, state)
self.assertIn("other-1", check.detail)
self.assertIn("other-2", check.detail)
def test_requester_ack_never_satisfies_another_sessions_obligation(self):
"""The requester is excluded from the required set and stays excluded."""
report = _identity_report(requester="req", others=("other-0",))
requester_rows = [s for s in report["affected_sessions"] if s["is_requester"]]
self.assertEqual([s["session_id"] for s in requester_rows], ["req"])
self.assertNotIn("req", report["ack_state"])
check = self.assertAcksFailClosed(report, self._state(acks={"req": "ack"}))
self.assertIn("other-0", check.detail)
def test_fabricated_ids_do_not_count_toward_coverage(self):
report = _identity_report(requester="req", others=("other-0",))
for bogus in ("", " ", "other-0 extra", "OTHER-0", "other-01", "0"):
with self.subTest(bogus=bogus):
self.assertAcksFailClosed(report, self._state(acks={bogus: "ack"}))
# --- per-session state must be explicitly valid ------------------------
def test_unproven_per_session_states_fail_closed(self):
"""A required id present but not explicitly acknowledged fails closed."""
report = _identity_report(requester="req", others=("other-0", "other-1"))
for value in ("pending", "stale", "unknown", "", None, True, 1, NOW):
with self.subTest(value=value):
self.assertAcksFailClosed(
report,
self._state(acks={"other-0": "ack", "other-1": value}),
)
def test_report_ack_state_placeholder_is_never_read_as_an_ack(self):
"""``ack_state`` values are the report's own placeholders, not evidence."""
report = _identity_report(requester="req", others=("other-0",))
report["ack_state"] = {"other-0": "ack"}
self.assertAcksFailClosed(report, self._state())
# --- missing / malformed / contradictory identity evidence -------------
def test_missing_identity_evidence_fails_closed(self):
report = _identity_report(requester="req", others=("other-0",))
report.pop("ack_state", None)
report.pop("affected_sessions", None)
check = self.assertAcksFailClosed(report, self._state(acks={"other-0": "ack"}))
self.assertIn("no session-identity evidence", check.detail)
def test_malformed_ack_state_fails_closed(self):
for malformed in ([], "other-0", 7, None, ("other-0",)):
with self.subTest(malformed=malformed):
report = _identity_report(requester="req", others=("other-0",))
report["ack_state"] = malformed
self.assertAcksFailClosed(
report, self._state(acks={"other-0": "ack"})
)
def test_non_string_ack_state_key_fails_closed(self):
report = _identity_report(requester="req", others=("other-0",))
report["ack_state"] = {7: "pending"}
self.assertAcksFailClosed(report, self._state(acks={"other-0": "ack"}))
def test_malformed_affected_sessions_fails_closed(self):
for malformed in ("sessions", 7, {"session_id": "other-0"}, [None], [7]):
with self.subTest(malformed=malformed):
report = _identity_report(requester="req", others=("other-0",))
report.pop("ack_state", None)
report["affected_sessions"] = malformed
self.assertAcksFailClosed(
report, self._state(acks={"other-0": "ack"})
)
def test_affected_sessions_without_explicit_booleans_fails_closed(self):
"""``live``/``is_requester`` must be real booleans, never inferred."""
report = _identity_report(requester="req", others=("other-0",))
report.pop("ack_state", None)
for row in report["affected_sessions"]:
if row["session_id"] == "other-0":
row["is_requester"] = "false"
self.assertAcksFailClosed(report, self._state(acks={"other-0": "ack"}))
def test_affected_sessions_missing_live_flag_fails_closed(self):
report = _identity_report(requester="req", others=("other-0",))
report.pop("ack_state", None)
for row in report["affected_sessions"]:
row.pop("live", None)
self.assertAcksFailClosed(report, self._state(acks={"other-0": "ack"}))
def test_contradictory_ack_state_and_affected_sessions_fails_closed(self):
"""Both views present and disagreeing is unresolvable, not a tie-break."""
report = _identity_report(requester="req", others=("other-0", "other-1"))
report["ack_state"] = {"other-0": "pending", "other-9": "pending"}
check = self.assertAcksFailClosed(
report, self._state(acks={"other-0": "ack", "other-9": "ack"})
)
self.assertIn("contradicts itself", check.detail)
def test_identity_count_mismatch_fails_closed(self):
"""Identity evidence that cannot be reconciled with the count denies."""
report = _identity_report(requester="req", others=("other-0", "other-1"))
report["counts"] = dict(report["counts"], sessions_live_other=1)
check = self.assertAcksFailClosed(
report, self._state(acks={"other-0": "ack", "other-1": "ack"})
)
self.assertIn("cannot be reconciled", check.detail)
def test_broken_identity_evidence_outranks_timeout_policy(self):
"""The sanctioned timeout path cannot paper over an unreadable report."""
report = _identity_report(requester="req", others=("other-0",))
report["ack_state"] = "not-a-mapping"
self.assertAcksFailClosed(report, self._state(ack_timeout_policy_applied=True))
# --- legitimate success is preserved -----------------------------------
def test_every_required_session_acknowledged_passes(self):
report = _identity_report(
requester="req", others=("other-0", "other-1", "other-2")
)
state = self._state(
acks={
"other-0": "ack",
"other-1": "acked",
"other-2": "acknowledged",
}
)
proof = self._build(report, state)
check = self._acks_check(proof)
self.assertTrue(check.passed)
self.assertTrue(proof.clean)
self.assertEqual(proof.failed_checks, [])
self.assertIn("acknowledged by identity", check.detail)
decision = dp.gate_apply_restart(
proof=proof.as_dict(),
now=NOW,
secret=SECRET,
expected_impact_fingerprint=dp.impact_fingerprint(report),
)
self.assertEqual(decision.verdict, dp.GATE_ALLOW)
self.assertTrue(decision.allow)
def test_required_session_ack_tolerates_surrounding_whitespace(self):
report = _identity_report(requester="req", others=("other-0",))
proof = self._build(report, self._state(acks={" other-0 ": " ACK "}))
self.assertTrue(self._acks_check(proof).passed)
self.assertTrue(proof.clean)
def test_no_other_live_sessions_still_passes_with_identity_evidence(self):
report = _identity_report(requester="req", others=())
self.assertEqual(report["counts"]["sessions_live_other"], 0)
self.assertEqual(report["ack_state"], {})
proof = self._build(report, self._state())
check = self._acks_check(proof)
self.assertTrue(check.passed)
self.assertIn("sessions_live_other=0", check.detail)
self.assertTrue(proof.clean)
def test_explicit_timeout_policy_retains_intended_behavior(self):
"""Valid, correctly typed timeout evidence still permits the check."""
report = _identity_report(requester="req", others=("other-0", "other-1"))
proof = self._build(report, self._state(ack_timeout_policy_applied=True))
check = self._acks_check(proof)
self.assertTrue(check.passed)
self.assertIn("timeout policy", check.detail)
self.assertTrue(proof.clean)
def test_timeout_policy_still_strictly_typed_under_identity_binding(self):
report = _identity_report(requester="req", others=("other-0",))
for value in (None, "true", "True", 1, [], {}):
with self.subTest(value=value):
self.assertAcksFailClosed(
report, self._state(ack_timeout_policy_applied=value)
)
if __name__ == "__main__":
unittest.main()
@@ -0,0 +1,381 @@
"""``apply_authorized`` requires BOTH authorizations (#886 review blocker B1).
The #663 restart-class matrix and the #661 drain-proof hard gate are independent
authorizations that first coexisted when PR #882 landed on master and PR #886
merged it into the restart-class branch. The union preserved both, but the apply
decision consulted only the drain gate::
payload["apply_authorized"] = gate.allow # pre-fix
so a clean drain proof or an authorized break-glass, which needs no proof at
all reported ``apply_authorized: True`` for a restart class the least-privilege
matrix had just denied, in the same payload that carried
``allow_restart: False`` and "role 'author' may not request full_mcp_restart".
These tests pin the conjunction and the properties that must survive it. They
exercise the real MCP tool, which previously had no test coverage at all that
absence is why the defect shipped.
"""
from __future__ import annotations
import json
import os
import unittest
from unittest.mock import patch
import drain_proof
import gitea_mcp_server as srv
CONTROLLER_APPROVAL_ENV = "GITEA_CONTROLLER_RESTART_APPROVAL_AUTHORIZATION"
BREAK_GLASS_ENV = "GITEA_BREAKGLASS_RESTART_AUTHORIZATION"
# A quiet control plane: nothing live, so the blast radius never masks the
# authorization outcome under test.
QUIET_SESSIONS: list[dict] = []
QUIET_LEASES: list[dict] = []
class _FakeDB:
"""Minimal control-plane DB stand-in for the restart inventory."""
def __init__(self, sessions=QUIET_SESSIONS, terminal=None):
self._sessions = list(sessions)
self._terminal = terminal
def list_sessions(self, statuses=None, limit=None):
return list(self._sessions)
def get_active_terminal_lock(self, remote=None, org=None, repo=None):
return self._terminal
def _profile(role: str) -> dict:
return {"profile_name": f"prgs-{role}", "role_kind": role, "role": role}
class _RestartToolHarness(unittest.TestCase):
"""Drives the real ``gitea_request_mcp_restart`` with a stubbed inventory."""
def _call(self, *, role: str, env: dict | None = None, **kwargs) -> dict:
environ = {k: v for k, v in os.environ.items()
if k not in (CONTROLLER_APPROVAL_ENV, BREAK_GLASS_ENV)}
environ.update(env or {})
with patch.object(srv, "_profile_operation_gate", return_value=None), \
patch.object(srv, "_resolve",
return_value=("gitea.prgs.cc",
"Scaled-Tech-Consulting",
"Gitea-Tools")), \
patch.object(srv, "get_profile", return_value=_profile(role)), \
patch.object(srv, "_control_plane_db_or_error",
return_value=(_FakeDB(), [])), \
patch.object(srv.lease_lifecycle, "list_active_leases",
return_value={"leases": list(QUIET_LEASES)}), \
patch.dict(os.environ, environ, clear=True):
return srv.gitea_request_mcp_restart(
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
session_id="probe-session",
**kwargs,
)
def _clean_proof_for(self, preview: dict) -> str:
"""Mint a genuinely clean, signature-valid proof bound to *preview*.
Built from the tool's own dry-run report, so the fingerprint matches and
the proof is rejected for authorization reasons only never because it
was stale or forged.
"""
proof = drain_proof.build_drain_proof(
impact_report=preview,
drain_state={
"assignments_stopped": True,
"checkpoints_complete": True,
"handoffs_verified": True,
"leases_handled": True,
"acks": {},
},
requesting_session_id="probe-session",
)
self.assertTrue(proof.clean, "harness must mint a clean proof")
return json.dumps(proof.as_dict())
class TestConjunction(_RestartToolHarness):
"""AC1/AC2 — the two authorizations are ANDed, in both directions."""
def test_gate_allow_with_class_denied_yields_apply_authorized_false(self):
# An author may not request full_mcp_restart (CONTROL_ROLES only).
preview = self._call(role="author", restart_class="full_mcp_restart")
self.assertFalse(preview["allow_restart"])
result = self._call(
role="author",
restart_class="full_mcp_restart",
dry_run=False,
drain_proof_json=self._clean_proof_for(preview),
)
self.assertTrue(result["apply_gate"]["drain_gate_allow"],
"drain gate itself should have allowed this proof")
self.assertFalse(result["apply_gate"]["restart_class_authorized"])
self.assertFalse(result["apply_authorized"],
"a clean proof must not authorize a denied class")
self.assertFalse(result["allow_restart"])
def test_gate_allow_with_class_allowed_can_yield_apply_authorized_true(self):
preview = self._call(
role="operator",
restart_class="full_mcp_restart",
env={CONTROLLER_APPROVAL_ENV: "operator-approved"},
)
self.assertTrue(preview["allow_restart"],
"operator + controller approval must authorize the class")
result = self._call(
role="operator",
restart_class="full_mcp_restart",
dry_run=False,
drain_proof_json=self._clean_proof_for(preview),
env={CONTROLLER_APPROVAL_ENV: "operator-approved"},
)
self.assertTrue(result["apply_gate"]["drain_gate_allow"])
self.assertTrue(result["apply_gate"]["restart_class_authorized"])
self.assertTrue(result["apply_authorized"],
"both authorizations pass; apply must be authorized")
def test_denial_is_attributable_to_the_authorization_that_caused_it(self):
preview = self._call(role="author", restart_class="full_mcp_restart")
result = self._call(
role="author",
restart_class="full_mcp_restart",
dry_run=False,
drain_proof_json=self._clean_proof_for(preview),
)
blob = " ".join(result["apply_gate"]["reasons"]).lower()
self.assertIn("restart class authorization denied", blob)
self.assertIn("full_mcp_restart", blob)
class TestProofCannotOverrideAuthorization(_RestartToolHarness):
"""AC3 — a clean proof never overrides a class or requester-role denial."""
def test_clean_proof_cannot_override_role_denial(self):
for role in ("author", "reviewer", "merger", "reconciler"):
with self.subTest(role=role):
preview = self._call(role=role, restart_class="full_mcp_restart")
result = self._call(
role=role,
restart_class="full_mcp_restart",
dry_run=False,
drain_proof_json=self._clean_proof_for(preview),
)
self.assertFalse(result["apply_authorized"])
def test_clean_proof_cannot_override_missing_controller_approval(self):
# Correct role, but the class demands controller approval and the
# environment carries none.
preview = self._call(role="operator", restart_class="full_mcp_restart")
self.assertFalse(preview["allow_restart"])
result = self._call(
role="operator",
restart_class="full_mcp_restart",
dry_run=False,
drain_proof_json=self._clean_proof_for(preview),
)
self.assertFalse(result["apply_authorized"])
def test_clean_proof_cannot_override_unknown_class(self):
preview = self._call(role="operator", restart_class="not_a_real_class",
env={CONTROLLER_APPROVAL_ENV: "yes"})
self.assertFalse(preview["allow_restart"])
result = self._call(
role="operator",
restart_class="not_a_real_class",
dry_run=False,
drain_proof_json=self._clean_proof_for(preview),
env={CONTROLLER_APPROVAL_ENV: "yes"},
)
self.assertFalse(result["apply_authorized"])
def test_clean_proof_cannot_override_missing_scope_target(self):
# worker_restart without target_session_id fails closed on scoping.
preview = self._call(role="operator", restart_class="worker_restart",
env={CONTROLLER_APPROVAL_ENV: "yes"})
self.assertFalse(preview["allow_restart"])
result = self._call(
role="operator",
restart_class="worker_restart",
dry_run=False,
drain_proof_json=self._clean_proof_for(preview),
env={CONTROLLER_APPROVAL_ENV: "yes"},
)
self.assertFalse(result["apply_authorized"])
class TestBreakGlassDoesNotCollapseTheMatrix(_RestartToolHarness):
"""AC4 — break-glass bypasses the drain proof only, never the class matrix."""
def test_break_glass_does_not_authorize_a_denied_class(self):
result = self._call(
role="author",
restart_class="host_restart",
dry_run=False,
request_break_glass=True,
env={BREAK_GLASS_ENV: "operator-issued"},
)
self.assertTrue(result["break_glass_authorized"])
self.assertTrue(result["apply_gate"]["drain_gate_allow"],
"break-glass does satisfy the drain gate")
self.assertFalse(result["apply_gate"]["restart_class_authorized"])
self.assertFalse(result["apply_authorized"],
"break-glass must not collapse the class matrix")
def test_break_glass_across_every_worker_role_and_restricted_class(self):
for role in ("author", "reviewer", "merger", "reconciler"):
for klass in ("rolling_mcp_restart", "full_mcp_restart",
"host_restart"):
with self.subTest(role=role, restart_class=klass):
result = self._call(
role=role,
restart_class=klass,
dry_run=False,
request_break_glass=True,
env={BREAK_GLASS_ENV: "operator-issued"},
)
self.assertFalse(result["apply_authorized"])
def test_break_glass_still_works_when_the_class_is_authorized(self):
# Break-glass keeps its purpose: skipping the drain proof for a caller
# the matrix does allow.
result = self._call(
role="operator",
restart_class="full_mcp_restart",
dry_run=False,
request_break_glass=True,
env={BREAK_GLASS_ENV: "operator-issued",
CONTROLLER_APPROVAL_ENV: "operator-approved"},
)
self.assertTrue(result["apply_authorized"])
self.assertEqual(result["apply_gate"]["verdict"], "break_glass")
def test_break_glass_is_not_self_assertable(self):
# Requested but no environment authorization -> no bypass, and the
# unproven apply is denied.
result = self._call(
role="operator",
restart_class="full_mcp_restart",
dry_run=False,
request_break_glass=True,
env={CONTROLLER_APPROVAL_ENV: "operator-approved"},
)
self.assertTrue(result["break_glass_requested"])
self.assertFalse(result["break_glass_authorized"])
self.assertFalse(result["apply_authorized"])
self.assertIn("incident", result)
class TestRestrictedClassesStayDenied(_RestartToolHarness):
"""AC5 — restricted classes remain denied to unauthorized requesters."""
def test_restricted_classes_denied_for_worker_roles(self):
for role in ("author", "reviewer", "merger", "reconciler"):
for klass in ("rolling_mcp_restart", "full_mcp_restart",
"host_restart"):
with self.subTest(role=role, restart_class=klass):
preview = self._call(
role=role,
restart_class=klass,
env={CONTROLLER_APPROVAL_ENV: "yes"},
)
self.assertFalse(preview["allow_restart"])
self.assertFalse(preview["permission_authorized"])
self.assertFalse(preview["role_authorized"])
def test_host_restart_needs_controller_and_infrastructure_operator(self):
# controller approval alone is not enough for host_restart.
preview = self._call(role="controller", restart_class="host_restart",
env={CONTROLLER_APPROVAL_ENV: "yes"})
self.assertFalse(preview["approval_satisfied"])
self.assertFalse(preview["allow_restart"])
class TestExistingPathsStillWork(_RestartToolHarness):
"""AC6 — valid scoped and unscoped restart paths are unaffected."""
def test_dry_run_never_reports_apply_authorization(self):
result = self._call(role="operator", restart_class="full_mcp_restart",
env={CONTROLLER_APPROVAL_ENV: "yes"})
self.assertNotIn("apply_authorized", result)
self.assertNotIn("apply_gate", result)
self.assertFalse(result["apply_supported"])
self.assertFalse(result["restart_performed"])
def test_self_service_unscoped_classes_authorize_for_every_role(self):
for role in ("author", "reviewer", "merger", "reconciler",
"controller", "operator", "admin"):
for klass in ("client_reconnect", "session_reconnect"):
with self.subTest(role=role, restart_class=klass):
preview = self._call(role=role, restart_class=klass)
self.assertTrue(preview["allow_restart"])
def test_scoped_class_with_target_authorizes_and_applies(self):
env = {CONTROLLER_APPROVAL_ENV: "operator-approved"}
preview = self._call(role="operator", restart_class="worker_restart",
target_session_id="worker-1", env=env)
self.assertTrue(preview["allow_restart"])
result = self._call(
role="operator",
restart_class="worker_restart",
target_session_id="worker-1",
dry_run=False,
drain_proof_json=self._clean_proof_for(preview),
env=env,
)
self.assertTrue(result["apply_authorized"])
def test_apply_still_denies_without_any_proof(self):
# The #661 hard gate is untouched by the conjunction.
result = self._call(
role="operator",
restart_class="full_mcp_restart",
dry_run=False,
env={CONTROLLER_APPROVAL_ENV: "operator-approved"},
)
self.assertFalse(result["apply_gate"]["drain_gate_allow"])
self.assertTrue(result["apply_gate"]["restart_class_authorized"])
self.assertFalse(result["apply_authorized"])
self.assertEqual(result["incident"]["kind"], "restart_drain_gate_denied")
def test_apply_denies_on_malformed_proof(self):
result = self._call(
role="operator",
restart_class="full_mcp_restart",
dry_run=False,
drain_proof_json="{not valid json",
env={CONTROLLER_APPROVAL_ENV: "operator-approved"},
)
self.assertFalse(result["apply_authorized"])
self.assertTrue(any("invalid drain_proof_json" in reason
for reason in result["apply_gate"]["reasons"]))
def test_tool_never_restarts_on_any_path(self):
for kwargs in (
{"restart_class": "client_reconnect"},
{"restart_class": "full_mcp_restart", "dry_run": False},
{"restart_class": "host_restart", "dry_run": False,
"request_break_glass": True},
):
with self.subTest(**kwargs):
result = self._call(role="operator", env={
CONTROLLER_APPROVAL_ENV: "yes", BREAK_GLASS_ENV: "yes"},
**kwargs)
self.assertFalse(result["restart_performed"])
self.assertFalse(result["apply_supported"])
if __name__ == "__main__":
unittest.main()
+117
View File
@@ -105,3 +105,120 @@ def test_cross_links_do_not_embed_secrets():
text = _read(path)
for marker in ("ghp_", "BEGIN PRIVATE KEY", "Authorization: Bearer"):
assert marker not in text, f"{path} contains {marker!r}"
# --- Coordinator doc stays in lock-step with the tool (#886 review blocker B2) --
#
# PR #882 moved the #661 drain-proof hard gate *into* gitea_request_mcp_restart,
# but the coordinator document still described the proof as "a separate child"
# and omitted both new parameters. Nothing referenced that document, so nothing
# caught the drift. These tests bind the prose to the real signature.
COORDINATOR_DOC = REPO_ROOT / "docs" / "mcp-restart-coordinator.md"
# Affirmative claims that were accurate before #661 landed and are now false.
# Matched against whitespace-normalized text so re-wrapping cannot hide them.
# Deliberately not the bare phrase "a separate child": the corrected prose uses
# it in a negation ("no longer a separate child operation"), and a guard that
# forbids naming the old behaviour would block explaining that it changed.
STALE_PRE_661_PHRASES = (
"gated by a drain proof (a separate child)",
"is a later child gated by a drain proof",
"mutative apply path is explicitly out of scope",
"apply is gated by a drain proof (a separate child)",
)
def _documented_signature_block() -> str:
"""The fenced signature block for the tool, as published in the doc."""
text = _read(COORDINATOR_DOC)
marker = "gitea_request_mcp_restart("
start = text.index(marker)
end = text.index("```", start)
return text[start:end]
def test_documented_signature_matches_the_real_tool_signature():
import inspect
import gitea_mcp_server
block = _documented_signature_block()
real = inspect.signature(gitea_mcp_server.gitea_request_mcp_restart)
for name in real.parameters:
assert name in block, (
f"docs/mcp-restart-coordinator.md documents no {name!r} parameter; "
"the published signature has drifted from the tool"
)
def test_drain_proof_and_break_glass_parameters_are_documented():
block = _documented_signature_block()
for name in ("drain_proof_json", "request_break_glass"):
assert name in block, f"signature block missing {name}"
def test_restart_class_and_target_scoping_parameters_survive():
block = _documented_signature_block()
for name in ("restart_class", "target_session_id", "target_role",
"target_connector"):
assert name in block, f"signature block lost #663 parameter {name}"
def test_gate_is_documented_as_executing_inside_this_tool():
lower = _read(COORDINATOR_DOC).lower()
assert "inside this tool" in lower, (
"the coordinator doc must state that the drain-proof gate executes in "
"gitea_request_mcp_restart, not in a later child"
)
assert "no longer a separate child operation" in lower
def test_stale_pre_661_wording_cannot_return():
normalized = " ".join(_read(COORDINATOR_DOC).split()).lower()
for phrase in STALE_PRE_661_PHRASES:
assert phrase not in normalized, (
f"stale pre-#661 wording returned to the coordinator doc: {phrase!r}"
)
def test_dry_run_versus_apply_behavior_is_documented():
lower = _read(COORDINATOR_DOC).lower()
assert "dry_run=true" in lower and "dry_run=false" in lower
assert "apply_supported" in lower and "restart_performed" in lower
assert "never restarts anything" in lower
def test_authorization_ordering_and_conjunction_are_documented():
text = _read(COORDINATOR_DOC)
lower = text.lower()
assert "authorization ordering" in lower
assert "allow_restart" in text
assert "apply_authorized" in text
# The conjunction itself, and the attribution fields behind it.
assert "gate.allow and allow_restart" in text
for field in ("drain_gate_allow", "restart_class_authorized"):
assert field in text, f"doc omits apply_gate.{field}"
def test_break_glass_scope_is_documented_as_drain_proof_only():
text = _read(COORDINATOR_DOC)
lower = text.lower()
assert "break-glass" in lower
assert "drain proof only" in lower, (
"doc must state break-glass never bypasses the restart-class matrix"
)
assert "GITEA_BREAKGLASS_RESTART_AUTHORIZATION" in text
def test_fail_closed_on_apply_is_documented():
lower = _read(COORDINATOR_DOC).lower()
assert "fail closed" in lower
for condition in ("expired", "unclean", "tampered", "stale"):
assert condition in lower, f"fail-closed list omits {condition!r}"
def test_coordinator_doc_embeds_no_secrets():
text = _read(COORDINATOR_DOC)
for marker in ("ghp_", "BEGIN PRIVATE KEY", "Authorization: Bearer"):
assert marker not in text, f"{COORDINATOR_DOC} contains {marker!r}"
+232
View File
@@ -0,0 +1,232 @@
"""Permission, drain, routing, and audit matrix for restart classes (#663)."""
from __future__ import annotations
import os
from datetime import datetime, timezone
import restart_coordinator as rc
NOW = datetime(2026, 7, 24, 20, 0, tzinfo=timezone.utc)
def _inventory() -> dict:
return {
"inventory_complete": True,
"sessions": [
{
"session_id": "requester",
"role": "author",
"profile": "prgs-author",
"pid": os.getpid(),
"status": "active",
"last_heartbeat_at": NOW.isoformat(),
},
{
"session_id": "reviewer",
"role": "reviewer",
"profile": "prgs-reviewer",
"pid": os.getpid(),
"status": "active",
"last_heartbeat_at": NOW.isoformat(),
},
],
"leases": [
{
"lease_id": "review-lease",
"session_id": "reviewer",
"role": "reviewer",
"phase": "reviewing",
"work_kind": "pr",
"work_number": 900,
"worktree_path": "/tmp/review-900",
"freshness": {"freshness": "active"},
}
],
}
def _evaluate(
restart_class: rc.RestartClass,
*,
role: str = "controller",
permissions: tuple[str, ...] | None = None,
approved: bool = True,
operator: bool = True,
**targets,
):
return rc.evaluate_restart_impact(
_inventory(),
now=NOW,
requesting_session_id="requester",
restart_class=restart_class,
requester_role=role,
requester_permissions=(
permissions if permissions is not None
else rc.permissions_for_role(role)
),
controller_approved=approved,
operator_authorized=operator,
**targets,
)
def test_policy_table_covers_exactly_all_nine_classes():
assert set(rc.RESTART_CLASS_POLICIES) == set(rc.RestartClass)
assert len(rc.RESTART_CLASS_POLICIES) == 9
for restart_class, policy in rc.RESTART_CLASS_POLICIES.items():
assert policy.restart_class is restart_class
assert policy.required_permission
assert policy.expected_blast_radius in {
rc.BLAST_NONE, rc.BLAST_LOW, rc.BLAST_MEDIUM, rc.BLAST_HIGH
}
assert policy.drain_requirement
assert policy.approval_requirement
assert policy.audit_requirement
assert policy.recovery_behavior
def test_permission_matrix_allows_each_class_with_exact_permission():
targets = {
rc.RestartClass.WORKER_RESTART: {"target_session_id": "reviewer"},
rc.RestartClass.ROLE_RUNTIME_RESTART: {"target_role": "reviewer"},
rc.RestartClass.CONNECTOR_RESTART: {"target_connector": "github"},
}
for restart_class, policy in rc.RESTART_CLASS_POLICIES.items():
report = _evaluate(
restart_class,
permissions=(policy.required_permission,),
**targets.get(restart_class, {}),
)
assert report.permission_authorized, restart_class
assert report.role_authorized, restart_class
assert report.approval_satisfied, restart_class
assert report.audit_record["restart_class"] == restart_class.value
assert (
report.audit_record["required_permission"]
== policy.required_permission
)
def test_missing_or_nearby_permission_denies():
report = _evaluate(
rc.RestartClass.ROLE_RUNTIME_RESTART,
permissions=("mcp.restart.worker.request",),
target_role="reviewer",
)
assert report.verdict == rc.VERDICT_UNSAFE
assert not report.allow_restart
assert not report.permission_authorized
assert any("missing required permission" in r for r in report.reasons)
def test_unknown_restart_class_denies_fail_closed():
report = rc.evaluate_restart_impact(
_inventory(),
now=NOW,
restart_class="surprise_reboot",
requester_role="admin",
requester_permissions=("mcp.restart.host.request",),
controller_approved=True,
operator_authorized=True,
)
assert report.verdict == rc.VERDICT_UNSAFE
assert not report.allow_restart
assert report.restart_policy == {}
assert any("unknown restart class" in r for r in report.reasons)
def test_worker_roles_cannot_request_full_or_host_restart():
for role in rc.WORKER_ROLES:
granted = rc.permissions_for_role(role)
assert "mcp.restart.full.request" not in granted
assert "mcp.restart.host.request" not in granted
report = _evaluate(
rc.RestartClass.FULL_MCP_RESTART,
role=role,
permissions=granted,
)
assert not report.role_authorized
assert not report.allow_restart
def test_controller_approval_is_independent_of_permission():
report = _evaluate(
rc.RestartClass.WORKER_RESTART,
approved=False,
target_session_id="reviewer",
)
assert report.permission_authorized
assert not report.approval_satisfied
assert not report.allow_restart
def test_narrow_classes_do_not_inherit_full_drain_or_peer_lease_block():
for restart_class in (
rc.RestartClass.CLIENT_RECONNECT,
rc.RestartClass.SESSION_RECONNECT,
rc.RestartClass.CONFIGURATION_RELOAD,
):
report = _evaluate(restart_class)
assert not report.restart_policy["full_drain_required"]
assert report.counts["leases_disruptive"] == 0
assert report.counts["sessions_live_other"] == 0
assert report.counts["critical_sections"] == 0
assert report.counts["mutations"] == 0
assert report.allow_restart, (restart_class, report.reasons)
def test_client_reconnect_does_not_wait_for_unrelated_terminal_lock():
inventory = _inventory()
inventory["terminal_lock"] = {"terminal_pr": 901}
report = rc.evaluate_restart_impact(
inventory,
now=NOW,
requesting_session_id="requester",
restart_class=rc.RestartClass.CLIENT_RECONNECT,
requester_role="author",
requester_permissions=rc.permissions_for_role("author"),
)
assert report.allow_restart
assert report.terminal_lock is None
def test_scoped_restart_only_counts_named_target():
report = _evaluate(
rc.RestartClass.ROLE_RUNTIME_RESTART,
target_role="author",
)
assert report.counts["leases_disruptive"] == 0
assert report.affected_prs == []
assert report.allow_restart
reviewer = _evaluate(
rc.RestartClass.ROLE_RUNTIME_RESTART,
target_role="reviewer",
)
assert reviewer.counts["leases_disruptive"] == 1
assert reviewer.affected_prs == [900]
assert not reviewer.allow_restart
def test_missing_scoped_target_denies_instead_of_widening():
for restart_class in (
rc.RestartClass.WORKER_RESTART,
rc.RestartClass.ROLE_RUNTIME_RESTART,
rc.RestartClass.CONNECTOR_RESTART,
):
report = _evaluate(restart_class)
assert not report.allow_restart
assert any("target required" in r for r in report.reasons)
def test_only_full_and_host_classes_require_full_drain():
requiring_full = {
restart_class
for restart_class, policy in rc.RESTART_CLASS_POLICIES.items()
if policy.full_drain_required
}
assert requiring_full == {
rc.RestartClass.FULL_MCP_RESTART,
rc.RestartClass.HOST_RESTART,
}
File diff suppressed because it is too large Load Diff
+6
View File
@@ -444,6 +444,12 @@ class TestAuditEmission(unittest.TestCase):
)
self.assertEqual(record["target"]["namespace"], NAMESPACE)
self.assertEqual(record["target"]["mode"], "restart")
self.assertEqual(
record["target"]["restart_class"], "role_runtime_restart"
)
self.assertEqual(
record["metadata"]["restart_class"], "role_runtime_restart"
)
self.assertEqual(record["result"], console_audit.RESULT_ALLOWED)
self.assertEqual(record["actor"]["subject"], "[email protected]")
self.assertFalse(record["metadata"]["process_kill_executed"])
+116
View File
@@ -67,6 +67,8 @@ from webui.system_health import (
snapshot_to_dict as system_health_to_dict,
)
from webui.system_health_views import render_system_health_page
from webui import request_service
from webui.request_views import render_requests_page
_READ_ONLY_METHODS = frozenset({"GET", "HEAD", "OPTIONS"})
_AUDIT_MUTATION_PATHS = frozenset({"/audit", "/api/audit"})
@@ -722,6 +724,109 @@ async def api_v1_analytics_ingest(request: Request) -> JSONResponse:
)
def _default_request_scope() -> dict[str, str]:
"""Resolve remote/org/repo from the project registry for request forms.
Returns an empty mapping when the registry cannot be read, which makes
``parse_request`` reject a request that did not name its own scope rather
than letting it default to some other repository.
"""
from webui.queue_loader import _host_from_url # host normalisation helper
registry, error = _load_project_registry()
if error is not None or not registry.projects:
return {}
project = registry.projects[0]
host = _host_from_url(project.remote_host)
return {
"remote": _derive_remote(host),
"org": project.gitea_owner or "",
"repo": project.repo_name or "",
}
async def _request_payload(request: Request) -> dict[str, object]:
"""Read a request body as JSON or form-encoded. Never raises."""
content_type = (request.headers.get("content-type") or "").lower()
if "application/json" in content_type:
try:
body = await request.json()
except Exception:
return {}
return dict(body) if isinstance(body, dict) else {}
try:
form = await request.form()
except Exception:
return {}
return {key: form[key] for key in form}
async def requests_page(request: Request) -> HTMLResponse:
"""Operator request form and intent preview (#643).
POST here only ever *previews*. Initiation is a separate confirmed call to
``/api/v1/requests/apply`` so that submitting this form cannot reserve
work as a side effect.
"""
submitted: dict[str, object] = {}
preview = None
error = None
if request.method == "POST":
submitted = await _request_payload(request)
work_request, error = request_service.parse_request(
submitted, default_scope=_default_request_scope()
)
if work_request is not None:
preview = request_service.preview_request(
work_request,
principal=resolve_principal(headers=dict(request.headers)),
)
return HTMLResponse(
render_requests_page(
preview=preview, error=error, submitted=submitted
)
)
async def api_v1_request_preview(request: Request) -> JSONResponse:
"""Dry-run authorization and intent preview for a work request (#643)."""
payload = await _request_payload(request)
work_request, error = request_service.parse_request(
payload, default_scope=_default_request_scope()
)
if work_request is None:
return JSONResponse(error.to_dict(), status_code=400)
preview = request_service.preview_request(
work_request,
principal=resolve_principal(headers=dict(request.headers)),
)
return JSONResponse(
preview.to_dict(), status_code=200 if preview.authorized else 403
)
async def api_v1_request_apply(request: Request) -> JSONResponse:
"""Initiate a previewed work request through the allocator (#643).
Fail-closed at every step: unauthorized, unconfirmed, not-next-safe, and
already-claimed all return without attempting an assignment.
"""
payload = await _request_payload(request)
work_request, error = request_service.parse_request(
payload, default_scope=_default_request_scope()
)
if work_request is None:
return JSONResponse(error.to_dict(), status_code=400)
confirm = _truthy_flag(str(payload.get("confirm") or ""))
result = request_service.apply_request(
work_request,
principal=resolve_principal(headers=dict(request.headers)),
confirm=confirm,
)
status = int(result.pop("status_code", 403))
return JSONResponse(result, status_code=status)
async def method_not_allowed(request: Request, _exc: Exception) -> Response:
path = request.url.path
if path in _AUDIT_MUTATION_PATHS and request.method == "POST":
@@ -786,6 +891,17 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
api_action_attempt,
methods=["POST"],
),
Route("/requests", requests_page, methods=["GET", "POST"]),
Route(
"/api/v1/requests/preview",
api_v1_request_preview,
methods=["POST"],
),
Route(
"/api/v1/requests/apply",
api_v1_request_apply,
methods=["POST"],
),
Route("/api/leases", api_leases, methods=["GET"]),
Route("/api/v1/inventory", api_inventory, methods=["GET"]),
Route(
+71 -8
View File
@@ -115,6 +115,12 @@ class ConsoleAction:
break_glass: bool
phase: int
summary: str
# Opt-in switch for an action whose execution path is genuinely wired
# ahead of its phase becoming globally active (#643). Naming a variable
# here enables nothing on its own: the variable must also be set in the
# environment. An action that leaves this ``None`` can only execute once
# ACTIVE_PHASE reaches its phase, exactly as before.
execution_env_flag: str | None = None
@property
def mcp_permission(self) -> str:
@@ -277,6 +283,27 @@ _ACTION_SPECS: tuple[ConsoleAction, ...] = (
phase=2,
summary="Restart one MCP namespace via the host supervisor.",
),
# #643: submit a work request — desired role, issue/PR, intent — and let
# the allocator reserve it. This is the one Phase 2 action whose execution
# path is actually implemented (``webui.request_service``), so it carries
# the opt-in flag; it stays denied until an operator sets that variable.
# Authority is operator-class because the outcome is a claim, not a Gitea
# verdict: initiating reviewer or merger *work* does not grant the right
# to approve or merge, which stays with the MCP role profile.
ConsoleAction(
action_id="initiate_workflow",
task_key="allocate_next_work",
action_class=CLASS_WRITE,
minimum_role=OPERATOR,
requires_confirmation=True,
dual_control=False,
break_glass=False,
phase=2,
summary=(
"Preview and initiate allocator-owned workflow work for a role."
),
execution_env_flag="WEBUI_REQUESTS_EXECUTION",
),
)
ACTIONS: dict[str, ConsoleAction] = {a.action_id: a for a in _ACTION_SPECS}
@@ -430,6 +457,33 @@ ALLOW_PREVIEW = "allowed_preview_only"
# gated on this model landing; nothing here enables it.
ACTIVE_PHASE = 1
_TRUTHY = frozenset({"1", "true", "yes", "on"})
def execution_wired(
action: ConsoleAction | None, env: dict[str, str] | None = None
) -> bool:
"""Whether *action* has a live execution path right now.
Two ways to be wired, and only two. The action's phase is active, or the
action declares an opt-in environment variable *and* that variable is set.
Everything else including every action that never declares a flag is
unwired, so the default across the registry stays deny.
Bumping ``ACTIVE_PHASE`` would enable execution for every action of that
phase at once. The per-action flag exists so a single implemented action
can go live without dragging its unimplemented phase-mates with it.
"""
if action is None:
return False
if action.phase <= ACTIVE_PHASE:
return True
flag = (action.execution_env_flag or "").strip()
if not flag:
return False
source = env if env is not None else os.environ
return (source.get(flag) or "").strip().lower() in _TRUTHY
@dataclass(frozen=True)
class AuthorizationDecision:
@@ -469,16 +523,19 @@ def authorize(
principal: Principal | None = None,
*,
for_execution: bool = False,
env: dict[str, str] | None = None,
) -> AuthorizationDecision:
"""Decide whether *principal* may invoke *action_id*. Deny by default.
``for_execution`` distinguishes a read-only preview from a real invocation.
Even an allowed decision reports ``execution_enabled=False`` while the
console is in Phase 1, so no caller can read an allow as permission to
mutate.
``execution_enabled`` reports whether the action has a live execution path
at all (:func:`execution_wired`) for every action without an explicit
opt-in flag that stays ``False`` while the console is in Phase 1, so no
caller can read an allow as permission to mutate.
"""
who = principal if principal is not None else ANONYMOUS
action = get_action(action_id)
wired = execution_wired(action, env)
if action is None:
return AuthorizationDecision(
@@ -497,7 +554,7 @@ def authorize(
"requires_confirmation": action.requires_confirmation,
"dual_control": action.dual_control,
"break_glass": action.break_glass,
"execution_enabled": False,
"execution_enabled": wired,
}
if not who.authenticated:
@@ -530,13 +587,19 @@ def authorize(
**base,
)
if for_execution and action.phase > ACTIVE_PHASE:
if for_execution and not wired:
return AuthorizationDecision(
allowed=False,
reason_code=DENY_PHASE_NOT_ACTIVE,
detail=(
f"Action {action_id!r} belongs to phase {action.phase}; the "
f"console is in phase {ACTIVE_PHASE}. Execution is not wired."
f"console is in phase {ACTIVE_PHASE}"
+ (
f" and {action.execution_env_flag} is not set"
if action.execution_env_flag
else ""
)
+ ". Execution is not wired."
),
**base,
)
@@ -545,8 +608,8 @@ def authorize(
allowed=True,
reason_code=ALLOW_PREVIEW,
detail=(
"Principal holds the required role. Preview only — execution "
"remains disabled until the Phase 2 action framework ships."
"Principal holds the required role. Execution proceeds only for an "
"action with a wired execution path; everything else is preview."
),
**base,
)
+1
View File
@@ -45,6 +45,7 @@ NAV_GROUPS: tuple[NavGroup, ...] = (
NavItem("/queue", "Queue"),
NavItem("/leases", "Leases"),
NavItem("/actions", "Actions"),
NavItem("/requests", "Requests"),
)),
NavGroup("Runtime/Sessions", (
NavItem("/runtime", "Runtime health"),
File diff suppressed because it is too large Load Diff
+164
View File
@@ -0,0 +1,164 @@
"""HTML views for the operator request surface (#643).
The form is deliberately a *preview* form. It has no initiate button, because
initiating requires a confirmed POST to ``/api/v1/requests/apply`` and a stray
form submission must not be able to produce one by accident.
Nothing rendered here is trusted input: every interpolated value is escaped,
and the page renders only values the service already produced rather than
echoing a raw request body back.
"""
from __future__ import annotations
import html
import json
from typing import Any
from webui.layout import render_page
from webui.request_service import (
REQUESTABLE_ROLES,
WORK_KINDS,
RequestError,
RequestPreview,
)
REQUESTS_PATH = "/requests"
PREVIEW_API_PATH = "/api/v1/requests/preview"
APPLY_API_PATH = "/api/v1/requests/apply"
def _escape(text: Any) -> str:
return html.escape(str(text if text is not None else ""), quote=True)
REQUEST_PAGE_STYLES = """
<style>
.request-form { display: grid; gap: 0.75rem; max-width: 44rem; }
.request-form label { display: grid; gap: 0.25rem; font-size: 0.9rem; }
.request-check { margin: 0.35rem 0; }
.request-check .verdict-ok { color: var(--accent); }
.request-check .verdict-fail { color: #d14; }
.request-prohibited code { margin-right: 0.4rem; }
</style>
"""
def _options(values: tuple[str, ...], selected: Any) -> str:
return "".join(
f"<option value='{_escape(value)}'"
+ (" selected" if selected == value else "")
+ f">{_escape(value)}</option>"
for value in values
)
def _form(values: dict[str, Any] | None = None) -> str:
current = dict(values or {})
number = current.get("work_number")
return (
f"<form class='request-form' method='post' action='{REQUESTS_PATH}'>"
"<label>Desired role<select name='desired_role'>"
f"{_options(REQUESTABLE_ROLES, current.get('desired_role'))}"
"</select></label>"
"<label>Work kind<select name='work_kind'>"
f"{_options(WORK_KINDS, current.get('work_kind'))}"
"</select></label>"
"<label>Issue or PR number"
"<input type='number' name='work_number' min='1' required "
f"value='{_escape(number) if number else ''}'></label>"
"<label>Intent summary"
"<input type='text' name='intent_summary' maxlength='500' required "
f"value='{_escape(current.get('intent_summary'))}'></label>"
"<label>Expected head SHA <span class='muted'>(PR work only)</span>"
"<input type='text' name='expected_head_sha' "
f"value='{_escape(current.get('expected_head_sha'))}'></label>"
"<button type='submit' class='copy-btn'>Preview request</button>"
"<p class='muted meta'>Preview is read-only and creates no assignment. "
f"Initiating requires a confirmed POST to <code>{APPLY_API_PATH}</code>."
"</p>"
"</form>"
)
def _checks_block(preview: RequestPreview) -> str:
rows = []
for check in preview.checks:
verdict = "PASS" if check.ok else "FAIL"
css = "verdict-ok" if check.ok else "verdict-fail"
rows.append(
"<li class='request-check'>"
f"<span class='{css}'><strong>{verdict}</strong></span> "
f"<code>{_escape(check.name)}</code> — {_escape(check.detail)} "
f"<span class='muted meta'>({_escape(check.reason_code)})</span>"
"</li>"
)
return "<ul>" + "".join(rows) + "</ul>"
def _preview_block(preview: RequestPreview) -> str:
verdict = "AUTHORIZED" if preview.authorized else "DENIED"
prohibited = "".join(
f"<code>{_escape(action)}</code>" for action in preview.prohibited_actions
)
request = preview.request
evidence = json.dumps(preview.allocator_evidence, indent=2, default=str)
return (
"<h3>Intent preview</h3>"
f"<p><strong>{verdict}</strong> — {_escape(preview.detail)}</p>"
"<p class='meta'>"
f"Role <code>{_escape(request.desired_role)}</code> · "
f"{_escape(request.work_kind)} <code>{_escape(request.display_ref)}</code>"
f" · profile <code>{_escape(preview.required_profile)}</code> · "
f"namespace <code>{_escape(preview.required_namespace)}</code> · "
f"permission <code>{_escape(preview.required_permission)}</code>"
"</p>"
f"<p>Intent: {_escape(request.intent_summary)}</p>"
f"{_checks_block(preview)}"
f"<p><strong>Next safe action:</strong> "
f"{_escape(preview.next_safe_action)}</p>"
"<p class='request-prohibited'><strong>Prohibited for this role:</strong> "
+ (prohibited or "<span class='muted'>none declared</span>")
+ "</p>"
"<p class='muted meta'>Correlation id "
f"<code>{_escape(preview.correlation_id)}</code></p>"
"<details><summary>Allocator evidence</summary>"
f"<pre class='prompt-text'>{_escape(evidence)}</pre>"
"</details>"
)
def _error_block(error: RequestError) -> str:
field = (
f"<p class='meta'>Field: <code>{_escape(error.field_name)}</code></p>"
if error.field_name
else ""
)
return (
"<h3>Request rejected</h3>"
f"<p><strong>{_escape(error.reason_code)}</strong> — "
f"{_escape(error.detail)}</p>{field}"
)
def render_requests_page(
*,
preview: RequestPreview | None = None,
error: RequestError | None = None,
submitted: dict[str, Any] | None = None,
) -> str:
"""Render the request form, plus a preview or rejection when one exists."""
body = (
"<h2>Requests</h2>"
"<p>Submit a work request — desired role, issue or PR, and intent — "
"and see whether it would be authorized before anything is reserved. "
"Initiation goes through the allocator (#600/#613); this console never "
"self-selects work, never approves, and never merges.</p>"
+ _form(submitted)
+ (_error_block(error) if error is not None else "")
+ (_preview_block(preview) if preview is not None else "")
+ f"<p class='meta'><a href='{PREVIEW_API_PATH}'>Preview API</a> · "
"<a href='/api/console/security-model'>RBAC model</a></p>"
+ REQUEST_PAGE_STYLES
)
return render_page(title="Requests", body_html=body)
+19 -1
View File
@@ -38,6 +38,7 @@ from dataclasses import asdict, dataclass
from typing import Any
import mcp_namespace_health
import restart_coordinator
import runtime_recovery_guard
from webui import console_audit, console_authz
@@ -99,6 +100,14 @@ def _clean(value: Any) -> str:
return str(value or "").strip()
def restart_class_for_mode(mode: str) -> str:
"""Map the existing namespace controls onto the #663 class taxonomy."""
if _clean(mode) == MODE_RELOAD:
return restart_coordinator.RestartClass.CONFIGURATION_RELOAD.value
return restart_coordinator.RestartClass.ROLE_RUNTIME_RESTART.value
# --- Mutation ledger --------------------------------------------------------
@@ -256,6 +265,7 @@ def build_restart_preview(
return {
"action_id": action_id,
"restart_class": restart_class_for_mode(md),
"namespace": ns,
"mode": md,
"scope_valid": scope_error is None,
@@ -309,6 +319,7 @@ def assess_restart_request(
"reason_code": reason_code,
"detail": detail,
"action_id": action_id,
"restart_class": restart_class_for_mode(md),
"namespace": ns,
"mode": md,
"preview": preview,
@@ -393,6 +404,7 @@ def assess_restart_request(
"process."
),
"action_id": action_id,
"restart_class": restart_class_for_mode(md),
"namespace": ns,
"mode": md,
"preview": preview,
@@ -441,7 +453,11 @@ def execute_restart(
else console_audit.RESULT_DENIED
),
principal=principal,
target={"namespace": assessment["namespace"], "mode": assessment["mode"]},
target={
"namespace": assessment["namespace"],
"mode": assessment["mode"],
"restart_class": assessment["restart_class"],
},
reason_code=assessment["reason_code"],
detail=assessment["detail"],
request_id=request_id,
@@ -450,6 +466,7 @@ def execute_restart(
"gates_passed": assessment["gates_passed"],
"process_kill_executed": False,
"post_restart_verification_required": True,
"restart_class": assessment["restart_class"],
},
)
@@ -463,6 +480,7 @@ def execute_restart(
"namespace": assessment["namespace"],
"mode": assessment["mode"],
"action_id": action_id,
"restart_class": assessment["restart_class"],
"process_kill_executed": False,
"host_hook": assessment["preview"]["restart_hook"],
"next_action": (
+10
View File
@@ -201,6 +201,16 @@ def _candidates_from_queue_snapshot(q_snap: QueueSnapshot) -> list[WorkCandidate
return candidates
def candidates_from_queue_snapshot(q_snap: QueueSnapshot) -> list[WorkCandidate]:
"""Public alias for :func:`_candidates_from_queue_snapshot` (#643).
The request-initiation service ranks the same candidate set this view
renders, so both must agree on how a queue row becomes a candidate. One
construction, two callers not two that can drift apart.
"""
return _candidates_from_queue_snapshot(q_snap)
def _claim_lease_records(inventory: dict[str, Any] | None) -> list[dict[str, Any]]:
"""Normalize ``build_claim_inventory`` entries into lease records.