795 lines
32 KiB
Python
795 lines
32 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_VERIFYING_TOLARIA, ST_UPDATING_SEARCH),
|
|
(ST_VERIFYING_TOLARIA, ST_RETRY_PENDING), # DRIFT != 0 / Tolaria-Fehler
|
|
(ST_VERIFYING_TOLARIA, ST_HUMAN_REVIEW_REQUIRED), # unexpected drift
|
|
(ST_UPDATING_SEARCH, ST_VERIFYING_SEARCH),
|
|
(ST_UPDATING_SEARCH, ST_RETRY_PENDING), # Search-Rebuild-Fehler
|
|
(ST_UPDATING_SEARCH, ST_HUMAN_REVIEW_REQUIRED),
|
|
(ST_VERIFYING_SEARCH, ST_APPLIED),
|
|
(ST_VERIFYING_SEARCH, ST_RETRY_PENDING), # Search-Health-Fehler
|
|
(ST_VERIFYING_SEARCH, ST_HUMAN_REVIEW_REQUIRED),
|
|
(ST_RETRY_PENDING, ST_READY), # Retry -> erneut propagieren
|
|
(ST_RETRY_PENDING, ST_RETRY_PENDING), # idempotenter erneuter Retry-Versuch
|
|
(ST_RETRY_PENDING, ST_FAILED), # Retry-Limit erreicht
|
|
(ST_RETRY_PENDING, ST_DEAD), # Max-Retry -> DEAD
|
|
(ST_RETRY_PENDING, ST_HUMAN_REVIEW_REQUIRED),# Human-Gate-Fehler
|
|
(ST_FAILED, ST_HUMAN_REVIEW_REQUIRED), # FAILED -> Human Review
|
|
(ST_DEAD, ST_HUMAN_REVIEW_REQUIRED), # DEAD -> Human Review
|
|
(ST_HUMAN_REVIEW_REQUIRED, ST_READY), # Human entscheidet -> erneut
|
|
(ST_HUMAN_REVIEW_REQUIRED, ST_APPLIED), # Human bestätigt als angewendet
|
|
(ST_HUMAN_REVIEW_REQUIRED, ST_DEAD), # Human verwirft
|
|
}
|
|
|
|
# Bootstrap-Übergänge
|
|
_ALLOWED_BOOTSTRAP = {
|
|
(BS_UNINITIALIZED, BS_RECONCILING),
|
|
(BS_RECONCILING, BS_BASELINE_READY),
|
|
(BS_RECONCILING, BS_UNINITIALIZED), # Reconciliation fehlgeschlagen -> zurück
|
|
(BS_BASELINE_READY, BS_ACTIVE),
|
|
}
|
|
|
|
|
|
class C5AError(Exception):
|
|
"""Basis-Fehlerklasse für C5A."""
|
|
|
|
def __init__(self, message: str, reason_code: Optional[str] = None):
|
|
super().__init__(message)
|
|
self.message = message
|
|
self.reason_code = reason_code
|
|
|
|
def to_dict(self) -> Dict[str, Any]:
|
|
return {"error": self.message, "reason_code": self.reason_code}
|
|
|
|
|
|
class InvalidTransitionError(C5AError):
|
|
"""Ungültiger State-Machine-Übergang."""
|
|
|
|
|
|
class InvalidReasonCodeError(C5AError):
|
|
"""Unbekannter Reason Code."""
|
|
|
|
|
|
class InvalidOperationError(C5AError):
|
|
"""Unbekannte Operation."""
|
|
|
|
|
|
class NoWriteGuaranteeError(C5AError):
|
|
"""C5A darf keine externen Writes ausführen."""
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Persistenz (SQLite, restart-fest, atomic, inspectable, backupbar)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
class C5AStore:
|
|
"""
|
|
Persistente Progress-State-Struktur für C5A.
|
|
|
|
SQLite-basiert, isoliert (eigene c5a.db), restart-fest. Keine externe DB,
|
|
kein Trading-/Forgejo-DB-Coupling. Atomic via SQLite-Transaktionen.
|
|
"""
|
|
|
|
SCHEMA_VERSION = 1
|
|
|
|
def __init__(self, db_path: str):
|
|
self.db_path = str(db_path)
|
|
self._conn = sqlite3.connect(self.db_path)
|
|
self._conn.row_factory = sqlite3.Row
|
|
self._conn.execute("PRAGMA journal_mode=WAL;")
|
|
self._init_schema()
|
|
|
|
def _init_schema(self) -> None:
|
|
with self._conn:
|
|
self._conn.execute("""
|
|
CREATE TABLE IF NOT EXISTS meta (
|
|
key TEXT PRIMARY KEY,
|
|
value TEXT NOT NULL
|
|
)
|
|
""")
|
|
self._conn.execute("""
|
|
CREATE TABLE IF NOT EXISTS commits (
|
|
commit_sha TEXT PRIMARY KEY,
|
|
parent_sha TEXT,
|
|
discovered_at INTEGER,
|
|
sequence INTEGER,
|
|
status TEXT NOT NULL,
|
|
retry_count INTEGER NOT NULL DEFAULT 0,
|
|
last_error TEXT,
|
|
last_error_code TEXT,
|
|
created_at INTEGER,
|
|
updated_at INTEGER
|
|
)
|
|
""")
|
|
self._conn.execute("""
|
|
CREATE TABLE IF NOT EXISTS objects (
|
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
commit_sha TEXT NOT NULL,
|
|
object_id TEXT NOT NULL,
|
|
path_before TEXT,
|
|
path_after TEXT,
|
|
operation TEXT NOT NULL,
|
|
content_hash_before TEXT,
|
|
content_hash_after TEXT,
|
|
metadata_hash_before TEXT,
|
|
metadata_hash_after TEXT,
|
|
representation TEXT,
|
|
state TEXT,
|
|
UNIQUE(commit_sha, object_id, operation)
|
|
)
|
|
""")
|
|
self._conn.execute("""
|
|
CREATE TABLE IF NOT EXISTS bootstrap (
|
|
id INTEGER PRIMARY KEY CHECK (id = 1),
|
|
state TEXT NOT NULL,
|
|
baseline_commit TEXT,
|
|
updated_at INTEGER
|
|
)
|
|
""")
|
|
self._conn.execute("""
|
|
CREATE TABLE IF NOT EXISTS health (
|
|
id INTEGER PRIMARY KEY CHECK (id = 1),
|
|
last_seen_commit TEXT,
|
|
last_applied_commit TEXT,
|
|
last_success_at INTEGER,
|
|
last_error_code TEXT,
|
|
updated_at INTEGER
|
|
)
|
|
""")
|
|
# Schema-Version setzen
|
|
self._conn.execute(
|
|
"INSERT OR REPLACE INTO meta (key, value) VALUES ('schema_version', ?)",
|
|
(str(self.SCHEMA_VERSION),),
|
|
)
|
|
# Bootstrap initialisieren
|
|
row = self._conn.execute("SELECT state FROM bootstrap WHERE id = 1").fetchone()
|
|
if row is None:
|
|
self._conn.execute(
|
|
"INSERT INTO bootstrap (id, state, updated_at) VALUES (1, ?, ?)",
|
|
(BS_UNINITIALIZED, int(time.time() * 1000)),
|
|
)
|
|
# Health initialisieren
|
|
row = self._conn.execute("SELECT id FROM health WHERE id = 1").fetchone()
|
|
if row is None:
|
|
self._conn.execute(
|
|
"INSERT INTO health (id, updated_at) VALUES (1, ?)",
|
|
(int(time.time() * 1000),),
|
|
)
|
|
|
|
# -- Bootstrap ----------------------------------------------------------
|
|
|
|
def bootstrap_state(self) -> str:
|
|
row = self._conn.execute("SELECT state FROM bootstrap WHERE id = 1").fetchone()
|
|
return row["state"] if row else BS_UNINITIALIZED
|
|
|
|
def bootstrap_transition(self, to_state: str) -> str:
|
|
"""Deterministischer Bootstrap-Übergang. Gibt den neuen Zustand zurück."""
|
|
if to_state not in BOOTSTRAP_STATES:
|
|
raise InvalidTransitionError(f"Unbekannter Bootstrap-Zustand: {to_state}")
|
|
cur = self.bootstrap_state()
|
|
if (cur, to_state) not in _ALLOWED_BOOTSTRAP:
|
|
raise InvalidTransitionError(
|
|
f"Ungültiger Bootstrap-Übergang: {cur} -> {to_state}"
|
|
)
|
|
with self._conn:
|
|
self._conn.execute(
|
|
"UPDATE bootstrap SET state = ?, updated_at = ? WHERE id = 1",
|
|
(to_state, int(time.time() * 1000)),
|
|
)
|
|
return to_state
|
|
|
|
def set_baseline(self, commit_sha: str) -> None:
|
|
"""Setzt die Baseline (nur aus BASELINE_READY)."""
|
|
if self.bootstrap_state() != BS_BASELINE_READY:
|
|
raise C5AError(
|
|
"Baseline darf nur im Zustand BASELINE_READY gesetzt werden",
|
|
reason_code=RC_OUT_OF_ORDER_COMMIT,
|
|
)
|
|
with self._conn:
|
|
self._conn.execute(
|
|
"UPDATE bootstrap SET baseline_commit = ?, updated_at = ? WHERE id = 1",
|
|
(commit_sha, int(time.time() * 1000)),
|
|
)
|
|
self._conn.execute(
|
|
"UPDATE health SET last_seen_commit = ?, last_applied_commit = ?, updated_at = ? WHERE id = 1",
|
|
(commit_sha, commit_sha, int(time.time() * 1000)),
|
|
)
|
|
|
|
def baseline_commit(self) -> Optional[str]:
|
|
row = self._conn.execute("SELECT baseline_commit FROM bootstrap WHERE id = 1").fetchone()
|
|
return row["baseline_commit"] if row else None
|
|
|
|
# -- Commits ------------------------------------------------------------
|
|
|
|
def upsert_commit(self, commit: Dict[str, Any]) -> Dict[str, Any]:
|
|
"""Legt einen Commit an oder aktualisiert ihn (idempotent per commit_sha)."""
|
|
sha = commit["commit_sha"]
|
|
parent = commit.get("parent_sha")
|
|
discovered = commit.get("discovered_at", int(time.time() * 1000))
|
|
seq = commit.get("sequence")
|
|
status = commit.get("status", ST_DISCOVERED)
|
|
retry = commit.get("retry_count", 0)
|
|
last_err = commit.get("last_error")
|
|
last_code = commit.get("last_error_code")
|
|
now = int(time.time() * 1000)
|
|
with self._conn:
|
|
self._conn.execute(
|
|
"""
|
|
INSERT INTO commits (commit_sha, parent_sha, discovered_at, sequence,
|
|
status, retry_count, last_error, last_error_code,
|
|
created_at, updated_at)
|
|
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
|
ON CONFLICT(commit_sha) DO UPDATE SET
|
|
parent_sha = excluded.parent_sha,
|
|
sequence = excluded.sequence,
|
|
status = excluded.status,
|
|
retry_count = excluded.retry_count,
|
|
last_error = excluded.last_error,
|
|
last_error_code = excluded.last_error_code,
|
|
updated_at = excluded.updated_at
|
|
""",
|
|
(sha, parent, discovered, seq, status, retry, last_err, last_code, now, now),
|
|
)
|
|
return self.get_commit(sha)
|
|
|
|
def get_commit(self, commit_sha: str) -> Optional[Dict[str, Any]]:
|
|
row = self._conn.execute(
|
|
"SELECT * FROM commits WHERE commit_sha = ?", (commit_sha,)
|
|
).fetchone()
|
|
return dict(row) if row else None
|
|
|
|
def commit_status(self, commit_sha: str) -> Optional[str]:
|
|
row = self._conn.execute(
|
|
"SELECT status FROM commits WHERE commit_sha = ?", (commit_sha,)
|
|
).fetchone()
|
|
return row["status"] if row else None
|
|
|
|
def transition_commit(self, commit_sha: str, to_state: str) -> str:
|
|
"""Deterministischer Sync-State-Übergang für einen Commit."""
|
|
if to_state not in SYNC_STATES:
|
|
raise InvalidTransitionError(f"Unbekannter Sync-Zustand: {to_state}")
|
|
cur = self.commit_status(commit_sha)
|
|
if cur is None:
|
|
raise C5AError(f"Commit nicht gefunden: {commit_sha}")
|
|
if (cur, to_state) not in _ALLOWED_TRANSITIONS:
|
|
raise InvalidTransitionError(
|
|
f"Ungültiger Sync-Übergang für {commit_sha}: {cur} -> {to_state}"
|
|
)
|
|
with self._conn:
|
|
self._conn.execute(
|
|
"UPDATE commits SET status = ?, updated_at = ? WHERE commit_sha = ?",
|
|
(to_state, int(time.time() * 1000), commit_sha),
|
|
)
|
|
return to_state
|
|
|
|
def set_commit_error(self, commit_sha: str, reason_code: str, message: str) -> None:
|
|
"""Setzt Fehler + Reason Code auf einen Commit."""
|
|
if reason_code not in REASON_CODES:
|
|
raise InvalidReasonCodeError(f"Unbekannter Reason Code: {reason_code}")
|
|
with self._conn:
|
|
self._conn.execute(
|
|
"UPDATE commits SET last_error = ?, last_error_code = ?, updated_at = ? WHERE commit_sha = ?",
|
|
(message, reason_code, int(time.time() * 1000), commit_sha),
|
|
)
|
|
|
|
def increment_retry(self, commit_sha: str) -> int:
|
|
"""Erhöht retry_count. Gibt den neuen Wert zurück."""
|
|
with self._conn:
|
|
self._conn.execute(
|
|
"UPDATE commits SET retry_count = retry_count + 1, updated_at = ? WHERE commit_sha = ?",
|
|
(int(time.time() * 1000), commit_sha),
|
|
)
|
|
row = self._conn.execute(
|
|
"SELECT retry_count FROM commits WHERE commit_sha = ?", (commit_sha,)
|
|
).fetchone()
|
|
return row["retry_count"] if row else 0
|
|
|
|
def list_commits(self, status: Optional[str] = None) -> List[Dict[str, Any]]:
|
|
if status:
|
|
rows = self._conn.execute(
|
|
"SELECT * FROM commits WHERE status = ? ORDER BY sequence ASC", (status,)
|
|
).fetchall()
|
|
else:
|
|
rows = self._conn.execute(
|
|
"SELECT * FROM commits ORDER BY sequence ASC"
|
|
).fetchall()
|
|
return [dict(r) for r in rows]
|
|
|
|
# -- Objects ------------------------------------------------------------
|
|
|
|
def add_object_change(self, obj: Dict[str, Any]) -> Dict[str, Any]:
|
|
"""Fügt eine Object-Change hinzu (idempotent per commit_sha+object_id+operation)."""
|
|
if obj.get("operation") not in OPERATIONS:
|
|
raise InvalidOperationError(f"Unbekannte Operation: {obj.get('operation')}")
|
|
with self._conn:
|
|
self._conn.execute(
|
|
"""
|
|
INSERT OR IGNORE INTO objects (commit_sha, object_id, path_before, path_after,
|
|
operation, content_hash_before, content_hash_after,
|
|
metadata_hash_before, metadata_hash_after, representation, state)
|
|
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
|
""",
|
|
(
|
|
obj["commit_sha"], obj["object_id"], obj.get("path_before"),
|
|
obj.get("path_after"), obj["operation"],
|
|
obj.get("content_hash_before"), obj.get("content_hash_after"),
|
|
obj.get("metadata_hash_before"), obj.get("metadata_hash_after"),
|
|
obj.get("representation"), obj.get("state"),
|
|
),
|
|
)
|
|
return self.get_object_change(obj["commit_sha"], obj["object_id"], obj["operation"])
|
|
|
|
def get_object_change(self, commit_sha: str, object_id: str, operation: str) -> Optional[Dict[str, Any]]:
|
|
row = self._conn.execute(
|
|
"SELECT * FROM objects WHERE commit_sha = ? AND object_id = ? AND operation = ?",
|
|
(commit_sha, object_id, operation),
|
|
).fetchone()
|
|
return dict(row) if row else None
|
|
|
|
def list_object_changes(self, commit_sha: str) -> List[Dict[str, Any]]:
|
|
rows = self._conn.execute(
|
|
"SELECT * FROM objects WHERE commit_sha = ? ORDER BY id ASC", (commit_sha,)
|
|
).fetchall()
|
|
return [dict(r) for r in rows]
|
|
|
|
# -- Health -------------------------------------------------------------
|
|
|
|
def health(self) -> Dict[str, Any]:
|
|
"""Health Contract (finalisiert)."""
|
|
h = self._conn.execute("SELECT * FROM health WHERE id = 1").fetchone()
|
|
h = dict(h) if h else {}
|
|
commits = self._conn.execute("SELECT status, COUNT(*) as c FROM commits GROUP BY status").fetchall()
|
|
counts = {r["status"]: r["c"] for r in commits}
|
|
pending = sum(counts.get(s, 0) for s in (
|
|
ST_DISCOVERED, ST_VALIDATING, ST_READY, ST_PROPAGATING_TOLARIA,
|
|
ST_VERIFYING_TOLARIA, ST_UPDATING_SEARCH, ST_VERIFYING_SEARCH,
|
|
ST_RETRY_PENDING, ST_WAITING_FOR_PREDECESSOR,
|
|
))
|
|
failed = counts.get(ST_FAILED, 0)
|
|
dead = counts.get(ST_DEAD, 0)
|
|
human = counts.get(ST_HUMAN_REVIEW_REQUIRED, 0)
|
|
drift = self._conn.execute(
|
|
"SELECT COUNT(*) as c FROM objects WHERE operation = ? AND state = ?",
|
|
(OP_CONTENT_UPDATE, "drift"),
|
|
).fetchone()["c"]
|
|
# Status ableiten
|
|
if dead > 0 or failed > 0:
|
|
status = HEALTH_ERROR
|
|
elif pending > 0 or human > 0:
|
|
status = HEALTH_DEGRADED
|
|
else:
|
|
status = HEALTH_OK
|
|
return {
|
|
"status": status,
|
|
"bootstrap_state": self.bootstrap_state(),
|
|
"last_seen_commit": h.get("last_seen_commit"),
|
|
"last_applied_commit": h.get("last_applied_commit"),
|
|
"pending_commits": pending,
|
|
"failed_commits": failed,
|
|
"dead_commits": dead,
|
|
"human_review_required": human,
|
|
"forgejo_status": "UNKNOWN", # C5A führt keine echten Probes aus
|
|
"tolaria_status": "UNKNOWN",
|
|
"search_status": "UNKNOWN",
|
|
"drift_count": drift,
|
|
"last_success_at": h.get("last_success_at"),
|
|
"last_error_code": h.get("last_error_code"),
|
|
}
|
|
|
|
def mark_applied(self, commit_sha: str) -> None:
|
|
"""Fortschreiben von last_applied_commit NUR nach vollständigem PASS."""
|
|
with self._conn:
|
|
self._conn.execute(
|
|
"UPDATE health SET last_applied_commit = ?, last_success_at = ?, updated_at = ? WHERE id = 1",
|
|
(commit_sha, int(time.time() * 1000), int(time.time() * 1000)),
|
|
)
|
|
|
|
def mark_seen(self, commit_sha: str) -> None:
|
|
with self._conn:
|
|
self._conn.execute(
|
|
"UPDATE health SET last_seen_commit = ?, updated_at = ? WHERE id = 1",
|
|
(commit_sha, int(time.time() * 1000)),
|
|
)
|
|
|
|
def close(self) -> None:
|
|
self._conn.close()
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Sync State Machine (deterministische Logik)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
class SyncStateMachine:
|
|
"""
|
|
Deterministische Sync State Machine.
|
|
|
|
Verarbeitet Commits in korrekter Reihenfolge (parent_sha / last_applied_commit),
|
|
mit Idempotenz, Retry/Backoff, Human-Gates und Reason Codes. Führt KEINE
|
|
externen Writes aus — nur Zustandsübergänge in der C5AStore.
|
|
"""
|
|
|
|
def __init__(self, store: C5AStore, max_retries: int = DEFAULT_MAX_RETRIES,
|
|
backoff_seconds: Optional[List[int]] = None):
|
|
self.store = store
|
|
self.max_retries = max_retries
|
|
self.backoff_seconds = backoff_seconds or DEFAULT_BACKOFF_SECONDS
|
|
|
|
# -- Idempotenz ---------------------------------------------------------
|
|
|
|
def idempotency_check(self, commit_sha: str, object_id: str, operation: str,
|
|
content_hash: Optional[str] = None) -> str:
|
|
"""
|
|
Deterministische Idempotenz-Prüfung.
|
|
|
|
Basis: commit_sha + object_id + operation. Zusätzlich Hash-Checks.
|
|
Rückgabe: ALREADY_APPLIED | ALREADY_AT_TARGET | RETRY_SAFE | CONFLICT
|
|
"""
|
|
existing = self.store.get_object_change(commit_sha, object_id, operation)
|
|
if existing is None:
|
|
return IDEM_RETRY_SAFE
|
|
# Commit bereits vollständig angewendet
|
|
if self.store.commit_status(commit_sha) == ST_APPLIED:
|
|
return IDEM_ALREADY_APPLIED
|
|
# Ziel-Hash bereits erreicht (gleicher Inhalt)
|
|
if content_hash is not None:
|
|
if existing.get("content_hash_after") == content_hash:
|
|
return IDEM_ALREADY_AT_TARGET
|
|
# Existiert, aber nicht angewendet und Hash weicht ab -> Konflikt
|
|
return IDEM_CONFLICT
|
|
|
|
# -- Ordering -----------------------------------------------------------
|
|
|
|
def predecessor_applied(self, commit_sha: str) -> bool:
|
|
"""Prüft, ob der Vorgänger (parent_sha) bereits angewendet ist."""
|
|
commit = self.store.get_commit(commit_sha)
|
|
if commit is None:
|
|
return False
|
|
parent = commit.get("parent_sha")
|
|
if not parent:
|
|
return True # Root-Commit, kein Vorgänger
|
|
parent_status = self.store.commit_status(parent)
|
|
return parent_status == ST_APPLIED
|
|
|
|
def check_ordering(self, commit_sha: str) -> str:
|
|
"""
|
|
Prüft die Commit-Reihenfolge. Gibt ST_READY oder ST_WAITING_FOR_PREDECESSOR zurück.
|
|
"""
|
|
if not self.predecessor_applied(commit_sha):
|
|
return ST_WAITING_FOR_PREDECESSOR
|
|
return ST_READY
|
|
|
|
# -- Retry / Backoff ----------------------------------------------------
|
|
|
|
def backoff_for(self, retry_count: int) -> int:
|
|
"""Backoff-Sekunden für den aktuellen Retry-Versuch."""
|
|
if retry_count <= 0:
|
|
return 0
|
|
idx = min(retry_count - 1, len(self.backoff_seconds) - 1)
|
|
return self.backoff_seconds[idx]
|
|
|
|
def retry_available(self, commit_sha: str) -> bool:
|
|
commit = self.store.get_commit(commit_sha)
|
|
if commit is None:
|
|
return False
|
|
return commit.get("retry_count", 0) < self.max_retries
|
|
|
|
# -- Human Gate ---------------------------------------------------------
|
|
|
|
def human_gate(self, commit_sha: str, reason_code: str, message: str) -> Dict[str, Any]:
|
|
"""
|
|
Überführt einen Commit in HUMAN_REVIEW_REQUIRED (FAIL CLOSED).
|
|
Gibt ein Ergebnis-Dict zurück.
|
|
"""
|
|
if reason_code not in REASON_CODES:
|
|
raise InvalidReasonCodeError(f"Unbekannter Reason Code: {reason_code}")
|
|
self.store.set_commit_error(commit_sha, reason_code, message)
|
|
self.store.transition_commit(commit_sha, ST_HUMAN_REVIEW_REQUIRED)
|
|
return {"commit_sha": commit_sha, "status": ST_HUMAN_REVIEW_REQUIRED,
|
|
"reason_code": reason_code}
|
|
|
|
# -- Hauptverarbeitung (deterministisch, NO-WRITE) ----------------------
|
|
|
|
def process_commit(self, commit: Dict[str, Any]) -> Dict[str, Any]:
|
|
"""
|
|
Verarbeitet einen Commit durch die State Machine.
|
|
|
|
NO-WRITE: Diese Methode führt KEINE externen Writes aus. Sie modelliert
|
|
nur die Zustandsübergänge. Die eigentliche Tolaria-Propagation und der
|
|
Search-Rebuild werden in späteren Phasen (C5C/C5D) als injizierte
|
|
Callables ergänzt — C5A ruft sie NICHT auf.
|
|
|
|
Rückgabe: Dict mit commit_sha, status, reason_code (optional), idempotency.
|
|
"""
|
|
sha = commit["commit_sha"]
|
|
# Idempotenz: bereits angewendet?
|
|
if self.store.commit_status(sha) == ST_APPLIED:
|
|
return {"commit_sha": sha, "status": ST_APPLIED, "idempotency": IDEM_ALREADY_APPLIED}
|
|
|
|
# Commit anlegen/aktualisieren
|
|
self.store.upsert_commit(commit)
|
|
|
|
# Ordering prüfen
|
|
order = self.check_ordering(sha)
|
|
if order == ST_WAITING_FOR_PREDECESSOR:
|
|
self.store.transition_commit(sha, ST_WAITING_FOR_PREDECESSOR)
|
|
return {"commit_sha": sha, "status": ST_WAITING_FOR_PREDECESSOR,
|
|
"reason_code": RC_OUT_OF_ORDER_COMMIT}
|
|
|
|
# DISCOVERED -> VALIDATING
|
|
self.store.transition_commit(sha, ST_VALIDATING)
|
|
|
|
# Validierung der Object-Changes
|
|
for obj in commit.get("changed_objects", []):
|
|
op = obj.get("operation")
|
|
if op not in OPERATIONS:
|
|
return self.human_gate(sha, RC_INVALID_SCHEMA,
|
|
f"Unbekannte Operation: {op}")
|
|
if not obj.get("object_id"):
|
|
return self.human_gate(sha, RC_UNKNOWN_OBJECT_ID,
|
|
"object_id fehlt")
|
|
# DELETE_REQUEST -> Human Gate (kein automatisches Hard Delete)
|
|
if op == OP_DELETE_REQUEST:
|
|
return self.human_gate(sha, RC_AMBIGUOUS_DELETE,
|
|
"DELETE_REQUEST erfordert Human Gate")
|
|
# DANGLING_DERIVED_FROM: wenn derived_from gesetzt, muss Ziel existieren
|
|
if obj.get("derived_from") and not self._object_exists(obj["derived_from"]):
|
|
return self.human_gate(sha, RC_DANGLING_DERIVED_FROM,
|
|
f"dangling derived_from: {obj['derived_from']}")
|
|
|
|
# VALIDATING -> READY
|
|
self.store.transition_commit(sha, ST_READY)
|
|
|
|
# READY -> PROPAGATING_TOLARIA (Modellierung; kein echter Write)
|
|
self.store.transition_commit(sha, ST_PROPAGATING_TOLARIA)
|
|
self.store.transition_commit(sha, ST_VERIFYING_TOLARIA)
|
|
|
|
# DRIFT-Prüfung (modelliert): wenn ein Objekt als drift markiert ist -> Human Gate
|
|
for obj in commit.get("changed_objects", []):
|
|
if obj.get("drift"):
|
|
return self.human_gate(sha, RC_UNEXPECTED_TOLARIA_DRIFT,
|
|
f"unexpected drift auf {obj['object_id']}")
|
|
|
|
# VERIFYING_TOLARIA -> UPDATING_SEARCH -> VERIFYING_SEARCH
|
|
self.store.transition_commit(sha, ST_UPDATING_SEARCH)
|
|
self.store.transition_commit(sha, ST_VERIFYING_SEARCH)
|
|
|
|
# VERIFYING_SEARCH -> APPLIED
|
|
self.store.transition_commit(sha, ST_APPLIED)
|
|
self.store.mark_applied(sha)
|
|
|
|
return {"commit_sha": sha, "status": ST_APPLIED, "idempotency": IDEM_RETRY_SAFE}
|
|
|
|
def _object_exists(self, object_id: str) -> bool:
|
|
"""Prüft, ob eine object_id bereits in irgendeinem Commit existiert."""
|
|
row = self.store._conn.execute(
|
|
"SELECT 1 FROM objects WHERE object_id = ? LIMIT 1", (object_id,)
|
|
).fetchone()
|
|
return row is not None
|
|
|
|
# -- Retry-Verarbeitung -------------------------------------------------
|
|
|
|
def retry_commit(self, commit_sha: str, reason_code: str, message: str) -> Dict[str, Any]:
|
|
"""
|
|
Behandelt einen Fehler: Retry, FAILED oder DEAD je nach retry_count.
|
|
"""
|
|
if reason_code not in REASON_CODES:
|
|
raise InvalidReasonCodeError(f"Unbekannter Reason Code: {reason_code}")
|
|
self.store.set_commit_error(commit_sha, reason_code, message)
|
|
retry = self.store.increment_retry(commit_sha)
|
|
if retry >= self.max_retries:
|
|
self.store.transition_commit(commit_sha, ST_DEAD)
|
|
return {"commit_sha": commit_sha, "status": ST_DEAD, "reason_code": reason_code,
|
|
"retry_count": retry}
|
|
self.store.transition_commit(commit_sha, ST_RETRY_PENDING)
|
|
return {"commit_sha": commit_sha, "status": ST_RETRY_PENDING, "reason_code": reason_code,
|
|
"retry_count": retry, "backoff_seconds": self.backoff_for(retry)}
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# No-Write-Guarantee (statische Prüfung)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
def assert_no_write_guarantee() -> Dict[str, Any]:
|
|
"""
|
|
Beweist, dass C5A keine externen Writes ausführen kann.
|
|
|
|
Prüft, dass keine Netzwerk-/HTTP-/Socket-/Subprocess-Mutationsfunktionen
|
|
importiert oder verwendet werden.
|
|
"""
|
|
import ast
|
|
import sys
|
|
this_file = Path(__file__).resolve()
|
|
tree = ast.parse(this_file.read_text(encoding="utf-8"))
|
|
banned_imports = {"requests", "urllib", "http", "socket", "subprocess", "os.system"}
|
|
found = set()
|
|
for node in ast.walk(tree):
|
|
if isinstance(node, ast.Import):
|
|
for alias in node.names:
|
|
root = alias.name.split(".")[0]
|
|
if root in banned_imports:
|
|
found.add(root)
|
|
elif isinstance(node, ast.ImportFrom):
|
|
if node.module:
|
|
root = node.module.split(".")[0]
|
|
if root in banned_imports:
|
|
found.add(root)
|
|
return {
|
|
"no_network_imports": len(found) == 0,
|
|
"banned_imports_found": sorted(found),
|
|
"no_write_guarantee": len(found) == 0,
|
|
}
|