feat(mcp): implement graceful maintenance-drain mode (Closes #659)

Add durable per-repo drain state, allocator assignment stop, mutation
deferral with a safety allowlist, observable status, and capability-gated
enter/exit tools. Drain proof/restart gate remain #661.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
This commit is contained in:
2026-07-25 17:02:12 -04:00
co-authored by Claude Opus 4.8
parent 76f293eb28
commit e91b94db56
7 changed files with 1056 additions and 1 deletions
+194 -1
View File
@@ -31,8 +31,9 @@ from typing import Any, Iterator, Sequence
import dependency_graph
import gitea_audit
import maintenance_drain
SCHEMA_VERSION = 5
SCHEMA_VERSION = 6
# Assignable work kinds only — raw monitoring incidents are never work items.
WORK_KINDS = frozenset({"issue", "pr"})
@@ -239,6 +240,31 @@ 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,
@@ -3025,3 +3051,170 @@ 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,
}