FRESH CHECKER-VERDICT (deleg_3ea8ed5a): FAIL. Defekt: rq_c5e.py Replay 'alle Objekte bereits propagated' rief transition_commit(ST_UPDATING_SEARCH) aus PROPAGATING_TOLARIA/READY direkt auf, was InvalidTransitionError warf (Transition nicht in _ALLOWED_TRANSITIONS). Reparatur (invarianten-treu): statt den VERIFYING_TOLARIA-Schritt zu ueberspringen (wuerde Read-Back/DRIFT-Check verletzen), wird der formale State-Pfad READY->PROPAGATING->VERIFYING->UPDATING_SEARCH durchlaufen. Nutzt nur bereits erlaubte Transitions; keine _ALLOWED_TRANSITIONS-Aenderung, kein C5A/C5C-Risiko. Kein Tolaria-Doppel-Write (prop_calls=0), kein verfrühter Search (search_calls=0). + 2 Regressionstests (test_c5e.py): crash-window + ready-edge-case. Volle Suite: C5A 25/0 + B/C/D/E 170/0 = 195 OK. Guarantees true.
740 lines
32 KiB
Python
740 lines
32 KiB
Python
#!/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_crash_window_propagating_tolaria_all_objects_propagated(self):
|
|
# Crash-Window (§3/§4): alle Objekte bereits nach Tolaria propagiert, aber die
|
|
# Transition PROPAGATING_TOLARIA -> VERIFYING_TOLARIA ist noch nicht erfolgt.
|
|
# Replay MUSS deterministisch zum Suchschritt weiterlaufen (kein
|
|
# InvalidTransitionError, kein Tolaria-Doppel-Write, kein verfrühter Search).
|
|
# Regression für FRESH CHECKER-VERDICT (rq_c5e.py Zeile 424).
|
|
add_commit(self.s, "c1", status=ST_PROPAGATING_TOLARIA, sequence=1)
|
|
add_object(self.s, "c1", "obj1")
|
|
add_object(self.s, "c1", "obj2")
|
|
for o in self.s.list_object_changes("c1"):
|
|
self.s.set_object_progress("c1", o["object_id"], o["operation"], 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_UPDATING_SEARCH,
|
|
"Replay aus PROPAGATING_TOLARIA (alle propagiert) -> UPDATING_SEARCH")
|
|
self.assertEqual(prop.calls, 0, "kein Tolaria-Doppel-Write im Crash-Window")
|
|
# Formal über VERIFYING_TOLARIA gelaufen (Read-Back-Invariante), nicht übersprungen.
|
|
self.assertEqual(self.s.commit_status("c1"), ST_UPDATING_SEARCH)
|
|
|
|
def test_replay_ready_with_all_objects_propagated(self):
|
|
# Edge-Case: Commit noch in READY, aber Objekte bereits alle propagiert
|
|
# (unwahrscheinlich, aber deterministisch auflösbar -> UPDATING_SEARCH).
|
|
add_commit(self.s, "c1", status=ST_READY, 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_UPDATING_SEARCH)
|
|
self.assertEqual(prop.calls, 0, "kein Tolaria-Doppel-Write")
|
|
|
|
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)
|