feat(mcp): post-restart reconciliation and completion proof (Closes #662)
Add pure post_restart_reconcile.reconcile_after_restart classifier with a machine-readable completion proof covering service health, sessions, leases, capabilities, worktrees, interrupted mutations (never auto-resumed), duplicates, and queue state. Soft-depends on #660 checkpoints (skipped with reason when the schema module is absent). Wire read-only MCP tool gitea_reconcile_after_restart, boot-once hook via gitea_assess_master_parity, log_only/enforce modes (mutation_hold), and docs. Closes #662 Related: #655 #652 #653 #660 #661 Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
This commit is contained in:
@@ -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
|
||||
@@ -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,
|
||||
|
||||
@@ -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
|
||||
@@ -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": {
|
||||
|
||||
@@ -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