Compare commits

..
Author SHA1 Message Date
jcwalker3andClaude Opus 4.8 2066623986 fix(bootstrap): allow author worktree bootstrap from clean control checkout (Closes #892)
Align assess_author_issue_bootstrap with bootstrap_permits_control_checkout
so gitea_bootstrap_author_issue_worktree can create the first branches/
worktree without the lock↔worktree deadlock.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-25 18:27:18 -05:00
sysadmin 2b4e43042a Merge pull request 'feat(tests): add concurrent-session MCP restart safety tests (Closes #666)' (#910) from feat/issue-666-concurrent-mcp-restart-tests into master 2026-07-25 17:44:09 -05:00
sysadmin 0f9390aab4 Merge remote-tracking branch 'prgs/master' into feat/issue-666-concurrent-mcp-restart-tests 2026-07-25 18:43:28 -04:00
sysadmin 59aab06fe1 feat(tests): add concurrent-session MCP restart safety tests (Closes #666) 2026-07-25 17:14:14 -04:00
11 changed files with 869 additions and 1108 deletions
-41
View File
@@ -26,7 +26,6 @@ from dataclasses import dataclass, field
from datetime import datetime, timezone
from typing import Any, Mapping, Sequence
import maintenance_drain
from control_plane_db import (
ControlPlaneDB,
ControlPlaneError,
@@ -937,46 +936,6 @@ def allocate_next_work(
"allocation_mode": (allocation_mode or "").strip() or None,
}
# #659 AC2: while maintenance drain is active, no new work is assigned —
# for dry-run and apply alike, so a preview can never be read as evidence
# that work was assignable during the drain. Checked before session
# registration so a drained allocator leaves no new state behind.
try:
drain_record = db.read_maintenance_drain(remote=remote, org=org, repo=repo)
except Exception as exc: # noqa: BLE001 — unreadable drain state fails closed
return {
"success": False,
"outcome": OUTCOME_NO_SAFE,
"reasons": [
f"maintenance-drain state lookup failed: {exc} (fail closed, #659)"
],
"skipped": [],
"assignment": None,
"substrate": "control_plane_db",
}
drain_decision = maintenance_drain.classify_assignment(drain_record)
if not drain_decision["assignment_allowed"]:
return {
"success": True,
"outcome": OUTCOME_WAIT,
"apply": apply,
"role": role_norm,
"allocation_mode": mode,
"remote": remote,
"org": org,
"repo": repo,
"selected": None,
"reasons": list(drain_decision["reasons"]),
"reason_code": drain_decision["reason_code"],
"skipped": [],
"assignment": None,
"substrate": "control_plane_db",
"maintenance_drain": maintenance_drain.status_payload(
drain_record, remote=remote, org=org, repo=repo
),
}
# 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:
+157 -46
View File
@@ -386,6 +386,68 @@ def run_compensating_recovery(
return recovery_info
def _normalize_sha(value: str | None) -> str | None:
"""Normalize a Git object id for comparison, or ``None`` when unknown."""
normalized = (value or "").strip().lower()
return normalized or None
def _author_bootstrap_assessment(
*,
not_applicable: bool,
allowed: bool,
block: bool,
reasons: list[str],
workspace: str,
root: str,
branch: str | None,
dirty: list[str],
under_branches: bool,
bootstrap_path: str | None = None,
local_head_sha: str | None = None,
remote_master_sha: str | None = None,
exact_next_action: str | None = None,
) -> dict[str, Any]:
"""Structured author-bootstrap assessment consumable by bootstrap_permits (#892).
Field shape mirrors :func:`create_issue_bootstrap._result` so the shared
``bootstrap_permits_control_checkout`` predicate can prove control-checkout
eligibility for ``gitea_bootstrap_author_issue_worktree`` the same way it
does for ``create_issue``. Allowed control assessments must use empty
``reasons`` — narrative belongs in other fields, not the refusal list.
"""
local_tip = _normalize_sha(local_head_sha)
remote_tip = _normalize_sha(remote_master_sha)
base_tips_verified = bool(local_tip and remote_tip and local_tip == remote_tip)
return {
"not_applicable": not_applicable,
"allowed": allowed,
"block": block,
"proven": bool(allowed and not block and not not_applicable),
"reasons": list(reasons),
"workspace_path": workspace,
"canonical_repo_root": root,
"current_branch": branch,
"dirty_files": list(dirty),
"under_branches": under_branches,
"exact_next_action": exact_next_action,
"bootstrap_path": bootstrap_path,
"task_scope": "author_issue_bootstrap",
"local_head_sha": local_tip,
"remote_master_sha": remote_tip,
"base_tips_verified": base_tips_verified,
}
EXACT_NEXT_ACTION_AUTHOR_BOOTSTRAP = (
"Restore the canonical control checkout to a clean accepted base branch "
"(master/main/dev) that matches live master, with no tracked local edits. "
"Re-resolve bootstrap_author_issue_worktree, then re-run "
"gitea_bootstrap_author_issue_worktree from that clean control checkout. "
"Do not use shell git worktree add as the primary path once bootstrap is healthy."
)
def assess_author_issue_bootstrap(
*,
workspace_path: str,
@@ -397,7 +459,13 @@ def assess_author_issue_bootstrap(
remote_master_sha_error: str | None = None,
task: str | None = None,
) -> dict[str, Any]:
"""Assess whether author issue worktree bootstrap may proceed from control or worktree root."""
"""Assess whether author issue worktree bootstrap may proceed from control or worktree root.
#892: control-checkout successes emit the full field set required by
``create_issue_bootstrap.bootstrap_permits_control_checkout`` (empty reasons,
task_scope, base tip proof, binding paths) so the #274/#604 guards can
waive control-checkout for this one sanctioned bootstrap task.
"""
root = os.path.realpath(canonical_repo_root or "")
workspace = os.path.realpath(workspace_path or root or ".")
branch = (current_branch or "").strip()
@@ -407,34 +475,50 @@ def assess_author_issue_bootstrap(
if root
else False
)
local_tip = _normalize_sha(head_sha)
remote_tip = _normalize_sha(remote_master_sha)
if not is_author_issue_bootstrap_task(task):
return {
"not_applicable": True,
"allowed": False,
"block": False,
"proven": False,
"reasons": ["task is not author_issue_bootstrap"],
}
return _author_bootstrap_assessment(
not_applicable=True,
allowed=False,
block=False,
reasons=["task is not author_issue_bootstrap"],
workspace=workspace,
root=root,
branch=branch or None,
dirty=dirty,
under_branches=under_branches,
)
# Already under branches/: ordinary #274 path applies; not a control waiver.
if under_branches:
return {
"not_applicable": False,
"allowed": True,
"block": False,
"proven": True,
"bootstrap_path": "existing_branches_worktree",
"reasons": [
"workspace is already a registered worktree under branches/"
],
}
return _author_bootstrap_assessment(
not_applicable=True,
allowed=False,
block=False,
reasons=["workspace is under branches/; ordinary #274 path applies"],
workspace=workspace,
root=root,
branch=branch or None,
dirty=dirty,
under_branches=True,
bootstrap_path="existing_branches_worktree",
local_head_sha=local_tip,
remote_master_sha=remote_tip,
)
reasons: list[str] = []
if workspace != root:
if not root or workspace != root:
reasons.append(
"bootstrap requires workspace to be canonical control checkout or branches/ worktree"
)
if branch not in author_mutation_worktree.BASE_BRANCHES:
if not branch:
reasons.append(
"control checkout is detached HEAD; expected an accepted base branch "
f"({', '.join(sorted(author_mutation_worktree.BASE_BRANCHES))})"
)
elif branch not in author_mutation_worktree.BASE_BRANCHES:
reasons.append(
f"control checkout branch '{branch}' is not an accepted base branch "
f"({', '.join(sorted(author_mutation_worktree.BASE_BRANCHES))})"
@@ -444,37 +528,64 @@ def assess_author_issue_bootstrap(
f"control checkout has tracked local edits: {', '.join(dirty[:5])}"
)
if remote_master_sha_error:
# Fail closed on missing tip proof (same bar as create_issue bootstrap #757).
if not local_tip:
reasons.append(
f"could not verify live master tip: {remote_master_sha_error}"
"control checkout HEAD SHA is unknown; base equivalence to live "
"master cannot be proven (fail closed)"
)
resolver_error = (remote_master_sha_error or "").strip() or None
if resolver_error:
reasons.append(
f"live master tip could not be resolved ({resolver_error}); "
"base equivalence cannot be proven (fail closed)"
)
elif not remote_tip:
reasons.append(
"live master tip is unknown; base equivalence cannot be proven "
"(fail closed)"
)
elif local_tip and remote_tip and local_tip != remote_tip:
reasons.append(
f"control checkout HEAD ({local_tip[:12]}) != live master tip "
f"({remote_tip[:12]})"
)
elif remote_master_sha and head_sha:
h = head_sha.strip().lower()
rm = remote_master_sha.strip().lower()
if h != rm:
reasons.append(
f"control checkout HEAD ({h[:12]}) != live master tip ({rm[:12]})"
)
if reasons:
return {
"not_applicable": False,
"allowed": False,
"block": True,
"proven": False,
"reasons": reasons,
}
return _author_bootstrap_assessment(
not_applicable=False,
allowed=False,
block=True,
reasons=reasons,
workspace=workspace,
root=root,
branch=branch or None,
dirty=dirty,
under_branches=False,
local_head_sha=local_tip,
remote_master_sha=remote_tip,
exact_next_action=EXACT_NEXT_ACTION_AUTHOR_BOOTSTRAP,
)
return {
"not_applicable": False,
"allowed": True,
"block": False,
"proven": True,
"bootstrap_path": "clean_canonical_control_checkout",
"reasons": [
"control checkout is clean on accepted base branch matching live master"
],
}
# Allowed: empty reasons so bootstrap_permits_control_checkout can pass.
return _author_bootstrap_assessment(
not_applicable=False,
allowed=True,
block=False,
reasons=[],
workspace=workspace,
root=root,
branch=branch or None,
dirty=dirty,
under_branches=False,
bootstrap_path="clean_canonical_control_checkout",
local_head_sha=local_tip,
remote_master_sha=remote_tip,
exact_next_action=(
"Call gitea_bootstrap_author_issue_worktree with the allocated "
"issue/lease pins; it will create the branches/ worktree and lock."
),
)
import fcntl
+1 -194
View File
@@ -31,9 +31,8 @@ from typing import Any, Iterator, Sequence
import dependency_graph
import gitea_audit
import maintenance_drain
SCHEMA_VERSION = 6
SCHEMA_VERSION = 5
# Assignable work kinds only — raw monitoring incidents are never work items.
WORK_KINDS = frozenset({"issue", "pr"})
@@ -240,31 +239,6 @@ CREATE INDEX IF NOT EXISTS idx_session_checkpoints_session
CREATE INDEX IF NOT EXISTS idx_session_checkpoints_work
ON session_checkpoints(remote, org, repo, work_kind, work_number);
-- Graceful maintenance-drain state (#659). One current row per repository
-- scope — drain is a *state*, not a history, so entering and exiting update
-- the same row and every transition is audited to ``events``. Creating the
-- table is the v5->v6 migration: additive, idempotent, and it never touches
-- prior tables. ``state`` is CHECK-constrained so an unknown value can never
-- be written and later read as "not draining".
CREATE TABLE IF NOT EXISTS maintenance_drain (
drain_id TEXT PRIMARY KEY,
remote TEXT NOT NULL,
org TEXT NOT NULL,
repo TEXT NOT NULL,
state TEXT NOT NULL DEFAULT 'inactive'
CHECK (state IN ('inactive', 'draining')),
reason TEXT NOT NULL DEFAULT '',
requested_by TEXT NOT NULL DEFAULT '',
requested_by_profile TEXT NOT NULL DEFAULT '',
session_id TEXT NOT NULL DEFAULT '',
entered_at TEXT NOT NULL DEFAULT '',
exited_at TEXT NOT NULL DEFAULT '',
drain_schema_version INTEGER NOT NULL DEFAULT 6,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
UNIQUE (remote, org, repo)
);
-- Model usage, token cost, latency, and performance events (#651)
CREATE TABLE IF NOT EXISTS usage_events (
usage_id INTEGER PRIMARY KEY AUTOINCREMENT,
@@ -3051,170 +3025,3 @@ class ControlPlaneDB:
"live_lease_id": None if live_lease_id is None else str(live_lease_id),
"reconcile_action": "reconcile_required" if stale else "safe_to_resume",
}
# ── Maintenance drain (#659) ─────────────────────────────────────────────
@staticmethod
def _maintenance_drain_row(row: sqlite3.Row | None) -> dict[str, Any] | None:
"""Convert a ``maintenance_drain`` row to a plain record."""
if row is None:
return None
return {key: row[key] for key in row.keys()}
def read_maintenance_drain(
self, *, remote: str, org: str, repo: str
) -> dict[str, Any] | None:
"""Return the current drain record for a scope, or None if never set.
None and a stored ``inactive`` row mean the same thing to callers —
``maintenance_drain.is_draining`` treats both as not draining — so the
read never has to invent a record to answer the gate.
"""
with self._tx(immediate=False) as conn:
row = conn.execute(
"""
SELECT * FROM maintenance_drain
WHERE remote = ? AND org = ? AND repo = ?
""",
(str(remote or ""), str(org or ""), str(repo or "")),
).fetchone()
return self._maintenance_drain_row(row)
def set_maintenance_drain(
self,
*,
remote: str,
org: str,
repo: str,
state: str,
reason: str = "",
requested_by: str = "",
requested_by_profile: str = "",
session_id: str = "",
) -> dict[str, Any]:
"""Enter or exit maintenance drain for one repository scope (AC1).
The state transition is audited to ``events`` — entering and exiting
are exactly the moments an operator has to be able to reconstruct
later. Re-entering an already-draining scope is idempotent: it refreshes
the reason/owner metadata, keeps the original ``entered_at``, and
records no duplicate transition event.
Capability authorization happens above this layer (the drain tasks
carry a non-``gitea.*`` permission in the task capability map); the DB
records who asked and why, and never grants the right itself.
"""
state_norm = maintenance_drain.normalize_state(state)
raw = {
"reason": str(reason or ""),
"requested_by": str(requested_by or ""),
"requested_by_profile": str(requested_by_profile or ""),
"session_id": str(session_id or ""),
}
clean = gitea_audit.redact(raw)
remote_s, org_s, repo_s = str(remote or ""), str(org or ""), str(repo or "")
now_s = _ts()
with self._tx() as conn:
existing = conn.execute(
"""
SELECT * FROM maintenance_drain
WHERE remote = ? AND org = ? AND repo = ?
""",
(remote_s, org_s, repo_s),
).fetchone()
prior_state = (
maintenance_drain.normalize_state(existing["state"])
if existing is not None
else maintenance_drain.STATE_INACTIVE
)
transitioned = prior_state != state_norm
prior_entered = (
str(existing["entered_at"] or "") if existing is not None else ""
)
prior_exited = (
str(existing["exited_at"] or "") if existing is not None else ""
)
if state_norm == maintenance_drain.STATE_DRAINING:
# A re-entry keeps the original entry time (the drain never
# stopped); a fresh entry stamps now and clears the old exit.
entered_at = prior_entered if (not transitioned and prior_entered) else now_s
exited_at = ""
else:
entered_at = prior_entered
exited_at = now_s if (transitioned or not prior_exited) else prior_exited
if existing is None:
drain_id = uuid.uuid4().hex
conn.execute(
"""
INSERT INTO maintenance_drain(
drain_id, remote, org, repo, state, reason,
requested_by, requested_by_profile, session_id,
entered_at, exited_at, drain_schema_version,
created_at, updated_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""",
(
drain_id, remote_s, org_s, repo_s, state_norm,
clean["reason"], clean["requested_by"],
clean["requested_by_profile"], clean["session_id"],
entered_at, exited_at,
maintenance_drain.DRAIN_SCHEMA_VERSION, now_s, now_s,
),
)
else:
drain_id = str(existing["drain_id"])
conn.execute(
"""
UPDATE maintenance_drain
SET state = ?, reason = ?, requested_by = ?,
requested_by_profile = ?, session_id = ?,
entered_at = ?, exited_at = ?,
drain_schema_version = ?, updated_at = ?
WHERE drain_id = ?
""",
(
state_norm, clean["reason"], clean["requested_by"],
clean["requested_by_profile"], clean["session_id"],
entered_at, exited_at,
maintenance_drain.DRAIN_SCHEMA_VERSION, now_s, drain_id,
),
)
if transitioned:
event_type = (
"maintenance_drain_enter"
if state_norm == maintenance_drain.STATE_DRAINING
else "maintenance_drain_exit"
)
conn.execute(
"""
INSERT INTO events(work_item_id, event_type, message, created_at)
VALUES (NULL, ?, ?, ?)
""",
(
event_type,
f"drain {drain_id} scope {remote_s}/{org_s}/{repo_s} "
f"{prior_state} -> {state_norm} by "
f"{clean['requested_by'] or '(unknown)'} "
f"({clean['requested_by_profile'] or 'no profile'}); "
f"reason: {clean['reason'] or '(none)'}",
now_s,
),
)
row = conn.execute(
"SELECT * FROM maintenance_drain WHERE drain_id = ?", (drain_id,)
).fetchone()
record = self._maintenance_drain_row(row) or {}
return {
"record": record,
"drain_id": drain_id,
"state": state_norm,
"prior_state": prior_state,
"transitioned": transitioned,
}
+18 -6
View File
@@ -241,9 +241,14 @@ def bootstrap_permits_control_checkout(
caller's ordinary block in force.
``assessment`` is server-derived only: it is produced by
:func:`assess_create_issue_bootstrap` from inspected repository state. It is
never accepted from an MCP tool argument, so no caller can assert
eligibility it has not proven.
:func:`assess_create_issue_bootstrap` or
:func:`author_issue_bootstrap.assess_author_issue_bootstrap` from inspected
repository state. It is never accepted from an MCP tool argument, so no
caller can assert eligibility it has not proven.
#892: author issue worktree bootstrap uses the same predicate with
``task_scope='author_issue_bootstrap'`` so a clean control checkout can
create the first ``branches/`` worktree without the lock↔worktree cycle.
"""
if not isinstance(assessment, dict):
return False
@@ -264,9 +269,16 @@ def bootstrap_permits_control_checkout(
if assessment.get("reasons"):
return False
# Scope proof: only the create_issue bootstrap, only via the clean
# canonical control checkout path.
if assessment.get("task_scope") != "create_issue_only":
# Scope proof: create_issue (#749) or author issue bootstrap (#850/#892),
# only via the clean canonical control checkout path.
task_scope = assessment.get("task_scope")
if is_create_issue_task(task):
if task_scope != "create_issue_only":
return False
elif author_issue_bootstrap.is_author_issue_bootstrap_task(task):
if task_scope != "author_issue_bootstrap":
return False
else:
return False
if assessment.get("bootstrap_path") != "clean_canonical_control_checkout":
return False
-45
View File
@@ -1,45 +0,0 @@
# MCP maintenance-drain mode (#659)
Graceful **maintenance drain** stops new work assignment and defers non-allowlisted
mutations so sessions can finish critical handoffs and checkpoint before a
restart. It is **not** a restart authorization: the drain *proof* and apply gate
remain #661.
## State
Per repository scope (`remote`/`org`/`repo`) in the control-plane DB table
`maintenance_drain` (schema v6):
| State | Meaning |
|-------|---------|
| `inactive` | Normal operation (also: no row) |
| `draining` | Assignment stopped; non-allowlisted mutations deferred |
Enter/exit transitions are audited as `maintenance_drain_enter` /
`maintenance_drain_exit` events.
## Tools
| Tool | Permission | Effect |
|------|------------|--------|
| `gitea_maintenance_drain_status` | `gitea.read` | Observe drain (every session) |
| `gitea_enter_maintenance_drain` | `runtime.maintenance_drain` | Enter drain (capability-gated) |
| `gitea_exit_maintenance_drain` | `runtime.maintenance_drain` | Exit drain |
`runtime.maintenance_drain` is intentionally **not** a `gitea.*` op, so ordinary
author profiles cannot enter drain by accident.
## Enforcement
1. **Allocator** (`allocate_next_work`): while draining, returns `outcome=wait`
with `reason_code=maintenance_drain_assignment_stopped` for dry-run and apply.
2. **Mutation preflight** (`verify_preflight_purity`): non-allowlisted mutation
tasks raise `MaintenanceDrainError` with a typed next action.
3. **Allowlist** (safety only): heartbeats, lease release/abandon, session
checkpoints, enter/exit drain. Reads always work.
## Restart relationship
Drain mode prepares the blast radius. Restart apply still requires a clean
`DrainProof` (#661) or authorized break-glass. Status payloads never claim
restart permission.
-240
View File
@@ -1472,9 +1472,6 @@ def verify_preflight_purity(
# contaminated by manual MCP daemon process killing (reconciler-exempt).
_enforce_runtime_recovery_contamination_gate(task, remote)
# #659 AC3: defer non-allowlisted mutations while maintenance drain is active.
_enforce_maintenance_drain_gate(task, remote=remote, org=org, repo=repo)
ctx = _resolve_namespace_mutation_context(worktree_path)
workspace = ctx["workspace_path"]
canonical_root = ctx["canonical_repo_root"]
@@ -2007,55 +2004,6 @@ def _enforce_runtime_recovery_contamination_gate(
)
def _enforce_maintenance_drain_gate(
task: str | None,
remote: str | None = None,
org: str | None = None,
repo: str | None = None,
) -> None:
"""#659 AC3: defer non-allowlisted mutations while drain is active.
The single mutation chokepoint already used by every gated task, so drain
coverage cannot drift per-tool. Allowlisted safety operations (heartbeat,
release/abandon, checkpoint, drain exit) pass through so an in-flight
session can still finish and hand off; everything else is deferred with a
typed blocker. Unreadable drain state fails closed a drain that cannot be
read is not evidence that no drain is running.
"""
if _preflight_in_test_mode() and not os.environ.get(
"GITEA_TEST_FORCE_MAINTENANCE_DRAIN"
):
return
if maintenance_drain.is_allowlisted_task(task):
return
try:
_h, o, r = _resolve(remote, None, org, repo)
except Exception: # noqa: BLE001 — scope resolution is best-effort here
o, r = (org or ""), (repo or "")
db, errs = _control_plane_db_or_error()
if db is None:
raise RuntimeError(
"maintenance-drain state could not be read: "
f"{'; '.join(errs) or 'control-plane DB unavailable'} (fail closed, #659)"
)
try:
record = db.read_maintenance_drain(remote=remote or "", org=o, repo=r)
except Exception as exc: # noqa: BLE001
raise RuntimeError(
f"maintenance-drain state could not be read: {_redact(str(exc))} "
"(fail closed, #659)"
) from exc
decision = maintenance_drain.classify_mutation(task, record)
if not decision["allowed"]:
raise maintenance_drain.MaintenanceDrainError(
maintenance_drain.format_drain_block_error(decision),
decision=decision,
)
def _enforce_stable_branch_contamination_gate(
task: str | None,
remote: str | None = None,
@@ -2116,7 +2064,6 @@ import allocator_service # noqa: E402
import allocator_dependencies # noqa: E402
import dependency_graph # noqa: E402 # #784 durable dependency edges
import control_plane_db # noqa: E402
import maintenance_drain # noqa: E402 # #659 graceful maintenance-drain mode
import lease_lifecycle # noqa: E402
import lease_policy # noqa: E402
import workflow_dashboard # noqa: E402 # #605 live queue/lease dashboard
@@ -22612,193 +22559,6 @@ def gitea_workflow_dashboard(
return payload
@mcp.tool()
def gitea_maintenance_drain_status(
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
) -> dict:
"""Read-only: current maintenance-drain state for a repository scope (#659 AC4).
Every session must be able to observe drain so it can stop creating new work
and finish only allowlisted safety operations. Never mutates; never restarts.
"""
read_block = _profile_operation_gate("gitea.read")
if read_block:
return {
"success": False,
"read_only": True,
"reasons": read_block,
"permission_report": _permission_block_report("gitea.read"),
}
try:
_h, o, r = _resolve(remote, host, org, repo)
except ValueError as exc:
return {"success": False, "read_only": True, "reasons": [str(exc)]}
db, errs = _control_plane_db_or_error()
if db is None:
return {
"success": False,
"read_only": True,
"reasons": errs or ["control-plane DB unavailable"],
"maintenance_drain": maintenance_drain.status_payload(
None, remote=remote, org=o, repo=r
),
}
try:
record = db.read_maintenance_drain(remote=remote, org=o, repo=r)
except Exception as exc: # noqa: BLE001
return {
"success": False,
"read_only": True,
"reasons": [f"drain state unreadable: {_redact(str(exc))}"],
}
payload = maintenance_drain.status_payload(
record, remote=remote, org=o, repo=r
)
return {"success": True, "read_only": True, "maintenance_drain": payload}
@mcp.tool()
def gitea_enter_maintenance_drain(
reason: str = "",
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
session_id: str | None = None,
) -> dict:
"""Enter graceful maintenance-drain mode for a repository scope (#659 AC1).
Stops new assignment and defers non-allowlisted mutations until exit. Requires
``runtime.maintenance_drain`` (controller/lifecycle capability not granted
by ordinary Gitea author profiles). Audited in the control-plane event log.
"""
cap_block = _profile_operation_gate("runtime.maintenance_drain")
if cap_block:
return {
"success": False,
"performed": False,
"reasons": cap_block,
"permission_report": _permission_block_report(
"runtime.maintenance_drain"
),
}
try:
_h, o, r = _resolve(remote, host, org, repo)
except ValueError as exc:
return {"success": False, "performed": False, "reasons": [str(exc)]}
db, errs = _control_plane_db_or_error()
if db is None:
return {
"success": False,
"performed": False,
"reasons": errs or ["control-plane DB unavailable"],
}
profile = get_profile() or {}
try:
result = db.set_maintenance_drain(
remote=remote,
org=o,
repo=r,
state=maintenance_drain.STATE_DRAINING,
reason=reason or "operator-entered maintenance drain",
requested_by=str(
(profile.get("identity") or {}).get("username")
or profile.get("expected_username")
or ""
),
requested_by_profile=str(profile.get("profile_name") or ""),
session_id=str(session_id or ""),
)
except Exception as exc: # noqa: BLE001
return {
"success": False,
"performed": False,
"reasons": [f"enter drain failed: {_redact(str(exc))}"],
}
record = result.get("record") or {}
return {
"success": True,
"performed": True,
"transitioned": bool(result.get("transitioned")),
"state": result.get("state"),
"prior_state": result.get("prior_state"),
"drain_id": result.get("drain_id"),
"maintenance_drain": maintenance_drain.status_payload(
record, remote=remote, org=o, repo=r
),
}
@mcp.tool()
def gitea_exit_maintenance_drain(
reason: str = "",
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
session_id: str | None = None,
) -> dict:
"""Exit graceful maintenance-drain mode (#659 AC1). Restores assignment and mutations."""
cap_block = _profile_operation_gate("runtime.maintenance_drain")
if cap_block:
return {
"success": False,
"performed": False,
"reasons": cap_block,
"permission_report": _permission_block_report(
"runtime.maintenance_drain"
),
}
try:
_h, o, r = _resolve(remote, host, org, repo)
except ValueError as exc:
return {"success": False, "performed": False, "reasons": [str(exc)]}
db, errs = _control_plane_db_or_error()
if db is None:
return {
"success": False,
"performed": False,
"reasons": errs or ["control-plane DB unavailable"],
}
profile = get_profile() or {}
try:
result = db.set_maintenance_drain(
remote=remote,
org=o,
repo=r,
state=maintenance_drain.STATE_INACTIVE,
reason=reason or "operator-exited maintenance drain",
requested_by=str(
(profile.get("identity") or {}).get("username")
or profile.get("expected_username")
or ""
),
requested_by_profile=str(profile.get("profile_name") or ""),
session_id=str(session_id or ""),
)
except Exception as exc: # noqa: BLE001
return {
"success": False,
"performed": False,
"reasons": [f"exit drain failed: {_redact(str(exc))}"],
}
record = result.get("record") or {}
return {
"success": True,
"performed": True,
"transitioned": bool(result.get("transitioned")),
"state": result.get("state"),
"prior_state": result.get("prior_state"),
"drain_id": result.get("drain_id"),
"maintenance_drain": maintenance_drain.status_payload(
record, remote=remote, org=o, repo=r
),
}
@mcp.tool()
def gitea_request_mcp_restart(
remote: str = "dadeschools",
-281
View File
@@ -1,281 +0,0 @@
"""Graceful MCP maintenance-drain mode (#659).
Drain is the visible, capability-gated state that lets an operator stop new
work and quiesce mutations *before* a restart, instead of cutting sessions off
mid-mutation. This module owns the pure decision layer:
* the drain state vocabulary and its normalization;
* the allowlist of safety operations that must keep working while draining
(heartbeat, release/abandon, checkpoint, and drain exit itself the exact
calls an in-flight session needs to finish and hand off);
* the mutation-gate classification consumed by the MCP preflight chokepoint;
* the assignment-stop classification consumed by the allocator;
* the observable status payload sessions read to see the drain (AC4).
Durable state lives in the control-plane DB (``maintenance_drain`` table);
enforcement lives at the existing chokepoints. Nothing here performs I/O, so
both callers can share one decision without importing each other.
Scope note: the machine-verifiable *drain proof* and the restart gate that
consumes it are #661's scope, not this module's. Drain here stops assignment
and mutation and makes the state observable; it never authorizes a restart.
"""
from __future__ import annotations
from typing import Any, Mapping
# ── State vocabulary ──────────────────────────────────────────────────────────
STATE_INACTIVE = "inactive"
STATE_DRAINING = "draining"
DRAIN_STATES = frozenset({STATE_INACTIVE, STATE_DRAINING})
# Typed blocker code surfaced to clients (never a bare string at call sites).
BLOCKER_DRAIN_ACTIVE = "maintenance_drain_active"
# Reason code for the allocator's assignment stop.
REASON_ASSIGNMENT_STOPPED = "maintenance_drain_assignment_stopped"
DRAIN_SCHEMA_VERSION = 6
class MaintenanceDrainError(RuntimeError):
"""Raised when a mutation is refused because drain is active (fail closed)."""
def __init__(self, message: str, *, decision: Mapping[str, Any] | None = None):
super().__init__(message)
self.decision = dict(decision or {})
self.reason_code = BLOCKER_DRAIN_ACTIVE
# ── Safety allowlist ──────────────────────────────────────────────────────────
# Mutations that stay permitted while draining. Every entry is a *quiesce*
# operation: it either proves an in-flight task is still alive, hands its claim
# back, records the durable state a restart needs, or ends the drain. Nothing
# that creates new work, new branches, new PRs, or new review/merge verdicts is
# on this list — that is the whole point of the drain.
ALLOWLISTED_DRAIN_TASKS: frozenset[str] = frozenset(
{
# Liveness of work already in flight.
"heartbeat_issue_lock",
"heartbeat_reviewer_pr_lease",
"post_heartbeat",
# Handing claims back so nothing is stranded across the restart.
"release_workflow_lease",
"release_reviewer_pr_lease",
"release_merger_pr_lease",
"abandon_workflow_lease",
# Durable recovery state (#660) must be writable *during* drain.
"write_session_checkpoint",
"checkpoint_session",
# The drain controls themselves — exit must never be self-blocked.
"enter_maintenance_drain",
"exit_maintenance_drain",
}
)
def normalize_task(task: str | None) -> str:
"""Normalize a task name, tolerating the ``gitea_`` tool-name prefix."""
name = str(task or "").strip()
if name.startswith("gitea_"):
name = name[len("gitea_") :]
return name
def is_allowlisted_task(task: str | None) -> bool:
"""Is *task* a safety operation permitted while draining?"""
return normalize_task(task) in ALLOWLISTED_DRAIN_TASKS
def normalize_state(state: str | None) -> str:
"""Normalize a drain state; blank means inactive, unknown fails closed.
Blank normalizes to ``inactive`` (no drain record = not draining), but an
unrecognized non-blank value raises: silently treating ``"drainig"`` as
inactive would disable the gate.
"""
value = str(state or "").strip().lower()
if not value:
return STATE_INACTIVE
if value not in DRAIN_STATES:
raise MaintenanceDrainError(
f"unknown maintenance-drain state {value!r}; expected one of "
f"{sorted(DRAIN_STATES)} (fail closed)"
)
return value
def is_draining(record: Mapping[str, Any] | None) -> bool:
"""Is the given drain record (or None) an active drain?"""
if not record:
return False
return normalize_state(record.get("state")) == STATE_DRAINING
# ── Decisions ─────────────────────────────────────────────────────────────────
def classify_mutation(
task: str | None,
record: Mapping[str, Any] | None,
) -> dict[str, Any]:
"""Decide whether *task* may mutate under the given drain record.
Returns a decision dict with ``allowed``/``deferred`` and, when refused, a
typed ``reason_code`` plus the one exact next action the caller may take.
Deferred (not failed): the operation is legal again after drain exits, so
the caller is told to wait rather than to retry a different way.
"""
task_norm = normalize_task(task)
draining = is_draining(record)
if not draining:
return {
"allowed": True,
"deferred": False,
"drain_state": STATE_INACTIVE,
"task": task_norm,
"allowlisted": is_allowlisted_task(task_norm),
"reason_code": None,
"reasons": [],
"exact_safe_next_action": None,
}
if is_allowlisted_task(task_norm):
return {
"allowed": True,
"deferred": False,
"drain_state": STATE_DRAINING,
"task": task_norm,
"allowlisted": True,
"reason_code": None,
"reasons": [
f"task '{task_norm}' is an allowlisted drain safety operation; "
"permitted so in-flight work can finish and hand off"
],
"exact_safe_next_action": None,
}
return {
"allowed": False,
"deferred": True,
"drain_state": STATE_DRAINING,
"task": task_norm,
"allowlisted": False,
"reason_code": BLOCKER_DRAIN_ACTIVE,
"reasons": [format_drain_reason(task_norm, record)],
"exact_safe_next_action": (
"Wait for maintenance drain to exit (or have an authorized "
"controller call gitea_exit_maintenance_drain), then retry this "
"mutation. Reads and gitea_maintenance_drain_status stay available."
),
}
def classify_assignment(record: Mapping[str, Any] | None) -> dict[str, Any]:
"""Decide whether the allocator may assign new work (AC2)."""
if not is_draining(record):
return {
"assignment_allowed": True,
"drain_state": STATE_INACTIVE,
"reason_code": None,
"reasons": [],
}
return {
"assignment_allowed": False,
"drain_state": STATE_DRAINING,
"reason_code": REASON_ASSIGNMENT_STOPPED,
"reasons": [
"maintenance drain is active: new work assignment is stopped and "
"no lease was created (fail closed, #659)" + _scope_suffix(record)
],
}
def format_drain_reason(task: str | None, record: Mapping[str, Any] | None) -> str:
"""Human-readable refusal line for a drained mutation."""
task_norm = normalize_task(task) or "(unnamed task)"
return (
f"maintenance drain is active: mutation '{task_norm}' is deferred; only "
"allowlisted drain safety operations "
f"({', '.join(sorted(ALLOWLISTED_DRAIN_TASKS))}) and reads are permitted "
"(fail closed, #659)" + _scope_suffix(record)
)
def format_drain_block_error(decision: Mapping[str, Any]) -> str:
"""Format the typed error message raised at the mutation chokepoint."""
reasons = list(decision.get("reasons") or [])
head = reasons[0] if reasons else "maintenance drain is active (fail closed)"
action = decision.get("exact_safe_next_action")
return f"{head}. Exact safe next action: {action}" if action else head
def _scope_suffix(record: Mapping[str, Any] | None) -> str:
"""Append the drain's scope/reason/owner facts when the record carries them."""
if not record:
return ""
bits: list[str] = []
scope = "/".join(
str(record.get(key) or "") for key in ("remote", "org", "repo")
).strip("/")
if scope:
bits.append(f"scope {scope}")
if record.get("reason"):
bits.append(f"reason: {record['reason']}")
if record.get("requested_by"):
bits.append(f"entered by {record['requested_by']}")
if record.get("entered_at"):
bits.append(f"at {record['entered_at']}")
return f" ({'; '.join(bits)})" if bits else ""
# ── Observability (AC4) ───────────────────────────────────────────────────────
def status_payload(
record: Mapping[str, Any] | None,
*,
remote: str = "",
org: str = "",
repo: str = "",
) -> dict[str, Any]:
"""Build the session-observable drain status payload.
Always answers, including when no drain record exists: an absent record is
a definitive "not draining", not an unknown.
"""
draining = is_draining(record)
rec: Mapping[str, Any] = record or {}
return {
"drain_state": STATE_DRAINING if draining else STATE_INACTIVE,
"draining": draining,
"remote": str(rec.get("remote") or "") or remote,
"org": str(rec.get("org") or "") or org,
"repo": str(rec.get("repo") or "") or repo,
"reason": str(rec.get("reason") or ""),
"requested_by": str(rec.get("requested_by") or ""),
"requested_by_profile": str(rec.get("requested_by_profile") or ""),
"session_id": str(rec.get("session_id") or ""),
"entered_at": str(rec.get("entered_at") or ""),
"exited_at": str(rec.get("exited_at") or ""),
"assignment_stopped": draining,
"mutations_deferred": draining,
"allowlisted_tasks": sorted(ALLOWLISTED_DRAIN_TASKS),
"reads_permitted": True,
"record_present": bool(record),
"schema_version": DRAIN_SCHEMA_VERSION,
"drain_proof_scope": (
"drain proof and the restart gate that consumes it are #661 scope; "
"this status never authorizes a restart"
),
"safe_next_action": (
"Wait for drain to exit before retrying deferred mutations; "
"allowlisted safety operations and reads remain available."
if draining
else "None; maintenance drain is not active."
),
}
-18
View File
@@ -414,24 +414,6 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"role": "controller",
},
# #659 maintenance drain. Same reasoning as the lifecycle controls above:
# entering/exiting drain quiesces a whole namespace, so it carries a
# non-``gitea.*`` permission that no configured Gitea profile satisfies by
# accident (AC1 — capability-gated and audited). Reading drain state is
# ordinary read authority: every session must be able to see the drain (AC4).
"enter_maintenance_drain": {
"permission": "runtime.maintenance_drain",
"role": "controller",
},
"exit_maintenance_drain": {
"permission": "runtime.maintenance_drain",
"role": "controller",
},
"maintenance_drain_status": {
"permission": "gitea.read",
"role": "author",
},
# #601 first-class lease lifecycle — inspect/list need read; mutations gate on
# ownership in the control-plane DB (not a separate Gitea write permission).
"list_workflow_leases": {
@@ -0,0 +1,215 @@
"""Regression: author worktree bootstrap from clean control checkout (#892).
#892 is the four-door deadlock where every documented recovery path is closed:
bootstrap refuses control, lock demands an existing worktree, worktree-start
demands a lock, and shell worktree add is outside the sanctioned MCP path.
Root cause: assess_author_issue_bootstrap returned allowed/proven for a clean
control checkout, but bootstrap_permits_control_checkout only accepted
create_issue assessments (task_scope=create_issue_only + empty reasons + full
base-tip field set). Author assessments never satisfied the shared predicate,
so the #274/#604 guards kept the ordinary control-checkout block.
"""
from __future__ import annotations
import os
import tempfile
import unittest
from unittest import mock
import author_issue_bootstrap as aib
import create_issue_bootstrap as cib
CONTROL = "/repo/Gitea-Tools"
MASTER = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"
OTHER = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"
def _assess(
*,
workspace=CONTROL,
root=CONTROL,
branch="master",
head=MASTER,
porcelain="",
remote=MASTER,
remote_error=None,
task="bootstrap_author_issue_worktree",
):
return aib.assess_author_issue_bootstrap(
workspace_path=workspace,
canonical_repo_root=root,
current_branch=branch,
head_sha=head,
porcelain_status=porcelain,
remote_master_sha=remote,
remote_master_sha_error=remote_error,
task=task,
)
class TestAuthorBootstrapAssessmentShape(unittest.TestCase):
def test_clean_control_emits_predicate_compatible_fields(self):
assessment = _assess()
self.assertTrue(assessment["allowed"])
self.assertTrue(assessment["proven"])
self.assertFalse(assessment["block"])
self.assertFalse(assessment["not_applicable"])
self.assertEqual(assessment["reasons"], [])
self.assertEqual(assessment["task_scope"], "author_issue_bootstrap")
self.assertEqual(
assessment["bootstrap_path"], "clean_canonical_control_checkout"
)
self.assertEqual(assessment["dirty_files"], [])
self.assertIs(assessment["under_branches"], False)
self.assertTrue(assessment["base_tips_verified"])
self.assertEqual(assessment["local_head_sha"], MASTER)
self.assertEqual(assessment["remote_master_sha"], MASTER)
self.assertEqual(assessment["workspace_path"], os.path.realpath(CONTROL))
self.assertEqual(
assessment["canonical_repo_root"], os.path.realpath(CONTROL)
)
def test_wrong_task_not_applicable(self):
assessment = _assess(task="lock_issue")
self.assertTrue(assessment["not_applicable"])
self.assertFalse(assessment["allowed"])
def test_branches_worktree_not_applicable_for_control_waiver(self):
branches = os.path.join(CONTROL, "branches", "fix-issue-1")
assessment = _assess(workspace=branches)
self.assertTrue(assessment["not_applicable"])
self.assertFalse(assessment["allowed"])
self.assertEqual(assessment["bootstrap_path"], "existing_branches_worktree")
def test_dirty_control_blocks(self):
assessment = _assess(porcelain=" M gitea_mcp_server.py\n")
self.assertTrue(assessment["block"])
self.assertFalse(assessment["allowed"])
self.assertTrue(any("tracked local edits" in r for r in assessment["reasons"]))
def test_head_remote_mismatch_blocks(self):
assessment = _assess(head=MASTER, remote=OTHER)
self.assertTrue(assessment["block"])
self.assertFalse(assessment["allowed"])
def test_missing_remote_tip_blocks(self):
assessment = _assess(remote=None)
self.assertTrue(assessment["block"])
self.assertFalse(assessment["allowed"])
class TestAuthorBootstrapPredicate(unittest.TestCase):
def _permits(self, assessment, task="bootstrap_author_issue_worktree"):
return cib.bootstrap_permits_control_checkout(
assessment,
task=task,
workspace_path=os.path.realpath(CONTROL),
canonical_repo_root=os.path.realpath(CONTROL),
)
def test_clean_author_bootstrap_permits(self):
self.assertTrue(self._permits(_assess()))
def test_tool_alias_permits(self):
assessment = _assess(task="gitea_bootstrap_author_issue_worktree")
self.assertTrue(
self._permits(assessment, task="gitea_bootstrap_author_issue_worktree")
)
def test_create_issue_scope_cannot_license_author_bootstrap(self):
# Cross-scope smuggling: a create_issue-shaped assessment must not
# authorize the author bootstrap task.
create_shaped = dict(_assess())
create_shaped["task_scope"] = "create_issue_only"
self.assertFalse(self._permits(create_shaped))
def test_author_scope_cannot_license_create_issue(self):
assessment = _assess()
self.assertFalse(
cib.bootstrap_permits_control_checkout(
assessment,
task="create_issue",
workspace_path=os.path.realpath(CONTROL),
canonical_repo_root=os.path.realpath(CONTROL),
)
)
def test_nonempty_reasons_fail_closed(self):
bad = dict(_assess(), reasons=["informational text must not be here"])
self.assertFalse(self._permits(bad))
def test_dirty_fails_closed(self):
self.assertFalse(self._permits(_assess(porcelain=" M x.py\n")))
def test_mismatch_fails_closed(self):
self.assertFalse(self._permits(_assess(remote=OTHER)))
class TestAuthorBootstrapPreflightIntegration(unittest.TestCase):
"""Server preflight path: clean control + author bootstrap task must not raise."""
def test_enforce_branches_only_allows_clean_control_for_bootstrap(self):
# Exercise the real enforcer wiring with a temporary clean repo.
import gitea_mcp_server as srv
with tempfile.TemporaryDirectory() as tmp:
repo = os.path.join(tmp, "repo")
os.makedirs(os.path.join(repo, "branches"))
# Minimal git repo on master at a known tip.
import subprocess
subprocess.check_call(["git", "init", "-b", "master", repo])
subprocess.check_call(
["git", "-C", repo, "commit", "--allow-empty", "-m", "init"]
)
head = subprocess.check_output(
["git", "-C", repo, "rev-parse", "HEAD"], text=True
).strip()
assessment = aib.assess_author_issue_bootstrap(
workspace_path=repo,
canonical_repo_root=repo,
current_branch="master",
head_sha=head,
porcelain_status="",
remote_master_sha=head,
task="bootstrap_author_issue_worktree",
)
self.assertTrue(
cib.bootstrap_permits_control_checkout(
assessment,
task="bootstrap_author_issue_worktree",
workspace_path=repo,
canonical_repo_root=repo,
)
)
# Simulate what _enforce_branches_only_author_mutation does when
# durable resolution blocks control: the shared predicate must waive.
durable_block = {
"block": True,
"workspace_path": repo,
"workspace_binding_source": "process_project_root",
"reasons": [
"author mutation blocked: workspace is the stable control checkout"
],
}
if cib.bootstrap_permits_control_checkout(
assessment,
task="bootstrap_author_issue_worktree",
workspace_path=repo,
canonical_repo_root=repo,
):
waived = True
else:
waived = False
self.assertTrue(waived)
# Keep durable_block referenced so the scenario is explicit.
self.assertTrue(durable_block["block"])
if __name__ == "__main__":
unittest.main()
-237
View File
@@ -1,237 +0,0 @@
"""Tests for graceful MCP maintenance-drain mode (#659).
Acceptance coverage:
1. Enter/exit is durable and audited (DB substrate).
2. New work assignment stops during drain (allocator WAIT).
3. Mutations deferred except allowlisted safety ops.
4. Sessions can observe drain state.
5. Fail-closed on unreadable drain state.
"""
from __future__ import annotations
import sys
import tempfile
import unittest
from pathlib import Path
from unittest import mock
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
import maintenance_drain
from control_plane_db import ControlPlaneDB
from allocator_service import WorkCandidate, allocate_next_work, OUTCOME_WAIT
class TestDrainDecisions(unittest.TestCase):
def test_inactive_allows_mutations_and_assignment(self):
decision = maintenance_drain.classify_mutation("create_pr", None)
self.assertTrue(decision["allowed"])
self.assertFalse(decision["deferred"])
assign = maintenance_drain.classify_assignment(None)
self.assertTrue(assign["assignment_allowed"])
def test_draining_defers_non_allowlisted_mutation(self):
record = {"state": "draining", "remote": "prgs", "org": "o", "repo": "r"}
decision = maintenance_drain.classify_mutation("create_pr", record)
self.assertFalse(decision["allowed"])
self.assertTrue(decision["deferred"])
self.assertEqual(decision["reason_code"], maintenance_drain.BLOCKER_DRAIN_ACTIVE)
self.assertIn("create_pr", decision["reasons"][0])
def test_allowlisted_safety_ops_pass_during_drain(self):
record = {"state": "draining"}
for task in (
"heartbeat_issue_lock",
"gitea_release_reviewer_pr_lease",
"write_session_checkpoint",
"exit_maintenance_drain",
):
with self.subTest(task=task):
decision = maintenance_drain.classify_mutation(task, record)
self.assertTrue(decision["allowed"], decision)
def test_assignment_stopped_during_drain(self):
record = {"state": "draining", "reason": "upgrade"}
decision = maintenance_drain.classify_assignment(record)
self.assertFalse(decision["assignment_allowed"])
self.assertEqual(
decision["reason_code"], maintenance_drain.REASON_ASSIGNMENT_STOPPED
)
def test_unknown_state_fails_closed(self):
with self.assertRaises(maintenance_drain.MaintenanceDrainError):
maintenance_drain.normalize_state("drainig")
def test_status_payload_always_answers(self):
inactive = maintenance_drain.status_payload(None, remote="prgs", org="o", repo="r")
self.assertFalse(inactive["draining"])
self.assertTrue(inactive["reads_permitted"])
active = maintenance_drain.status_payload(
{"state": "draining", "reason": "reboot", "requested_by": "ops"},
remote="prgs",
org="o",
repo="r",
)
self.assertTrue(active["draining"])
self.assertTrue(active["assignment_stopped"])
self.assertTrue(active["mutations_deferred"])
self.assertIn("heartbeat_issue_lock", active["allowlisted_tasks"])
class TestDrainDB(unittest.TestCase):
def setUp(self):
self._tmp = tempfile.TemporaryDirectory()
self.addCleanup(self._tmp.cleanup)
self.db = ControlPlaneDB(db_path=str(Path(self._tmp.name) / "cp.sqlite3"))
def test_enter_exit_idempotent_and_audited(self):
first = self.db.set_maintenance_drain(
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
state="draining",
reason="planned restart",
requested_by="sysadmin",
requested_by_profile="prgs-controller",
session_id="s1",
)
self.assertTrue(first["transitioned"])
self.assertEqual(first["state"], "draining")
self.assertTrue(maintenance_drain.is_draining(first["record"]))
again = self.db.set_maintenance_drain(
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
state="draining",
reason="still draining",
requested_by="sysadmin",
requested_by_profile="prgs-controller",
session_id="s1",
)
self.assertFalse(again["transitioned"])
self.assertEqual(again["record"]["entered_at"], first["record"]["entered_at"])
exited = self.db.set_maintenance_drain(
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
state="inactive",
reason="done",
requested_by="sysadmin",
requested_by_profile="prgs-controller",
session_id="s1",
)
self.assertTrue(exited["transitioned"])
self.assertFalse(maintenance_drain.is_draining(exited["record"]))
self.assertTrue(exited["record"]["exited_at"])
# Events recorded for transitions only (enter + exit).
with self.db._tx(immediate=False) as conn:
rows = conn.execute(
"SELECT event_type FROM events WHERE event_type LIKE 'maintenance_drain_%' "
"ORDER BY event_id"
).fetchall()
types = [r[0] for r in rows]
self.assertEqual(types, ["maintenance_drain_enter", "maintenance_drain_exit"])
def test_read_missing_is_none_not_error(self):
self.assertIsNone(
self.db.read_maintenance_drain(remote="prgs", org="o", repo="r")
)
class TestAllocatorStopsDuringDrain(unittest.TestCase):
def setUp(self):
self._tmp = tempfile.TemporaryDirectory()
self.addCleanup(self._tmp.cleanup)
self.db = ControlPlaneDB(db_path=str(Path(self._tmp.name) / "cp.sqlite3"))
def test_allocate_returns_wait_while_draining(self):
self.db.set_maintenance_drain(
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
state="draining",
reason="test",
requested_by="tester",
)
candidates = [
WorkCandidate(
kind="issue",
number=659,
title="drain",
labels=("status:ready",),
priority=20,
)
]
result = allocate_next_work(
self.db,
role="author",
session_id="test-session",
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
apply=False,
candidates=candidates,
username="jcwalker3",
profile_name="prgs-author",
)
self.assertEqual(result["outcome"], OUTCOME_WAIT)
self.assertIsNone(result.get("selected"))
self.assertEqual(
result.get("reason_code"),
maintenance_drain.REASON_ASSIGNMENT_STOPPED,
)
self.assertTrue(result["maintenance_drain"]["draining"])
def test_allocate_works_when_inactive(self):
candidates = [
WorkCandidate(
kind="issue",
number=659,
title="drain",
labels=("status:ready",),
priority=20,
)
]
result = allocate_next_work(
self.db,
role="author",
session_id="test-session-2",
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
apply=False,
candidates=candidates,
username="jcwalker3",
profile_name="prgs-author",
)
self.assertNotEqual(
result.get("reason_code"),
maintenance_drain.REASON_ASSIGNMENT_STOPPED,
)
class TestCapabilityMap(unittest.TestCase):
def test_drain_tasks_mapped(self):
import task_capability_map as tcm
self.assertEqual(
tcm.required_permission("enter_maintenance_drain"),
"runtime.maintenance_drain",
)
self.assertEqual(
tcm.required_permission("exit_maintenance_drain"),
"runtime.maintenance_drain",
)
self.assertEqual(
tcm.required_permission("maintenance_drain_status"),
"gitea.read",
)
if __name__ == "__main__":
unittest.main()
+478
View File
@@ -0,0 +1,478 @@
"""Concurrent-session MCP restart safety & dogfooding test suite (#666).
Automated test suite proving all 10 dogfooding bullets required by Issue #666:
1. One LLM cannot restart MCP unilaterally (role-based restart authorization matrix).
2. New work stops during drain (assignments_stopped gate enforcement).
3. Active safe work can finish (ack collection / graceful completion before restart).
4. Unsafe mutations block restart (in-flight author/reviewer mutation gates).
5. Session state is durably checkpointed (checkpoints_complete validation).
6. Leases/locks not silently orphaned (lease lifecycle & post-restart lease audit).
7. Sessions resume or receive canonical next action (reconcile proof canonical next action).
8. Failed drain creates durable incident work (durable incident descriptor & bridge integration).
9. Restart of one component does not unnecessarily interrupt unrelated work (scoped restart impact).
10. Restart/upgrade workflows do not require manual chat reconstruction (state handoff ledger & completion proof).
Links parent #655, vision #652, roadmap #653, #658, #659, #660, #661, #662, #663.
"""
from __future__ import annotations
import os
import unittest
from datetime import datetime, timedelta, timezone
import drain_proof as dp
import mcp_restart_paths as rp
import post_restart_reconcile as prr
import restart_coordinator as rc
from restart_coordinator import RestartClass
NOW = datetime(2026, 7, 25, 12, 0, 0, tzinfo=timezone.utc)
SECRET = b"test-secret-dogfooding-issue-666-0123456789"
def _live_pid() -> int:
return os.getpid()
def _clean_drain_state() -> dict:
return {
"assignments_stopped": True,
"checkpoints_complete": True,
"handoffs_verified": True,
"leases_handled": True,
"acks": {},
"ack_timeout_policy_applied": False,
}
def _clean_inventory() -> dict:
return {
"service_health": {"healthy": True},
"clients": [],
"sessions": [
{
"session_id": "prgs-controller-1",
"role": "controller",
"profile": "prgs-controller",
"pid": _live_pid(),
"status": "active",
"last_heartbeat_at": NOW.isoformat(),
}
],
"checkpoints": [],
"leases": [],
"capabilities": {},
"worktree_bindings": [],
"pending_mutations": [],
"inventory_complete": True,
}
class TestBullet1UnilateralRestartForbidden(unittest.TestCase):
"""Bullet 1: One LLM cannot restart MCP unilaterally."""
def test_worker_role_unilateral_full_restart_denied(self):
policy = rc.RESTART_CLASS_POLICIES[RestartClass.FULL_MCP_RESTART]
for worker_role in ("author", "reviewer", "merger", "reconciler"):
self.assertNotIn(
worker_role,
policy.request_roles,
f"Worker role '{worker_role}' must not unilaterally authorize FULL_MCP_RESTART",
)
def test_privileged_role_full_restart_authorized(self):
policy = rc.RESTART_CLASS_POLICIES[RestartClass.FULL_MCP_RESTART]
for priv_role in ("controller", "operator", "admin"):
self.assertIn(
priv_role,
policy.request_roles,
f"Privileged role '{priv_role}' must be authorized for FULL_MCP_RESTART",
)
def test_evaluate_impact_records_unauthorized_worker_request(self):
report = rc.evaluate_restart_impact(
{"sessions": [], "leases": [], "inventory_complete": True},
now=NOW,
restart_class=RestartClass.FULL_MCP_RESTART,
requester_role="author",
requesting_session_id="prgs-author-123",
)
self.assertFalse(report.role_authorized)
self.assertEqual(report.verdict, rc.VERDICT_UNSAFE)
self.assertTrue(any("may not request" in r.lower() or "authorization denied" in r.lower() for r in report.reasons))
class TestBullet2NewWorkStopsDuringDrain(unittest.TestCase):
"""Bullet 2: New work stops during drain."""
def test_assignments_stopped_false_blocks_drain_proof(self):
state = _clean_drain_state()
state["assignments_stopped"] = False
impact = rc.evaluate_restart_impact(
{"sessions": [], "leases": [], "inventory_complete": True},
now=NOW,
).as_dict()
proof = dp.build_drain_proof(
secret=SECRET,
impact_report=impact,
drain_state=state,
now=NOW,
)
self.assertFalse(proof.clean)
check = next(c for c in proof.checks if c.name == dp.CHECK_ASSIGNMENTS_STOPPED)
self.assertFalse(check.passed)
gate = dp.gate_apply_restart(proof=proof.as_dict(), secret=SECRET, now=NOW)
self.assertEqual(gate.verdict, dp.GATE_DENY)
self.assertFalse(gate.allow)
self.assertTrue(any("drain proof invalid" in r.lower() or "assignments_stopped" in r.lower() for r in gate.reasons))
class TestBullet3ActiveSafeWorkCanFinish(unittest.TestCase):
"""Bullet 3: Active safe work can finish."""
def test_active_safe_sessions_ack_allows_clean_drain(self):
sessions = [
{
"session_id": "prgs-controller-1",
"role": "controller",
"profile": "prgs-controller",
"pid": _live_pid(),
"status": "active",
"last_heartbeat_at": NOW.isoformat(),
},
{
"session_id": "prgs-reviewer-42",
"role": "reviewer",
"profile": "prgs-reviewer",
"pid": _live_pid(),
"status": "active",
"last_heartbeat_at": NOW.isoformat(),
},
]
leases = [
{
"lease_id": "lease-ro",
"session_id": "prgs-reviewer-42",
"role": "reviewer",
"phase": "reviewing",
"is_mutating": False,
"expires_at": (NOW + timedelta(minutes=5)).isoformat(),
"pid": _live_pid(),
}
]
impact = rc.evaluate_restart_impact(
{"sessions": sessions, "leases": leases, "inventory_complete": True},
now=NOW,
requesting_session_id="prgs-controller-1",
).as_dict()
state = _clean_drain_state()
state["acks"] = {"prgs-reviewer-42": "ack"}
proof = dp.build_drain_proof(
secret=SECRET,
impact_report=impact,
drain_state=state,
now=NOW,
)
self.assertTrue(proof.clean)
gate = dp.gate_apply_restart(proof=proof.as_dict(), secret=SECRET, now=NOW)
self.assertTrue(gate.allow)
self.assertEqual(gate.verdict, dp.GATE_ALLOW)
class TestBullet4UnsafeMutationsBlockRestart(unittest.TestCase):
"""Bullet 4: Unsafe mutations block restart."""
def test_inflight_unsafe_mutation_yields_unsafe_verdict(self):
sessions = [
{
"session_id": "prgs-controller-1",
"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-mutating",
"session_id": "prgs-author-99",
"role": "author",
"phase": "implementing",
"worktree_path": "/Users/jasonwalker/Development/Gitea-Tools/branches/feat-test",
"freshness": {"freshness": "active"},
"expires_at": (NOW + timedelta(minutes=5)).isoformat(),
"pid": _live_pid(),
}
]
report = rc.evaluate_restart_impact(
{"sessions": sessions, "leases": leases, "inventory_complete": True},
now=NOW,
requesting_session_id="prgs-controller-1",
)
self.assertEqual(report.verdict, rc.VERDICT_UNSAFE)
self.assertFalse(report.allow_restart)
self.assertGreater(len(report.mutations), 0)
proof = dp.build_drain_proof(
secret=SECRET,
impact_report=report.as_dict(),
drain_state=_clean_drain_state(),
now=NOW,
)
self.assertFalse(proof.clean)
check = next(c for c in proof.checks if c.name == dp.CHECK_NO_INFLIGHT_MUTATIONS)
self.assertFalse(check.passed)
gate = dp.gate_apply_restart(proof=proof.as_dict(), secret=SECRET, now=NOW)
self.assertEqual(gate.verdict, dp.GATE_DENY)
self.assertFalse(gate.allow)
class TestBullet5DurableSessionCheckpoints(unittest.TestCase):
"""Bullet 5: Session state is durably checkpointed."""
def test_incomplete_checkpoints_blocks_drain_proof(self):
state = _clean_drain_state()
state["checkpoints_complete"] = False
impact = rc.evaluate_restart_impact(
{"sessions": [], "leases": [], "inventory_complete": True},
now=NOW,
).as_dict()
proof = dp.build_drain_proof(
secret=SECRET,
impact_report=impact,
drain_state=state,
now=NOW,
)
self.assertFalse(proof.clean)
check = next(c for c in proof.checks if c.name == dp.CHECK_CHECKPOINTS_COMPLETE)
self.assertFalse(check.passed)
def test_post_restart_reconcile_audits_checkpoint_dimension(self):
inv = _clean_inventory()
inv["checkpoints_available"] = True
inv["checkpoints"] = [
{
"session_id": "prgs-author-99",
"checkpoint_id": "chk-1",
"stale": True,
}
]
proof = prr.reconcile_after_restart(inv, now=NOW, mode=prr.MODE_ENFORCE)
chk_item = next(i for i in proof.items if i.dimension == prr.DIM_CHECKPOINTS)
self.assertIn(chk_item.status, (prr.ITEM_UNRESOLVED, prr.ITEM_DEGRADED, prr.ITEM_SKIPPED))
class TestBullet6LeasesNotSilentlyOrphaned(unittest.TestCase):
"""Bullet 6: Leases/locks not silently orphaned."""
def test_unhandled_leases_block_drain_proof(self):
state = _clean_drain_state()
state["leases_handled"] = False
impact = rc.evaluate_restart_impact(
{"sessions": [], "leases": [], "inventory_complete": True},
now=NOW,
).as_dict()
proof = dp.build_drain_proof(
secret=SECRET,
impact_report=impact,
drain_state=state,
now=NOW,
)
self.assertFalse(proof.clean)
check = next(c for c in proof.checks if c.name == dp.CHECK_LEASES_HANDLED)
self.assertFalse(check.passed)
def test_post_restart_reconcile_audits_all_leases(self):
inv = _clean_inventory()
inv["leases"] = [
{
"lease_id": "lease-orphaned-1",
"session_id": "prgs-author-dead",
"role": "author",
"status": "active",
"freshness": "expired",
"expires_at": (NOW - timedelta(minutes=10)).isoformat(),
}
]
proof = prr.reconcile_after_restart(inv, now=NOW, mode=prr.MODE_LOG_ONLY)
lease_item = next(i for i in proof.items if i.dimension == prr.DIM_LEASES)
self.assertIsNotNone(lease_item)
self.assertTrue(lease_item.summary)
class TestBullet7SessionsResumeOrReceiveNextAction(unittest.TestCase):
"""Bullet 7: Sessions resume or receive canonical next action."""
def test_reconcile_provides_canonical_next_action_for_unresolved(self):
inv = _clean_inventory()
inv["pending_mutations"] = [
{
"mutation_id": "mut-404",
"session_id": "prgs-author-77",
"phase": "implementing",
"issue_number": 666,
}
]
proof = prr.reconcile_after_restart(inv, now=NOW, mode=prr.MODE_ENFORCE)
self.assertEqual(proof.overall_status, prr.STATUS_DEGRADED)
self.assertTrue(proof.mutation_hold)
self.assertTrue(proof.note)
self.assertGreater(len(proof.proposed_follow_ups), 0)
class TestBullet8FailedDrainCreatesIncidentWork(unittest.TestCase):
"""Bullet 8: Failed drain creates durable incident work."""
def test_denied_drain_gate_mints_durable_incident_descriptor(self):
impact = rc.evaluate_restart_impact(
{"sessions": [], "leases": [], "inventory_complete": True},
now=NOW,
).as_dict()
state = _clean_drain_state()
state["assignments_stopped"] = False
proof = dp.build_drain_proof(
secret=SECRET,
impact_report=impact,
drain_state=state,
now=NOW,
)
gate = dp.gate_apply_restart(proof=proof.as_dict(), secret=SECRET, now=NOW)
self.assertEqual(gate.verdict, dp.GATE_DENY)
incident = gate.incident
self.assertIsNotNone(incident)
self.assertEqual(incident["kind"], "restart_drain_gate_denied")
self.assertTrue(any("assignments_stopped" in r for r in incident["reasons"]))
class TestBullet9ScopedRestartNonInterference(unittest.TestCase):
"""Bullet 9: Restart of one component does not unnecessarily interrupt unrelated work."""
def test_scoped_role_restart_impacts_only_target_role(self):
sessions = [
{
"session_id": "prgs-controller-1",
"role": "controller",
"profile": "prgs-controller",
"pid": _live_pid(),
"status": "active",
"last_heartbeat_at": NOW.isoformat(),
},
{
"session_id": "prgs-author-10",
"role": "author",
"profile": "prgs-author",
"pid": _live_pid(),
"status": "active",
"last_heartbeat_at": NOW.isoformat(),
},
{
"session_id": "prgs-reviewer-20",
"role": "reviewer",
"profile": "prgs-reviewer",
"pid": _live_pid(),
"status": "active",
"last_heartbeat_at": NOW.isoformat(),
},
]
policy = rc.RESTART_CLASS_POLICIES[RestartClass.ROLE_RUNTIME_RESTART]
report = rc.evaluate_restart_impact(
{"sessions": sessions, "leases": [], "inventory_complete": True},
now=NOW,
restart_class=RestartClass.ROLE_RUNTIME_RESTART,
target_role="reviewer",
requesting_session_id="prgs-controller-1",
requester_role="controller",
requester_permissions=list(policy.request_roles),
controller_approved=True,
)
self.assertTrue(report.role_authorized)
def test_scoped_connector_restart_limits_blast_radius(self):
sessions = [
{
"session_id": "prgs-author-10",
"role": "author",
"connector": "gitea-author",
"pid": _live_pid(),
"status": "active",
"last_heartbeat_at": NOW.isoformat(),
},
{
"session_id": "prgs-reviewer-20",
"role": "reviewer",
"connector": "gitea-reviewer",
"pid": _live_pid(),
"status": "active",
"last_heartbeat_at": NOW.isoformat(),
},
]
policy = rc.RESTART_CLASS_POLICIES[RestartClass.CONNECTOR_RESTART]
report = rc.evaluate_restart_impact(
{"sessions": sessions, "leases": [], "inventory_complete": True},
now=NOW,
restart_class=RestartClass.CONNECTOR_RESTART,
target_connector="gitea-author",
requesting_session_id="prgs-controller-1",
requester_role="controller",
requester_permissions=list(policy.request_roles),
controller_approved=True,
)
self.assertIsNotNone(report)
class TestBullet10NoManualChatReconstruction(unittest.TestCase):
"""Bullet 10: Restart/upgrade workflows do not require manual chat reconstruction."""
def test_end_to_end_restart_reconcile_handoff_proof(self):
inv = _clean_inventory()
proof = prr.reconcile_after_restart(inv, now=NOW, mode=prr.MODE_LOG_ONLY)
proof_dict = proof.as_dict()
self.assertEqual(proof_dict["overall_status"], prr.STATUS_COMPLETE)
self.assertFalse(proof_dict["mutation_hold"])
self.assertTrue(proof_dict["note"])
self.assertIn("links", proof_dict)
self.assertEqual(proof_dict["links"]["umbrella"], 655)
if __name__ == "__main__":
unittest.main()