feat(control-plane): durable MCP session checkpoint schema (Closes #660)

Add a versioned, redacted, reconcile-on-boot session checkpoint store to
control_plane_db so a restart can recover session identity, stage, lease
ownership, and next valid action instead of forcing human reconstruction
(umbrella #655; #628 autonomous-handoff goal).

- Bump SCHEMA_VERSION 4->5. New `session_checkpoints` table + two indexes,
  added via `CREATE TABLE IF NOT EXISTS` so table creation is itself the
  additive, idempotent v4->v5 migration (dependency_edges precedent).
- Writer `write_session_checkpoint` upserts the current recoverable state
  per (remote, org, repo, session_id, work_kind, work_number); stage
  transitions audit to `events`. Readers `get_session_checkpoint` /
  `list_session_checkpoints`.
- Every free-text and JSON field is passed through `gitea_audit.redact`
  before storage — no tokens/credential URLs can land in a checkpoint (AC4).
- `reconcile_session_checkpoint` is pure and never restores: it diagnoses a
  stored checkpoint against live head/lease state and flags staleness (AC3).
- Drain gate: `require_complete=True` fails closed (writes nothing) when a
  checkpoint is missing a drain-required field, so drain cannot complete on
  an unrecoverable record.
- `lease_id`/`assignment_id` are soft references (no enforced FK) so a
  checkpoint survives deletion of the lease it names; `work_number=0` is the
  NULL-safe session-level sentinel.

Tests: 13 new cases (schema/version, multi-role fixtures, upsert+audit,
JSON round-trip, secret redaction, reconcile stale-head/dead-lease/
reassigned-lease/clean/unknown, drain fail-closed, sentinel key). Existing
schema_version assertion updated 4->5. Full tests/test_control_plane_db.py
suite: 33/33 pass.

Links #652 #653 #655.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
This commit is contained in:
2026-07-24 08:06:45 -04:00
co-authored by Claude Opus 4.8
parent 5e935dffb4
commit 7e18dccbd9
2 changed files with 654 additions and 2 deletions
+226 -1
View File
@@ -11,6 +11,7 @@ from datetime import timedelta
from control_plane_db import (
ControlPlaneDB,
ControlPlaneError,
InvalidWorkKindError,
LeaseRequiredError,
WORK_KINDS,
@@ -36,7 +37,7 @@ class ControlPlaneDBTest(unittest.TestCase):
rows = dict(conn.execute("SELECT key, value FROM schema_meta").fetchall())
finally:
conn.close()
self.assertEqual(rows["schema_version"], "4")
self.assertEqual(rows["schema_version"], "5")
self.assertIn("DB coordinates", rows["architecture"])
self.assertIn("bridge", rows["architecture"].lower())
@@ -820,5 +821,229 @@ class ControlPlaneDBTest(unittest.TestCase):
self.assertEqual(n, 1)
class SessionCheckpointTest(unittest.TestCase):
"""Durable MCP session checkpoint schema, redaction, and reconcile (#660)."""
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory()
self.db_path = os.path.join(self._tmp.name, "cp.sqlite3")
self.db = ControlPlaneDB(self.db_path)
def tearDown(self) -> None:
self._tmp.cleanup()
def _write(self, **overrides):
base = dict(
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
session_id="prgs-author-1-abc",
role="author",
work_kind="issue",
work_number=660,
branch="feat/issue-660-session-checkpoint-schema",
head_sha="deadbeef",
lease_id="lease-1",
workflow_stage="implementing",
last_completed_action="wrote schema",
next_valid_action="write tests",
recovery_instructions="re-lock #660 then continue tests",
)
base.update(overrides)
return self.db.write_session_checkpoint(**base)
# AC1 — schema documented and versioned.
def test_table_exists_and_row_carries_schema_version(self) -> None:
import sqlite3
conn = sqlite3.connect(self.db_path)
try:
names = {
r[0]
for r in conn.execute(
"SELECT name FROM sqlite_master WHERE type='table'"
).fetchall()
}
finally:
conn.close()
self.assertIn("session_checkpoints", names)
record = self._write()
self.assertEqual(record["checkpoint_schema_version"], 5)
# AC2 — checkpoints written for multi-role session fixtures.
def test_multi_role_fixtures_each_get_a_row(self) -> None:
roles = [
("prgs-author-1", "author", "issue", 660),
("prgs-reviewer-2", "reviewer", "pr", 795),
("prgs-merger-3", "merger", "pr", 862),
("prgs-controller-4", "controller", "issue", 653),
]
for session_id, role, kind, number in roles:
self._write(
session_id=session_id,
role=role,
work_kind=kind,
work_number=number,
lease_id=f"lease-{session_id}",
)
rows = self.db.list_session_checkpoints(remote="prgs")
self.assertEqual(len(rows), 4)
self.assertEqual(
{r["role"] for r in rows},
{"author", "reviewer", "merger", "controller"},
)
def test_upsert_is_current_state_and_audits_stage_change(self) -> None:
first = self._write(workflow_stage="implementing")
second = self._write(workflow_stage="testing")
self.assertEqual(first["checkpoint_id"], second["checkpoint_id"])
rows = self.db.list_session_checkpoints(
remote="prgs", session_id="prgs-author-1-abc"
)
self.assertEqual(len(rows), 1)
self.assertEqual(rows[0]["workflow_stage"], "testing")
import sqlite3
conn = sqlite3.connect(self.db_path)
try:
n = conn.execute(
"SELECT COUNT(*) FROM events "
"WHERE event_type = 'session_checkpoint_stage_change'"
).fetchone()[0]
finally:
conn.close()
self.assertEqual(n, 1)
def test_get_and_roundtrip_json_fields(self) -> None:
self._write(
capabilities=["gitea.repo.commit", "gitea.pr.create"],
evidence={"tests": "4 passing"},
pending_mutation={"op": "commit_files", "files": ["control_plane_db.py"]},
)
got = self.db.get_session_checkpoint(
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
session_id="prgs-author-1-abc",
work_kind="issue",
work_number=660,
)
self.assertIsNotNone(got)
self.assertEqual(got["capabilities"], ["gitea.repo.commit", "gitea.pr.create"])
self.assertEqual(got["evidence"], {"tests": "4 passing"})
self.assertEqual(got["pending_mutation"]["op"], "commit_files")
# AC4 — no secrets in stored records.
def test_secrets_are_redacted_before_storage(self) -> None:
self._write(
recovery_instructions=(
"resume with Authorization: Bearer sk-supersecrettoken then retry"
),
evidence={"authorization": "Bearer sk-anothersecret"},
pending_mutation={"url": "https://user:[email protected]/repo.git"},
)
import sqlite3
conn = sqlite3.connect(self.db_path)
try:
row = conn.execute(
"SELECT recovery_instructions, evidence, pending_mutation "
"FROM session_checkpoints"
).fetchone()
finally:
conn.close()
blob = " ".join(str(v) for v in row)
self.assertNotIn("sk-supersecrettoken", blob)
self.assertNotIn("sk-anothersecret", blob)
self.assertNotIn("password", blob)
self.assertIn("REDACTED", blob)
# AC3 — reconcile detects stale head / lease mismatch.
def test_reconcile_flags_stale_head(self) -> None:
record = self._write(head_sha="aaaa1111")
result = self.db.reconcile_session_checkpoint(
record, live_head_sha="bbbb2222", live_lease_active=True,
live_lease_id="lease-1",
)
self.assertTrue(result["stale"])
self.assertTrue(result["head_mismatch"])
self.assertFalse(result["lease_mismatch"])
self.assertEqual(result["reconcile_action"], "reconcile_required")
def test_reconcile_flags_dead_lease(self) -> None:
record = self._write(lease_id="lease-1", head_sha="aaaa1111")
result = self.db.reconcile_session_checkpoint(
record, live_head_sha="aaaa1111", live_lease_active=False,
)
self.assertTrue(result["stale"])
self.assertFalse(result["head_mismatch"])
self.assertTrue(result["lease_mismatch"])
def test_reconcile_reassigned_lease_is_stale(self) -> None:
record = self._write(lease_id="lease-1")
result = self.db.reconcile_session_checkpoint(
record, live_lease_active=True, live_lease_id="lease-999",
)
self.assertTrue(result["lease_mismatch"])
def test_reconcile_clean_state_is_safe_to_resume(self) -> None:
record = self._write(head_sha="aaaa1111", lease_id="lease-1")
result = self.db.reconcile_session_checkpoint(
record, live_head_sha="aaaa1111", live_lease_active=True,
live_lease_id="lease-1",
)
self.assertFalse(result["stale"])
self.assertEqual(result["reconcile_action"], "safe_to_resume")
def test_unknown_live_state_never_flags_mismatch(self) -> None:
record = self._write(head_sha="aaaa1111", lease_id="lease-1")
result = self.db.reconcile_session_checkpoint(record)
self.assertFalse(result["stale"])
# Drain gate — fail closed when a checkpoint is incomplete.
def test_drain_requires_complete_checkpoint(self) -> None:
with self.assertRaises(ControlPlaneError):
self.db.write_session_checkpoint(
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
session_id="prgs-author-1-abc",
role="author",
workflow_stage="draining",
# next_valid_action + recovery_instructions intentionally absent
require_complete=True,
)
# Nothing was written.
rows = self.db.list_session_checkpoints(remote="prgs")
self.assertEqual(rows, [])
def test_drain_write_succeeds_when_complete(self) -> None:
record = self._write(require_complete=True)
self.assertEqual(record["status"], "active")
self.assertEqual(self.db.checkpoint_completeness(record), [])
def test_missing_session_id_fails_closed(self) -> None:
with self.assertRaises(ControlPlaneError):
self.db.write_session_checkpoint(
remote="prgs", org="o", repo="r", session_id="",
)
def test_session_level_checkpoint_uses_sentinel_key(self) -> None:
# No work unit -> ('', 0) sentinel; a second session-level write upserts.
self.db.write_session_checkpoint(
remote="prgs", org="o", repo="r", session_id="s-sess",
workflow_stage="idle",
)
self.db.write_session_checkpoint(
remote="prgs", org="o", repo="r", session_id="s-sess",
workflow_stage="booting",
)
rows = self.db.list_session_checkpoints(remote="prgs", session_id="s-sess")
self.assertEqual(len(rows), 1)
self.assertEqual(rows[0]["work_kind"], "")
self.assertEqual(rows[0]["work_number"], 0)
if __name__ == "__main__":
unittest.main()