trading-system-docs/tolaria/c5-sync-service/rq_c5a.py

797 lines
33 KiB
Python

#!/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_PROPAGATING_TOLARIA, ST_HUMAN_REVIEW_REQUIRED), # C5C: Drift/Secret/dangling waehrend Propagation
(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,
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,
reason_code 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]) -> Optional[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, reason_code)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""",
(
obj["commit_sha"], obj.get("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"), obj.get("reason_code"),
),
)
return self.get_object_change(obj["commit_sha"], obj.get("object_id"), obj["operation"])
def get_object_change(self, commit_sha: str, object_id: Optional[str], operation: str) -> Optional[Dict[str, Any]]:
row = self._conn.execute(
"SELECT * FROM objects WHERE commit_sha = ? AND object_id IS ? 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,
}