Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1232789b41 | ||
|
|
2976c21ee6 | ||
|
|
a7a283f449 |
@@ -0,0 +1,56 @@
|
|||||||
|
# Post-restart MCP reconciliation (#662)
|
||||||
|
|
||||||
|
After an MCP process restart, sessions, leases, capabilities, worktrees, and
|
||||||
|
interrupted mutations must be reconciled before operators claim a clean runtime.
|
||||||
|
This document describes the #662 completion-proof path.
|
||||||
|
|
||||||
|
## Components
|
||||||
|
|
||||||
|
| Piece | Where | Responsibility |
|
||||||
|
|-------|-------|----------------|
|
||||||
|
| `post_restart_reconcile.reconcile_after_restart` | `post_restart_reconcile.py` | Pure classification: inventory → completion proof DTO. No I/O. |
|
||||||
|
| `RestartCompletionProof` | `post_restart_reconcile.py` | Machine-readable proof (`.as_dict()` is JSON-serializable). |
|
||||||
|
| `gitea_reconcile_after_restart` | `gitea_mcp_server.py` | MCP tool: gathers inventory from the #613 control-plane DB + master-parity, classifies, returns the proof. Read-only. |
|
||||||
|
| Boot hook | `gitea_assess_master_parity` | First post-restart parity probe also runs reconcile once (log-only by default). |
|
||||||
|
|
||||||
|
## Dimensions
|
||||||
|
|
||||||
|
The assessor classifies:
|
||||||
|
|
||||||
|
- **service_health** — process healthy / parity mutation-safe
|
||||||
|
- **clients** — connected client descriptors (optional inventory)
|
||||||
|
- **sessions** — active session rows with dead owner pids are unresolved
|
||||||
|
- **checkpoints** — soft-depends on #660; skipped with reason when schema absent
|
||||||
|
- **leases** — live control-plane leases after restart
|
||||||
|
- **capabilities** — master-parity / stale-runtime (#610)
|
||||||
|
- **worktrees** — lease-bound paths missing on disk
|
||||||
|
- **interrupted_mutations** — mutating lease phases or explicit pending inventory; **never auto-resumed**
|
||||||
|
- **duplicates** — multiple live claims on the same work item
|
||||||
|
- **queue** — allocator resume safety
|
||||||
|
|
||||||
|
## Modes
|
||||||
|
|
||||||
|
| Mode | Env / arg | Behavior |
|
||||||
|
|------|-----------|----------|
|
||||||
|
| `log_only` (default) | unset or `GITEA_POST_RESTART_RECONCILE_MODE=log_only` | Proof only; `mutation_hold=false` |
|
||||||
|
| `enforce` | `GITEA_POST_RESTART_RECONCILE_MODE=enforce` or `mode=enforce` | Sets `mutation_hold=true` when overall status is degraded/failed or interrupted mutations remain |
|
||||||
|
|
||||||
|
## Follow-up issues
|
||||||
|
|
||||||
|
Unresolved dimensions produce `proposed_follow_ups` entries suitable for durable
|
||||||
|
Gitea issues. The MCP tool **does not create** those issues in v1 (rollout is
|
||||||
|
log-only first). Controllers may file them from the proof payload.
|
||||||
|
|
||||||
|
## Links
|
||||||
|
|
||||||
|
- Umbrella: #655
|
||||||
|
- Vision: #652 · Roadmap: #653
|
||||||
|
- Checkpoint schema: #660 (soft dependency)
|
||||||
|
- Drain proof: #661 (soft)
|
||||||
|
- This issue: #662
|
||||||
|
|
||||||
|
## Non-goals
|
||||||
|
|
||||||
|
- HA multi-instance failover
|
||||||
|
- Automatic silent mutation replay
|
||||||
|
- Implementing the #660 checkpoint schema itself
|
||||||
@@ -18022,6 +18022,17 @@ def gitea_assess_master_parity(
|
|||||||
}
|
}
|
||||||
if parity["restart_required"] and enforced:
|
if parity["restart_required"] and enforced:
|
||||||
out["report"] = master_parity_gate.parity_report(parity)
|
out["report"] = master_parity_gate.parity_report(parity)
|
||||||
|
# #662 AC1: first post-restart master-parity probe also runs the boot
|
||||||
|
# reconcile once (log-only by default; never raises).
|
||||||
|
boot_proof = _ensure_boot_post_restart_reconcile()
|
||||||
|
if boot_proof is not None:
|
||||||
|
out["post_restart_reconcile"] = {
|
||||||
|
"reconcile_id": boot_proof.get("reconcile_id"),
|
||||||
|
"overall_status": boot_proof.get("overall_status"),
|
||||||
|
"mutation_hold": boot_proof.get("mutation_hold"),
|
||||||
|
"unresolved_count": boot_proof.get("unresolved_count"),
|
||||||
|
"mode": boot_proof.get("mode"),
|
||||||
|
}
|
||||||
return out
|
return out
|
||||||
|
|
||||||
|
|
||||||
@@ -22470,6 +22481,249 @@ def gitea_request_mcp_restart(
|
|||||||
return payload
|
return payload
|
||||||
|
|
||||||
|
|
||||||
|
# --- #662 post-restart reconciliation ---------------------------------------
|
||||||
|
|
||||||
|
_POST_RESTART_LAST_PROOF: dict | None = None
|
||||||
|
_POST_RESTART_BOOT_RAN = False
|
||||||
|
|
||||||
|
|
||||||
|
def _post_restart_reconcile_mode() -> str:
|
||||||
|
"""Return log_only (default) or enforce from process environment."""
|
||||||
|
raw = (os.environ.get("GITEA_POST_RESTART_RECONCILE_MODE") or "").strip().lower()
|
||||||
|
if raw in {"enforce", "enforced", "hold"}:
|
||||||
|
return "enforce"
|
||||||
|
return "log_only"
|
||||||
|
|
||||||
|
|
||||||
|
def _gather_post_restart_inventory(
|
||||||
|
*,
|
||||||
|
remote: str,
|
||||||
|
org: str,
|
||||||
|
repo: str,
|
||||||
|
limit: int = 200,
|
||||||
|
) -> dict:
|
||||||
|
"""Gather live control-plane facts for post-restart reconcile (#662).
|
||||||
|
|
||||||
|
Best-effort and fail-closed: any inventory section that cannot be read is
|
||||||
|
recorded on ``incomplete_reasons`` and ``inventory_complete`` is cleared.
|
||||||
|
Never mutates Gitea or the control-plane DB.
|
||||||
|
"""
|
||||||
|
import master_parity_gate
|
||||||
|
import post_restart_reconcile as prr
|
||||||
|
|
||||||
|
inventory_complete = True
|
||||||
|
incomplete_reasons: list[str] = []
|
||||||
|
sessions: list[dict] = []
|
||||||
|
leases: list[dict] = []
|
||||||
|
worktree_bindings: list[dict] = []
|
||||||
|
|
||||||
|
db, db_errs = _control_plane_db_or_error()
|
||||||
|
if db is None:
|
||||||
|
inventory_complete = False
|
||||||
|
incomplete_reasons.extend(
|
||||||
|
db_errs or ["control-plane DB unavailable; cannot reconcile after restart"]
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
try:
|
||||||
|
sessions = db.list_sessions(statuses=("active",), limit=max(1, int(limit)))
|
||||||
|
except Exception as exc: # noqa: BLE001
|
||||||
|
inventory_complete = False
|
||||||
|
incomplete_reasons.append(
|
||||||
|
f"session inventory failed: {_redact(str(exc))}"
|
||||||
|
)
|
||||||
|
try:
|
||||||
|
lease_result = lease_lifecycle.list_active_leases(
|
||||||
|
db,
|
||||||
|
remote=remote if remote in REMOTES else remote,
|
||||||
|
org=org,
|
||||||
|
repo=repo,
|
||||||
|
role=None,
|
||||||
|
include_non_active=True,
|
||||||
|
limit=max(1, int(limit)),
|
||||||
|
)
|
||||||
|
leases = list(lease_result.get("leases") or [])
|
||||||
|
except Exception as exc: # noqa: BLE001
|
||||||
|
inventory_complete = False
|
||||||
|
incomplete_reasons.append(
|
||||||
|
f"lease inventory failed: {_redact(str(exc))}"
|
||||||
|
)
|
||||||
|
leases = []
|
||||||
|
|
||||||
|
# Derive worktree bindings from live leases that carry a path.
|
||||||
|
for row in leases:
|
||||||
|
wt = (row.get("worktree_path") or "").strip()
|
||||||
|
if not wt:
|
||||||
|
continue
|
||||||
|
exists = bool(lease_lifecycle.worktree_exists(wt))
|
||||||
|
worktree_bindings.append(
|
||||||
|
{
|
||||||
|
"path": wt,
|
||||||
|
"exists": exists,
|
||||||
|
"missing": not exists,
|
||||||
|
"lease_id": row.get("lease_id"),
|
||||||
|
"session_id": row.get("session_id"),
|
||||||
|
"work_kind": row.get("work_kind"),
|
||||||
|
"work_number": row.get("work_number"),
|
||||||
|
}
|
||||||
|
)
|
||||||
|
|
||||||
|
# Capability / master-parity dimension (code parity after restart).
|
||||||
|
try:
|
||||||
|
parity = _current_master_parity()
|
||||||
|
capabilities = {
|
||||||
|
"stale": bool(parity.get("stale")),
|
||||||
|
"startup_head": parity.get("startup_head"),
|
||||||
|
"current_head": parity.get("current_head"),
|
||||||
|
"in_parity": parity.get("in_parity"),
|
||||||
|
"mutation_safe": parity.get("mutation_safe"),
|
||||||
|
}
|
||||||
|
service_health = {
|
||||||
|
"healthy": bool(parity.get("mutation_safe", not parity.get("stale"))),
|
||||||
|
"parity_summary": master_parity_gate.format_parity(parity),
|
||||||
|
}
|
||||||
|
boot_head = parity.get("startup_head") or _process_boot_head_sha
|
||||||
|
current_head = parity.get("current_head")
|
||||||
|
except Exception as exc: # noqa: BLE001
|
||||||
|
inventory_complete = False
|
||||||
|
incomplete_reasons.append(
|
||||||
|
f"master-parity inventory failed: {_redact(str(exc))}"
|
||||||
|
)
|
||||||
|
capabilities = {"stale": None}
|
||||||
|
service_health = {"healthy": None, "error": "parity assessment failed"}
|
||||||
|
boot_head = _process_boot_head_sha
|
||||||
|
current_head = None
|
||||||
|
|
||||||
|
# #660 soft dependency: checkpoint schema not landed → skip dimension.
|
||||||
|
checkpoints_available = False
|
||||||
|
try:
|
||||||
|
import importlib
|
||||||
|
|
||||||
|
importlib.import_module("session_checkpoint_schema")
|
||||||
|
checkpoints_available = True
|
||||||
|
except Exception: # noqa: BLE001
|
||||||
|
checkpoints_available = False
|
||||||
|
|
||||||
|
return {
|
||||||
|
"inventory_complete": inventory_complete,
|
||||||
|
"incomplete_reasons": incomplete_reasons,
|
||||||
|
"service_health": service_health,
|
||||||
|
"clients": [], # client transport inventory is host/IDE-owned
|
||||||
|
"sessions": sessions,
|
||||||
|
"leases": leases,
|
||||||
|
"checkpoints_available": checkpoints_available,
|
||||||
|
"checkpoints": [] if checkpoints_available else None,
|
||||||
|
"worktree_bindings": worktree_bindings,
|
||||||
|
"pending_mutations": [],
|
||||||
|
"capabilities": capabilities,
|
||||||
|
"boot_head_sha": boot_head,
|
||||||
|
"current_head_sha": current_head,
|
||||||
|
"queue_state": {"safe_to_resume": inventory_complete},
|
||||||
|
"reconcile_version": prr.RECONCILE_VERSION,
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def _run_post_restart_reconcile(
|
||||||
|
*,
|
||||||
|
remote: str = "prgs",
|
||||||
|
org: str | None = None,
|
||||||
|
repo: str | None = None,
|
||||||
|
mode: str | None = None,
|
||||||
|
limit: int = 200,
|
||||||
|
) -> dict:
|
||||||
|
"""Gather + classify post-restart state; cache the latest proof (#662)."""
|
||||||
|
global _POST_RESTART_LAST_PROOF, _POST_RESTART_BOOT_RAN
|
||||||
|
import post_restart_reconcile as prr
|
||||||
|
|
||||||
|
try:
|
||||||
|
_h, o, r = _resolve(remote, None, org, repo)
|
||||||
|
except ValueError as exc:
|
||||||
|
return {
|
||||||
|
"success": False,
|
||||||
|
"read_only": True,
|
||||||
|
"reasons": [str(exc)],
|
||||||
|
}
|
||||||
|
|
||||||
|
inventory = _gather_post_restart_inventory(
|
||||||
|
remote=remote, org=o, repo=r, limit=limit
|
||||||
|
)
|
||||||
|
proof = prr.reconcile_after_restart(
|
||||||
|
inventory,
|
||||||
|
mode=mode or _post_restart_reconcile_mode(),
|
||||||
|
)
|
||||||
|
payload = proof.as_dict()
|
||||||
|
payload["success"] = True
|
||||||
|
payload["read_only"] = True
|
||||||
|
payload["remote"] = remote
|
||||||
|
payload["org"] = o
|
||||||
|
payload["repo"] = r
|
||||||
|
payload["follow_up_create_supported"] = False
|
||||||
|
payload["follow_up_create_note"] = (
|
||||||
|
"proposed_follow_ups lists durable issues the apply path may create; "
|
||||||
|
"this tool never creates them (log-only by default, #662 rollout)"
|
||||||
|
)
|
||||||
|
_POST_RESTART_LAST_PROOF = payload
|
||||||
|
_POST_RESTART_BOOT_RAN = True
|
||||||
|
return payload
|
||||||
|
|
||||||
|
|
||||||
|
def _ensure_boot_post_restart_reconcile() -> dict | None:
|
||||||
|
"""Run post-restart reconcile once per process (boot hook, #662 AC1)."""
|
||||||
|
global _POST_RESTART_BOOT_RAN
|
||||||
|
if _POST_RESTART_BOOT_RAN:
|
||||||
|
return _POST_RESTART_LAST_PROOF
|
||||||
|
# Best-effort: never raise from the boot hook.
|
||||||
|
try:
|
||||||
|
return _run_post_restart_reconcile()
|
||||||
|
except Exception: # noqa: BLE001
|
||||||
|
_POST_RESTART_BOOT_RAN = True
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
@mcp.tool()
|
||||||
|
def gitea_reconcile_after_restart(
|
||||||
|
remote: str = "dadeschools",
|
||||||
|
host: str | None = None,
|
||||||
|
org: str | None = None,
|
||||||
|
repo: str | None = None,
|
||||||
|
mode: str | None = None,
|
||||||
|
limit: int = 200,
|
||||||
|
) -> dict:
|
||||||
|
"""Run post-restart MCP reconciliation and return a completion proof (#662).
|
||||||
|
|
||||||
|
Gathers live control-plane sessions, leases, worktree bindings, and
|
||||||
|
master-parity evidence, then classifies them with the pure
|
||||||
|
``post_restart_reconcile.reconcile_after_restart`` assessor. Returns a
|
||||||
|
machine-readable completion proof listing resolved / unresolved dimensions
|
||||||
|
and proposed durable follow-up issues.
|
||||||
|
|
||||||
|
Read-only by design: never restarts MCP, never auto-resumes write
|
||||||
|
mutations, and never creates Gitea issues (those are a separate apply
|
||||||
|
path). Default mode is ``log_only``; set
|
||||||
|
``GITEA_POST_RESTART_RECONCILE_MODE=enforce`` (or pass ``mode='enforce'``)
|
||||||
|
to set ``mutation_hold`` when anything remains unresolved.
|
||||||
|
|
||||||
|
Soft-depends on #660 for session checkpoints: when the checkpoint schema
|
||||||
|
module is absent the checkpoints dimension is ``skipped`` with an explicit
|
||||||
|
reason rather than inventing a schema.
|
||||||
|
"""
|
||||||
|
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"),
|
||||||
|
}
|
||||||
|
|
||||||
|
return _run_post_restart_reconcile(
|
||||||
|
remote=remote,
|
||||||
|
org=org,
|
||||||
|
repo=repo,
|
||||||
|
mode=mode,
|
||||||
|
limit=limit,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
@mcp.tool()
|
@mcp.tool()
|
||||||
def gitea_inspect_workflow_lease(
|
def gitea_inspect_workflow_lease(
|
||||||
lease_id: str,
|
lease_id: str,
|
||||||
|
|||||||
@@ -0,0 +1,791 @@
|
|||||||
|
"""Post-restart MCP reconciliation and completion proof (#662).
|
||||||
|
|
||||||
|
After an MCP process restart, sessions, leases, capabilities, worktrees, and
|
||||||
|
interrupted mutations are not systematically reconciled; operators rebuild
|
||||||
|
context from chat. This module is the pure classification core of the
|
||||||
|
post-restart reconcile path.
|
||||||
|
|
||||||
|
Design rules (mirrors ``restart_coordinator`` / ``workflow_dashboard``):
|
||||||
|
|
||||||
|
* **Pure classification.** :func:`reconcile_after_restart` takes an already
|
||||||
|
gathered inventory and returns a structured *completion proof*. It never
|
||||||
|
touches the network, the filesystem, or a live process, so multi-session
|
||||||
|
fixtures can drive every branch in unit tests.
|
||||||
|
* **Fail closed.** Incomplete inventory never reports overall ``complete``.
|
||||||
|
Ambiguous interrupted mutations are ``unresolved`` (never silently resumed).
|
||||||
|
* **No blind write resume.** The proof never authorizes replaying a mutation;
|
||||||
|
it only classifies evidence and names follow-up work.
|
||||||
|
* **#660 soft dependency.** When durable session checkpoints are not present
|
||||||
|
in the inventory, the checkpoint dimension is ``skipped`` with an explicit
|
||||||
|
reason rather than inventing a schema (#660 lands separately).
|
||||||
|
* **Log-only then enforce.** Default mode is ``log_only``. ``enforce`` sets
|
||||||
|
``mutation_hold`` when anything remains unresolved so callers can block
|
||||||
|
write ops until reconcile is complete or degraded mode is documented.
|
||||||
|
|
||||||
|
The single sanctioned gather+classify entry point is the MCP tool
|
||||||
|
``gitea_reconcile_after_restart`` (read-only inventory gather + pure classify).
|
||||||
|
Creating durable follow-up Gitea issues from unresolved items is an explicit
|
||||||
|
apply step outside this pure module.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from dataclasses import dataclass, field
|
||||||
|
from datetime import datetime, timezone
|
||||||
|
from typing import Any, Mapping, Sequence
|
||||||
|
from uuid import uuid4
|
||||||
|
|
||||||
|
import lease_lifecycle
|
||||||
|
|
||||||
|
RECONCILE_VERSION = "1.0.0-issue-662"
|
||||||
|
SCHEMA_VERSION = 1
|
||||||
|
|
||||||
|
# Overall proof statuses.
|
||||||
|
STATUS_COMPLETE = "complete"
|
||||||
|
STATUS_DEGRADED = "degraded"
|
||||||
|
STATUS_FAILED = "failed"
|
||||||
|
|
||||||
|
# Per-dimension item statuses.
|
||||||
|
ITEM_RESOLVED = "resolved"
|
||||||
|
ITEM_UNRESOLVED = "unresolved"
|
||||||
|
ITEM_DEGRADED = "degraded"
|
||||||
|
ITEM_SKIPPED = "skipped"
|
||||||
|
|
||||||
|
# Modes.
|
||||||
|
MODE_LOG_ONLY = "log_only"
|
||||||
|
MODE_ENFORCE = "enforce"
|
||||||
|
|
||||||
|
# Lease / session phases that imply a write critical section was in flight.
|
||||||
|
MUTATING_PHASES = frozenset(
|
||||||
|
{
|
||||||
|
"implementing",
|
||||||
|
"publishing",
|
||||||
|
"merging",
|
||||||
|
"reviewing",
|
||||||
|
"committing",
|
||||||
|
"pushing",
|
||||||
|
"closing",
|
||||||
|
"mutating",
|
||||||
|
"critical_section",
|
||||||
|
}
|
||||||
|
)
|
||||||
|
|
||||||
|
# Dimensions the acceptance criteria require.
|
||||||
|
DIM_SERVICE_HEALTH = "service_health"
|
||||||
|
DIM_CLIENTS = "clients"
|
||||||
|
DIM_SESSIONS = "sessions"
|
||||||
|
DIM_CHECKPOINTS = "checkpoints"
|
||||||
|
DIM_LEASES = "leases"
|
||||||
|
DIM_CAPABILITIES = "capabilities"
|
||||||
|
DIM_WORKTREES = "worktrees"
|
||||||
|
DIM_MUTATIONS = "interrupted_mutations"
|
||||||
|
DIM_DUPLICATES = "duplicates"
|
||||||
|
DIM_QUEUE = "queue"
|
||||||
|
|
||||||
|
|
||||||
|
def _utc_now() -> datetime:
|
||||||
|
return datetime.now(timezone.utc)
|
||||||
|
|
||||||
|
|
||||||
|
def _ts(dt: datetime) -> str:
|
||||||
|
return dt.astimezone(timezone.utc).isoformat().replace("+00:00", "Z")
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass(frozen=True)
|
||||||
|
class ReconcileItem:
|
||||||
|
"""One dimension of the post-restart reconcile report."""
|
||||||
|
|
||||||
|
dimension: str
|
||||||
|
status: str
|
||||||
|
summary: str
|
||||||
|
details: dict[str, Any] = field(default_factory=dict)
|
||||||
|
follow_up_required: bool = False
|
||||||
|
|
||||||
|
def as_dict(self) -> dict[str, Any]:
|
||||||
|
return {
|
||||||
|
"dimension": self.dimension,
|
||||||
|
"status": self.status,
|
||||||
|
"summary": self.summary,
|
||||||
|
"details": dict(self.details),
|
||||||
|
"follow_up_required": self.follow_up_required,
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass(frozen=True)
|
||||||
|
class FollowUpIssue:
|
||||||
|
"""A durable follow-up issue the apply path may create for unresolved work."""
|
||||||
|
|
||||||
|
title: str
|
||||||
|
body: str
|
||||||
|
dimension: str
|
||||||
|
severity: str = "high"
|
||||||
|
|
||||||
|
def as_dict(self) -> dict[str, Any]:
|
||||||
|
return {
|
||||||
|
"title": self.title,
|
||||||
|
"body": self.body,
|
||||||
|
"dimension": self.dimension,
|
||||||
|
"severity": self.severity,
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass(frozen=True)
|
||||||
|
class RestartCompletionProof:
|
||||||
|
"""Machine-readable post-restart completion proof (#662 AC2)."""
|
||||||
|
|
||||||
|
schema_version: int
|
||||||
|
reconcile_version: str
|
||||||
|
reconcile_id: str
|
||||||
|
started_at: str
|
||||||
|
finished_at: str
|
||||||
|
boot_head_sha: str | None
|
||||||
|
current_head_sha: str | None
|
||||||
|
inventory_complete: bool
|
||||||
|
incomplete_reasons: tuple[str, ...]
|
||||||
|
mode: str
|
||||||
|
mutation_hold: bool
|
||||||
|
overall_status: str
|
||||||
|
items: tuple[ReconcileItem, ...]
|
||||||
|
proposed_follow_ups: tuple[FollowUpIssue, ...]
|
||||||
|
resolved_count: int
|
||||||
|
unresolved_count: int
|
||||||
|
skipped_count: int
|
||||||
|
note: str
|
||||||
|
|
||||||
|
def as_dict(self) -> dict[str, Any]:
|
||||||
|
return {
|
||||||
|
"schema_version": self.schema_version,
|
||||||
|
"reconcile_version": self.reconcile_version,
|
||||||
|
"reconcile_id": self.reconcile_id,
|
||||||
|
"started_at": self.started_at,
|
||||||
|
"finished_at": self.finished_at,
|
||||||
|
"boot_head_sha": self.boot_head_sha,
|
||||||
|
"current_head_sha": self.current_head_sha,
|
||||||
|
"inventory_complete": self.inventory_complete,
|
||||||
|
"incomplete_reasons": list(self.incomplete_reasons),
|
||||||
|
"mode": self.mode,
|
||||||
|
"mutation_hold": self.mutation_hold,
|
||||||
|
"overall_status": self.overall_status,
|
||||||
|
"items": [i.as_dict() for i in self.items],
|
||||||
|
"proposed_follow_ups": [f.as_dict() for f in self.proposed_follow_ups],
|
||||||
|
"resolved_count": self.resolved_count,
|
||||||
|
"unresolved_count": self.unresolved_count,
|
||||||
|
"skipped_count": self.skipped_count,
|
||||||
|
"note": self.note,
|
||||||
|
"links": {
|
||||||
|
"umbrella": 655,
|
||||||
|
"vision": 652,
|
||||||
|
"roadmap": 653,
|
||||||
|
"issue": 662,
|
||||||
|
"checkpoint_schema": 660,
|
||||||
|
"drain_proof": 661,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def _item(
|
||||||
|
dimension: str,
|
||||||
|
status: str,
|
||||||
|
summary: str,
|
||||||
|
*,
|
||||||
|
details: dict[str, Any] | None = None,
|
||||||
|
follow_up: bool = False,
|
||||||
|
) -> ReconcileItem:
|
||||||
|
return ReconcileItem(
|
||||||
|
dimension=dimension,
|
||||||
|
status=status,
|
||||||
|
summary=summary,
|
||||||
|
details=dict(details or {}),
|
||||||
|
follow_up_required=follow_up,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _lease_freshness(lease: Mapping[str, Any]) -> str:
|
||||||
|
fr = lease.get("freshness")
|
||||||
|
if isinstance(fr, Mapping):
|
||||||
|
return str(fr.get("freshness") or fr.get("status") or "unknown")
|
||||||
|
if isinstance(fr, str):
|
||||||
|
return fr
|
||||||
|
# Fall back to pure classifier when raw lease rows are supplied.
|
||||||
|
try:
|
||||||
|
return str(lease_lifecycle.classify_lease_freshness(dict(lease)).get("freshness") or "unknown")
|
||||||
|
except Exception: # noqa: BLE001 - pure path must not raise on bad rows
|
||||||
|
return "unknown"
|
||||||
|
|
||||||
|
|
||||||
|
def _is_live_freshness(freshness: str) -> bool:
|
||||||
|
return freshness in {"active", "live", "fresh"}
|
||||||
|
|
||||||
|
|
||||||
|
def _is_mutating_phase(phase: str | None) -> bool:
|
||||||
|
p = (phase or "").strip().lower()
|
||||||
|
if not p:
|
||||||
|
return False
|
||||||
|
if p in MUTATING_PHASES:
|
||||||
|
return True
|
||||||
|
# Soft match for compound phases like "author_implementing".
|
||||||
|
return any(token in p for token in MUTATING_PHASES)
|
||||||
|
|
||||||
|
|
||||||
|
def _detect_interrupted_mutations(
|
||||||
|
leases: Sequence[Mapping[str, Any]],
|
||||||
|
pending_mutations: Sequence[Mapping[str, Any]],
|
||||||
|
) -> list[dict[str, Any]]:
|
||||||
|
"""Return interrupted-mutation evidence (never auto-resumes writes)."""
|
||||||
|
found: list[dict[str, Any]] = []
|
||||||
|
|
||||||
|
for raw in pending_mutations or ():
|
||||||
|
if not isinstance(raw, Mapping):
|
||||||
|
continue
|
||||||
|
found.append(
|
||||||
|
{
|
||||||
|
"source": "pending_mutation_inventory",
|
||||||
|
"status": "unresolved",
|
||||||
|
"phase": raw.get("phase"),
|
||||||
|
"session_id": raw.get("session_id"),
|
||||||
|
"work_kind": raw.get("work_kind") or raw.get("kind"),
|
||||||
|
"work_number": raw.get("work_number") or raw.get("number"),
|
||||||
|
"reason": raw.get("reason")
|
||||||
|
or "pending mutation recorded across process restart",
|
||||||
|
"resume_allowed": False,
|
||||||
|
}
|
||||||
|
)
|
||||||
|
|
||||||
|
for lease in leases or ():
|
||||||
|
if not isinstance(lease, Mapping):
|
||||||
|
continue
|
||||||
|
phase = lease.get("phase")
|
||||||
|
freshness = _lease_freshness(lease)
|
||||||
|
if not _is_mutating_phase(str(phase) if phase is not None else None):
|
||||||
|
continue
|
||||||
|
# A mutating phase whose owner is not live is interrupted.
|
||||||
|
if _is_live_freshness(freshness):
|
||||||
|
# Still live after restart is itself surprising — flag for review.
|
||||||
|
found.append(
|
||||||
|
{
|
||||||
|
"source": "lease_mutating_phase",
|
||||||
|
"status": "unresolved",
|
||||||
|
"phase": phase,
|
||||||
|
"freshness": freshness,
|
||||||
|
"lease_id": lease.get("lease_id"),
|
||||||
|
"session_id": lease.get("session_id"),
|
||||||
|
"work_kind": lease.get("work_kind"),
|
||||||
|
"work_number": lease.get("work_number"),
|
||||||
|
"worktree_path": lease.get("worktree_path"),
|
||||||
|
"reason": (
|
||||||
|
"mutating lease phase still classified live after restart; "
|
||||||
|
"do not auto-resume writes"
|
||||||
|
),
|
||||||
|
"resume_allowed": False,
|
||||||
|
}
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
found.append(
|
||||||
|
{
|
||||||
|
"source": "lease_mutating_phase",
|
||||||
|
"status": "unresolved",
|
||||||
|
"phase": phase,
|
||||||
|
"freshness": freshness,
|
||||||
|
"lease_id": lease.get("lease_id"),
|
||||||
|
"session_id": lease.get("session_id"),
|
||||||
|
"work_kind": lease.get("work_kind"),
|
||||||
|
"work_number": lease.get("work_number"),
|
||||||
|
"worktree_path": lease.get("worktree_path"),
|
||||||
|
"reason": (
|
||||||
|
f"mutating lease phase '{phase}' with non-live freshness "
|
||||||
|
f"'{freshness}' — interrupted by restart"
|
||||||
|
),
|
||||||
|
"resume_allowed": False,
|
||||||
|
}
|
||||||
|
)
|
||||||
|
return found
|
||||||
|
|
||||||
|
|
||||||
|
def _detect_duplicate_work(
|
||||||
|
leases: Sequence[Mapping[str, Any]],
|
||||||
|
) -> list[dict[str, Any]]:
|
||||||
|
"""Surface duplicate live claims on the same work item."""
|
||||||
|
by_work: dict[tuple[Any, Any], list[Mapping[str, Any]]] = {}
|
||||||
|
for lease in leases or ():
|
||||||
|
if not isinstance(lease, Mapping):
|
||||||
|
continue
|
||||||
|
if not _is_live_freshness(_lease_freshness(lease)):
|
||||||
|
continue
|
||||||
|
key = (lease.get("work_kind"), lease.get("work_number"))
|
||||||
|
if key[0] is None or key[1] is None:
|
||||||
|
continue
|
||||||
|
by_work.setdefault(key, []).append(lease)
|
||||||
|
|
||||||
|
dups: list[dict[str, Any]] = []
|
||||||
|
for (kind, number), rows in sorted(by_work.items(), key=lambda kv: str(kv[0])):
|
||||||
|
if len(rows) < 2:
|
||||||
|
continue
|
||||||
|
dups.append(
|
||||||
|
{
|
||||||
|
"work_kind": kind,
|
||||||
|
"work_number": number,
|
||||||
|
"claim_count": len(rows),
|
||||||
|
"session_ids": [r.get("session_id") for r in rows],
|
||||||
|
"lease_ids": [r.get("lease_id") for r in rows],
|
||||||
|
}
|
||||||
|
)
|
||||||
|
return dups
|
||||||
|
|
||||||
|
|
||||||
|
def _follow_up_for_item(item: ReconcileItem) -> FollowUpIssue | None:
|
||||||
|
if not item.follow_up_required:
|
||||||
|
return None
|
||||||
|
title = f"[post-restart] unresolved {item.dimension} after MCP restart"
|
||||||
|
body = (
|
||||||
|
f"## Post-restart reconcile follow-up (#662)\n\n"
|
||||||
|
f"**Dimension:** `{item.dimension}`\n"
|
||||||
|
f"**Status:** `{item.status}`\n"
|
||||||
|
f"**Summary:** {item.summary}\n\n"
|
||||||
|
f"```json\n{item.details!r}\n```\n\n"
|
||||||
|
f"Parent umbrella: #655 · Vision: #652 · Roadmap: #653 · Reconcile: #662\n"
|
||||||
|
f"Do **not** auto-resume write mutations; reconcile evidence first.\n"
|
||||||
|
)
|
||||||
|
return FollowUpIssue(
|
||||||
|
title=title,
|
||||||
|
body=body,
|
||||||
|
dimension=item.dimension,
|
||||||
|
severity="high" if item.dimension == DIM_MUTATIONS else "medium",
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def reconcile_after_restart(
|
||||||
|
inventory: Mapping[str, Any],
|
||||||
|
*,
|
||||||
|
now: datetime | None = None,
|
||||||
|
mode: str = MODE_LOG_ONLY,
|
||||||
|
reconcile_id: str | None = None,
|
||||||
|
) -> RestartCompletionProof:
|
||||||
|
"""Classify a post-restart inventory into a completion proof (#662).
|
||||||
|
|
||||||
|
Parameters
|
||||||
|
----------
|
||||||
|
inventory:
|
||||||
|
Gathered facts. Expected keys (all optional except completeness):
|
||||||
|
|
||||||
|
* ``inventory_complete`` (bool) — fail closed when false
|
||||||
|
* ``incomplete_reasons`` (list[str])
|
||||||
|
* ``service_health`` (dict with ``healthy`` bool)
|
||||||
|
* ``clients`` (list) — connected client descriptors
|
||||||
|
* ``sessions`` (list)
|
||||||
|
* ``leases`` (list, optionally with ``freshness``)
|
||||||
|
* ``checkpoints`` (list | None) — durable session checkpoints (#660)
|
||||||
|
* ``checkpoints_available`` (bool) — False when #660 schema absent
|
||||||
|
* ``worktree_bindings`` (list)
|
||||||
|
* ``pending_mutations`` (list) — explicit interrupted-mutation evidence
|
||||||
|
* ``capabilities`` (dict with optional ``stale`` / heads)
|
||||||
|
* ``boot_head_sha`` / ``current_head_sha``
|
||||||
|
* ``queue_state`` (dict)
|
||||||
|
mode:
|
||||||
|
``log_only`` (default) or ``enforce`` (sets mutation_hold on unresolved).
|
||||||
|
"""
|
||||||
|
started = now or _utc_now()
|
||||||
|
mode_norm = (mode or MODE_LOG_ONLY).strip().lower()
|
||||||
|
if mode_norm not in {MODE_LOG_ONLY, MODE_ENFORCE}:
|
||||||
|
mode_norm = MODE_LOG_ONLY
|
||||||
|
|
||||||
|
inventory_complete = bool(inventory.get("inventory_complete", False))
|
||||||
|
incomplete_reasons = tuple(
|
||||||
|
str(r) for r in (inventory.get("incomplete_reasons") or []) if str(r).strip()
|
||||||
|
)
|
||||||
|
|
||||||
|
items: list[ReconcileItem] = []
|
||||||
|
|
||||||
|
# --- service health -------------------------------------------------
|
||||||
|
health = inventory.get("service_health") or {}
|
||||||
|
if not isinstance(health, Mapping):
|
||||||
|
health = {}
|
||||||
|
if not inventory_complete and "service_health" not in inventory:
|
||||||
|
items.append(
|
||||||
|
_item(
|
||||||
|
DIM_SERVICE_HEALTH,
|
||||||
|
ITEM_UNRESOLVED,
|
||||||
|
"service health unknown because inventory is incomplete",
|
||||||
|
details={"inventory_complete": False},
|
||||||
|
follow_up=True,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
elif health.get("healthy") is True:
|
||||||
|
items.append(
|
||||||
|
_item(
|
||||||
|
DIM_SERVICE_HEALTH,
|
||||||
|
ITEM_RESOLVED,
|
||||||
|
"service health verified",
|
||||||
|
details=dict(health),
|
||||||
|
)
|
||||||
|
)
|
||||||
|
elif health.get("healthy") is False:
|
||||||
|
items.append(
|
||||||
|
_item(
|
||||||
|
DIM_SERVICE_HEALTH,
|
||||||
|
ITEM_UNRESOLVED,
|
||||||
|
"service health check failed",
|
||||||
|
details=dict(health),
|
||||||
|
follow_up=True,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
items.append(
|
||||||
|
_item(
|
||||||
|
DIM_SERVICE_HEALTH,
|
||||||
|
ITEM_DEGRADED,
|
||||||
|
"service health not reported; treating as degraded",
|
||||||
|
details=dict(health),
|
||||||
|
follow_up=True,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
|
||||||
|
# --- clients --------------------------------------------------------
|
||||||
|
clients = list(inventory.get("clients") or [])
|
||||||
|
disconnected = [
|
||||||
|
c
|
||||||
|
for c in clients
|
||||||
|
if isinstance(c, Mapping) and c.get("connected") is False
|
||||||
|
]
|
||||||
|
if "clients" not in inventory:
|
||||||
|
items.append(
|
||||||
|
_item(
|
||||||
|
DIM_CLIENTS,
|
||||||
|
ITEM_SKIPPED,
|
||||||
|
"client inventory not supplied",
|
||||||
|
details={},
|
||||||
|
)
|
||||||
|
)
|
||||||
|
elif disconnected:
|
||||||
|
items.append(
|
||||||
|
_item(
|
||||||
|
DIM_CLIENTS,
|
||||||
|
ITEM_UNRESOLVED,
|
||||||
|
f"{len(disconnected)} disconnected client(s) need reconnect",
|
||||||
|
details={"disconnected": disconnected, "total": len(clients)},
|
||||||
|
follow_up=True,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
items.append(
|
||||||
|
_item(
|
||||||
|
DIM_CLIENTS,
|
||||||
|
ITEM_RESOLVED,
|
||||||
|
f"{len(clients)} client(s) accounted for",
|
||||||
|
details={"total": len(clients)},
|
||||||
|
)
|
||||||
|
)
|
||||||
|
|
||||||
|
# --- sessions -------------------------------------------------------
|
||||||
|
sessions = [s for s in (inventory.get("sessions") or []) if isinstance(s, Mapping)]
|
||||||
|
orphan_sessions = [
|
||||||
|
s
|
||||||
|
for s in sessions
|
||||||
|
if str(s.get("status") or "").lower() == "active"
|
||||||
|
and s.get("pid") is not None
|
||||||
|
and not lease_lifecycle.is_process_alive(s.get("pid"))
|
||||||
|
]
|
||||||
|
if orphan_sessions:
|
||||||
|
items.append(
|
||||||
|
_item(
|
||||||
|
DIM_SESSIONS,
|
||||||
|
ITEM_UNRESOLVED,
|
||||||
|
f"{len(orphan_sessions)} active session row(s) with dead owner pid",
|
||||||
|
details={
|
||||||
|
"orphan_session_ids": [s.get("session_id") for s in orphan_sessions],
|
||||||
|
"total_sessions": len(sessions),
|
||||||
|
},
|
||||||
|
follow_up=True,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
items.append(
|
||||||
|
_item(
|
||||||
|
DIM_SESSIONS,
|
||||||
|
ITEM_RESOLVED,
|
||||||
|
f"{len(sessions)} session row(s) reconciled (no dead-pid orphans)",
|
||||||
|
details={"total_sessions": len(sessions)},
|
||||||
|
)
|
||||||
|
)
|
||||||
|
|
||||||
|
# --- checkpoints (#660 soft) ----------------------------------------
|
||||||
|
checkpoints_available = inventory.get("checkpoints_available")
|
||||||
|
checkpoints = inventory.get("checkpoints")
|
||||||
|
if checkpoints_available is False or (
|
||||||
|
checkpoints is None and "checkpoints" not in inventory
|
||||||
|
):
|
||||||
|
items.append(
|
||||||
|
_item(
|
||||||
|
DIM_CHECKPOINTS,
|
||||||
|
ITEM_SKIPPED,
|
||||||
|
"durable session checkpoint schema not available yet (#660)",
|
||||||
|
details={"depends_on": 660},
|
||||||
|
)
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
cp_list = [c for c in (checkpoints or []) if isinstance(c, Mapping)]
|
||||||
|
stale_cp = [c for c in cp_list if c.get("stale") or c.get("invalid")]
|
||||||
|
if stale_cp:
|
||||||
|
items.append(
|
||||||
|
_item(
|
||||||
|
DIM_CHECKPOINTS,
|
||||||
|
ITEM_UNRESOLVED,
|
||||||
|
f"{len(stale_cp)} checkpoint(s) invalid or stale vs live state",
|
||||||
|
details={"stale_count": len(stale_cp), "total": len(cp_list)},
|
||||||
|
follow_up=True,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
items.append(
|
||||||
|
_item(
|
||||||
|
DIM_CHECKPOINTS,
|
||||||
|
ITEM_RESOLVED,
|
||||||
|
f"{len(cp_list)} checkpoint(s) consistent with live state",
|
||||||
|
details={"total": len(cp_list)},
|
||||||
|
)
|
||||||
|
)
|
||||||
|
|
||||||
|
# --- leases / locks -------------------------------------------------
|
||||||
|
leases = [L for L in (inventory.get("leases") or []) if isinstance(L, Mapping)]
|
||||||
|
live_leases = [L for L in leases if _is_live_freshness(_lease_freshness(L))]
|
||||||
|
items.append(
|
||||||
|
_item(
|
||||||
|
DIM_LEASES,
|
||||||
|
ITEM_RESOLVED if inventory_complete else ITEM_DEGRADED,
|
||||||
|
f"{len(live_leases)} live lease(s) of {len(leases)} inventoried",
|
||||||
|
details={
|
||||||
|
"live_count": len(live_leases),
|
||||||
|
"total": len(leases),
|
||||||
|
"live_lease_ids": [L.get("lease_id") for L in live_leases],
|
||||||
|
},
|
||||||
|
follow_up=not inventory_complete,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
|
||||||
|
# --- capabilities / stale runtime -----------------------------------
|
||||||
|
caps = inventory.get("capabilities") or {}
|
||||||
|
if not isinstance(caps, Mapping):
|
||||||
|
caps = {}
|
||||||
|
if caps.get("stale") is True:
|
||||||
|
items.append(
|
||||||
|
_item(
|
||||||
|
DIM_CAPABILITIES,
|
||||||
|
ITEM_UNRESOLVED,
|
||||||
|
"runtime code is stale vs on-disk master; restart did not reach parity",
|
||||||
|
details=dict(caps),
|
||||||
|
follow_up=True,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
items.append(
|
||||||
|
_item(
|
||||||
|
DIM_CAPABILITIES,
|
||||||
|
ITEM_RESOLVED,
|
||||||
|
"capability/runtime parity acceptable",
|
||||||
|
details=dict(caps) if caps else {"stale": False},
|
||||||
|
)
|
||||||
|
)
|
||||||
|
|
||||||
|
# --- worktrees ------------------------------------------------------
|
||||||
|
bindings = [
|
||||||
|
b for b in (inventory.get("worktree_bindings") or []) if isinstance(b, Mapping)
|
||||||
|
]
|
||||||
|
missing_wt = [
|
||||||
|
b
|
||||||
|
for b in bindings
|
||||||
|
if b.get("missing") is True or b.get("exists") is False
|
||||||
|
]
|
||||||
|
if "worktree_bindings" not in inventory:
|
||||||
|
items.append(
|
||||||
|
_item(
|
||||||
|
DIM_WORKTREES,
|
||||||
|
ITEM_SKIPPED,
|
||||||
|
"worktree binding inventory not supplied",
|
||||||
|
)
|
||||||
|
)
|
||||||
|
elif missing_wt:
|
||||||
|
items.append(
|
||||||
|
_item(
|
||||||
|
DIM_WORKTREES,
|
||||||
|
ITEM_UNRESOLVED,
|
||||||
|
f"{len(missing_wt)} worktree binding(s) missing on disk",
|
||||||
|
details={"missing": missing_wt, "total": len(bindings)},
|
||||||
|
follow_up=True,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
items.append(
|
||||||
|
_item(
|
||||||
|
DIM_WORKTREES,
|
||||||
|
ITEM_RESOLVED,
|
||||||
|
f"{len(bindings)} worktree binding(s) present",
|
||||||
|
details={"total": len(bindings)},
|
||||||
|
)
|
||||||
|
)
|
||||||
|
|
||||||
|
# --- interrupted mutations (AC4) ------------------------------------
|
||||||
|
pending = [
|
||||||
|
m
|
||||||
|
for m in (inventory.get("pending_mutations") or [])
|
||||||
|
if isinstance(m, Mapping)
|
||||||
|
]
|
||||||
|
interrupted = _detect_interrupted_mutations(leases, pending)
|
||||||
|
if interrupted:
|
||||||
|
items.append(
|
||||||
|
_item(
|
||||||
|
DIM_MUTATIONS,
|
||||||
|
ITEM_UNRESOLVED,
|
||||||
|
f"{len(interrupted)} interrupted mutation(s); write resume forbidden",
|
||||||
|
details={"interrupted": interrupted},
|
||||||
|
follow_up=True,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
items.append(
|
||||||
|
_item(
|
||||||
|
DIM_MUTATIONS,
|
||||||
|
ITEM_RESOLVED,
|
||||||
|
"no interrupted mutations detected",
|
||||||
|
details={"interrupted": []},
|
||||||
|
)
|
||||||
|
)
|
||||||
|
|
||||||
|
# --- duplicates -----------------------------------------------------
|
||||||
|
dups = _detect_duplicate_work(leases)
|
||||||
|
if dups:
|
||||||
|
items.append(
|
||||||
|
_item(
|
||||||
|
DIM_DUPLICATES,
|
||||||
|
ITEM_UNRESOLVED,
|
||||||
|
f"{len(dups)} work item(s) have multiple live claims",
|
||||||
|
details={"duplicates": dups},
|
||||||
|
follow_up=True,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
items.append(
|
||||||
|
_item(
|
||||||
|
DIM_DUPLICATES,
|
||||||
|
ITEM_RESOLVED,
|
||||||
|
"no duplicate live claims detected",
|
||||||
|
details={"duplicates": []},
|
||||||
|
)
|
||||||
|
)
|
||||||
|
|
||||||
|
# --- queue ----------------------------------------------------------
|
||||||
|
queue = inventory.get("queue_state")
|
||||||
|
if queue is None:
|
||||||
|
items.append(
|
||||||
|
_item(
|
||||||
|
DIM_QUEUE,
|
||||||
|
ITEM_SKIPPED,
|
||||||
|
"allocator queue state not supplied",
|
||||||
|
)
|
||||||
|
)
|
||||||
|
elif isinstance(queue, Mapping) and queue.get("safe_to_resume") is False:
|
||||||
|
items.append(
|
||||||
|
_item(
|
||||||
|
DIM_QUEUE,
|
||||||
|
ITEM_UNRESOLVED,
|
||||||
|
"allocator queue not safe to resume",
|
||||||
|
details=dict(queue),
|
||||||
|
follow_up=True,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
items.append(
|
||||||
|
_item(
|
||||||
|
DIM_QUEUE,
|
||||||
|
ITEM_RESOLVED,
|
||||||
|
"allocator queue state acceptable",
|
||||||
|
details=dict(queue) if isinstance(queue, Mapping) else {},
|
||||||
|
)
|
||||||
|
)
|
||||||
|
|
||||||
|
# Incomplete inventory always degrades the whole proof.
|
||||||
|
if not inventory_complete:
|
||||||
|
# Ensure at least one follow-up names the incomplete inventory.
|
||||||
|
items.append(
|
||||||
|
_item(
|
||||||
|
"inventory",
|
||||||
|
ITEM_UNRESOLVED,
|
||||||
|
"control-plane inventory incomplete; reconcile cannot claim success",
|
||||||
|
details={"reasons": list(incomplete_reasons)},
|
||||||
|
follow_up=True,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
|
||||||
|
resolved = sum(1 for i in items if i.status == ITEM_RESOLVED)
|
||||||
|
unresolved = sum(1 for i in items if i.status in {ITEM_UNRESOLVED, ITEM_DEGRADED})
|
||||||
|
skipped = sum(1 for i in items if i.status == ITEM_SKIPPED)
|
||||||
|
|
||||||
|
if not inventory_complete or any(i.status == ITEM_UNRESOLVED for i in items):
|
||||||
|
if any(i.status == ITEM_UNRESOLVED for i in items) and inventory_complete:
|
||||||
|
overall = STATUS_DEGRADED
|
||||||
|
elif not inventory_complete:
|
||||||
|
overall = STATUS_FAILED
|
||||||
|
else:
|
||||||
|
overall = STATUS_DEGRADED
|
||||||
|
elif any(i.status == ITEM_DEGRADED for i in items):
|
||||||
|
overall = STATUS_DEGRADED
|
||||||
|
else:
|
||||||
|
overall = STATUS_COMPLETE
|
||||||
|
|
||||||
|
# Enforce mode holds mutations whenever anything is unresolved/failed.
|
||||||
|
mutation_hold = False
|
||||||
|
if mode_norm == MODE_ENFORCE and overall in {STATUS_DEGRADED, STATUS_FAILED}:
|
||||||
|
mutation_hold = True
|
||||||
|
if mode_norm == MODE_ENFORCE and any(
|
||||||
|
i.dimension == DIM_MUTATIONS and i.status == ITEM_UNRESOLVED for i in items
|
||||||
|
):
|
||||||
|
mutation_hold = True
|
||||||
|
|
||||||
|
follow_ups = tuple(
|
||||||
|
fu for i in items if (fu := _follow_up_for_item(i)) is not None
|
||||||
|
)
|
||||||
|
|
||||||
|
finished = _utc_now() if now is None else now
|
||||||
|
note = (
|
||||||
|
"Read-only completion proof. Never auto-resumes write mutations. "
|
||||||
|
"Unresolved items require durable follow-up before claiming clean restart. "
|
||||||
|
f"Mode={mode_norm}."
|
||||||
|
)
|
||||||
|
|
||||||
|
return RestartCompletionProof(
|
||||||
|
schema_version=SCHEMA_VERSION,
|
||||||
|
reconcile_version=RECONCILE_VERSION,
|
||||||
|
reconcile_id=(reconcile_id or f"reconcile-{uuid4().hex[:12]}"),
|
||||||
|
started_at=_ts(started),
|
||||||
|
finished_at=_ts(finished),
|
||||||
|
boot_head_sha=(
|
||||||
|
str(inventory.get("boot_head_sha")).strip()
|
||||||
|
if inventory.get("boot_head_sha")
|
||||||
|
else None
|
||||||
|
),
|
||||||
|
current_head_sha=(
|
||||||
|
str(inventory.get("current_head_sha")).strip()
|
||||||
|
if inventory.get("current_head_sha")
|
||||||
|
else None
|
||||||
|
),
|
||||||
|
inventory_complete=inventory_complete,
|
||||||
|
incomplete_reasons=incomplete_reasons,
|
||||||
|
mode=mode_norm,
|
||||||
|
mutation_hold=mutation_hold,
|
||||||
|
overall_status=overall,
|
||||||
|
items=tuple(items),
|
||||||
|
proposed_follow_ups=follow_ups,
|
||||||
|
resolved_count=resolved,
|
||||||
|
unresolved_count=unresolved,
|
||||||
|
skipped_count=skipped,
|
||||||
|
note=note,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def mutations_allowed(proof: RestartCompletionProof | Mapping[str, Any] | None) -> bool:
|
||||||
|
"""Return whether write mutations may proceed under the given proof."""
|
||||||
|
if proof is None:
|
||||||
|
return True # no proof yet → caller decides; enforce path sets hold
|
||||||
|
if isinstance(proof, RestartCompletionProof):
|
||||||
|
return not proof.mutation_hold
|
||||||
|
if isinstance(proof, Mapping):
|
||||||
|
return not bool(proof.get("mutation_hold"))
|
||||||
|
return True
|
||||||
@@ -132,6 +132,16 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
|
|||||||
"permission": "gitea.branch.push",
|
"permission": "gitea.branch.push",
|
||||||
"role": "author",
|
"role": "author",
|
||||||
},
|
},
|
||||||
|
# #662: post-restart reconcile is read-only inventory + pure classification.
|
||||||
|
# Durable follow-up issue creation is a separate apply path (not this task).
|
||||||
|
"reconcile_after_restart": {
|
||||||
|
"permission": "gitea.read",
|
||||||
|
"role": "author",
|
||||||
|
},
|
||||||
|
"gitea_reconcile_after_restart": {
|
||||||
|
"permission": "gitea.read",
|
||||||
|
"role": "author",
|
||||||
|
},
|
||||||
# PR synchronization lifecycle: assess is read-only (any role with gitea.read);
|
# PR synchronization lifecycle: assess is read-only (any role with gitea.read);
|
||||||
# update-by-merge is author-only and mutates the PR head via Gitea API.
|
# update-by-merge is author-only and mutates the PR head via Gitea API.
|
||||||
"assess_pr_sync_status": {
|
"assess_pr_sync_status": {
|
||||||
|
|||||||
@@ -0,0 +1,249 @@
|
|||||||
|
"""Tests for post-restart MCP reconciliation and completion proof (#662)."""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import json
|
||||||
|
import os
|
||||||
|
import unittest
|
||||||
|
from datetime import datetime, timezone
|
||||||
|
|
||||||
|
import post_restart_reconcile as prr
|
||||||
|
|
||||||
|
NOW = datetime(2026, 7, 24, 12, 0, 0, tzinfo=timezone.utc)
|
||||||
|
|
||||||
|
|
||||||
|
def _base_inventory(**overrides):
|
||||||
|
inv = {
|
||||||
|
"inventory_complete": True,
|
||||||
|
"incomplete_reasons": [],
|
||||||
|
"service_health": {"healthy": True},
|
||||||
|
"clients": [{"session_id": "c1", "connected": True}],
|
||||||
|
"sessions": [
|
||||||
|
{
|
||||||
|
"session_id": "s-live",
|
||||||
|
"status": "active",
|
||||||
|
"pid": os.getpid(),
|
||||||
|
"role": "author",
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"leases": [],
|
||||||
|
"checkpoints_available": False,
|
||||||
|
"worktree_bindings": [{"path": "/tmp/wt", "exists": True}],
|
||||||
|
"pending_mutations": [],
|
||||||
|
"capabilities": {"stale": False},
|
||||||
|
"boot_head_sha": "a" * 40,
|
||||||
|
"current_head_sha": "a" * 40,
|
||||||
|
"queue_state": {"safe_to_resume": True},
|
||||||
|
}
|
||||||
|
inv.update(overrides)
|
||||||
|
return inv
|
||||||
|
|
||||||
|
|
||||||
|
class IncompleteInventoryTest(unittest.TestCase):
|
||||||
|
def test_incomplete_inventory_fails_closed(self) -> None:
|
||||||
|
proof = prr.reconcile_after_restart(
|
||||||
|
{
|
||||||
|
"inventory_complete": False,
|
||||||
|
"incomplete_reasons": ["control-plane DB unavailable"],
|
||||||
|
},
|
||||||
|
now=NOW,
|
||||||
|
mode=prr.MODE_ENFORCE,
|
||||||
|
reconcile_id="test-incomplete",
|
||||||
|
)
|
||||||
|
self.assertEqual(proof.overall_status, prr.STATUS_FAILED)
|
||||||
|
self.assertTrue(proof.mutation_hold)
|
||||||
|
self.assertFalse(proof.inventory_complete)
|
||||||
|
self.assertTrue(proof.proposed_follow_ups)
|
||||||
|
self.assertIn("control-plane DB unavailable", proof.incomplete_reasons)
|
||||||
|
|
||||||
|
|
||||||
|
class HappyPathTest(unittest.TestCase):
|
||||||
|
def test_clean_restart_is_complete_without_mutation_hold(self) -> None:
|
||||||
|
proof = prr.reconcile_after_restart(
|
||||||
|
_base_inventory(),
|
||||||
|
now=NOW,
|
||||||
|
mode=prr.MODE_ENFORCE,
|
||||||
|
reconcile_id="test-clean",
|
||||||
|
)
|
||||||
|
self.assertEqual(proof.overall_status, prr.STATUS_COMPLETE)
|
||||||
|
self.assertFalse(proof.mutation_hold)
|
||||||
|
self.assertEqual(proof.unresolved_count, 0)
|
||||||
|
cp = next(i for i in proof.items if i.dimension == prr.DIM_CHECKPOINTS)
|
||||||
|
self.assertEqual(cp.status, prr.ITEM_SKIPPED)
|
||||||
|
links = proof.as_dict()["links"]
|
||||||
|
self.assertEqual(links["umbrella"], 655)
|
||||||
|
self.assertEqual(links["issue"], 662)
|
||||||
|
self.assertEqual(links["vision"], 652)
|
||||||
|
self.assertEqual(links["roadmap"], 653)
|
||||||
|
|
||||||
|
|
||||||
|
class InterruptedMutationTest(unittest.TestCase):
|
||||||
|
def test_mutating_lease_with_dead_owner_is_unresolved(self) -> None:
|
||||||
|
proof = prr.reconcile_after_restart(
|
||||||
|
_base_inventory(
|
||||||
|
leases=[
|
||||||
|
{
|
||||||
|
"lease_id": "lease-mut",
|
||||||
|
"session_id": "s-dead",
|
||||||
|
"phase": "implementing",
|
||||||
|
"work_kind": "issue",
|
||||||
|
"work_number": 662,
|
||||||
|
"worktree_path": "/tmp/wt-662",
|
||||||
|
"freshness": {"freshness": "stale_dead_process"},
|
||||||
|
}
|
||||||
|
]
|
||||||
|
),
|
||||||
|
now=NOW,
|
||||||
|
mode=prr.MODE_ENFORCE,
|
||||||
|
)
|
||||||
|
mut = next(i for i in proof.items if i.dimension == prr.DIM_MUTATIONS)
|
||||||
|
self.assertEqual(mut.status, prr.ITEM_UNRESOLVED)
|
||||||
|
self.assertTrue(mut.follow_up_required)
|
||||||
|
interrupted = mut.details["interrupted"]
|
||||||
|
self.assertEqual(len(interrupted), 1)
|
||||||
|
self.assertFalse(interrupted[0]["resume_allowed"])
|
||||||
|
self.assertTrue(proof.mutation_hold)
|
||||||
|
self.assertTrue(
|
||||||
|
any(f.dimension == prr.DIM_MUTATIONS for f in proof.proposed_follow_ups)
|
||||||
|
)
|
||||||
|
|
||||||
|
def test_explicit_pending_mutation_inventory(self) -> None:
|
||||||
|
proof = prr.reconcile_after_restart(
|
||||||
|
_base_inventory(
|
||||||
|
pending_mutations=[
|
||||||
|
{
|
||||||
|
"session_id": "s1",
|
||||||
|
"phase": "publishing",
|
||||||
|
"work_kind": "pr",
|
||||||
|
"work_number": 856,
|
||||||
|
"reason": "push interrupted mid-flight",
|
||||||
|
}
|
||||||
|
]
|
||||||
|
),
|
||||||
|
now=NOW,
|
||||||
|
mode=prr.MODE_LOG_ONLY,
|
||||||
|
)
|
||||||
|
mut = next(i for i in proof.items if i.dimension == prr.DIM_MUTATIONS)
|
||||||
|
self.assertEqual(mut.status, prr.ITEM_UNRESOLVED)
|
||||||
|
# log_only never holds mutations even when unresolved
|
||||||
|
self.assertFalse(proof.mutation_hold)
|
||||||
|
self.assertEqual(proof.overall_status, prr.STATUS_DEGRADED)
|
||||||
|
|
||||||
|
|
||||||
|
class DuplicateClaimsTest(unittest.TestCase):
|
||||||
|
def test_duplicate_live_claims_flagged(self) -> None:
|
||||||
|
proof = prr.reconcile_after_restart(
|
||||||
|
_base_inventory(
|
||||||
|
leases=[
|
||||||
|
{
|
||||||
|
"lease_id": "l1",
|
||||||
|
"session_id": "s1",
|
||||||
|
"phase": "allocated",
|
||||||
|
"work_kind": "issue",
|
||||||
|
"work_number": 100,
|
||||||
|
"freshness": {"freshness": "active"},
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"lease_id": "l2",
|
||||||
|
"session_id": "s2",
|
||||||
|
"phase": "allocated",
|
||||||
|
"work_kind": "issue",
|
||||||
|
"work_number": 100,
|
||||||
|
"freshness": {"freshness": "active"},
|
||||||
|
},
|
||||||
|
]
|
||||||
|
),
|
||||||
|
now=NOW,
|
||||||
|
mode=prr.MODE_ENFORCE,
|
||||||
|
)
|
||||||
|
dups = next(i for i in proof.items if i.dimension == prr.DIM_DUPLICATES)
|
||||||
|
self.assertEqual(dups.status, prr.ITEM_UNRESOLVED)
|
||||||
|
self.assertEqual(dups.details["duplicates"][0]["claim_count"], 2)
|
||||||
|
self.assertTrue(proof.mutation_hold)
|
||||||
|
|
||||||
|
|
||||||
|
class OrphanSessionTest(unittest.TestCase):
|
||||||
|
def test_active_session_dead_pid_is_unresolved(self) -> None:
|
||||||
|
proof = prr.reconcile_after_restart(
|
||||||
|
_base_inventory(
|
||||||
|
sessions=[
|
||||||
|
{
|
||||||
|
"session_id": "ghost",
|
||||||
|
"status": "active",
|
||||||
|
"pid": 2_000_000_000,
|
||||||
|
"role": "author",
|
||||||
|
}
|
||||||
|
]
|
||||||
|
),
|
||||||
|
now=NOW,
|
||||||
|
mode=prr.MODE_ENFORCE,
|
||||||
|
)
|
||||||
|
sess = next(i for i in proof.items if i.dimension == prr.DIM_SESSIONS)
|
||||||
|
self.assertEqual(sess.status, prr.ITEM_UNRESOLVED)
|
||||||
|
self.assertIn("ghost", sess.details["orphan_session_ids"])
|
||||||
|
|
||||||
|
|
||||||
|
class CapabilityStaleTest(unittest.TestCase):
|
||||||
|
def test_stale_runtime_unresolved(self) -> None:
|
||||||
|
proof = prr.reconcile_after_restart(
|
||||||
|
_base_inventory(capabilities={"stale": True, "startup_head": "aaa"}),
|
||||||
|
now=NOW,
|
||||||
|
mode=prr.MODE_ENFORCE,
|
||||||
|
)
|
||||||
|
caps = next(i for i in proof.items if i.dimension == prr.DIM_CAPABILITIES)
|
||||||
|
self.assertEqual(caps.status, prr.ITEM_UNRESOLVED)
|
||||||
|
self.assertTrue(proof.mutation_hold)
|
||||||
|
|
||||||
|
|
||||||
|
class MutationsAllowedHelperTest(unittest.TestCase):
|
||||||
|
def test_mutations_allowed_respects_hold(self) -> None:
|
||||||
|
held = prr.reconcile_after_restart(
|
||||||
|
_base_inventory(
|
||||||
|
pending_mutations=[{"phase": "merging", "session_id": "x"}]
|
||||||
|
),
|
||||||
|
now=NOW,
|
||||||
|
mode=prr.MODE_ENFORCE,
|
||||||
|
)
|
||||||
|
self.assertFalse(prr.mutations_allowed(held))
|
||||||
|
self.assertFalse(prr.mutations_allowed(held.as_dict()))
|
||||||
|
clean = prr.reconcile_after_restart(
|
||||||
|
_base_inventory(), now=NOW, mode=prr.MODE_ENFORCE
|
||||||
|
)
|
||||||
|
self.assertTrue(prr.mutations_allowed(clean))
|
||||||
|
|
||||||
|
|
||||||
|
class CheckpointSoftDependencyTest(unittest.TestCase):
|
||||||
|
def test_checkpoints_when_schema_present(self) -> None:
|
||||||
|
proof = prr.reconcile_after_restart(
|
||||||
|
_base_inventory(
|
||||||
|
checkpoints_available=True,
|
||||||
|
checkpoints=[{"session_id": "s1", "stale": False}],
|
||||||
|
),
|
||||||
|
now=NOW,
|
||||||
|
)
|
||||||
|
cp = next(i for i in proof.items if i.dimension == prr.DIM_CHECKPOINTS)
|
||||||
|
self.assertEqual(cp.status, prr.ITEM_RESOLVED)
|
||||||
|
|
||||||
|
def test_stale_checkpoints_unresolved(self) -> None:
|
||||||
|
proof = prr.reconcile_after_restart(
|
||||||
|
_base_inventory(
|
||||||
|
checkpoints_available=True,
|
||||||
|
checkpoints=[{"session_id": "s1", "stale": True}],
|
||||||
|
),
|
||||||
|
now=NOW,
|
||||||
|
mode=prr.MODE_ENFORCE,
|
||||||
|
)
|
||||||
|
cp = next(i for i in proof.items if i.dimension == prr.DIM_CHECKPOINTS)
|
||||||
|
self.assertEqual(cp.status, prr.ITEM_UNRESOLVED)
|
||||||
|
|
||||||
|
|
||||||
|
class ProofSerializationTest(unittest.TestCase):
|
||||||
|
def test_as_dict_is_json_friendly(self) -> None:
|
||||||
|
proof = prr.reconcile_after_restart(_base_inventory(), now=NOW)
|
||||||
|
blob = json.dumps(proof.as_dict())
|
||||||
|
self.assertIn("reconcile_id", blob)
|
||||||
|
self.assertIn("proposed_follow_ups", blob)
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
unittest.main()
|
||||||
Reference in New Issue
Block a user