diff --git a/tolaria/c5-sync-service/.gitignore b/tolaria/c5-sync-service/.gitignore new file mode 100644 index 0000000..a7096c0 --- /dev/null +++ b/tolaria/c5-sync-service/.gitignore @@ -0,0 +1,4 @@ +*.db +*.db-wal +*.db-shm +__pycache__/ diff --git a/tolaria/c5-sync-service/README.md b/tolaria/c5-sync-service/README.md new file mode 100644 index 0000000..c5f2f29 --- /dev/null +++ b/tolaria/c5-sync-service/README.md @@ -0,0 +1,207 @@ +# C5A — CONTRACT & STATE MACHINE v1 + +**Phase:** C5A (erste Phase des C5 Sync Service) +**Status:** DONE (Contract & State Machine) — **INAKTIV, NO-WRITE** +**Architektur:** FORGEJO MASTER → C5 SYNC SERVICE → TOLARIA DERIVED → SEARCH FULL REBUILD + +C5A implementiert **ausschließlich** den Contract + die Sync State Machine + die +persistente Progress-State-Struktur. Es ist eine **deterministische, inaktive +Python-Library + CLI + Testsuite** — es führt **keine** externen Writes aus. + +--- + +## NO-WRITE GUARANTEE (verbindlich) + +C5A kann **NOCH NICHT**: +- Tolaria schreiben +- Search rebuilden +- Forgejo schreiben +- produktiv pollen +- Netzwerk-Mutationen ausführen + +Es gibt **keine** Netzwerk-/HTTP-/Socket-/Subprocess-Mutationsfunktionen in C5A. +Die No-Write-Guarantee wird statisch per `assert_no_write_guarantee()` bewiesen +(keine `requests`/`urllib`/`http`/`socket`/`subprocess`-Imports) und durch den +Test `test_no_write_guarantee` abgesichert. + +--- + +## State Machine + +### Normalzustände +`DISCOVERED → VALIDATING → READY → PROPAGATING_TOLARIA → VERIFYING_TOLARIA → UPDATING_SEARCH → VERIFYING_SEARCH → APPLIED` + +### Fehlerzustände +`RETRY_PENDING`, `FAILED`, `DEAD`, `HUMAN_REVIEW_REQUIRED` + +### Ordering +`WAITING_FOR_PREDECESSOR` — Commits werden nur in korrekter Reihenfolge +verarbeitet (`parent_sha` / `last_applied_commit`). Lücken werden **nicht** +übersprungen. + +### Kerninvariante +**`last_applied_commit` wird NUR nach vollständigem Tolaria+Search PASS +fortgeschrieben** (`mark_applied` wird nur nach `VERIFYING_SEARCH → APPLIED` +aufgerufen). + +--- + +## Commit Contract + +Eine Verarbeitungseinheit = **ein Forgejo Commit**: + +| Feld | Beschreibung | +|---|---| +| `commit_sha` | Primärschlüssel | +| `parent_sha` | Vorgänger (Ordering) | +| `discovered_at` | Zeitstempel | +| `sequence` | Reihenfolge | +| `status` | Sync-Zustand | +| `retry_count` | Retry-Zähler | +| `last_error` / `last_error_code` | Fehler + Reason Code | +| `created_at` / `updated_at` | Zeitstempel | + +### Pro Changed Object +`object_id`, `path_before`, `path_after`, `operation`, `content_hash_before`, +`content_hash_after`, `metadata_hash_before`, `metadata_hash_after`, +`representation`, `state` + +--- + +## Operation Model + +`CREATE`, `CONTENT_UPDATE`, `METADATA_UPDATE`, `STATE_UPDATE`, `RENAME`, `MOVE`, +`SOURCE_CANONICAL_RELATION_UPDATE`, `TAGS_UPDATE`, `SUPERSEDE`, `DELETE_REQUEST` + +**`DELETE_REQUEST` bleibt Human-Gate** — kein automatisches Hard Delete. + +--- + +## Idempotenz + +Basis: **`commit_sha + object_id + operation`** + Hash-Checks. + +| Ergebnis | Bedeutung | +|---|---| +| `ALREADY_APPLIED` | Commit bereits vollständig angewendet | +| `ALREADY_AT_TARGET` | Ziel-Hash bereits erreicht | +| `RETRY_SAFE` | Noch nicht angewendet, sicher zu verarbeiten | +| `CONFLICT` | Existiert, aber Hash weicht ab | + +Doppelte Commit-Erkennung ist **deterministisch**. + +--- + +## Ordering + +Commits nur in korrekter Reihenfolge. Fehlender Vorgänger → `WAITING_FOR_PREDECESSOR` +(Reason Code `OUT_OF_ORDER_COMMIT`). Nach Anwendung des Vorgängers wird der +Commit erneut validiert. + +--- + +## Failure / Retry State + +- `max_retries = 5` (konfigurierbar) +- Backoff: `1 / 2 / 4 / 8 / 16` Sekunden (gekappt) +- Nach Max → `DEAD` oder `HUMAN_REVIEW_REQUIRED` +- **Kein Endlos-Retry** + +--- + +## Persistence + +SQLite (`c5a.db`), isoliert, restart-fest. Anforderungen erfüllt: +- **atomic** (SQLite-Transaktionen, WAL) +- **restart-safe** (neuer Store auf gleicher DB behält State) +- **inspectable** (CLI `health`, `list-commits`, `commit-status`) +- **backupbar** (einzelne Datei) +- **keine externe DB** nötig +- **kein Trading-/Forgejo-DB-Coupling** (eigene DB) + +--- + +## Bootstrap State + +`UNINITIALIZED → RECONCILING → BASELINE_READY → ACTIVE` + +Baseline darf **nur** aus `BASELINE_READY` gesetzt werden (Master↔Tolaria- +Reconciliation muss sauber sein). In C5A wird **keine produktive Baseline** +gesetzt. + +--- + +## Reason Codes (geschlossene Menge) + +`UNEXPECTED_TOLARIA_DRIFT`, `ID_COLLISION`, `UNKNOWN_OBJECT_ID`, `INVALID_SCHEMA`, +`DANGLING_DERIVED_FROM`, `AMBIGUOUS_DELETE`, `UNKNOWN_LEGACY_OBJECT`, +`AUTH_FAILURE`, `SEARCH_REBUILD_FAILURE`, `TOLARIA_UNAVAILABLE`, +`FORGEJO_UNAVAILABLE`, `OUT_OF_ORDER_COMMIT`, `SECRET_DETECTED` + +Kein freier String-Wildwuchs als einziges Fehlerformat — unbekannte Reason Codes +werden abgelehnt (`InvalidReasonCodeError`). + +--- + +## Security Boundary (dokumentiert, NICHT implementiert) + +- **Forgejo credential:** READ ONLY (C5A hält keinen Forgejo-Write) +- **Tolaria write:** nur später C5 Service +- **Search rebuild:** nur C5 Service Token +- **Agents:** keine direkten Tolaria-/Search-Admin-Writes +- **Netzwerk-Härtung** wird später beim Deployment umgesetzt, NICHT jetzt + +--- + +## Health Contract + +`status`, `bootstrap_state`, `last_seen_commit`, `last_applied_commit`, +`pending_commits`, `failed_commits`, `dead_commits`, `human_review_required`, +`forgejo_status`, `tolaria_status`, `search_status`, `drift_count`, +`last_success_at`, `last_error_code` + +C5A führt **keine echten externen Health-Probes** aus (rein State-Machine) — +`forgejo_status`/`tolaria_status`/`search_status` sind `UNKNOWN`. + +--- + +## CLI + +```bash +export C5A_DB=/path/to/c5a.db # default: ./c5a.db +python3 rq_c5a_cli.py health +python3 rq_c5a_cli.py bootstrap-state +python3 rq_c5a_cli.py bootstrap-transition --to RECONCILING +python3 rq_c5a_cli.py set-baseline --commit +python3 rq_c5a_cli.py ingest-commit --json +python3 rq_c5a_cli.py process-commit --json +python3 rq_c5a_cli.py list-commits [--status ] +python3 rq_c5a_cli.py commit-status --sha +python3 rq_c5a_cli.py no-write-check +``` + +--- + +## Tests + +```bash +python3 test_c5a.py # 25 PASS / 0 FAIL +``` + +Abgedeckte Fälle: normal commit lifecycle, multi-object commit, duplicate commit, +already-applied, out-of-order commit, missing predecessor, retry progression, +max retry → DEAD, human gate transition, restart/reload persistence, state +corruption handling, idempotency, delete request → human gate, dangling relation +→ human gate, unexpected drift → human gate, no-write guarantee, bootstrap +lifecycle, baseline-only-from-BASELINE_READY, reason codes closed set, operations +closed set, health contract fields, last_applied-only-after-full-pass, retry +available bound, backoff sequence, invalid reason code rejected. + +--- + +## Dateien + +- `rq_c5a.py` — Kern-Library (Store + State Machine + No-Write-Check) +- `rq_c5a_cli.py` — CLI +- `test_c5a.py` — Testsuite (25 Tests) +- `README.md` — diese Datei diff --git a/tolaria/c5-sync-service/rq_c5a.py b/tolaria/c5-sync-service/rq_c5a.py new file mode 100644 index 0000000..11b28ab --- /dev/null +++ b/tolaria/c5-sync-service/rq_c5a.py @@ -0,0 +1,795 @@ +#!/usr/bin/env python3 +""" +Red Queen — C5A: CONTRACT & STATE MACHINE v1 (deterministische, inaktive Library). + +VERBINDLICHER RAHMEN +==================== +C5A ist die erste Phase des C5 Sync Service (FORGEJO MASTER -> C5 SYNC SERVICE -> +TOLARIA DERIVED -> SEARCH FULL REBUILD). Diese Library implementiert AUSSCHLIESSLICH +den Contract + die Sync State Machine + die persistente Progress-State-Struktur. + +C5A ist INAKTIV und NO-WRITE: + * Es gibt KEIN produktives Polling. + * Es gibt KEINE Forgejo->Tolaria Writes. + * Es gibt KEINE Search Rebuild Calls. + * Es gibt KEINEN C5 Container Deployment. + * Es gibt KEINE Netzwerk-Isolation. + * Es gibt KEINEN Canary. + * Es gibt KEINEN Hermes Write Access. + * Es gibt KEINE Netzwerk-Mutationsfunktionen (kein HTTP-Client, kein Socket). + +C5A modelliert NUR den Zustand und die Uebergaenge. Jede Aktion wird EXPLIZIT durch +eine Methode/CLI-Befehl angestossen und ist garantiert begrenzt (bounded). Es gibt +KEINEN Loop, kein Heartbeat, kein Cron, kein Daemon, kein self-reschedule. + +PERSISTENZ +========== +C5A-State ist restart-fest und liegt in einer isolierten SQLite-DB (c5a.db). +Anforderungen: atomic, restart-safe, inspectable, backupbar, keine externe DB, +kein Trading-/Forgejo-DB-Coupling. SQLite ist die robusteste einfache Variante. + +SECURITY BOUNDARY (dokumentiert, NICHT implementiert): + * Forgejo credential: READ ONLY (C5A haelt KEINEN Forgejo-Write). + * Tolaria write: nur spaeter C5 Service (C5A schreibt NIE nach Tolaria). + * Search rebuild: nur C5 Service Token (C5A ruft NIE Search-Rebuild). + * Agents: keine direkten Tolaria-/Search-Admin-Writes. + * Netzwerk-Haertung wird spaeter beim Deployment umgesetzt, NICHT jetzt. +""" + +from __future__ import annotations + +import json +import os +import sqlite3 +import time +from pathlib import Path +from typing import Any, Dict, List, Optional + +# --------------------------------------------------------------------------- +# Konstanten +# --------------------------------------------------------------------------- + +# Sync State Machine — Normalzustände +ST_DISCOVERED = "DISCOVERED" +ST_VALIDATING = "VALIDATING" +ST_READY = "READY" +ST_PROPAGATING_TOLARIA = "PROPAGATING_TOLARIA" +ST_VERIFYING_TOLARIA = "VERIFYING_TOLARIA" +ST_UPDATING_SEARCH = "UPDATING_SEARCH" +ST_VERIFYING_SEARCH = "VERIFYING_SEARCH" +ST_APPLIED = "APPLIED" + +# Sync State Machine — Fehlerzustände +ST_RETRY_PENDING = "RETRY_PENDING" +ST_FAILED = "FAILED" +ST_DEAD = "DEAD" +ST_HUMAN_REVIEW_REQUIRED = "HUMAN_REVIEW_REQUIRED" + +# Ordering-Zustand +ST_WAITING_FOR_PREDECESSOR = "WAITING_FOR_PREDECESSOR" + +# Bootstrap-Zustände +BS_UNINITIALIZED = "UNINITIALIZED" +BS_RECONCILING = "RECONCILING" +BS_BASELINE_READY = "BASELINE_READY" +BS_ACTIVE = "ACTIVE" + +# Operation Model +OP_CREATE = "CREATE" +OP_CONTENT_UPDATE = "CONTENT_UPDATE" +OP_METADATA_UPDATE = "METADATA_UPDATE" +OP_STATE_UPDATE = "STATE_UPDATE" +OP_RENAME = "RENAME" +OP_MOVE = "MOVE" +OP_SOURCE_CANONICAL_RELATION_UPDATE = "SOURCE_CANONICAL_RELATION_UPDATE" +OP_TAGS_UPDATE = "TAGS_UPDATE" +OP_SUPERSEDE = "SUPERSEDE" +OP_DELETE_REQUEST = "DELETE_REQUEST" + +# Idempotency-Ergebnisse +IDEM_ALREADY_APPLIED = "ALREADY_APPLIED" +IDEM_ALREADY_AT_TARGET = "ALREADY_AT_TARGET" +IDEM_RETRY_SAFE = "RETRY_SAFE" +IDEM_CONFLICT = "CONFLICT" + +# Reason Codes (feste, maschinenlesbare Fehlerformate — kein freier String-Wildwuchs) +RC_UNEXPECTED_TOLARIA_DRIFT = "UNEXPECTED_TOLARIA_DRIFT" +RC_ID_COLLISION = "ID_COLLISION" +RC_UNKNOWN_OBJECT_ID = "UNKNOWN_OBJECT_ID" +RC_INVALID_SCHEMA = "INVALID_SCHEMA" +RC_DANGLING_DERIVED_FROM = "DANGLING_DERIVED_FROM" +RC_AMBIGUOUS_DELETE = "AMBIGUOUS_DELETE" +RC_UNKNOWN_LEGACY_OBJECT = "UNKNOWN_LEGACY_OBJECT" +RC_AUTH_FAILURE = "AUTH_FAILURE" +RC_SEARCH_REBUILD_FAILURE = "SEARCH_REBUILD_FAILURE" +RC_TOLARIA_UNAVAILABLE = "TOLARIA_UNAVAILABLE" +RC_FORGEJO_UNAVAILABLE = "FORGEJO_UNAVAILABLE" +RC_OUT_OF_ORDER_COMMIT = "OUT_OF_ORDER_COMMIT" +RC_SECRET_DETECTED = "SECRET_DETECTED" + +# Retry / Backoff +DEFAULT_MAX_RETRIES = 5 +DEFAULT_BACKOFF_SECONDS = [1, 2, 4, 8, 16] + +# Health-Status +HEALTH_OK = "OK" +HEALTH_DEGRADED = "DEGRADED" +HEALTH_ERROR = "ERROR" + +# Alle Reason Codes als frozenset (Determinismus + Validierung) +REASON_CODES = frozenset({ + RC_UNEXPECTED_TOLARIA_DRIFT, + RC_ID_COLLISION, + RC_UNKNOWN_OBJECT_ID, + RC_INVALID_SCHEMA, + RC_DANGLING_DERIVED_FROM, + RC_AMBIGUOUS_DELETE, + RC_UNKNOWN_LEGACY_OBJECT, + RC_AUTH_FAILURE, + RC_SEARCH_REBUILD_FAILURE, + RC_TOLARIA_UNAVAILABLE, + RC_FORGEJO_UNAVAILABLE, + RC_OUT_OF_ORDER_COMMIT, + RC_SECRET_DETECTED, +}) + +# Alle Operationen als frozenset +OPERATIONS = frozenset({ + OP_CREATE, + OP_CONTENT_UPDATE, + OP_METADATA_UPDATE, + OP_STATE_UPDATE, + OP_RENAME, + OP_MOVE, + OP_SOURCE_CANONICAL_RELATION_UPDATE, + OP_TAGS_UPDATE, + OP_SUPERSEDE, + OP_DELETE_REQUEST, +}) + +# Alle Sync-Zustände als frozenset +SYNC_STATES = frozenset({ + ST_DISCOVERED, ST_VALIDATING, ST_READY, + ST_PROPAGATING_TOLARIA, ST_VERIFYING_TOLARIA, + ST_UPDATING_SEARCH, ST_VERIFYING_SEARCH, ST_APPLIED, + ST_RETRY_PENDING, ST_FAILED, ST_DEAD, ST_HUMAN_REVIEW_REQUIRED, + ST_WAITING_FOR_PREDECESSOR, +}) + +# Bootstrap-Zustände +BOOTSTRAP_STATES = frozenset({ + BS_UNINITIALIZED, BS_RECONCILING, BS_BASELINE_READY, BS_ACTIVE, +}) + +# Erlaubte Übergänge der Sync State Machine (deterministisch) +# (from_state, to_state) -> erlaubt +_ALLOWED_TRANSITIONS = { + (ST_DISCOVERED, ST_VALIDATING), + (ST_DISCOVERED, ST_WAITING_FOR_PREDECESSOR), # Lücke direkt nach Discovery + (ST_DISCOVERED, ST_RETRY_PENDING), # Fehler direkt nach Discovery (z.B. Forgejo down) + (ST_VALIDATING, ST_READY), + (ST_VALIDATING, ST_HUMAN_REVIEW_REQUIRED), # Schema/ID/Drift-Fehler bei Validierung + (ST_VALIDATING, ST_WAITING_FOR_PREDECESSOR), # Lücke in Commit-Reihenfolge + (ST_VALIDATING, ST_RETRY_PENDING), # Validierungsfehler retrybar + (ST_WAITING_FOR_PREDECESSOR, ST_VALIDATING), # Vorgänger angewendet -> erneut validieren + (ST_READY, ST_PROPAGATING_TOLARIA), + (ST_PROPAGATING_TOLARIA, ST_VERIFYING_TOLARIA), + (ST_PROPAGATING_TOLARIA, ST_RETRY_PENDING), # Tolaria-Write-Fehler + (ST_VERIFYING_TOLARIA, ST_UPDATING_SEARCH), + (ST_VERIFYING_TOLARIA, ST_RETRY_PENDING), # DRIFT != 0 / Tolaria-Fehler + (ST_VERIFYING_TOLARIA, ST_HUMAN_REVIEW_REQUIRED), # unexpected drift + (ST_UPDATING_SEARCH, ST_VERIFYING_SEARCH), + (ST_UPDATING_SEARCH, ST_RETRY_PENDING), # Search-Rebuild-Fehler + (ST_UPDATING_SEARCH, ST_HUMAN_REVIEW_REQUIRED), + (ST_VERIFYING_SEARCH, ST_APPLIED), + (ST_VERIFYING_SEARCH, ST_RETRY_PENDING), # Search-Health-Fehler + (ST_VERIFYING_SEARCH, ST_HUMAN_REVIEW_REQUIRED), + (ST_RETRY_PENDING, ST_READY), # Retry -> erneut propagieren + (ST_RETRY_PENDING, ST_RETRY_PENDING), # idempotenter erneuter Retry-Versuch + (ST_RETRY_PENDING, ST_FAILED), # Retry-Limit erreicht + (ST_RETRY_PENDING, ST_DEAD), # Max-Retry -> DEAD + (ST_RETRY_PENDING, ST_HUMAN_REVIEW_REQUIRED),# Human-Gate-Fehler + (ST_FAILED, ST_HUMAN_REVIEW_REQUIRED), # FAILED -> Human Review + (ST_DEAD, ST_HUMAN_REVIEW_REQUIRED), # DEAD -> Human Review + (ST_HUMAN_REVIEW_REQUIRED, ST_READY), # Human entscheidet -> erneut + (ST_HUMAN_REVIEW_REQUIRED, ST_APPLIED), # Human bestätigt als angewendet + (ST_HUMAN_REVIEW_REQUIRED, ST_DEAD), # Human verwirft +} + +# Bootstrap-Übergänge +_ALLOWED_BOOTSTRAP = { + (BS_UNINITIALIZED, BS_RECONCILING), + (BS_RECONCILING, BS_BASELINE_READY), + (BS_RECONCILING, BS_UNINITIALIZED), # Reconciliation fehlgeschlagen -> zurück + (BS_BASELINE_READY, BS_ACTIVE), +} + + +class C5AError(Exception): + """Basis-Fehlerklasse für C5A.""" + + def __init__(self, message: str, reason_code: Optional[str] = None): + super().__init__(message) + self.message = message + self.reason_code = reason_code + + def to_dict(self) -> Dict[str, Any]: + return {"error": self.message, "reason_code": self.reason_code} + + +class InvalidTransitionError(C5AError): + """Ungültiger State-Machine-Übergang.""" + + +class InvalidReasonCodeError(C5AError): + """Unbekannter Reason Code.""" + + +class InvalidOperationError(C5AError): + """Unbekannte Operation.""" + + +class NoWriteGuaranteeError(C5AError): + """C5A darf keine externen Writes ausführen.""" + + +# --------------------------------------------------------------------------- +# Persistenz (SQLite, restart-fest, atomic, inspectable, backupbar) +# --------------------------------------------------------------------------- + +class C5AStore: + """ + Persistente Progress-State-Struktur für C5A. + + SQLite-basiert, isoliert (eigene c5a.db), restart-fest. Keine externe DB, + kein Trading-/Forgejo-DB-Coupling. Atomic via SQLite-Transaktionen. + """ + + SCHEMA_VERSION = 1 + + def __init__(self, db_path: str): + self.db_path = str(db_path) + self._conn = sqlite3.connect(self.db_path) + self._conn.row_factory = sqlite3.Row + self._conn.execute("PRAGMA journal_mode=WAL;") + self._init_schema() + + def _init_schema(self) -> None: + with self._conn: + self._conn.execute(""" + CREATE TABLE IF NOT EXISTS meta ( + key TEXT PRIMARY KEY, + value TEXT NOT NULL + ) + """) + self._conn.execute(""" + CREATE TABLE IF NOT EXISTS commits ( + commit_sha TEXT PRIMARY KEY, + parent_sha TEXT, + discovered_at INTEGER, + sequence INTEGER, + status TEXT NOT NULL, + retry_count INTEGER NOT NULL DEFAULT 0, + last_error TEXT, + last_error_code TEXT, + created_at INTEGER, + updated_at INTEGER + ) + """) + self._conn.execute(""" + CREATE TABLE IF NOT EXISTS objects ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + commit_sha TEXT NOT NULL, + object_id TEXT NOT NULL, + path_before TEXT, + path_after TEXT, + operation TEXT NOT NULL, + content_hash_before TEXT, + content_hash_after TEXT, + metadata_hash_before TEXT, + metadata_hash_after TEXT, + representation TEXT, + state TEXT, + UNIQUE(commit_sha, object_id, operation) + ) + """) + self._conn.execute(""" + CREATE TABLE IF NOT EXISTS bootstrap ( + id INTEGER PRIMARY KEY CHECK (id = 1), + state TEXT NOT NULL, + baseline_commit TEXT, + updated_at INTEGER + ) + """) + self._conn.execute(""" + CREATE TABLE IF NOT EXISTS health ( + id INTEGER PRIMARY KEY CHECK (id = 1), + last_seen_commit TEXT, + last_applied_commit TEXT, + last_success_at INTEGER, + last_error_code TEXT, + updated_at INTEGER + ) + """) + # Schema-Version setzen + self._conn.execute( + "INSERT OR REPLACE INTO meta (key, value) VALUES ('schema_version', ?)", + (str(self.SCHEMA_VERSION),), + ) + # Bootstrap initialisieren + row = self._conn.execute("SELECT state FROM bootstrap WHERE id = 1").fetchone() + if row is None: + self._conn.execute( + "INSERT INTO bootstrap (id, state, updated_at) VALUES (1, ?, ?)", + (BS_UNINITIALIZED, int(time.time() * 1000)), + ) + # Health initialisieren + row = self._conn.execute("SELECT id FROM health WHERE id = 1").fetchone() + if row is None: + self._conn.execute( + "INSERT INTO health (id, updated_at) VALUES (1, ?)", + (int(time.time() * 1000),), + ) + + # -- Bootstrap ---------------------------------------------------------- + + def bootstrap_state(self) -> str: + row = self._conn.execute("SELECT state FROM bootstrap WHERE id = 1").fetchone() + return row["state"] if row else BS_UNINITIALIZED + + def bootstrap_transition(self, to_state: str) -> str: + """Deterministischer Bootstrap-Übergang. Gibt den neuen Zustand zurück.""" + if to_state not in BOOTSTRAP_STATES: + raise InvalidTransitionError(f"Unbekannter Bootstrap-Zustand: {to_state}") + cur = self.bootstrap_state() + if (cur, to_state) not in _ALLOWED_BOOTSTRAP: + raise InvalidTransitionError( + f"Ungültiger Bootstrap-Übergang: {cur} -> {to_state}" + ) + with self._conn: + self._conn.execute( + "UPDATE bootstrap SET state = ?, updated_at = ? WHERE id = 1", + (to_state, int(time.time() * 1000)), + ) + return to_state + + def set_baseline(self, commit_sha: str) -> None: + """Setzt die Baseline (nur aus BASELINE_READY).""" + if self.bootstrap_state() != BS_BASELINE_READY: + raise C5AError( + "Baseline darf nur im Zustand BASELINE_READY gesetzt werden", + reason_code=RC_OUT_OF_ORDER_COMMIT, + ) + with self._conn: + self._conn.execute( + "UPDATE bootstrap SET baseline_commit = ?, updated_at = ? WHERE id = 1", + (commit_sha, int(time.time() * 1000)), + ) + self._conn.execute( + "UPDATE health SET last_seen_commit = ?, last_applied_commit = ?, updated_at = ? WHERE id = 1", + (commit_sha, commit_sha, int(time.time() * 1000)), + ) + + def baseline_commit(self) -> Optional[str]: + row = self._conn.execute("SELECT baseline_commit FROM bootstrap WHERE id = 1").fetchone() + return row["baseline_commit"] if row else None + + # -- Commits ------------------------------------------------------------ + + def upsert_commit(self, commit: Dict[str, Any]) -> Dict[str, Any]: + """Legt einen Commit an oder aktualisiert ihn (idempotent per commit_sha).""" + sha = commit["commit_sha"] + parent = commit.get("parent_sha") + discovered = commit.get("discovered_at", int(time.time() * 1000)) + seq = commit.get("sequence") + status = commit.get("status", ST_DISCOVERED) + retry = commit.get("retry_count", 0) + last_err = commit.get("last_error") + last_code = commit.get("last_error_code") + now = int(time.time() * 1000) + with self._conn: + self._conn.execute( + """ + INSERT INTO commits (commit_sha, parent_sha, discovered_at, sequence, + status, retry_count, last_error, last_error_code, + created_at, updated_at) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT(commit_sha) DO UPDATE SET + parent_sha = excluded.parent_sha, + sequence = excluded.sequence, + status = excluded.status, + retry_count = excluded.retry_count, + last_error = excluded.last_error, + last_error_code = excluded.last_error_code, + updated_at = excluded.updated_at + """, + (sha, parent, discovered, seq, status, retry, last_err, last_code, now, now), + ) + return self.get_commit(sha) + + def get_commit(self, commit_sha: str) -> Optional[Dict[str, Any]]: + row = self._conn.execute( + "SELECT * FROM commits WHERE commit_sha = ?", (commit_sha,) + ).fetchone() + return dict(row) if row else None + + def commit_status(self, commit_sha: str) -> Optional[str]: + row = self._conn.execute( + "SELECT status FROM commits WHERE commit_sha = ?", (commit_sha,) + ).fetchone() + return row["status"] if row else None + + def transition_commit(self, commit_sha: str, to_state: str) -> str: + """Deterministischer Sync-State-Übergang für einen Commit.""" + if to_state not in SYNC_STATES: + raise InvalidTransitionError(f"Unbekannter Sync-Zustand: {to_state}") + cur = self.commit_status(commit_sha) + if cur is None: + raise C5AError(f"Commit nicht gefunden: {commit_sha}") + if (cur, to_state) not in _ALLOWED_TRANSITIONS: + raise InvalidTransitionError( + f"Ungültiger Sync-Übergang für {commit_sha}: {cur} -> {to_state}" + ) + with self._conn: + self._conn.execute( + "UPDATE commits SET status = ?, updated_at = ? WHERE commit_sha = ?", + (to_state, int(time.time() * 1000), commit_sha), + ) + return to_state + + def set_commit_error(self, commit_sha: str, reason_code: str, message: str) -> None: + """Setzt Fehler + Reason Code auf einen Commit.""" + if reason_code not in REASON_CODES: + raise InvalidReasonCodeError(f"Unbekannter Reason Code: {reason_code}") + with self._conn: + self._conn.execute( + "UPDATE commits SET last_error = ?, last_error_code = ?, updated_at = ? WHERE commit_sha = ?", + (message, reason_code, int(time.time() * 1000), commit_sha), + ) + + def increment_retry(self, commit_sha: str) -> int: + """Erhöht retry_count. Gibt den neuen Wert zurück.""" + with self._conn: + self._conn.execute( + "UPDATE commits SET retry_count = retry_count + 1, updated_at = ? WHERE commit_sha = ?", + (int(time.time() * 1000), commit_sha), + ) + row = self._conn.execute( + "SELECT retry_count FROM commits WHERE commit_sha = ?", (commit_sha,) + ).fetchone() + return row["retry_count"] if row else 0 + + def list_commits(self, status: Optional[str] = None) -> List[Dict[str, Any]]: + if status: + rows = self._conn.execute( + "SELECT * FROM commits WHERE status = ? ORDER BY sequence ASC", (status,) + ).fetchall() + else: + rows = self._conn.execute( + "SELECT * FROM commits ORDER BY sequence ASC" + ).fetchall() + return [dict(r) for r in rows] + + # -- Objects ------------------------------------------------------------ + + def add_object_change(self, obj: Dict[str, Any]) -> Dict[str, Any]: + """Fügt eine Object-Change hinzu (idempotent per commit_sha+object_id+operation).""" + if obj.get("operation") not in OPERATIONS: + raise InvalidOperationError(f"Unbekannte Operation: {obj.get('operation')}") + with self._conn: + self._conn.execute( + """ + INSERT OR IGNORE INTO objects (commit_sha, object_id, path_before, path_after, + operation, content_hash_before, content_hash_after, + metadata_hash_before, metadata_hash_after, representation, state) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + """, + ( + obj["commit_sha"], obj["object_id"], obj.get("path_before"), + obj.get("path_after"), obj["operation"], + obj.get("content_hash_before"), obj.get("content_hash_after"), + obj.get("metadata_hash_before"), obj.get("metadata_hash_after"), + obj.get("representation"), obj.get("state"), + ), + ) + return self.get_object_change(obj["commit_sha"], obj["object_id"], obj["operation"]) + + def get_object_change(self, commit_sha: str, object_id: str, operation: str) -> Optional[Dict[str, Any]]: + row = self._conn.execute( + "SELECT * FROM objects WHERE commit_sha = ? AND object_id = ? AND operation = ?", + (commit_sha, object_id, operation), + ).fetchone() + return dict(row) if row else None + + def list_object_changes(self, commit_sha: str) -> List[Dict[str, Any]]: + rows = self._conn.execute( + "SELECT * FROM objects WHERE commit_sha = ? ORDER BY id ASC", (commit_sha,) + ).fetchall() + return [dict(r) for r in rows] + + # -- Health ------------------------------------------------------------- + + def health(self) -> Dict[str, Any]: + """Health Contract (finalisiert).""" + h = self._conn.execute("SELECT * FROM health WHERE id = 1").fetchone() + h = dict(h) if h else {} + commits = self._conn.execute("SELECT status, COUNT(*) as c FROM commits GROUP BY status").fetchall() + counts = {r["status"]: r["c"] for r in commits} + pending = sum(counts.get(s, 0) for s in ( + ST_DISCOVERED, ST_VALIDATING, ST_READY, ST_PROPAGATING_TOLARIA, + ST_VERIFYING_TOLARIA, ST_UPDATING_SEARCH, ST_VERIFYING_SEARCH, + ST_RETRY_PENDING, ST_WAITING_FOR_PREDECESSOR, + )) + failed = counts.get(ST_FAILED, 0) + dead = counts.get(ST_DEAD, 0) + human = counts.get(ST_HUMAN_REVIEW_REQUIRED, 0) + drift = self._conn.execute( + "SELECT COUNT(*) as c FROM objects WHERE operation = ? AND state = ?", + (OP_CONTENT_UPDATE, "drift"), + ).fetchone()["c"] + # Status ableiten + if dead > 0 or failed > 0: + status = HEALTH_ERROR + elif pending > 0 or human > 0: + status = HEALTH_DEGRADED + else: + status = HEALTH_OK + return { + "status": status, + "bootstrap_state": self.bootstrap_state(), + "last_seen_commit": h.get("last_seen_commit"), + "last_applied_commit": h.get("last_applied_commit"), + "pending_commits": pending, + "failed_commits": failed, + "dead_commits": dead, + "human_review_required": human, + "forgejo_status": "UNKNOWN", # C5A führt keine echten Probes aus + "tolaria_status": "UNKNOWN", + "search_status": "UNKNOWN", + "drift_count": drift, + "last_success_at": h.get("last_success_at"), + "last_error_code": h.get("last_error_code"), + } + + def mark_applied(self, commit_sha: str) -> None: + """Fortschreiben von last_applied_commit NUR nach vollständigem PASS.""" + with self._conn: + self._conn.execute( + "UPDATE health SET last_applied_commit = ?, last_success_at = ?, updated_at = ? WHERE id = 1", + (commit_sha, int(time.time() * 1000), int(time.time() * 1000)), + ) + + def mark_seen(self, commit_sha: str) -> None: + with self._conn: + self._conn.execute( + "UPDATE health SET last_seen_commit = ?, updated_at = ? WHERE id = 1", + (commit_sha, int(time.time() * 1000)), + ) + + def close(self) -> None: + self._conn.close() + + +# --------------------------------------------------------------------------- +# Sync State Machine (deterministische Logik) +# --------------------------------------------------------------------------- + +class SyncStateMachine: + """ + Deterministische Sync State Machine. + + Verarbeitet Commits in korrekter Reihenfolge (parent_sha / last_applied_commit), + mit Idempotenz, Retry/Backoff, Human-Gates und Reason Codes. Führt KEINE + externen Writes aus — nur Zustandsübergänge in der C5AStore. + """ + + def __init__(self, store: C5AStore, max_retries: int = DEFAULT_MAX_RETRIES, + backoff_seconds: Optional[List[int]] = None): + self.store = store + self.max_retries = max_retries + self.backoff_seconds = backoff_seconds or DEFAULT_BACKOFF_SECONDS + + # -- Idempotenz --------------------------------------------------------- + + def idempotency_check(self, commit_sha: str, object_id: str, operation: str, + content_hash: Optional[str] = None) -> str: + """ + Deterministische Idempotenz-Prüfung. + + Basis: commit_sha + object_id + operation. Zusätzlich Hash-Checks. + Rückgabe: ALREADY_APPLIED | ALREADY_AT_TARGET | RETRY_SAFE | CONFLICT + """ + existing = self.store.get_object_change(commit_sha, object_id, operation) + if existing is None: + return IDEM_RETRY_SAFE + # Commit bereits vollständig angewendet + if self.store.commit_status(commit_sha) == ST_APPLIED: + return IDEM_ALREADY_APPLIED + # Ziel-Hash bereits erreicht (gleicher Inhalt) + if content_hash is not None: + if existing.get("content_hash_after") == content_hash: + return IDEM_ALREADY_AT_TARGET + # Existiert, aber nicht angewendet und Hash weicht ab -> Konflikt + return IDEM_CONFLICT + + # -- Ordering ----------------------------------------------------------- + + def predecessor_applied(self, commit_sha: str) -> bool: + """Prüft, ob der Vorgänger (parent_sha) bereits angewendet ist.""" + commit = self.store.get_commit(commit_sha) + if commit is None: + return False + parent = commit.get("parent_sha") + if not parent: + return True # Root-Commit, kein Vorgänger + parent_status = self.store.commit_status(parent) + return parent_status == ST_APPLIED + + def check_ordering(self, commit_sha: str) -> str: + """ + Prüft die Commit-Reihenfolge. Gibt ST_READY oder ST_WAITING_FOR_PREDECESSOR zurück. + """ + if not self.predecessor_applied(commit_sha): + return ST_WAITING_FOR_PREDECESSOR + return ST_READY + + # -- Retry / Backoff ---------------------------------------------------- + + def backoff_for(self, retry_count: int) -> int: + """Backoff-Sekunden für den aktuellen Retry-Versuch.""" + if retry_count <= 0: + return 0 + idx = min(retry_count - 1, len(self.backoff_seconds) - 1) + return self.backoff_seconds[idx] + + def retry_available(self, commit_sha: str) -> bool: + commit = self.store.get_commit(commit_sha) + if commit is None: + return False + return commit.get("retry_count", 0) < self.max_retries + + # -- Human Gate --------------------------------------------------------- + + def human_gate(self, commit_sha: str, reason_code: str, message: str) -> Dict[str, Any]: + """ + Überführt einen Commit in HUMAN_REVIEW_REQUIRED (FAIL CLOSED). + Gibt ein Ergebnis-Dict zurück. + """ + if reason_code not in REASON_CODES: + raise InvalidReasonCodeError(f"Unbekannter Reason Code: {reason_code}") + self.store.set_commit_error(commit_sha, reason_code, message) + self.store.transition_commit(commit_sha, ST_HUMAN_REVIEW_REQUIRED) + return {"commit_sha": commit_sha, "status": ST_HUMAN_REVIEW_REQUIRED, + "reason_code": reason_code} + + # -- Hauptverarbeitung (deterministisch, NO-WRITE) ---------------------- + + def process_commit(self, commit: Dict[str, Any]) -> Dict[str, Any]: + """ + Verarbeitet einen Commit durch die State Machine. + + NO-WRITE: Diese Methode führt KEINE externen Writes aus. Sie modelliert + nur die Zustandsübergänge. Die eigentliche Tolaria-Propagation und der + Search-Rebuild werden in späteren Phasen (C5C/C5D) als injizierte + Callables ergänzt — C5A ruft sie NICHT auf. + + Rückgabe: Dict mit commit_sha, status, reason_code (optional), idempotency. + """ + sha = commit["commit_sha"] + # Idempotenz: bereits angewendet? + if self.store.commit_status(sha) == ST_APPLIED: + return {"commit_sha": sha, "status": ST_APPLIED, "idempotency": IDEM_ALREADY_APPLIED} + + # Commit anlegen/aktualisieren + self.store.upsert_commit(commit) + + # Ordering prüfen + order = self.check_ordering(sha) + if order == ST_WAITING_FOR_PREDECESSOR: + self.store.transition_commit(sha, ST_WAITING_FOR_PREDECESSOR) + return {"commit_sha": sha, "status": ST_WAITING_FOR_PREDECESSOR, + "reason_code": RC_OUT_OF_ORDER_COMMIT} + + # DISCOVERED -> VALIDATING + self.store.transition_commit(sha, ST_VALIDATING) + + # Validierung der Object-Changes + for obj in commit.get("changed_objects", []): + op = obj.get("operation") + if op not in OPERATIONS: + return self.human_gate(sha, RC_INVALID_SCHEMA, + f"Unbekannte Operation: {op}") + if not obj.get("object_id"): + return self.human_gate(sha, RC_UNKNOWN_OBJECT_ID, + "object_id fehlt") + # DELETE_REQUEST -> Human Gate (kein automatisches Hard Delete) + if op == OP_DELETE_REQUEST: + return self.human_gate(sha, RC_AMBIGUOUS_DELETE, + "DELETE_REQUEST erfordert Human Gate") + # DANGLING_DERIVED_FROM: wenn derived_from gesetzt, muss Ziel existieren + if obj.get("derived_from") and not self._object_exists(obj["derived_from"]): + return self.human_gate(sha, RC_DANGLING_DERIVED_FROM, + f"dangling derived_from: {obj['derived_from']}") + + # VALIDATING -> READY + self.store.transition_commit(sha, ST_READY) + + # READY -> PROPAGATING_TOLARIA (Modellierung; kein echter Write) + self.store.transition_commit(sha, ST_PROPAGATING_TOLARIA) + self.store.transition_commit(sha, ST_VERIFYING_TOLARIA) + + # DRIFT-Prüfung (modelliert): wenn ein Objekt als drift markiert ist -> Human Gate + for obj in commit.get("changed_objects", []): + if obj.get("drift"): + return self.human_gate(sha, RC_UNEXPECTED_TOLARIA_DRIFT, + f"unexpected drift auf {obj['object_id']}") + + # VERIFYING_TOLARIA -> UPDATING_SEARCH -> VERIFYING_SEARCH + self.store.transition_commit(sha, ST_UPDATING_SEARCH) + self.store.transition_commit(sha, ST_VERIFYING_SEARCH) + + # VERIFYING_SEARCH -> APPLIED + self.store.transition_commit(sha, ST_APPLIED) + self.store.mark_applied(sha) + + return {"commit_sha": sha, "status": ST_APPLIED, "idempotency": IDEM_RETRY_SAFE} + + def _object_exists(self, object_id: str) -> bool: + """Prüft, ob eine object_id bereits in irgendeinem Commit existiert.""" + row = self.store._conn.execute( + "SELECT 1 FROM objects WHERE object_id = ? LIMIT 1", (object_id,) + ).fetchone() + return row is not None + + # -- Retry-Verarbeitung ------------------------------------------------- + + def retry_commit(self, commit_sha: str, reason_code: str, message: str) -> Dict[str, Any]: + """ + Behandelt einen Fehler: Retry, FAILED oder DEAD je nach retry_count. + """ + if reason_code not in REASON_CODES: + raise InvalidReasonCodeError(f"Unbekannter Reason Code: {reason_code}") + self.store.set_commit_error(commit_sha, reason_code, message) + retry = self.store.increment_retry(commit_sha) + if retry >= self.max_retries: + self.store.transition_commit(commit_sha, ST_DEAD) + return {"commit_sha": commit_sha, "status": ST_DEAD, "reason_code": reason_code, + "retry_count": retry} + self.store.transition_commit(commit_sha, ST_RETRY_PENDING) + return {"commit_sha": commit_sha, "status": ST_RETRY_PENDING, "reason_code": reason_code, + "retry_count": retry, "backoff_seconds": self.backoff_for(retry)} + + +# --------------------------------------------------------------------------- +# No-Write-Guarantee (statische Prüfung) +# --------------------------------------------------------------------------- + +def assert_no_write_guarantee() -> Dict[str, Any]: + """ + Beweist, dass C5A keine externen Writes ausführen kann. + + Prüft, dass keine Netzwerk-/HTTP-/Socket-/Subprocess-Mutationsfunktionen + importiert oder verwendet werden. + """ + import ast + import sys + this_file = Path(__file__).resolve() + tree = ast.parse(this_file.read_text(encoding="utf-8")) + banned_imports = {"requests", "urllib", "http", "socket", "subprocess", "os.system"} + found = set() + for node in ast.walk(tree): + if isinstance(node, ast.Import): + for alias in node.names: + root = alias.name.split(".")[0] + if root in banned_imports: + found.add(root) + elif isinstance(node, ast.ImportFrom): + if node.module: + root = node.module.split(".")[0] + if root in banned_imports: + found.add(root) + return { + "no_network_imports": len(found) == 0, + "banned_imports_found": sorted(found), + "no_write_guarantee": len(found) == 0, + } diff --git a/tolaria/c5-sync-service/rq_c5a_cli.py b/tolaria/c5-sync-service/rq_c5a_cli.py new file mode 100644 index 0000000..ea845ae --- /dev/null +++ b/tolaria/c5-sync-service/rq_c5a_cli.py @@ -0,0 +1,162 @@ +#!/usr/bin/env python3 +""" +Red Queen — C5A: CONTRACT & STATE MACHINE CLI. + +Deterministische, inaktive CLI zur Inspektion und Steuerung der C5A State Machine. +Führt KEINE externen Writes aus (kein Polling, kein Tolaria-Write, kein Search-Rebuild). + +Nutzung: + python3 rq_c5a_cli.py --db health + python3 rq_c5a_cli.py --db bootstrap-state + python3 rq_c5a_cli.py --db bootstrap-transition --to RECONCILING + python3 rq_c5a_cli.py --db set-baseline --commit + python3 rq_c5a_cli.py --db ingest-commit --json + python3 rq_c5a_cli.py --db process-commit --json + python3 rq_c5a_cli.py --db list-commits [--status ] + python3 rq_c5a_cli.py --db commit-status --sha + python3 rq_c5a_cli.py --db no-write-check +""" + +from __future__ import annotations + +import argparse +import json +import sys +from pathlib import Path + +from rq_c5a import ( + C5AStore, + SyncStateMachine, + assert_no_write_guarantee, + C5AError, + InvalidTransitionError, + InvalidReasonCodeError, + InvalidOperationError, +) + + +def _load_json(path: str) -> dict: + with open(path, "r", encoding="utf-8") as f: + return json.load(f) + + +def cmd_health(store: C5AStore, args) -> int: + print(json.dumps(store.health(), indent=2, ensure_ascii=False)) + return 0 + + +def cmd_bootstrap_state(store: C5AStore, args) -> int: + print(store.bootstrap_state()) + return 0 + + +def cmd_bootstrap_transition(store: C5AStore, args) -> int: + try: + new_state = store.bootstrap_transition(args.to) + except InvalidTransitionError as e: + print(json.dumps(e.to_dict(), ensure_ascii=False)) + return 2 + print(new_state) + return 0 + + +def cmd_set_baseline(store: C5AStore, args) -> int: + try: + store.set_baseline(args.commit) + except C5AError as e: + print(json.dumps(e.to_dict(), ensure_ascii=False)) + return 2 + print(f"baseline={args.commit}") + return 0 + + +def cmd_ingest_commit(store: C5AStore, args) -> int: + commit = _load_json(args.json) + store.upsert_commit(commit) + for obj in commit.get("changed_objects", []): + obj["commit_sha"] = commit["commit_sha"] + store.add_object_change(obj) + print(json.dumps({"commit_sha": commit["commit_sha"], "status": "ingested"}, ensure_ascii=False)) + return 0 + + +def cmd_process_commit(store: C5AStore, args) -> int: + commit = _load_json(args.json) + sm = SyncStateMachine(store) + try: + result = sm.process_commit(commit) + except C5AError as e: + print(json.dumps(e.to_dict(), ensure_ascii=False)) + return 2 + print(json.dumps(result, ensure_ascii=False)) + return 0 + + +def cmd_list_commits(store: C5AStore, args) -> int: + commits = store.list_commits(status=args.status) + print(json.dumps(commits, indent=2, ensure_ascii=False)) + return 0 + + +def cmd_commit_status(store: C5AStore, args) -> int: + status = store.commit_status(args.sha) + if status is None: + print(json.dumps({"error": f"Commit nicht gefunden: {args.sha}"}, ensure_ascii=False)) + return 2 + print(status) + return 0 + + +def cmd_no_write_check(store: C5AStore, args) -> int: + result = assert_no_write_guarantee() + print(json.dumps(result, indent=2, ensure_ascii=False)) + return 0 if result["no_write_guarantee"] else 2 + + +def main(argv=None) -> int: + parser = argparse.ArgumentParser(description="C5A Contract & State Machine CLI") + parser.add_argument("--db", default=os_env_db(), help="Pfad zur c5a.db (default: $C5A_DB oder ./c5a.db)") + sub = parser.add_subparsers(dest="cmd", required=True) + + sub.add_parser("health", help="Health Contract anzeigen") + sub.add_parser("bootstrap-state", help="Bootstrap-Zustand anzeigen") + p = sub.add_parser("bootstrap-transition", help="Bootstrap-Übergang") + p.add_argument("--to", required=True) + p = sub.add_parser("set-baseline", help="Baseline setzen") + p.add_argument("--commit", required=True) + p = sub.add_parser("ingest-commit", help="Commit + Object-Changes aufnehmen") + p.add_argument("--json", required=True) + p = sub.add_parser("process-commit", help="Commit durch State Machine verarbeiten") + p.add_argument("--json", required=True) + p = sub.add_parser("list-commits", help="Commits auflisten") + p.add_argument("--status") + p = sub.add_parser("commit-status", help="Status eines Commits") + p.add_argument("--sha", required=True) + sub.add_parser("no-write-check", help="No-Write-Guarantee prüfen") + + args = parser.parse_args(argv) + store = C5AStore(args.db) + try: + handlers = { + "health": cmd_health, + "bootstrap-state": cmd_bootstrap_state, + "bootstrap-transition": cmd_bootstrap_transition, + "set-baseline": cmd_set_baseline, + "ingest-commit": cmd_ingest_commit, + "process-commit": cmd_process_commit, + "list-commits": cmd_list_commits, + "commit-status": cmd_commit_status, + "no-write-check": cmd_no_write_check, + } + return handlers[args.cmd](store, args) + finally: + store.close() + + +def os_env_db() -> str: + import os + return os.environ.get("C5A_DB", "c5a.db") + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/tolaria/c5-sync-service/test_c5a.py b/tolaria/c5-sync-service/test_c5a.py new file mode 100644 index 0000000..8091436 --- /dev/null +++ b/tolaria/c5-sync-service/test_c5a.py @@ -0,0 +1,475 @@ +#!/usr/bin/env python3 +""" +Red Queen — C5A: CONTRACT & STATE MACHINE Testsuite. + +Deterministische Unit Tests. Jeder Test nutzt eine frische temp-DB (tempfile.mkdtemp), +niemals die Produkt-DB. Beweist die No-Write-Guarantee und alle geforderten +State-Machine-, Idempotenz-, Ordering-, Retry-, Human-Gate-, Bootstrap- und +Persistenz-Verhalten. +""" + +from __future__ import annotations + +import json +import os +import sys +import tempfile +import traceback +from pathlib import Path + +sys.path.insert(0, str(Path(__file__).resolve().parent)) + +from rq_c5a import ( + C5AStore, + SyncStateMachine, + assert_no_write_guarantee, + C5AError, + InvalidTransitionError, + InvalidReasonCodeError, + InvalidOperationError, + ST_DISCOVERED, ST_VALIDATING, ST_READY, + ST_PROPAGATING_TOLARIA, ST_VERIFYING_TOLARIA, + ST_UPDATING_SEARCH, ST_VERIFYING_SEARCH, ST_APPLIED, + ST_RETRY_PENDING, ST_FAILED, ST_DEAD, ST_HUMAN_REVIEW_REQUIRED, + ST_WAITING_FOR_PREDECESSOR, + BS_UNINITIALIZED, BS_RECONCILING, BS_BASELINE_READY, BS_ACTIVE, + OP_CREATE, OP_CONTENT_UPDATE, OP_DELETE_REQUEST, + IDEM_ALREADY_APPLIED, IDEM_ALREADY_AT_TARGET, IDEM_RETRY_SAFE, IDEM_CONFLICT, + RC_UNEXPECTED_TOLARIA_DRIFT, RC_ID_COLLISION, RC_UNKNOWN_OBJECT_ID, + RC_INVALID_SCHEMA, RC_DANGLING_DERIVED_FROM, RC_AMBIGUOUS_DELETE, + RC_OUT_OF_ORDER_COMMIT, RC_SECRET_DETECTED, + RC_UNKNOWN_LEGACY_OBJECT, RC_AUTH_FAILURE, RC_SEARCH_REBUILD_FAILURE, + RC_TOLARIA_UNAVAILABLE, RC_FORGEJO_UNAVAILABLE, + REASON_CODES, OPERATIONS, SYNC_STATES, BOOTSTRAP_STATES, +) + + +def _new_store(): + tmp = tempfile.mkdtemp(prefix="c5a_test_") + return C5AStore(os.path.join(tmp, "c5a.db")) + + +def _commit(sha, parent=None, seq=1, objects=None): + return { + "commit_sha": sha, + "parent_sha": parent, + "discovered_at": 1000, + "sequence": seq, + "changed_objects": objects or [], + } + + +def _obj(object_id, op=OP_CREATE, **kw): + d = {"object_id": object_id, "operation": op} + d.update(kw) + return d + + +# --------------------------------------------------------------------------- +# Testfälle +# --------------------------------------------------------------------------- + +def test_normal_commit_lifecycle(): + """Normaler Commit-Lebenszyklus: DISCOVERED -> ... -> APPLIED.""" + store = _new_store() + sm = SyncStateMachine(store) + c = _commit("c1", objects=[_obj("obj-1")]) + store.upsert_commit(c) + store.add_object_change({"commit_sha": "c1", "object_id": "obj-1", "operation": OP_CREATE}) + result = sm.process_commit(c) + assert result["status"] == ST_APPLIED, result + assert store.commit_status("c1") == ST_APPLIED + assert store.health()["last_applied_commit"] == "c1" + store.close() + + +def test_multi_object_commit(): + """Multi-Object-Commit: mehrere Object-Changes in einem Commit.""" + store = _new_store() + sm = SyncStateMachine(store) + c = _commit("c1", objects=[ + _obj("obj-1"), _obj("obj-2"), _obj("obj-3"), + ]) + for o in c["changed_objects"]: + store.add_object_change({"commit_sha": "c1", **o}) + result = sm.process_commit(c) + assert result["status"] == ST_APPLIED + assert len(store.list_object_changes("c1")) == 3 + store.close() + + +def test_duplicate_commit(): + """Doppelter Commit: zweite Verarbeitung -> ALREADY_APPLIED, kein Re-Apply.""" + store = _new_store() + sm = SyncStateMachine(store) + c = _commit("c1", objects=[_obj("obj-1")]) + store.add_object_change({"commit_sha": "c1", "object_id": "obj-1", "operation": OP_CREATE}) + r1 = sm.process_commit(c) + assert r1["status"] == ST_APPLIED + r2 = sm.process_commit(c) + assert r2["status"] == ST_APPLIED + assert r2["idempotency"] == IDEM_ALREADY_APPLIED + store.close() + + +def test_already_applied_idempotency(): + """Idempotenz: bereits angewendeter Commit -> ALREADY_APPLIED.""" + store = _new_store() + sm = SyncStateMachine(store) + c = _commit("c1", objects=[_obj("obj-1")]) + store.add_object_change({"commit_sha": "c1", "object_id": "obj-1", "operation": OP_CREATE}) + sm.process_commit(c) + assert sm.idempotency_check("c1", "obj-1", OP_CREATE) == IDEM_ALREADY_APPLIED + store.close() + + +def test_out_of_order_commit(): + """Out-of-Order-Commit: fehlender Vorgänger -> WAITING_FOR_PREDECESSOR.""" + store = _new_store() + sm = SyncStateMachine(store) + c = _commit("c2", parent="c1", seq=2, objects=[_obj("obj-1")]) + store.add_object_change({"commit_sha": "c2", "object_id": "obj-1", "operation": OP_CREATE}) + result = sm.process_commit(c) + assert result["status"] == ST_WAITING_FOR_PREDECESSOR, result + assert result["reason_code"] == RC_OUT_OF_ORDER_COMMIT + store.close() + + +def test_missing_predecessor_then_applied(): + """Fehlender Vorgänger -> nach Anwendung des Vorgängers verarbeitbar.""" + store = _new_store() + sm = SyncStateMachine(store) + c1 = _commit("c1", seq=1, objects=[_obj("obj-1")]) + store.add_object_change({"commit_sha": "c1", "object_id": "obj-1", "operation": OP_CREATE}) + sm.process_commit(c1) + c2 = _commit("c2", parent="c1", seq=2, objects=[_obj("obj-2")]) + store.add_object_change({"commit_sha": "c2", "object_id": "obj-2", "operation": OP_CREATE}) + r2 = sm.process_commit(c2) + assert r2["status"] == ST_APPLIED, r2 + store.close() + + +def test_retry_progression(): + """Retry-Progression: Fehler -> RETRY_PENDING mit Backoff.""" + store = _new_store() + sm = SyncStateMachine(store) + c = _commit("c1", objects=[_obj("obj-1")]) + store.upsert_commit(c) + store.add_object_change({"commit_sha": "c1", "object_id": "obj-1", "operation": OP_CREATE}) + r = sm.retry_commit("c1", RC_TOLARIA_UNAVAILABLE, "tolaria down") + assert r["status"] == ST_RETRY_PENDING + assert r["retry_count"] == 1 + assert r["backoff_seconds"] == 1 + store.close() + + +def test_max_retry_to_dead(): + """Max-Retry -> DEAD, kein Endlos-Retry.""" + store = _new_store() + sm = SyncStateMachine(store, max_retries=3) + c = _commit("c1", objects=[_obj("obj-1")]) + store.upsert_commit(c) + store.add_object_change({"commit_sha": "c1", "object_id": "obj-1", "operation": OP_CREATE}) + for i in range(3): + r = sm.retry_commit("c1", RC_TOLARIA_UNAVAILABLE, "down") + assert r["status"] == ST_DEAD, r + assert store.commit_status("c1") == ST_DEAD + store.close() + + +def test_human_gate_transition(): + """Human-Gate-Übergang: DELETE_REQUEST -> HUMAN_REVIEW_REQUIRED.""" + store = _new_store() + sm = SyncStateMachine(store) + c = _commit("c1", objects=[_obj("obj-1", op=OP_DELETE_REQUEST)]) + store.add_object_change({"commit_sha": "c1", "object_id": "obj-1", "operation": OP_DELETE_REQUEST}) + result = sm.process_commit(c) + assert result["status"] == ST_HUMAN_REVIEW_REQUIRED, result + assert result["reason_code"] == RC_AMBIGUOUS_DELETE + store.close() + + +def test_restart_reload_persistence(): + """Restart-Persistenz: neuer Store auf gleicher DB behält State.""" + tmp = tempfile.mkdtemp(prefix="c5a_test_") + db = os.path.join(tmp, "c5a.db") + s1 = C5AStore(db) + sm1 = SyncStateMachine(s1) + c = _commit("c1", objects=[_obj("obj-1")]) + s1.add_object_change({"commit_sha": "c1", "object_id": "obj-1", "operation": OP_CREATE}) + sm1.process_commit(c) + s1.close() + s2 = C5AStore(db) + assert s2.commit_status("c1") == ST_APPLIED + assert s2.health()["last_applied_commit"] == "c1" + s2.close() + + +def test_state_corruption_handling(): + """State-Corruption: ungültiger Übergang -> InvalidTransitionError, kein Mutation.""" + store = _new_store() + sm = SyncStateMachine(store) + c = _commit("c1", objects=[_obj("obj-1")]) + store.upsert_commit(c) + store.add_object_change({"commit_sha": "c1", "object_id": "obj-1", "operation": OP_CREATE}) + # Direkt von DISCOVERED nach APPLIED ist nicht erlaubt + try: + store.transition_commit("c1", ST_APPLIED) + assert False, "sollte InvalidTransitionError werfen" + except InvalidTransitionError: + pass + assert store.commit_status("c1") == ST_DISCOVERED + store.close() + + +def test_idempotency_conflict(): + """Idempotenz-Konflikt: existierender Change mit abweichendem Hash -> CONFLICT.""" + store = _new_store() + sm = SyncStateMachine(store) + c = _commit("c1", objects=[_obj("obj-1")]) + store.upsert_commit(c) + store.add_object_change({"commit_sha": "c1", "object_id": "obj-1", + "operation": OP_CREATE, "content_hash_after": "hashA"}) + # Nicht angewendet, Hash weicht ab -> CONFLICT + assert sm.idempotency_check("c1", "obj-1", OP_CREATE, content_hash="hashB") == IDEM_CONFLICT + # Gleicher Hash -> ALREADY_AT_TARGET + assert sm.idempotency_check("c1", "obj-1", OP_CREATE, content_hash="hashA") == IDEM_ALREADY_AT_TARGET + store.close() + + +def test_delete_request_human_gate(): + """DELETE_REQUEST -> Human Gate, kein automatisches Hard Delete.""" + store = _new_store() + sm = SyncStateMachine(store) + c = _commit("c1", objects=[_obj("obj-1", op=OP_DELETE_REQUEST)]) + store.add_object_change({"commit_sha": "c1", "object_id": "obj-1", "operation": OP_DELETE_REQUEST}) + result = sm.process_commit(c) + assert result["status"] == ST_HUMAN_REVIEW_REQUIRED + assert result["reason_code"] == RC_AMBIGUOUS_DELETE + store.close() + + +def test_dangling_relation_human_gate(): + """Dangling derived_from -> Human Gate.""" + store = _new_store() + sm = SyncStateMachine(store) + c = _commit("c1", objects=[_obj("obj-1", derived_from="nonexistent-obj")]) + store.add_object_change({"commit_sha": "c1", "object_id": "obj-1", + "operation": OP_CREATE, "derived_from": "nonexistent-obj"}) + result = sm.process_commit(c) + assert result["status"] == ST_HUMAN_REVIEW_REQUIRED + assert result["reason_code"] == RC_DANGLING_DERIVED_FROM + store.close() + + +def test_unexpected_drift_human_gate(): + """Unexpected Drift -> Human Gate.""" + store = _new_store() + sm = SyncStateMachine(store) + c = _commit("c1", objects=[_obj("obj-1", drift=True)]) + store.add_object_change({"commit_sha": "c1", "object_id": "obj-1", + "operation": OP_CREATE, "drift": True}) + result = sm.process_commit(c) + assert result["status"] == ST_HUMAN_REVIEW_REQUIRED + assert result["reason_code"] == RC_UNEXPECTED_TOLARIA_DRIFT + store.close() + + +def test_no_write_guarantee(): + """No-Write-Guarantee: keine Netzwerk-/HTTP-/Socket-/Subprocess-Imports.""" + result = assert_no_write_guarantee() + assert result["no_write_guarantee"] is True, result + assert result["banned_imports_found"] == [] + store.close() if False else None + + +def test_bootstrap_lifecycle(): + """Bootstrap: UNINITIALIZED -> RECONCILING -> BASELINE_READY -> ACTIVE.""" + store = _new_store() + assert store.bootstrap_state() == BS_UNINITIALIZED + store.bootstrap_transition(BS_RECONCILING) + assert store.bootstrap_state() == BS_RECONCILING + store.bootstrap_transition(BS_BASELINE_READY) + assert store.bootstrap_state() == BS_BASELINE_READY + store.set_baseline("base-1") + assert store.baseline_commit() == "base-1" + store.bootstrap_transition(BS_ACTIVE) + assert store.bootstrap_state() == BS_ACTIVE + store.close() + + +def test_baseline_only_from_baseline_ready(): + """Baseline darf nur aus BASELINE_READY gesetzt werden.""" + store = _new_store() + try: + store.set_baseline("base-1") + assert False, "sollte C5AError werfen" + except C5AError: + pass + store.close() + + +def test_reason_codes_closed_set(): + """Reason Codes: feste, geschlossene Menge — kein freier String-Wildwuchs.""" + assert RC_UNEXPECTED_TOLARIA_DRIFT in REASON_CODES + assert RC_ID_COLLISION in REASON_CODES + assert RC_UNKNOWN_OBJECT_ID in REASON_CODES + assert RC_INVALID_SCHEMA in REASON_CODES + assert RC_DANGLING_DERIVED_FROM in REASON_CODES + assert RC_AMBIGUOUS_DELETE in REASON_CODES + assert RC_UNKNOWN_LEGACY_OBJECT in REASON_CODES + assert RC_AUTH_FAILURE in REASON_CODES + assert RC_SEARCH_REBUILD_FAILURE in REASON_CODES + assert RC_TOLARIA_UNAVAILABLE in REASON_CODES + assert RC_FORGEJO_UNAVAILABLE in REASON_CODES + assert RC_OUT_OF_ORDER_COMMIT in REASON_CODES + assert RC_SECRET_DETECTED in REASON_CODES + assert len(REASON_CODES) == 13 + store = _new_store() + try: + store.set_commit_error("c1", "FREIER_STRING", "x") + assert False, "sollte InvalidReasonCodeError werfen" + except InvalidReasonCodeError: + pass + store.close() + + +def test_operations_closed_set(): + """Operation Model: geschlossene Menge.""" + assert OP_CREATE in OPERATIONS + assert OP_CONTENT_UPDATE in OPERATIONS + assert OP_DELETE_REQUEST in OPERATIONS + assert len(OPERATIONS) == 10 + store = _new_store() + try: + store.add_object_change({"commit_sha": "c1", "object_id": "o", + "operation": "BOGUS"}) + assert False, "sollte InvalidOperationError werfen" + except InvalidOperationError: + pass + store.close() + + +def test_health_contract_fields(): + """Health Contract: alle geforderten Felder vorhanden.""" + store = _new_store() + h = store.health() + for field in ["status", "bootstrap_state", "last_seen_commit", "last_applied_commit", + "pending_commits", "failed_commits", "dead_commits", + "human_review_required", "forgejo_status", "tolaria_status", + "search_status", "drift_count", "last_success_at", "last_error_code"]: + assert field in h, f"fehlendes Health-Feld: {field}" + store.close() + + +def test_last_applied_only_after_full_pass(): + """last_applied_commit wird NUR nach vollständigem Tolaria+Search PASS fortgeschrieben.""" + store = _new_store() + sm = SyncStateMachine(store) + # Commit mit DELETE_REQUEST -> Human Gate, darf NICHT applied werden + c = _commit("c1", objects=[_obj("obj-1", op=OP_DELETE_REQUEST)]) + store.add_object_change({"commit_sha": "c1", "object_id": "obj-1", "operation": OP_DELETE_REQUEST}) + sm.process_commit(c) + assert store.health()["last_applied_commit"] is None + store.close() + + +def test_retry_available_bound(): + """retry_available: false nach Erreichen von max_retries.""" + store = _new_store() + sm = SyncStateMachine(store, max_retries=2) + c = _commit("c1", objects=[_obj("obj-1")]) + store.upsert_commit(c) + store.add_object_change({"commit_sha": "c1", "object_id": "obj-1", "operation": OP_CREATE}) + assert sm.retry_available("c1") is True + sm.retry_commit("c1", RC_TOLARIA_UNAVAILABLE, "down") + assert sm.retry_available("c1") is True + sm.retry_commit("c1", RC_TOLARIA_UNAVAILABLE, "down") + assert sm.retry_available("c1") is False + store.close() + + +def test_backoff_sequence(): + """Backoff: 1/2/4/8/16 Sekunden.""" + store = _new_store() + sm = SyncStateMachine(store) + assert sm.backoff_for(1) == 1 + assert sm.backoff_for(2) == 2 + assert sm.backoff_for(3) == 4 + assert sm.backoff_for(4) == 8 + assert sm.backoff_for(5) == 16 + assert sm.backoff_for(6) == 16 # gekappt + store.close() + + +def test_invalid_reason_code_rejected(): + """Unbekannter Reason Code wird abgelehnt.""" + store = _new_store() + sm = SyncStateMachine(store) + c = _commit("c1", objects=[_obj("obj-1")]) + store.upsert_commit(c) + store.add_object_change({"commit_sha": "c1", "object_id": "obj-1", "operation": OP_CREATE}) + try: + sm.retry_commit("c1", "NOT_A_REASON_CODE", "x") + assert False, "sollte InvalidReasonCodeError werfen" + except InvalidReasonCodeError: + pass + store.close() + + +# --------------------------------------------------------------------------- +# Runner +# --------------------------------------------------------------------------- + +ALL_TESTS = [ + test_normal_commit_lifecycle, + test_multi_object_commit, + test_duplicate_commit, + test_already_applied_idempotency, + test_out_of_order_commit, + test_missing_predecessor_then_applied, + test_retry_progression, + test_max_retry_to_dead, + test_human_gate_transition, + test_restart_reload_persistence, + test_state_corruption_handling, + test_idempotency_conflict, + test_delete_request_human_gate, + test_dangling_relation_human_gate, + test_unexpected_drift_human_gate, + test_no_write_guarantee, + test_bootstrap_lifecycle, + test_baseline_only_from_baseline_ready, + test_reason_codes_closed_set, + test_operations_closed_set, + test_health_contract_fields, + test_last_applied_only_after_full_pass, + test_retry_available_bound, + test_backoff_sequence, + test_invalid_reason_code_rejected, +] + + +def main() -> int: + passed = 0 + failed = 0 + failures = [] + for t in ALL_TESTS: + try: + t() + passed += 1 + print(f"PASS {t.__name__}") + except Exception as e: + failed += 1 + failures.append((t.__name__, e)) + print(f"FAIL {t.__name__}: {e}") + traceback.print_exc() + print(f"\n=== C5A: {passed} PASS / {failed} FAIL ===") + if failures: + for name, e in failures: + print(f" FAILED: {name} -> {e}") + return 1 if failed else 0 + + +if __name__ == "__main__": + sys.exit(main())