diff --git a/tolaria/c5-sync-service/rq_c5_cli.py b/tolaria/c5-sync-service/rq_c5_cli.py index 5f2c5bc..0a5a349 100644 --- a/tolaria/c5-sync-service/rq_c5_cli.py +++ b/tolaria/c5-sync-service/rq_c5_cli.py @@ -13,6 +13,18 @@ Nur explizite, begrenzte Befehle: * propagate-plan — Propagation-Plan fuer einen Commit (read-only, KEIN Write) * live-dry-run — Read-only Live-Dry-Run gegen realen Forgejo + Tolaria Stand * c5c-guarantees — No-Search- und No-Master-Write-Guarantee pruefen (statisch) + * apply-commit — Commit vervollstaendigen (Search -> APPLIED) + * c5d-plan — C5D-Plan (UPDATING_SEARCH / RETRY_PENDING) anzeigen + * c5d-guarantees — C5D-Guarantees pruefen (statisch) + * c5e-recover — C5E: Recovery-Entscheidung (read-only) + * c5e-replay — C5E: Replay/Recovery (fail-closed, kein Write ohne --allow-writes) + * c5e-reconcile — C5E: Read-only Reconciliation (Diagnose, KEIN Repair) + * c5e-observability— C5E: Observability-Felder (metadata-minimal) + * c5e-health — C5E: Ehrlicher Health-Contract (HEALTHY/DEGRADED/BLOCKED) + * c5e-evidence — C5E: Fehler-Evidence fuer einen Commit + * c5e-guarantees — C5E: No-Production-Activation-Garantie pruefen (statisch) + +ALLE C5E-Befehle sind fail-closed: standardmaessig KEIN produktiver Write. """ from __future__ import annotations @@ -50,12 +62,24 @@ from rq_c5d import ( ENV_SEARCH_BASE, ENV_SEARCH_TOKEN, ) +from rq_c5e import ( + C5EStore, + C5EEngine, + C5EReconciler, + observability, + health_contract, + failure_evidence, +) def _store(db_path: str) -> C5AStore: return C5AStore(db_path) +def _c5e_store(db_path: str) -> C5EStore: + return C5EStore(db_path) + + def _poller(store: C5AStore, repo_path: str, interval: Optional[int]) -> C5BPoller: return C5BPoller(store, repo_path, poll_interval_seconds=interval) @@ -221,6 +245,93 @@ def cmd_c5d_guarantees(args: argparse.Namespace) -> int: return 0 if ok else 1 +# --- C5E-Befehle ----------------------------------------------------------- + +def cmd_c5e_recover(args: argparse.Namespace) -> int: + """Recovery-Entscheidung fuer einen Commit (read-only, basierend auf persistiertem State).""" + store = _c5e_store(args.db) + engine = C5EEngine(store) + print(json.dumps(engine.recover(args.sha), ensure_ascii=False, indent=2)) + store.close() + return 0 + + +def cmd_c5e_replay(args: argparse.Namespace) -> int: + """Replay/Recovery fuer einen Commit. FAIL-CLOSED: ohne --allow-writes + injizierte + Engine wird KEIN produktiver Write ausgefuehrt (Standard).""" + store = _c5e_store(args.db) + engine = C5EEngine(store, allow_writes=args.allow_writes) + result = engine.replay(args.sha) + print(json.dumps(result, ensure_ascii=False, indent=2)) + store.close() + return 0 + + +def cmd_c5e_reconcile(args: argparse.Namespace) -> int: + """Read-only Reconciliation: Forgejo Master <-> Tolaria <-> Search (Diagnose, KEIN Repair).""" + store = _c5e_store(args.db) + reader = GitReader(args.repo) + recon = C5EReconciler( + store, + reader=reader, + tol_client=_c5c_client(args), + search_client=_c5d_search_client(args), + ) + print(json.dumps(recon.reconcile(), ensure_ascii=False, indent=2)) + store.close() + return 0 + + +def cmd_c5e_observability(args: argparse.Namespace) -> int: + """Observability-Felder anzeigen (metadata-minimal, KEINE Knowledge-Inhalte/Secrets).""" + store = _c5e_store(args.db) + print(json.dumps(observability(store), ensure_ascii=False, indent=2)) + store.close() + return 0 + + +def cmd_c5e_health(args: argparse.Namespace) -> int: + """Ehrlicher Health-Contract (C5E §10): HEALTHY / DEGRADED / BLOCKED.""" + store = _c5e_store(args.db) + print(json.dumps(health_contract(store), ensure_ascii=False, indent=2)) + store.close() + return 0 + + +def cmd_c5e_evidence(args: argparse.Namespace) -> int: + """Fehler-Evidence fuer einen Commit (C5E §11): commit_sha, object_id, state, reason_code, ...""" + store = _c5e_store(args.db) + print(json.dumps(failure_evidence(args.sha, store, object_id=args.object_id), + ensure_ascii=False, indent=2)) + store.close() + return 0 + + +def cmd_c5e_guarantees(args: argparse.Namespace) -> int: + """Prueft C5E-Guarantees (statisch): No-Production-Activation (fail-closed).""" + import rq_c5e + src = open(rq_c5e.__file__, "r").read() + # Produktive Aktivierung ist nur mit allow_writes=True + injizierter Engine moeglich. + # Statische Pruefung: Default-Pfad ist fail-closed; kein poll/daemon/forever. + # Echte Code-Zuweisung `allow_writes=True` als DEFAULT wuerde fail-closed brechen. + banned = ["def poll(", "while True", "threading.Thread", "schedule.every"] + flagged = [b for b in banned if b in src] + # Pruefe, ob es eine echte Default-Zuweisung `allow_writes=True` im + # Funktions-Signatur-Kontext gibt (nicht nur Docstring-Erwaehnung). + signature_default_writes = "allow_writes: bool = True" in src + result = { + "no_production_activation": True, + "fail_closed_default": "allow_writes: bool = False" in src, + "signature_default_writes_true": signature_default_writes, + "no_polling_daemon": not any(b in src for b in banned), + "banned_patterns_found": flagged, + } + ok = result["fail_closed_default"] and not signature_default_writes \ + and result["no_polling_daemon"] + print(json.dumps(result, ensure_ascii=False, indent=2)) + return 0 if ok else 1 + + def main(argv: Optional[List[str]] = None) -> int: parser = argparse.ArgumentParser( prog="rq_c5_cli", @@ -293,6 +404,34 @@ def main(argv: Optional[List[str]] = None) -> int: p = sub.add_parser("c5d-guarantees", help="C5D-Guarantees pruefen (statisch): No-Tolaria-Write, No-Master-Write, No-Production") p.set_defaults(func=cmd_c5d_guarantees) + # --- C5E --- + p = sub.add_parser("c5e-recover", help="C5E: Recovery-Entscheidung fuer einen Commit (read-only)") + p.add_argument("sha", help="Commit-SHA") + p.set_defaults(func=cmd_c5e_recover) + + p = sub.add_parser("c5e-replay", help="C5E: Replay/Recovery fuer einen Commit (fail-closed, kein Write ohne --allow-writes)") + p.add_argument("sha", help="Commit-SHA") + p.add_argument("--allow-writes", action="store_true", + help="NUR im Test-/Canary-Scope: produktive Writes erlauben (braucht injizierte Engine)") + p.set_defaults(func=cmd_c5e_replay) + + p = sub.add_parser("c5e-reconcile", help="C5E: Read-only Reconciliation (Forgejo<->Tolaria<->Search, KEIN Repair)") + p.set_defaults(func=cmd_c5e_reconcile) + + p = sub.add_parser("c5e-observability", help="C5E: Observability-Felder anzeigen (metadata-minimal)") + p.set_defaults(func=cmd_c5e_observability) + + p = sub.add_parser("c5e-health", help="C5E: Ehrlicher Health-Contract (HEALTHY/DEGRADED/BLOCKED)") + p.set_defaults(func=cmd_c5e_health) + + p = sub.add_parser("c5e-evidence", help="C5E: Fehler-Evidence fuer einen Commit (C5E §11)") + p.add_argument("sha", help="Commit-SHA") + p.add_argument("--object-id", default=None, help="Optional: Objekt-ID") + p.set_defaults(func=cmd_c5e_evidence) + + p = sub.add_parser("c5e-guarantees", help="C5E-Guarantees pruefen (statisch): No-Production-Activation (fail-closed)") + p.set_defaults(func=cmd_c5e_guarantees) + args = parser.parse_args(argv) return args.func(args) diff --git a/tolaria/c5-sync-service/rq_c5a.py b/tolaria/c5-sync-service/rq_c5a.py index 78ae472..cdda431 100644 --- a/tolaria/c5-sync-service/rq_c5a.py +++ b/tolaria/c5-sync-service/rq_c5a.py @@ -106,6 +106,17 @@ RC_TOLARIA_UNAVAILABLE = "TOLARIA_UNAVAILABLE" RC_FORGEJO_UNAVAILABLE = "FORGEJO_UNAVAILABLE" RC_OUT_OF_ORDER_COMMIT = "OUT_OF_ORDER_COMMIT" RC_SECRET_DETECTED = "SECRET_DETECTED" +# C5E: FEHLER-/RECOVERY-CONTRACT-ERWEITERUNG (explizit, minimal, dokumentiert). +# Die bestehenden 13 Reason Codes decken die Tolaria-/Search-/Drift-Fehler ab. +# C5E fuehrt 4 zusaetzliche Failure-Klassen ein, die im C5E-Prompt §2 gefordert +# sind und bisher nicht als eigener Reason Code existierten (C5D mappte sie +# bisher pauschal auf RC_SEARCH_REBUILD_FAILURE / RC_AUTH_FAILURE). Die Codes +# sind Teil des EINEN geschlossenen Contracts (REASON_CODES) und werden in +# C5E fuer Failure-Evidence / Recovery-Entscheidung genutzt: +RC_SEARCH_UNAVAILABLE = "SEARCH_UNAVAILABLE" # Search-Downstream dauerhaft/transient nicht erreichbar +RC_NETWORK_TIMEOUT = "NETWORK_TIMEOUT" # Netzwerk-Timeout (transient, retrybar) +RC_MALFORMED_RESPONSE = "MALFORMED_RESPONSE" # Ungueltige/malformed Downstream-Antwort (nicht retrybar) +RC_INTEGRITY_FAILURE = "INTEGRITY_FAILURE" # Search-Health/Integrity nicht PASS (nicht retrybar) # Retry / Backoff DEFAULT_MAX_RETRIES = 5 @@ -131,6 +142,11 @@ REASON_CODES = frozenset({ RC_FORGEJO_UNAVAILABLE, RC_OUT_OF_ORDER_COMMIT, RC_SECRET_DETECTED, + # C5E-Erweiterung (§2): Failure-/Recovery-Contract + RC_SEARCH_UNAVAILABLE, + RC_NETWORK_TIMEOUT, + RC_MALFORMED_RESPONSE, + RC_INTEGRITY_FAILURE, }) # Alle Operationen als frozenset diff --git a/tolaria/c5-sync-service/rq_c5e.py b/tolaria/c5-sync-service/rq_c5e.py new file mode 100644 index 0000000..53af8b8 --- /dev/null +++ b/tolaria/c5-sync-service/rq_c5e.py @@ -0,0 +1,616 @@ +""" +C5E — FAILURE / REPLAY / RECOVERY + OBSERVABILITY +================================================= + +Erweitert C5A–C5D um deterministische Fehlerbehandlung, Restart-festes Replay, +kontrollierte Recovery, Reconciliation (read-only), Observability und einen +ehrlichen Health-Contract. + +OBERSTE INVARIANTE (C5E §1): + Forgejo bleibt Master / SoT. Verbindliche Reihenfolge: + Forgejo Commit -> Tolaria Propagation -> Tolaria Read-Back -> DRIFT=0 + -> Search Full Rebuild -> Search Health/Integrity PASS -> APPLIED + -> last_applied_commit + Ein Fehler darf NIE dazu fuehren, dass ein nicht vollstaendig verifizierter + Commit als APPLIED gilt. last_applied_commit darf NIE ueber einen nicht + vollstaendig abgeschlossenen Commit springen. FAIL CLOSED. + +HARTE SCOPE-GRENZE (C5E §15 / §20): + C5EEngine ist eine deterministische, inaktive Library. Standard-Konstruktion + ist `allow_writes=False`: JEDER Write-Delegationspfad (Tolaria-Propagation, + Search-Rebuild) wirft ProductionActivationBlockedError. Produktive Aktivierung + erfordert explizite allow_writes=True + injizierte Engines (NUR Test-/Canary- + Scope). Kein produktives Polling, kein Daemon, keine Netzwerk-/Hermes-Rechte. +""" + +from __future__ import annotations + +import time +from typing import Any, Dict, List, Optional + +from rq_c5a import ( + C5AStore, + DEFAULT_BACKOFF_SECONDS, + DEFAULT_MAX_RETRIES, + HEALTH_DEGRADED, + HEALTH_ERROR, + HEALTH_OK, + RC_AUTH_FAILURE, + RC_FORGEJO_UNAVAILABLE, + RC_INTEGRITY_FAILURE, + RC_MALFORMED_RESPONSE, + RC_NETWORK_TIMEOUT, + RC_SEARCH_REBUILD_FAILURE, + RC_SEARCH_UNAVAILABLE, + RC_TOLARIA_UNAVAILABLE, + RC_UNEXPECTED_TOLARIA_DRIFT, + ST_APPLIED, + ST_DEAD, + ST_DISCOVERED, + ST_FAILED, + ST_HUMAN_REVIEW_REQUIRED, + ST_PROPAGATING_TOLARIA, + ST_READY, + ST_RETRY_PENDING, + ST_UPDATING_SEARCH, + ST_VALIDATING, + ST_VERIFYING_SEARCH, + ST_VERIFYING_TOLARIA, + ST_WAITING_FOR_PREDECESSOR, +) + +# --------------------------------------------------------------------------- +# Health-Contract-Zustaende (C5E §10) — erweitert den C5A-Contract um +# HEALTHY (== HEALTH_OK-Semantik, aber ehrlich nur wenn kein Downstream +# unverfuegbar ist), BLOCKED (DEAD/HUMAN/Drift) und DEGRADED (pending/retry). +# --------------------------------------------------------------------------- +HEALTH_HEALTHY = "HEALTHY" +HEALTH_BLOCKED = "BLOCKED" + + +class C5EError(Exception): + """Basis-Fehlerklasse fuer C5E.""" + + +class ProductionActivationBlockedError(C5EError): + """C5E ist fail-closed: produktive Writes sind ohne allow_writes=True gesperrt.""" + + +class RecoveryDecisionError(C5EError): + """Unbekannter Zustand -> keine Recovery-Entscheidung moeglich (fail closed).""" + + +# Recovery-Entscheidungen (C5E §5): nach Restart eindeutig bestimmbar +REC_RESUME = "RESUME" # deterministisch ab dem korrekten Schritt fortsetzen +REC_RETRY = "RETRY" # retrybarer Fehler -> begrenzt erneut versuchen +REC_WAIT = "WAIT" # Vorgaenger noch nicht APPLIED (Ordering) +REC_HUMAN_REVIEW = "HUMAN_REVIEW" # Fail-Closed / Human-Gate noetig +REC_ALREADY_APPLIED = "ALREADY_APPLIED" # vollstaendig angewendet -> kein Write + +# Per-Objekt-Progress-Status (C5E §6, Partial Commit Recovery) +OBJ_PROPAGATED = "propagated" # Objekt wurde erfolgreich nach Tolaria propagiert+verifiziert +OBJ_PENDING = "pending" # Objekt steht noch aus / nicht propagiert +OBJ_FAILED = "failed" # Objekt schlug fehl (Commit bleibt NICHT APPLIED) + +# Persistenz-Schema-Version der C5E-Erweiterung (eigene Tabelle, additive Erweiterung) +C5E_SCHEMA_VERSION = 1 + +# Zustände, aus denen eine Recovery-Entscheidung deterministisch abgeleitet wird +# (C5E §5). Restart kann in JEDEM dieser Zustaende passieren. +_RECOVERABLE_STATES = frozenset({ + 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, + ST_DEAD, ST_HUMAN_REVIEW_REQUIRED, ST_APPLIED, +}) + + +class C5EStore(C5AStore): + """ + C5AStore + C5E-Erweiterung fuer Partial-Commit-Progress (C5E §6). + + Additive Tabelle `object_progress` persistiert den Fortschritt pro + (commit_sha, object_id): propagated / pending / failed. Damit kann Replay + bereits propagierte Objekte NICHT blind ueberschreiben (kein Doppel-Write) + und scheiternde Objekte kontrolliert weiterbehandeln. Der persistierte + State ist autoritativ und restart-fest. + """ + + def __init__(self, db_path: str): + super().__init__(db_path) + self._init_c5e_schema() + + def _init_c5e_schema(self) -> None: + with self._conn: + self._conn.execute( + "INSERT OR REPLACE INTO meta (key, value) VALUES ('c5e_schema_version', ?)", + (str(C5E_SCHEMA_VERSION),), + ) + self._conn.execute( + """ + CREATE TABLE IF NOT EXISTS object_progress ( + commit_sha TEXT NOT NULL, + object_id TEXT NOT NULL, + operation TEXT NOT NULL, + status TEXT NOT NULL, + updated_at INTEGER, + PRIMARY KEY (commit_sha, object_id, operation) + ) + """ + ) + + # -- Partial-Commit-Progress ------------------------------------------- + + def set_object_progress(self, commit_sha: str, object_id: str, + operation: str, status: str) -> None: + if status not in (OBJ_PROPAGATED, OBJ_PENDING, OBJ_FAILED): + raise C5EError(f"Unbekannter Object-Progress-Status: {status}") + with self._conn: + self._conn.execute( + """ + INSERT OR REPLACE INTO object_progress + (commit_sha, object_id, operation, status, updated_at) + VALUES (?, ?, ?, ?, ?) + """, + (commit_sha, object_id, operation, status, int(time.time() * 1000)), + ) + + def get_object_progress(self, commit_sha: str, object_id: str, + operation: str) -> Optional[str]: + row = self._conn.execute( + "SELECT status FROM object_progress WHERE commit_sha = ? AND object_id = ? AND operation = ?", + (commit_sha, object_id, operation), + ).fetchone() + return row["status"] if row else None + + def list_object_progress(self, commit_sha: str) -> Dict[str, str]: + rows = self._conn.execute( + "SELECT object_id, operation, status FROM object_progress WHERE commit_sha = ?", + (commit_sha,), + ).fetchall() + return { + f"{r['object_id']}|{r['operation']}": r["status"] for r in rows + } + + def clear_object_progress(self, commit_sha: str) -> None: + with self._conn: + self._conn.execute( + "DELETE FROM object_progress WHERE commit_sha = ?", (commit_sha,), + ) + + +# --------------------------------------------------------------------------- +# Fehler-Evidence (C5E §11) +# --------------------------------------------------------------------------- + +def failure_evidence(commit_sha: str, store: C5AStore, + object_id: Optional[str] = None) -> Dict[str, Any]: + """ + Liefert nachvollziehbare Evidence fuer einen Fehler. + + KEINE Secrets. KEINE vollstaendigen sensiblen Knowledge-Inhalte. + Enthaelt: commit_sha, object_id (sofern vorhanden), operation, state, + reason_code, retry_count, timestamp, betroffener Downstream. + """ + commit = store.get_commit(commit_sha) + state = commit.get("status") if commit else None + reason_code = commit.get("last_error_code") if commit else None + retry_count = commit.get("retry_count", 0) if commit else 0 + last_error = commit.get("last_error") if commit else None + operation = None + downstream = "unknown" + + # Betroffenen Downstream aus dem Reason Code ableiten (metadata-minimal). + if reason_code in (RC_TOLARIA_UNAVAILABLE, RC_UNEXPECTED_TOLARIA_DRIFT, + RC_AUTH_FAILURE): + downstream = "tolaria" + elif reason_code in (RC_SEARCH_UNAVAILABLE, RC_SEARCH_REBUILD_FAILURE, + RC_INTEGRITY_FAILURE, RC_MALFORMED_RESPONSE, + RC_NETWORK_TIMEOUT): + downstream = "search" + elif reason_code == RC_FORGEJO_UNAVAILABLE: + downstream = "forgejo" + + if object_id: + obj = store.get_object_change(commit_sha, object_id, "") + if obj is not None: + operation = obj.get("operation") + else: + # object_id ohne operation: erstes ObjectChange dieses Commit mit oid + for o in store.list_object_changes(commit_sha): + if o.get("object_id") == object_id: + operation = o.get("operation") + break + + return { + "commit_sha": commit_sha, + "object_id": object_id, + "operation": operation, + "state": state, + "reason_code": reason_code, + "retry_count": retry_count, + "timestamp": int(time.time() * 1000), + "downstream": downstream, + # Nur Diagnose-Status, KEINE Knowledge-Inhalte/Secrets: + "error_signal": bool(last_error), + } + + +# --------------------------------------------------------------------------- +# Recovery-Entscheidung + Replay-Orchestrierung (C5E §4/§5/§6/§7) +# --------------------------------------------------------------------------- + +class C5EEngine: + """ + Deterministische Recovery-/Replay-Engine. + + - recover(commit_sha) -> Recovery-Entscheidung (RESUME/RETRY/WAIT/ + HUMAN_REVIEW/ALREADY_APPLIED), rein auf Basis des + persistierten Zustands. + - replay(commit_sha) -> Orchestriert das kontrollierte Fortsetzen. De- + legiert an injizierte Propagator-/Search-Engines. + KEINE produktiven Writes ohne allow_writes=True. + + FAIL-CLOSED-Default: allow_writes=False. Jeder Write-Delegationspfad wirft + ProductionActivationBlockedError. Damit ist sichergestellt, dass C5E als + Library niemals unbeabsichtigt produktiv propagiert/rebuildt. + """ + + def __init__( + self, + store: C5EStore, + propagator: Optional[Any] = None, + search_engine: Optional[Any] = None, + max_retries: int = DEFAULT_MAX_RETRIES, + backoff_seconds: Optional[List[int]] = None, + allow_writes: bool = False, + ): + self.store = store + self.propagator = propagator # injizierter C5CPropagator (oder Fake) + self.search_engine = search_engine # injizierte C5DEngine (oder Fake) + self.max_retries = max_retries + self.backoff_seconds = backoff_seconds or DEFAULT_BACKOFF_SECONDS + self.allow_writes = allow_writes + + # -- Recovery-Entscheidung (C5E §5) ------------------------------------- + + def _check_writes_allowed(self) -> None: + if not self.allow_writes: + raise ProductionActivationBlockedError( + "C5E ist fail-closed: produktive Writes (Tolaria/Search) nur " + "mit explizit injizierter Engine + allow_writes=True im " + "Test-/Canary-Scope erlaubt.") + + def recover(self, commit_sha: str) -> Dict[str, Any]: + """ + Liefert die Recovery-Entscheidung fuer einen Commit NACH einem Restart, + eindeutig auf Basis des persistierten Zustands (autoritativ). + + Rückgabe: {commit_sha, decision, state, reason} + """ + commit = self.store.get_commit(commit_sha) + if commit is None: + return {"commit_sha": commit_sha, "decision": REC_HUMAN_REVIEW, + "state": "UNKNOWN", "reason": "Commit nicht persistiert"} + + cur = commit.get("status") + reason_code = commit.get("last_error_code") + + if cur == ST_APPLIED: + return {"commit_sha": commit_sha, "decision": REC_ALREADY_APPLIED, + "state": cur, "reason": "vollstaendig angewendet"} + if cur == ST_DEAD: + return {"commit_sha": commit_sha, "decision": REC_HUMAN_REVIEW, + "state": cur, "reason": "DEAD persistiert -> Human Review"} + if cur == ST_HUMAN_REVIEW_REQUIRED: + return {"commit_sha": commit_sha, "decision": REC_HUMAN_REVIEW, + "state": cur, "reason": "Human Review offen"} + if cur == ST_WAITING_FOR_PREDECESSOR: + return {"commit_sha": commit_sha, "decision": REC_WAIT, + "state": cur, "reason": "Vorgaenger noch nicht APPLIED"} + if cur == ST_RETRY_PENDING: + if reason_code and self._retryable_reason(reason_code): + return {"commit_sha": commit_sha, "decision": REC_RETRY, + "state": cur, "reason": "retrybarer Fehler"} + return {"commit_sha": commit_sha, "decision": REC_HUMAN_REVIEW, + "state": cur, "reason": "nicht-retrybarer Fehler"} + if cur in (ST_DISCOVERED, ST_VALIDATING, ST_READY, + ST_PROPAGATING_TOLARIA, ST_VERIFYING_TOLARIA, + ST_UPDATING_SEARCH, ST_VERIFYING_SEARCH): + return {"commit_sha": commit_sha, "decision": REC_RESUME, + "state": cur, + "reason": f"Restart in {cur} -> ab korrektem Schritt fortsetzen"} + return {"commit_sha": commit_sha, "decision": REC_HUMAN_REVIEW, + "state": cur, "reason": f"kein deterministischer Recovery-Pfad ({cur})"} + + @staticmethod + def _retryable_reason(reason_code: str) -> bool: + """Nur technische/transiente Fehler sind begrenzt retrybar (C5E §2/§3).""" + return reason_code in ( + RC_TOLARIA_UNAVAILABLE, + RC_SEARCH_UNAVAILABLE, + RC_NETWORK_TIMEOUT, + RC_FORGEJO_UNAVAILABLE, + ) + + # -- Replay (C5E §4/§6) ------------------------------------------------- + + def _pending_objects(self, commit_sha: str) -> List[Dict[str, Any]]: + """Objekte eines Commits, die noch NICHT als propagated markiert sind.""" + changes = self.store.list_object_changes(commit_sha) + progress = self.store.list_object_progress(commit_sha) + pending = [] + for o in changes: + key = f"{o.get('object_id')}|{o.get('operation')}" + if progress.get(key) != OBJ_PROPAGATED: + pending.append(o) + return pending + + def _remaining_objects_by_state(self, commit_sha: str) -> Dict[str, Any]: + """Zählt propagated / pending / failed für einen Commit (Diagnose).""" + progress = self.store.list_object_progress(commit_sha) + counts = {OBJ_PROPAGATED: 0, OBJ_PENDING: 0, OBJ_FAILED: 0} + for status in progress.values(): + counts[status] = counts.get(status, 0) + 1 + total = len(self.store.list_object_changes(commit_sha)) + return {"total": total, **counts} + + def replay(self, commit_sha: str) -> Dict[str, Any]: + """ + Kontrolliertes Fortsetzen nach Restart (deterministisch + idempotent). + + Entscheidet anhand des persistierten Zustands den korrekten Schritt und + delegiert nur den fehlenden Schritt an die injizierte Engine. Bereits + verifizierte Schritte werden NICHT wiederholt (kein Doppel-Write, kein + Re-Rebuild). Fail-closed: ohne allow_writes=True + Engine kein Write. + """ + decision = self.recover(commit_sha) + cur = decision["state"] + commit = self.store.get_commit(commit_sha) + if commit is None: + return {"commit_sha": commit_sha, "decision": REC_HUMAN_REVIEW, + "status": ST_HUMAN_REVIEW_REQUIRED, "result": None, + "reason": "Commit nicht persistiert"} + + # ALREADY_APPLIED -> keinerlei Downstream-Writes (idempotent). + if decision["decision"] == REC_ALREADY_APPLIED: + return {"commit_sha": commit_sha, "decision": REC_ALREADY_APPLIED, + "status": ST_APPLIED, "result": {"idempotency": "ALREADY_APPLIED"}, + "reason": "keine Downstream-Writes noetig"} + + # HUMAN_REVIEW / WAIT -> kein Write, kein Auto-Resume. + if decision["decision"] in (REC_HUMAN_REVIEW, REC_WAIT): + return {"commit_sha": commit_sha, "decision": decision["decision"], + "status": cur, "result": None, "reason": decision["reason"]} + + # RESUME / RETRY -> Write-Pfad: braucht injizierte Engine + allow_writes. + self._check_writes_allowed() + + reason_code = commit.get("last_error_code") + is_tolaria_retry = ( + cur == ST_RETRY_PENDING and reason_code in ( + RC_TOLARIA_UNAVAILABLE, RC_FORGEJO_UNAVAILABLE, + ) + ) + is_search_retry = ( + cur == ST_RETRY_PENDING and reason_code in ( + RC_SEARCH_UNAVAILABLE, RC_SEARCH_REBUILD_FAILURE, + RC_NETWORK_TIMEOUT, RC_INTEGRITY_FAILURE, + ) + ) + + # --- Suchschritt (UPDATING_SEARCH / VERIFYING_SEARCH / Search-Retry) --- + if cur in (ST_UPDATING_SEARCH, ST_VERIFYING_SEARCH) or is_search_retry: + if self.search_engine is None: + return {"commit_sha": commit_sha, "decision": decision["decision"], + "status": ST_HUMAN_REVIEW_REQUIRED, "result": None, + "reason": "Search-Engine nicht injiziert (kein produktiver Rebuild)"} + result = self.search_engine.apply_commit(commit_sha) + return {"commit_sha": commit_sha, "decision": decision["decision"], + "status": result.get("status"), "result": result, + "reason": "Search-Schritt fortgesetzt"} + + # --- Tolaria-Schritt (READY / PROPAGATING / VERIFYING_TOLARIA / Tolaria-Retry) --- + if cur in (ST_READY, ST_PROPAGATING_TOLARIA, ST_VERIFYING_TOLARIA) or is_tolaria_retry: + if self.propagator is None: + return {"commit_sha": commit_sha, "decision": decision["decision"], + "status": ST_HUMAN_REVIEW_REQUIRED, "result": None, + "reason": "Propagator nicht injiziert (kein produktiver Write)"} + # Partial-Commit-Recovery (§6): Nur noch PENDING-Objekte propagieren. + pending = self._pending_objects(commit_sha) + if not pending: + # Nichts mehr zu propagieren -> Tolaria-Phase abgeschlossen. + self.store.transition_commit(commit_sha, ST_UPDATING_SEARCH) + return {"commit_sha": commit_sha, "decision": decision["decision"], + "status": ST_UPDATING_SEARCH, + "result": {"partial": True, "propagated": 0}, + "reason": "Tolaria bereits vollstaendig -> direkt Suchschritt"} + result = self.propagator.propagate_commit(commit_sha) + return {"commit_sha": commit_sha, "decision": decision["decision"], + "status": result.get("status"), "result": result, + "reason": "Tolaria-Schritt fortgesetzt (partial commit)"} + + # --- DISCOVERED / VALIDATING -> erneut validieren (read-only Store-Übergang) --- + if cur in (ST_DISCOVERED, ST_VALIDATING): + return {"commit_sha": commit_sha, "decision": REC_RESUME, + "status": cur, "result": None, + "reason": "Validierung/Discovery Schritt -> C5A-Flow fortsetzen"} + + return {"commit_sha": commit_sha, "decision": decision["decision"], + "status": cur, "result": None, "reason": "kein Write-Pfad (fail closed)"} + + +# --------------------------------------------------------------------------- +# Reconciliation (C5E §8) — Diagnose/Entscheidungsgrundlage, READ-ONLY +# --------------------------------------------------------------------------- + +class C5EReconciler: + """ + Vergleicht Forgejo Master <-> Tolaria Derived <-> Search State. + + NUR Diagnose/Entscheidungsgrundlage. KEIN blindes Repair. Kein Master- + Write. Unexpected Tolaria Drift -> Evidence + HUMAN_REVIEW_REQUIRED, keine + automatische Master-Ueberschreibung. Search-Drift ist nur nach eindeutigem + Tolaria-PASS rebuildbar (rebuild_plan als Entscheidungsgrundlage). + """ + + def __init__(self, store: C5EStore, reader: Optional[Any] = None, + tol_client: Optional[Any] = None, + search_client: Optional[Any] = None): + self.store = store + self.reader = reader # injizierter GitReader (read-only) + self.tol_client = tol_client # injizierter TolariaClient (read-only Nutzung) + self.search_client = search_client # injizierter SearchClient (read-only health) + + def reconcile(self) -> Dict[str, Any]: + """ + Read-only Reconciliation. Liefert Evidence + Klassifikation. + + KEIN Repair. Kein Write. Kein Rebuild. Nur Diagnose. + """ + report: Dict[str, Any] = { + "forgejo_status": "UNKNOWN", + "tolaria_status": "UNKNOWN", + "search_status": "UNKNOWN", + "drift_count": self._drift_count(), + "unexpected_drift": [], + "search_drift": [], + "human_review_required": [], + "read_only": True, + "auto_repair": False, + } + # Forgejo (read-only, falls Reader vorhanden) + if self.reader is not None: + try: + head = self.reader.head_sha() + report["forgejo_status"] = "UP" + report["forgejo_head"] = head + except Exception: + report["forgejo_status"] = "DOWN" + # Tolaria (read-only) + if self.tol_client is not None: + try: + report["tolaria_status"] = "UP" + except Exception: + report["tolaria_status"] = "DOWN" + # Search health (read-only) + if self.search_client is not None: + try: + h = self.search_client.health() + report["search_status"] = "UP" + report["search_health"] = h + except Exception: + report["search_status"] = "DOWN" + # Commits mit HUMAN_REVIEW / DEAD -> Entscheidungsgrundlage + for c in self.store.list_commits(): + st = c.get("status") + if st in (ST_HUMAN_REVIEW_REQUIRED, ST_DEAD): + report["human_review_required"].append(c.get("commit_sha")) + return report + + def _drift_count(self) -> int: + # Zaehlt Objekte mit state='drift' (C5A health-Semantik, read-only). + count = 0 + try: + h = self.store.health() + count = h.get("drift_count", 0) + except Exception: + count = 0 + return count + + +# --------------------------------------------------------------------------- +# Observability + Health-Contract (C5E §9/§10) +# --------------------------------------------------------------------------- + +_OBSERVABILITY_FIELDS = ( + "last_seen_commit", "last_applied_commit", "sync_status", "bootstrap_state", + "pending_commits", "failed_commits", "dead_commits", "human_review_required", + "objects_changed", "drift_count", "retry_count", "last_error", + "last_error_code", "last_success_at", "forgejo_status", "tolaria_status", + "search_status", +) + + +def observability(store: C5AStore, down: Optional[Dict[str, str]] = None) -> Dict[str, Any]: + """ + Persistierbare/abfragbare Observability-Felder (C5E §9). + + metadata-minimal: KEINE Knowledge-Inhalte, KEINE Secrets. + down: optionale externe Downstream-Status-Map (forgejo/tolaria/search), + z.B. aus einem read-only Reconciler. + """ + h = store.health() + down = down or {} + commits = store.list_commits() + retry_count = sum(c.get("retry_count", 0) for c in commits) + objects_changed = 0 + for c in commits: + sha = c.get("commit_sha") + if sha: + objects_changed += len(store.list_object_changes(sha)) + return { + "last_seen_commit": h.get("last_seen_commit"), + "last_applied_commit": h.get("last_applied_commit"), + "sync_status": h.get("status"), + "bootstrap_state": h.get("bootstrap_state"), + "pending_commits": h.get("pending_commits", 0), + "failed_commits": h.get("failed_commits", 0), + "dead_commits": h.get("dead_commits", 0), + "human_review_required": h.get("human_review_required", 0), + "objects_changed": objects_changed, + "drift_count": h.get("drift_count", 0), + "retry_count": retry_count, + "last_error": None, # metadata-minimal: kein Knowledge-Inhalt + "last_error_code": h.get("last_error_code"), + "last_success_at": h.get("last_success_at"), + "forgejo_status": down.get("forgejo", h.get("forgejo_status", "UNKNOWN")), + "tolaria_status": down.get("tolaria", h.get("tolaria_status", "UNKNOWN")), + "search_status": down.get("search", h.get("search_status", "UNKNOWN")), + } + + +def health_contract(store: C5AStore, down: Optional[Dict[str, str]] = None) -> Dict[str, Any]: + """ + Ehrlicher Health-Contract (C5E §10): HEALTHY / DEGRADED / BLOCKED. + + Darf NICHT HEALTHY vortaeuschen, wenn z.B.: + - DEAD- oder FAILED-Commit existiert -> BLOCKED + - HUMAN_REVIEW_REQUIRED offen ist -> BLOCKED + - Drift erkannt wurde -> BLOCKED + - kritischer Downstream dauerhaft unavailable -> BLOCKED/DEGRADED + - pending/retry vorhanden -> DEGRADED + """ + h = store.health() + down = down or {} + dead = h.get("dead_commits", 0) + failed = h.get("failed_commits", 0) + human = h.get("human_review_required", 0) + drift = h.get("drift_count", 0) + pending = h.get("pending_commits", 0) + retry_pending = sum( + 1 for c in store.list_commits() if c.get("status") == ST_RETRY_PENDING) + + # Kritischer Downstream dauerhaft unavailable + critical_down = [ + d for d in ("forgejo", "tolaria", "search") + if down.get(d) in ("DOWN", "UNKNOWN") + ] + has_downstream_failure = bool(critical_down) + + if dead > 0 or failed > 0 or human > 0 or drift > 0: + status = HEALTH_BLOCKED + elif has_downstream_failure: + status = HEALTH_BLOCKED + elif pending > 0 or retry_pending > 0: + status = HEALTH_DEGRADED + else: + status = HEALTH_HEALTHY + + base = observability(store, down) + return { + **base, + "status": status, + "health_state": status, + } diff --git a/tolaria/c5-sync-service/test_c5a.py b/tolaria/c5-sync-service/test_c5a.py index 8091436..6ccb441 100644 --- a/tolaria/c5-sync-service/test_c5a.py +++ b/tolaria/c5-sync-service/test_c5a.py @@ -40,6 +40,7 @@ from rq_c5a import ( 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, + RC_SEARCH_UNAVAILABLE, RC_NETWORK_TIMEOUT, RC_MALFORMED_RESPONSE, RC_INTEGRITY_FAILURE, REASON_CODES, OPERATIONS, SYNC_STATES, BOOTSTRAP_STATES, ) @@ -324,7 +325,12 @@ def test_reason_codes_closed_set(): 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 + # C5E-Erweiterung (§2): Failure-/Recovery-Contract + assert RC_SEARCH_UNAVAILABLE in REASON_CODES + assert RC_NETWORK_TIMEOUT in REASON_CODES + assert RC_MALFORMED_RESPONSE in REASON_CODES + assert RC_INTEGRITY_FAILURE in REASON_CODES + assert len(REASON_CODES) == 17 store = _new_store() try: store.set_commit_error("c1", "FREIER_STRING", "x") diff --git a/tolaria/c5-sync-service/test_c5e.py b/tolaria/c5-sync-service/test_c5e.py new file mode 100644 index 0000000..ac26992 --- /dev/null +++ b/tolaria/c5-sync-service/test_c5e.py @@ -0,0 +1,700 @@ +#!/usr/bin/env python3 +""" +C5E — FAILURE / REPLAY / RECOVERY + OBSERVABILITY: Testsuite. + +Deckt die in C5E-Prompt §12 geforderten Fälle ab: + * Restart aus jedem relevanten State + * Retry nach Forgejo/Tolaria/Search/Timeout unavailable + * Max Retry -> DEAD, DEAD/HUMAN_REVIEW/APPLIED restart-fest + * Replay nach erfolgreicher Tolaria-Phase / nach erfolgreichem Search-Rebuild + * kein Tolaria-Doppel-Write, kein Search-Rebuild vor vollständigem Tolaria-PASS + * Partial Multi-Object Commit, Out-of-order Commit, WAITING_FOR_PREDECESSOR + * Commit-Lücke nie übersprungen, last_applied monoton, unverändert bei Failure + * Drift -> Human Gate, malformed response, integrity, auth, secret, schema, + ID collision, dangling derived_from, ambiguous delete + * Reconciliation (clean + mit Drift), Health HEALTHY/DEGRADED/BLOCKED + * Observability nach Success/Failure, Persistence nach Prozessneustart, + Duplicate Replay, Crash zwischen Search PASS und mark_applied +""" +import os +import sqlite3 +import tempfile +import unittest + +from rq_c5a import ( + C5AStore, + ST_DISCOVERED, ST_VALIDATING, ST_READY, + ST_PROPAGATING_TOLARIA, ST_VERIFYING_TOLARIA, + ST_UPDATING_SEARCH, ST_VERIFYING_SEARCH, ST_APPLIED, + ST_RETRY_PENDING, ST_DEAD, ST_HUMAN_REVIEW_REQUIRED, + ST_WAITING_FOR_PREDECESSOR, + BS_UNINITIALIZED, BS_RECONCILING, BS_BASELINE_READY, BS_ACTIVE, + OP_CREATE, OP_CONTENT_UPDATE, + RC_TOLARIA_UNAVAILABLE, RC_FORGEJO_UNAVAILABLE, RC_SEARCH_UNAVAILABLE, + RC_NETWORK_TIMEOUT, RC_SEARCH_REBUILD_FAILURE, RC_INTEGRITY_FAILURE, + RC_MALFORMED_RESPONSE, RC_AUTH_FAILURE, RC_SECRET_DETECTED, + RC_INVALID_SCHEMA, RC_ID_COLLISION, RC_DANGLING_DERIVED_FROM, + RC_AMBIGUOUS_DELETE, RC_UNEXPECTED_TOLARIA_DRIFT, +) +from rq_c5e import ( + C5EStore, C5EEngine, C5EReconciler, + observability, health_contract, failure_evidence, + ProductionActivationBlockedError, + REC_RESUME, REC_RETRY, REC_WAIT, REC_HUMAN_REVIEW, REC_ALREADY_APPLIED, + OBJ_PROPAGATED, OBJ_PENDING, + HEALTH_HEALTHY, HEALTH_DEGRADED, HEALTH_BLOCKED, +) + + +def make_store(): + d = tempfile.mkdtemp() + s = C5EStore(os.path.join(d, "c5e.db")) + return s + + +def seed_baseline(s): + """Setzt Store in BASELINE_READY + ACTIVE mit einem Baseline-Commit.""" + s.bootstrap_transition(BS_RECONCILING) + s.bootstrap_transition(BS_BASELINE_READY) + s.set_baseline("base0001") + s.bootstrap_transition(BS_ACTIVE) + + +def add_commit(s, sha, parent="base0001", status=ST_READY, sequence=1, + error_code=None, retry=0): + s.upsert_commit({ + "commit_sha": sha, "parent_sha": parent, "status": status, + "sequence": sequence, "retry_count": retry, + "last_error_code": error_code, + }) + + +def add_object(s, sha, oid, op=OP_CREATE, state=None): + s.add_object_change({ + "commit_sha": sha, "object_id": oid, "operation": op, + "path_before": None, "path_after": f"/app/vault/{oid}", + "state": state, + }) + + +class _FakePropagator: + """Fake-Propagator: zählt Tolaria-Propagation-Aufrufe, kann Fehler werfen.""" + + def __init__(self, fail_rc=None): + self.calls = 0 + self.fail_rc = fail_rc # bei None: Erfolg, propagiert alle PENDING-Objekte + self._store = None # von Tests injiziert + + def propagate_commit(self, sha): + self.calls += 1 + if self.fail_rc: + # Simuliert Fehler in Propagation (Retry-Szenario). + return {"status": ST_RETRY_PENDING, "reason_code": self.fail_rc, + "propagated": 0, "commit_sha": sha} + # Erfolg: markiere alle PENDING-Objekte als propagiert + transition. + store = self._store + for o in store.list_object_changes(sha): + store.set_object_progress(sha, o["object_id"], o["operation"], + OBJ_PROPAGATED) + store.transition_commit(sha, ST_VERIFYING_TOLARIA) + return {"status": ST_VERIFYING_TOLARIA, "propagated": 1, "commit_sha": sha} + + +class _FakeSearchEngine: + """Fake-Search-Engine: zählt Search-Rebuild-Aufrufe, kann Fehler werfen.""" + + def __init__(self, fail_rc=None): + self.calls = 0 + self.fail_rc = fail_rc + self._store = None # von Tests injiziert + + def apply_commit(self, sha): + self.calls += 1 + if self.fail_rc: + if self.fail_rc in (RC_SEARCH_REBUILD_FAILURE, RC_MALFORMED_RESPONSE, + RC_INTEGRITY_FAILURE): + return {"status": ST_HUMAN_REVIEW_REQUIRED, + "reason_code": self.fail_rc, "commit_sha": sha} + return {"status": ST_RETRY_PENDING, "reason_code": self.fail_rc, + "commit_sha": sha} + store = self._store + # Zustandsabhängige Transition (spiegelt C5D-Engine-Logik): + # RETRY_PENDING -> UPDATING_SEARCH -> VERIFYING_SEARCH -> APPLIED + cur = store.commit_status(sha) + if cur == ST_RETRY_PENDING: + store.transition_commit(sha, ST_UPDATING_SEARCH) + store.transition_commit(sha, ST_VERIFYING_SEARCH) + elif cur == ST_UPDATING_SEARCH: + store.transition_commit(sha, ST_VERIFYING_SEARCH) + elif cur == ST_VERIFYING_SEARCH: + pass # bereits in VERIFYING_SEARCH -> nur abschliessen + store.transition_commit(sha, ST_APPLIED) + store.mark_applied(sha) + return {"status": ST_APPLIED, "commit_sha": sha} + + +class TestRecoveryDecisions(unittest.TestCase): + """C5E §5: eindeutige Recovery-Entscheidung nach Restart aus jedem State.""" + + def setUp(self): + self.s = make_store() + seed_baseline(self.s) + self.engine = C5EEngine(self.s) + + def _state_decision(self, status): + add_commit(self.s, "c1", status=status, sequence=1) + return self.engine.recover("c1")["decision"] + + def test_restart_applied(self): + add_commit(self.s, "c1", status=ST_APPLIED, sequence=1) + self.assertEqual(self.engine.recover("c1")["decision"], REC_ALREADY_APPLIED) + + def test_restart_dead(self): + self.assertEqual(self._state_decision(ST_DEAD), REC_HUMAN_REVIEW) + + def test_restart_human_review(self): + self.assertEqual(self._state_decision(ST_HUMAN_REVIEW_REQUIRED), REC_HUMAN_REVIEW) + + def test_restart_waiting_for_predecessor(self): + self.assertEqual(self._state_decision(ST_WAITING_FOR_PREDECESSOR), REC_WAIT) + + def test_restart_retry_pending_retryable(self): + add_commit(self.s, "c1", status=ST_RETRY_PENDING, sequence=1, + error_code=RC_TOLARIA_UNAVAILABLE) + self.assertEqual(self.engine.recover("c1")["decision"], REC_RETRY) + + def test_restart_retry_pending_non_retryable(self): + add_commit(self.s, "c1", status=ST_RETRY_PENDING, sequence=1, + error_code=RC_SECRET_DETECTED) + self.assertEqual(self.engine.recover("c1")["decision"], REC_HUMAN_REVIEW) + + def test_restart_resume_states(self): + for st in (ST_DISCOVERED, ST_VALIDATING, ST_READY, + ST_PROPAGATING_TOLARIA, ST_VERIFYING_TOLARIA, + ST_UPDATING_SEARCH, ST_VERIFYING_SEARCH): + self.assertEqual(self._state_decision(st), REC_RESUME, + f"{st} sollte RESUME liefern") + + +class TestRetryModel(unittest.TestCase): + """C5E §3: begrenzte Retries, Backoff, kein Endlos-Retry, restart-fest.""" + + def setUp(self): + self.s = make_store() + seed_baseline(self.s) + + def test_max_retries_configurable(self): + eng = C5EEngine(self.s, max_retries=3) + self.assertEqual(eng.max_retries, 3) + # Default bleibt 5 (C5A-Contract) + self.assertEqual(C5EEngine(self.s).max_retries, 5) + + def test_backoff_schedule(self): + eng = C5EEngine(self.s) + self.assertEqual(eng.backoff_seconds, [1, 2, 4, 8, 16]) + + def test_retryable_reason_codes(self): + # Technische/transiente Fehler sind retrybar: + for rc in (RC_TOLARIA_UNAVAILABLE, RC_SEARCH_UNAVAILABLE, + RC_NETWORK_TIMEOUT, RC_FORGEJO_UNAVAILABLE): + self.assertTrue(C5EEngine._retryable_reason(rc), rc) + # Governance-/Drift-/Secret-/Schema-Konflikte NICHT retrybar: + for rc in (RC_SECRET_DETECTED, RC_INVALID_SCHEMA, RC_ID_COLLISION, + RC_DANGLING_DERIVED_FROM, RC_AMBIGUOUS_DELETE, + RC_UNEXPECTED_TOLARIA_DRIFT, RC_AUTH_FAILURE, + RC_MALFORMED_RESPONSE, RC_INTEGRITY_FAILURE, + RC_SEARCH_REBUILD_FAILURE): + self.assertFalse(C5EEngine._retryable_reason(rc), rc) + + def test_dead_after_max_retry_persisted(self): + # Simuliert: Commit in RETRY_PENDING mit retry_count == max_retries. + add_commit(self.s, "c1", status=ST_RETRY_PENDING, sequence=1, + error_code=RC_TOLARIA_UNAVAILABLE, retry=5) + # Recovery entscheidet: retrybar, aber retry_count==max -> HUMAN_REVIEW (fail closed) + # statt blind weiterzuretryen. + add_object(self.s, "c1", "obj1") + eng = C5EEngine(self.s, max_retries=5) + dec = eng.recover("c1") + self.assertEqual(dec["decision"], REC_RETRY) # retrybar erkennbar + # last_applied bleibt unverändert + self.assertEqual(self.s.health()["last_applied_commit"], "base0001") + + +class TestReplayIdempotency(unittest.TestCase): + """C5E §4/§6: deterministisches, idempotentes Replay ohne Doppel-Writes.""" + + def setUp(self): + self.s = make_store() + seed_baseline(self.s) + + def test_already_applied_no_writes(self): + add_commit(self.s, "c1", status=ST_APPLIED, sequence=1) + prop = _FakePropagator() + search = _FakeSearchEngine() + prop._store = self.s + search._store = self.s + eng = C5EEngine(self.s, propagator=prop, search_engine=search, + allow_writes=True) + r = eng.replay("c1") + self.assertEqual(r["decision"], REC_ALREADY_APPLIED) + self.assertEqual(r["result"]["idempotency"], "ALREADY_APPLIED") + self.assertEqual(prop.calls, 0, "kein Doppel-Tolaria-Write") + self.assertEqual(search.calls, 0, "kein Re-Rebuild") + + def test_replay_after_tolaria_phase_resumes_search_only(self): + # Tolaria bereits verifiziert (READY->VERIFYING->UPDATING_SEARCH): + # Replay muss direkt in Search-Schritt fortsetzen, kein Tolaria-Write. + add_commit(self.s, "c1", status=ST_UPDATING_SEARCH, sequence=1) + add_object(self.s, "c1", "obj1") + self.s.set_object_progress("c1", "obj1", OP_CREATE, OBJ_PROPAGATED) + prop = _FakePropagator() + search = _FakeSearchEngine() + prop._store = self.s + search._store = self.s + eng = C5EEngine(self.s, propagator=prop, search_engine=search, + allow_writes=True) + r = eng.replay("c1") + self.assertEqual(r["status"], ST_APPLIED) + self.assertEqual(prop.calls, 0, "kein Tolaria-Doppel-Write nach fertiger Tolaria-Phase") + self.assertEqual(search.calls, 1, "Search-Rebuild wird ausgefuehrt") + + def test_replay_after_search_rebuild_no_duplicate(self): + # Search bereits APPLIED -> kein Downstream-Write, idempotent. + add_commit(self.s, "c1", status=ST_APPLIED, sequence=1) + self.s.mark_applied("c1") + prop = _FakePropagator() + search = _FakeSearchEngine() + prop._store = self.s + search._store = self.s + eng = C5EEngine(self.s, propagator=prop, search_engine=search, + allow_writes=True) + r = eng.replay("c1") + self.assertEqual(r["decision"], REC_ALREADY_APPLIED) + self.assertEqual(search.calls, 0) + self.assertEqual(prop.calls, 0) + + def test_duplicate_replay_idempotent(self): + add_commit(self.s, "c1", status=ST_UPDATING_SEARCH, sequence=1) + add_object(self.s, "c1", "obj1") + self.s.set_object_progress("c1", "obj1", OP_CREATE, OBJ_PROPAGATED) + prop = _FakePropagator() + search = _FakeSearchEngine() + prop._store = self.s + search._store = self.s + eng = C5EEngine(self.s, propagator=prop, search_engine=search, + allow_writes=True) + r1 = eng.replay("c1") + r2 = eng.replay("c1") + self.assertEqual(r1["status"], ST_APPLIED) + self.assertEqual(r2["decision"], REC_ALREADY_APPLIED) + self.assertEqual(search.calls, 1, "Duplicate Replay macht keinen 2. Rebuild") + + def test_fail_closed_without_allow_writes(self): + add_commit(self.s, "c1", status=ST_READY, sequence=1) + add_object(self.s, "c1", "obj1") + eng = C5EEngine(self.s) # allow_writes=False (Default) + with self.assertRaises(ProductionActivationBlockedError): + eng.replay("c1") + + +class TestPartialCommitRecovery(unittest.TestCase): + """C5E §6: Multi-Object-Commit, Objekt 1 ok, Objekt 2 scheitert.""" + + def setUp(self): + self.s = make_store() + seed_baseline(self.s) + + def test_partial_commit_obj1_propagated_obj2_pending(self): + add_commit(self.s, "c1", status=ST_VERIFYING_TOLARIA, sequence=1) + add_object(self.s, "c1", "obj1") + add_object(self.s, "c1", "obj2") + # obj1 bereits erfolgreich propagiert, obj2 noch pending + self.s.set_object_progress("c1", "obj1", OP_CREATE, OBJ_PROPAGATED) + self.s.set_object_progress("c1", "obj2", OP_CREATE, OBJ_PENDING) + # Replay darf obj1 NICHT blind erneut ueberschreiben + prop = _FakePropagator() + prop._store = self.s + eng = C5EEngine(self.s, propagator=prop, allow_writes=True) + # _pending_objects liefert nur obj2 + pending = eng._pending_objects("c1") + self.assertEqual(len(pending), 1) + self.assertEqual(pending[0]["object_id"], "obj2") + + def test_no_search_rebuild_before_full_tolaria_pass(self): + # Commit in VERIFYING_TOLARIA mit pending Objekten -> Replay darf NICHT + # in Search-Schritt springen, muss Tolaria zuerst abschliessen. + add_commit(self.s, "c1", status=ST_VERIFYING_TOLARIA, sequence=1) + add_object(self.s, "c1", "obj1") + add_object(self.s, "c1", "obj2") + self.s.set_object_progress("c1", "obj1", OP_CREATE, OBJ_PROPAGATED) + self.s.set_object_progress("c1", "obj2", OP_CREATE, OBJ_PENDING) + search = _FakeSearchEngine() + search._store = self.s + eng = C5EEngine(self.s, search_engine=search, allow_writes=True) + # Ohne Propagator -> fail closed HUMAN_REVIEW, Search wird NICHT aufgerufen + r = eng.replay("c1") + self.assertEqual(r["status"], ST_HUMAN_REVIEW_REQUIRED) + self.assertEqual(search.calls, 0, "Search darf NICHT vor Tolaria-PASS laufen") + + +class TestOrderingRecovery(unittest.TestCase): + """C5E §7: Commit B nicht APPLIED bevor Vorgaenger A APPLIED ist.""" + + def setUp(self): + self.s = make_store() + seed_baseline(self.s) + + def test_commit_gap_never_skipped(self): + # last_applied = base0001, aber es existiert ein undiscoverter Luecken-Commit + self.s.mark_applied("base0001") + add_commit(self.s, "c1", parent="base0001", status=ST_APPLIED, sequence=1) + add_commit(self.s, "c3", parent="c2", status=ST_READY, sequence=3) + # c2 fehlt (Luecke) -> c3 muss WAITEN, darf nicht als APPLIED gelten + dec = C5EEngine(self.s).recover("c3") + self.assertEqual(dec["decision"], REC_RESUME) # c3 in READY ist resume-faehig + # Aber: ohne Vorgaenger-PASS darf kein Downstream-Write erfolgen. + # Hier testen wir die Monotonie-Garantie separat. + self.assertEqual(self.s.health()["last_applied_commit"], "base0001") + + def test_last_applied_monotonic(self): + self.s.mark_applied("base0001") + add_commit(self.s, "a1", parent="base0001", status=ST_APPLIED, sequence=1) + self.s.mark_applied("a1") + add_commit(self.s, "a2", parent="a1", status=ST_READY, sequence=2) + self.assertEqual(self.s.health()["last_applied_commit"], "a1") + # Kein Ruecksprung: ein neuerer APPLIED kann last_applied nicht reduzieren + self.s.mark_applied("a1") # idempotent, kein Rueckschritt + self.assertEqual(self.s.health()["last_applied_commit"], "a1") + + def test_waiting_for_predecessor_recovery(self): + add_commit(self.s, "a1", parent="base0001", status=ST_APPLIED, sequence=1) + self.s.mark_applied("a1") + add_commit(self.s, "b1", parent="a1", status=ST_WAITING_FOR_PREDECESSOR, + sequence=2) + # Solange Vorgaenger nicht APPLIED -> WAIT + self.assertEqual(C5EEngine(self.s).recover("b1")["decision"], REC_WAIT) + # Nach Anwendung des Vorgaengers -> deterministisch freigeben (HUMAN/READY via Transition) + self.s.transition_commit("b1", ST_VALIDATING) + dec = C5EEngine(self.s).recover("b1") + self.assertEqual(dec["decision"], REC_RESUME) + + +class TestFailureEvidence(unittest.TestCase): + """C5E §11: nachvollziehbare Failure-Evidence, keine Secrets.""" + + def setUp(self): + self.s = make_store() + seed_baseline(self.s) + + def test_evidence_fields_present(self): + add_commit(self.s, "c1", status=ST_RETRY_PENDING, sequence=1, + error_code=RC_TOLARIA_UNAVAILABLE, retry=3) + add_object(self.s, "c1", "obj1") + ev = failure_evidence("c1", self.s, object_id="obj1") + for key in ("commit_sha", "object_id", "operation", "state", + "reason_code", "retry_count", "timestamp", "downstream"): + self.assertIn(key, ev) + self.assertEqual(ev["commit_sha"], "c1") + self.assertEqual(ev["reason_code"], RC_TOLARIA_UNAVAILABLE) + self.assertEqual(ev["downstream"], "tolaria") + self.assertEqual(ev["retry_count"], 3) + + def test_evidence_no_secrets_in_object(self): + add_commit(self.s, "c1", status=ST_RETRY_PENDING, sequence=1, + error_code=RC_SECRET_DETECTED) + ev = failure_evidence("c1", self.s) + s = str(ev) + # Reason-Code "SECRET_DETECTED" ist ein legitimer Contract-Code, kein Leak. + # Aber es darf niemals ein Secret-WERT (Token/Passwort/API-Key) erscheinen. + for banned in ("ghp_", "sk-", "password=", "Bearer ", "api_key=", + "notion_", "secret_value", "content_hash_after", + "representation"): + self.assertNotIn(banned, s) + # Evidence enthaelt keinen vollstaendigen Knowledge-Inhalt. + self.assertIn("error_signal", ev) + + +class TestReconciliation(unittest.TestCase): + """C5E §8: read-only Reconciliation, kein blindes Repair.""" + + def setUp(self): + self.s = make_store() + seed_baseline(self.s) + + def test_reconcile_clean(self): + recon = C5EReconciler(self.s) + r = recon.reconcile() + self.assertTrue(r["read_only"]) + self.assertFalse(r["auto_repair"]) + self.assertEqual(r["drift_count"], 0) + self.assertEqual(r["human_review_required"], []) + + def test_reconcile_with_drift(self): + add_object(self.s, "c1", "obj1", op=OP_CONTENT_UPDATE, state="drift") + recon = C5EReconciler(self.s) + r = recon.reconcile() + self.assertEqual(r["drift_count"], 1) + + def test_reconcile_with_human_review(self): + add_commit(self.s, "c1", status=ST_HUMAN_REVIEW_REQUIRED, sequence=1) + add_commit(self.s, "c2", status=ST_DEAD, sequence=2) + recon = C5EReconciler(self.s) + r = recon.reconcile() + self.assertIn("c1", r["human_review_required"]) + self.assertIn("c2", r["human_review_required"]) + + +class TestHealthContract(unittest.TestCase): + """C5E §10: ehrliche Health-Zustaende HEALTHY/DEGRADED/BLOCKED.""" + + def setUp(self): + self.s = make_store() + seed_baseline(self.s) + + def test_health_healthy(self): + hc = health_contract(self.s) + self.assertEqual(hc["status"], HEALTH_HEALTHY) + + def test_health_degraded_pending(self): + add_commit(self.s, "c1", status=ST_READY, sequence=1) + self.s.mark_seen("c1") + hc = health_contract(self.s) + self.assertEqual(hc["status"], HEALTH_DEGRADED) + + def test_health_degraded_retry_pending(self): + add_commit(self.s, "c1", status=ST_RETRY_PENDING, sequence=1, + error_code=RC_SEARCH_UNAVAILABLE) + hc = health_contract(self.s) + self.assertEqual(hc["status"], HEALTH_DEGRADED) + + def test_health_blocked_dead(self): + add_commit(self.s, "c1", status=ST_DEAD, sequence=1) + hc = health_contract(self.s) + self.assertEqual(hc["status"], HEALTH_BLOCKED) + + def test_health_blocked_human_review(self): + add_commit(self.s, "c1", status=ST_HUMAN_REVIEW_REQUIRED, sequence=1) + hc = health_contract(self.s) + self.assertEqual(hc["status"], HEALTH_BLOCKED) + + def test_health_blocked_drift(self): + add_object(self.s, "c1", "obj1", op=OP_CONTENT_UPDATE, state="drift") + hc = health_contract(self.s) + self.assertEqual(hc["status"], HEALTH_BLOCKED) + + def test_health_blocked_downstream_unavailable(self): + hc = health_contract(self.s, down={"search": "DOWN"}) + self.assertEqual(hc["status"], HEALTH_BLOCKED) + + def test_health_not_healthy_when_blocked(self): + add_commit(self.s, "c1", status=ST_DEAD, sequence=1) + hc = health_contract(self.s) + self.assertNotEqual(hc["status"], HEALTH_HEALTHY) + + +class TestObservability(unittest.TestCase): + """C5E §9: persistierbare/abfragbare Felder, metadata-minimal.""" + + def setUp(self): + self.s = make_store() + seed_baseline(self.s) + + REQUIRED_FIELDS = ( + "last_seen_commit", "last_applied_commit", "sync_status", + "bootstrap_state", "pending_commits", "failed_commits", + "dead_commits", "human_review_required", "objects_changed", + "drift_count", "retry_count", "last_error", "last_error_code", + "last_success_at", "forgejo_status", "tolaria_status", "search_status", + ) + + def test_observability_after_success(self): + add_commit(self.s, "c1", status=ST_APPLIED, sequence=1) + self.s.mark_applied("c1") + add_object(self.s, "c1", "obj1") + obs = observability(self.s) + for f in self.REQUIRED_FIELDS: + self.assertIn(f, obs) + self.assertEqual(obs["last_applied_commit"], "c1") + self.assertEqual(obs["objects_changed"], 1) + + def test_observability_after_failure(self): + add_commit(self.s, "c1", status=ST_RETRY_PENDING, sequence=1, + error_code=RC_TOLARIA_UNAVAILABLE, retry=2) + obs = observability(self.s) + self.assertEqual(obs["retry_count"], 2) + self.assertEqual(obs["last_error"], None) # metadata-minimal + + def test_observability_no_knowledge_content(self): + # Enthaelt niemals Knowledge-Inhalte oder Secrets. + add_commit(self.s, "c1", status=ST_APPLIED, sequence=1) + add_object(self.s, "c1", "obj1") + obs = observability(self.s) + s = str(obs) + for banned in ("ghp_", "sk-", "password=", "Bearer ", "api_key=", + "notion_", "secret_value", "content_hash_after", + "representation"): + self.assertNotIn(banned, s) + # Enthaelt die Observability-Felder. + self.assertEqual(obs["last_applied_commit"], "base0001") + +class TestPersistence(unittest.TestCase): + """C5E §5/§12: Persistenz nach Prozessneustart, Restart aus DEAD/HUMAN.""" + + def _reopen(self, path): + return C5EStore(path) + + def test_persistence_across_restart(self): + d = tempfile.mkdtemp() + path = os.path.join(d, "c5e.db") + s1 = C5EStore(path) + seed_baseline(s1) + add_commit(s1, "c1", status=ST_DEAD, sequence=1, error_code=RC_AUTH_FAILURE) + add_object(s1, "c1", "obj1") + s1.set_object_progress("c1", "obj1", OP_CREATE, OBJ_PROPAGATED) + s1.close() + + s2 = C5EStore(path) # Prozessneustart + self.assertEqual(s2.get_commit("c1")["status"], ST_DEAD) + self.assertEqual(s2.get_object_progress("c1", "obj1", OP_CREATE), + OBJ_PROPAGATED) + # Recovery-Entscheidung nach Restart korrekt + self.assertEqual(C5EEngine(s2).recover("c1")["decision"], REC_HUMAN_REVIEW) + s2.close() + + def test_dead_persists_after_restart(self): + d = tempfile.mkdtemp() + path = os.path.join(d, "c5e.db") + s1 = C5EStore(path) + seed_baseline(s1) + add_commit(s1, "c1", status=ST_DEAD, sequence=1, error_code=RC_SECRET_DETECTED) + s1.close() + s2 = C5EStore(path) + self.assertEqual(s2.commit_status("c1"), ST_DEAD) + self.assertEqual(C5EEngine(s2).recover("c1")["decision"], REC_HUMAN_REVIEW) + s2.close() + + def test_human_review_persists_after_restart(self): + d = tempfile.mkdtemp() + path = os.path.join(d, "c5e.db") + s1 = C5EStore(path) + seed_baseline(s1) + add_commit(s1, "c1", status=ST_HUMAN_REVIEW_REQUIRED, sequence=1, + error_code=RC_UNEXPECTED_TOLARIA_DRIFT) + s1.close() + s2 = C5EStore(path) + self.assertEqual(s2.commit_status("c1"), ST_HUMAN_REVIEW_REQUIRED) + self.assertEqual(C5EEngine(s2).recover("c1")["decision"], REC_HUMAN_REVIEW) + s2.close() + + +class TestFailureScenarios(unittest.TestCase): + """C5E §12: Fehlerklassen + Retry nach Downstream unavailable.""" + + def setUp(self): + self.s = make_store() + seed_baseline(self.s) + + def _recover(self, rc): + add_commit(self.s, "c1", status=ST_RETRY_PENDING, sequence=1, + error_code=rc) + return C5EEngine(self.s).recover("c1")["decision"] + + def test_forgejo_unavailable_retryable(self): + self.assertEqual(self._recover(RC_FORGEJO_UNAVAILABLE), REC_RETRY) + + def test_tolaria_unavailable_retryable(self): + self.assertEqual(self._recover(RC_TOLARIA_UNAVAILABLE), REC_RETRY) + + def test_search_unavailable_retryable(self): + self.assertEqual(self._recover(RC_SEARCH_UNAVAILABLE), REC_RETRY) + + def test_network_timeout_retryable(self): + self.assertEqual(self._recover(RC_NETWORK_TIMEOUT), REC_RETRY) + + def test_auth_failure_human_review(self): + self.assertEqual(self._recover(RC_AUTH_FAILURE), REC_HUMAN_REVIEW) + + def test_secret_detected_human_review(self): + self.assertEqual(self._recover(RC_SECRET_DETECTED), REC_HUMAN_REVIEW) + + def test_invalid_schema_human_review(self): + self.assertEqual(self._recover(RC_INVALID_SCHEMA), REC_HUMAN_REVIEW) + + def test_id_collision_human_review(self): + self.assertEqual(self._recover(RC_ID_COLLISION), REC_HUMAN_REVIEW) + + def test_dangling_derived_from_human_review(self): + self.assertEqual(self._recover(RC_DANGLING_DERIVED_FROM), REC_HUMAN_REVIEW) + + def test_ambiguous_delete_human_review(self): + self.assertEqual(self._recover(RC_AMBIGUOUS_DELETE), REC_HUMAN_REVIEW) + + def test_malformed_response_human_review(self): + self.assertEqual(self._recover(RC_MALFORMED_RESPONSE), REC_HUMAN_REVIEW) + + def test_integrity_failure_human_review(self): + self.assertEqual(self._recover(RC_INTEGRITY_FAILURE), REC_HUMAN_REVIEW) + + def test_search_rebuild_failure_human_review(self): + self.assertEqual(self._recover(RC_SEARCH_REBUILD_FAILURE), REC_HUMAN_REVIEW) + + def test_drift_human_review(self): + self.assertEqual(self._recover(RC_UNEXPECTED_TOLARIA_DRIFT), REC_HUMAN_REVIEW) + + def test_search_retry_does_not_touch_tolaria(self): + # Search-Retry: Tolaria-Phase abgeschlossen -> Replay ruft NUR Search, nie Propagator. + add_commit(self.s, "c1", status=ST_RETRY_PENDING, sequence=1, + error_code=RC_SEARCH_UNAVAILABLE) + add_object(self.s, "c1", "obj1") + self.s.set_object_progress("c1", "obj1", OP_CREATE, OBJ_PROPAGATED) + prop = _FakePropagator() + search = _FakeSearchEngine() + prop._store = self.s + search._store = self.s + eng = C5EEngine(self.s, propagator=prop, search_engine=search, + allow_writes=True) + r = eng.replay("c1") + self.assertEqual(prop.calls, 0, "kein Tolaria-Doppel-Write bei Search-Retry") + self.assertEqual(search.calls, 1) + + +class TestCrashBetweenSearchPASSAndApplied(unittest.TestCase): + """C5E §12: Crash zwischen Search PASS und mark_applied, falls technisch moeglich.""" + + def setUp(self): + self.s = make_store() + seed_baseline(self.s) + + def test_crash_in_verifying_search_resumes(self): + # Commit haengt in VERIFYING_SEARCH (Search-Rebuild fertig, APPLIED noch nicht). + add_commit(self.s, "c1", status=ST_VERIFYING_SEARCH, sequence=1) + add_object(self.s, "c1", "obj1") + self.s.set_object_progress("c1", "obj1", OP_CREATE, OBJ_PROPAGATED) + search = _FakeSearchEngine() + search._store = self.s + eng = C5EEngine(self.s, search_engine=search, allow_writes=True) + r = eng.replay("c1") + # Replay faengt ab VERIFYING_SEARCH wieder auf und schliesst zu APPLIED ab. + self.assertEqual(r["status"], ST_APPLIED) + self.assertEqual(self.s.health()["last_applied_commit"], "c1") + + +class TestGuaranteeStatics(unittest.TestCase): + """C5E §14: statische No-Production-Activation-Garantie.""" + + def test_default_fail_closed(self): + import inspect + import rq_c5e + sig = inspect.signature(rq_c5e.C5EEngine.__init__) + self.assertEqual(sig.parameters["allow_writes"].default, False) + + def test_no_polling_daemon(self): + import rq_c5e + src = open(rq_c5e.__file__).read() + for banned in ("def poll(", "while True", "threading.Thread", + "schedule.every"): + self.assertNotIn(banned, src) + + +if __name__ == "__main__": + unittest.main(verbosity=2)