diff --git a/docs/post-restart-reconcile.md b/docs/post-restart-reconcile.md new file mode 100644 index 0000000..42511b3 --- /dev/null +++ b/docs/post-restart-reconcile.md @@ -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 diff --git a/gitea_mcp_server.py b/gitea_mcp_server.py index a658eb1..c8a9461 100644 --- a/gitea_mcp_server.py +++ b/gitea_mcp_server.py @@ -17880,6 +17880,17 @@ def gitea_assess_master_parity( } if parity["restart_required"] and enforced: 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 @@ -22328,6 +22339,249 @@ def gitea_request_mcp_restart( 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() def gitea_inspect_workflow_lease( lease_id: str, diff --git a/post_restart_reconcile.py b/post_restart_reconcile.py new file mode 100644 index 0000000..6477197 --- /dev/null +++ b/post_restart_reconcile.py @@ -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 diff --git a/task_capability_map.py b/task_capability_map.py index 01309fc..9875a2d 100644 --- a/task_capability_map.py +++ b/task_capability_map.py @@ -123,6 +123,16 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = { "permission": "gitea.branch.push", "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); # update-by-merge is author-only and mutates the PR head via Gitea API. "assess_pr_sync_status": { diff --git a/tests/test_issue_662_post_restart_reconcile.py b/tests/test_issue_662_post_restart_reconcile.py new file mode 100644 index 0000000..0782397 --- /dev/null +++ b/tests/test_issue_662_post_restart_reconcile.py @@ -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()