C5F canary commit parked in PROPAGATING_TOLARIA with objects already written to Tolaria (real write before crash). C5EEngine.replay() delegates the pending case to propagate_commit, whose entry guard only accepted READY/RETRY_PENDING, so the commit could never resume past the crash window and the ALREADY_AT_TARGET idempotency (pre_write_drift_check) was never reached. Minimal fix: accept PROPAGATING_TOLARIA as a resume entry state (idempotency still determined per-object via pre_write_drift_check -> no double write; read-back verify() remains the mandatory gate) and skip the READY->PROPAGATING_TOLARIA transition on resume (no self-transition entry exists in _ALLOWED_TRANSITIONS). Adds 2 regression tests covering the crash-window resume (already-at-target and pending-create). Full C5A-E suite: 190 tests, 0 failures.
857 lines
35 KiB
Python
857 lines
35 KiB
Python
#!/usr/bin/env python3
|
|
"""
|
|
Red Queen — C5C Tolaria Propagation & Drift Verification Engine.
|
|
|
|
C5C implementiert ausschliesslich:
|
|
VALIDIERTE OBJECT CHANGES -> TOLARIA PROPAGATION -> READ-BACK
|
|
-> HASH-/METADATA-VERIFY -> DRIFT DETECTION -> STATE-MACHINE UPDATE
|
|
|
|
ARCHITEKTURREGEL (verbindlich):
|
|
* Forgejo bleibt MASTER / Source of Truth.
|
|
* Tolaria ist ausschliesslich DERIVED.
|
|
* C5C darf: Forgejo lesen, Tolaria kontrolliert schreiben.
|
|
* C5C darf NIEMALS: Tolaria -> Forgejo schreiben, Forgejo veraendern,
|
|
IDs neu vergeben, Knowledge Content eigenmaessig umformulieren.
|
|
|
|
NO-SEARCH-GUARANTEE: C5C ruft NICHT /api/search/rebuild auf.
|
|
NO-MASTER-WRITE-GUARANTEE: C5C fuehrt KEIN git push / Forgejo-Write aus.
|
|
NO-PRODUCTION-PROPAGATION: C5C fuehrt KEINE produktive Propagation aus
|
|
(nur Tests gegen Fake/Mock + read-only Live-Dry-Run).
|
|
|
|
Write-Komponente (TolariaClient) ist klar isoliert. Nur C5C-Writer darf den
|
|
Tolaria-Write-Pfad verwenden. Keine generische Agent-Write-Funktion.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import hashlib
|
|
import json
|
|
import os
|
|
import time
|
|
import urllib.error
|
|
import urllib.request
|
|
from pathlib import Path
|
|
from typing import Any, Dict, List, Optional, Tuple
|
|
|
|
# C5A / C5B wiederverwenden (keine konkurrierende State Machine)
|
|
from rq_c5a import (
|
|
C5AStore,
|
|
ST_READY,
|
|
ST_PROPAGATING_TOLARIA,
|
|
ST_VERIFYING_TOLARIA,
|
|
ST_UPDATING_SEARCH,
|
|
ST_HUMAN_REVIEW_REQUIRED,
|
|
ST_RETRY_PENDING,
|
|
ST_DEAD,
|
|
ST_APPLIED,
|
|
IDEM_ALREADY_AT_TARGET,
|
|
IDEM_ALREADY_APPLIED,
|
|
IDEM_RETRY_SAFE,
|
|
IDEM_CONFLICT,
|
|
RC_UNEXPECTED_TOLARIA_DRIFT,
|
|
RC_UNKNOWN_OBJECT_ID,
|
|
RC_INVALID_SCHEMA,
|
|
RC_DANGLING_DERIVED_FROM,
|
|
RC_AMBIGUOUS_DELETE,
|
|
RC_UNKNOWN_LEGACY_OBJECT,
|
|
RC_SECRET_DETECTED,
|
|
RC_TOLARIA_UNAVAILABLE,
|
|
RC_ID_COLLISION,
|
|
RC_AUTH_FAILURE,
|
|
DEFAULT_MAX_RETRIES,
|
|
DEFAULT_BACKOFF_SECONDS,
|
|
OP_CREATE,
|
|
OP_CONTENT_UPDATE,
|
|
OP_METADATA_UPDATE,
|
|
OP_STATE_UPDATE,
|
|
OP_TAGS_UPDATE,
|
|
OP_SOURCE_CANONICAL_RELATION_UPDATE,
|
|
OP_RENAME,
|
|
OP_MOVE,
|
|
OP_DELETE_REQUEST,
|
|
OP_SUPERSEDE,
|
|
)
|
|
from rq_c5b import (
|
|
GitReader,
|
|
parse_frontmatter,
|
|
extract_object_id,
|
|
content_hash,
|
|
metadata_hash,
|
|
detect_secret,
|
|
SCOPE_LEGACY_SPECIAL,
|
|
)
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Konstanten
|
|
# ---------------------------------------------------------------------------
|
|
|
|
# Tolaria-API-Basis (Default; per Env C5_TOLARIA_BASE ueberschreibbar)
|
|
DEFAULT_TOLARIA_BASE = "http://187.124.31.123:5173/api/vault"
|
|
ENV_TOLARIA_BASE = "C5_TOLARIA_BASE"
|
|
|
|
# Vault-Pfad-Praefix (verbindlich — ohne Praefix liest /content das Tolaria-eigene
|
|
# README statt des Vault-Objekts -> false drift)
|
|
VAULT_PREFIX = "/app/vault"
|
|
|
|
# Eindeutige, enge Not-Found-Semantik der Tolaria-Vault-API:
|
|
# Ein nicht existierendes Objekt wird als HTTP 400 mit dieser Fehlermeldung
|
|
# zurueckgemeldet. NUR dieser exakte Fall darf von read() als None (Objekt
|
|
# nicht vorhanden) interpretiert werden — kein pauschales 400-Schlucken.
|
|
TOLARIA_NOT_FOUND_MSG = "Invalid or missing path"
|
|
|
|
# Pre-Write-Drift-Ergebnisse
|
|
DRIFT_WRITE_ALLOWED = "WRITE_ALLOWED"
|
|
DRIFT_ALREADY_AT_TARGET = "ALREADY_AT_TARGET"
|
|
DRIFT_UNEXPECTED = "UNEXPECTED_TOLARIA_DRIFT"
|
|
|
|
# Propagation-Plan-Klassifikation (Live-Dry-Run, §20)
|
|
PLAN_WOULD_WRITE = "WOULD_WRITE"
|
|
PLAN_ALREADY_AT_TARGET = "ALREADY_AT_TARGET"
|
|
PLAN_HUMAN_REVIEW = "HUMAN_REVIEW"
|
|
PLAN_DRIFT = "DRIFT"
|
|
PLAN_LEGACY_SPECIAL = "LEGACY_SPECIAL"
|
|
|
|
# Operationen, die C5C automatisch propagieren darf (kein Human Gate)
|
|
_AUTO_PROPAGATE_OPS = frozenset({
|
|
OP_CREATE,
|
|
OP_CONTENT_UPDATE,
|
|
OP_METADATA_UPDATE,
|
|
OP_STATE_UPDATE,
|
|
OP_TAGS_UPDATE,
|
|
OP_SOURCE_CANONICAL_RELATION_UPDATE,
|
|
OP_SUPERSEDE,
|
|
})
|
|
|
|
# Operationen, die IMMER Human Gate erfordern (C5C darf sie nie automatisch anwenden)
|
|
_HUMAN_GATE_OPS = frozenset({
|
|
OP_DELETE_REQUEST, # §11: niemals automatisch
|
|
OP_RENAME, # §10: API-Semantik nicht sicher verifizierbar -> Human Gate
|
|
OP_MOVE, # §10: dito
|
|
})
|
|
|
|
# Reason Codes, die NICHT retrybar sind (direkt Human Gate / Fail Closed, §17)
|
|
_NON_RETRYABLE_REASONS = frozenset({
|
|
RC_UNEXPECTED_TOLARIA_DRIFT,
|
|
RC_INVALID_SCHEMA,
|
|
RC_UNKNOWN_OBJECT_ID,
|
|
RC_SECRET_DETECTED,
|
|
RC_DANGLING_DERIVED_FROM,
|
|
RC_AMBIGUOUS_DELETE,
|
|
RC_ID_COLLISION,
|
|
RC_UNKNOWN_LEGACY_OBJECT,
|
|
RC_AUTH_FAILURE,
|
|
})
|
|
|
|
|
|
class C5CError(Exception):
|
|
"""Basis-Fehler fuer C5C."""
|
|
|
|
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 TolariaUnavailableError(C5CError):
|
|
"""Tolaria nicht erreichbar / Timeout / transienter HTTP-Fehler (retrybar)."""
|
|
|
|
|
|
class TolariaWriteError(C5CError):
|
|
"""Tolaria-Write fehlgeschlagen (nicht retrybar, z.B. Auth/Schema).
|
|
|
|
Tragt optional den HTTP-Statuscode der Antwort (http_code), damit
|
|
Caller (z.B. read()) den engen 'Objekt existiert nicht'-Fall gezielt
|
|
von anderen 400-/Schema-/Auth-Fehlern unterscheiden koennen.
|
|
"""
|
|
|
|
def __init__(self, message: str, reason_code: Optional[str] = None,
|
|
http_code: Optional[int] = None):
|
|
super().__init__(message, reason_code)
|
|
self.http_code = http_code
|
|
|
|
|
|
class UnexpectedDriftError(C5CError):
|
|
"""Tolaria CURRENT entspricht weder BEFORE noch TARGET (fail closed)."""
|
|
|
|
|
|
class ReadBackMismatchError(C5CError):
|
|
"""Read-Back nach Write stimmt nicht mit dem erwarteten Ziel ueberein."""
|
|
|
|
|
|
class NoSearchGuaranteeError(C5CError):
|
|
"""C5C darf Search-Rebuild nicht aufrufen."""
|
|
|
|
|
|
class NoMasterWriteGuaranteeError(C5CError):
|
|
"""C5C darf nicht nach Forgejo schreiben."""
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Tolaria Client (isoliert: read / write / verify)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
class TolariaClient:
|
|
"""
|
|
Isolierter Tolaria-Client. Nur C5C-Writer darf write() verwenden.
|
|
|
|
read() — POST /content (Read-Back)
|
|
list() — POST /list (Metadaten)
|
|
write() — POST /save (NUR C5C-Writer; isolierte Write-Komponente)
|
|
verify() — Read-Back + Hash-/Metadata-Vergleich
|
|
|
|
Keine generische Agent-Write-Funktion. Kein Search-Rebuild-Aufruf.
|
|
"""
|
|
|
|
def __init__(self, base_url: Optional[str] = None, timeout: float = 15.0):
|
|
self.base_url = (base_url or os.environ.get(ENV_TOLARIA_BASE)
|
|
or DEFAULT_TOLARIA_BASE).rstrip("/")
|
|
self.timeout = timeout
|
|
|
|
# -- HTTP-Helfer --------------------------------------------------------
|
|
|
|
def _post(self, endpoint: str, payload: Dict[str, Any]) -> Dict[str, Any]:
|
|
url = f"{self.base_url}/{endpoint.lstrip('/')}"
|
|
data = json.dumps(payload).encode("utf-8")
|
|
req = urllib.request.Request(
|
|
url, data=data, headers={"Content-Type": "application/json"},
|
|
method="POST",
|
|
)
|
|
try:
|
|
with urllib.request.urlopen(req, timeout=self.timeout) as resp:
|
|
body = resp.read().decode("utf-8")
|
|
return json.loads(body) if body else {}
|
|
except urllib.error.HTTPError as e:
|
|
# 4xx/5xx: transient (5xx) vs. nicht-retrybar (4xx Auth/Schema)
|
|
if e.code >= 500:
|
|
raise TolariaUnavailableError(
|
|
f"Tolaria HTTP {e.code} auf {endpoint}", RC_TOLARIA_UNAVAILABLE)
|
|
raise TolariaWriteError(
|
|
f"Tolaria HTTP {e.code} auf {endpoint}: {e.read().decode('utf-8', 'replace')[:200]}",
|
|
RC_AUTH_FAILURE if e.code in (401, 403) else RC_INVALID_SCHEMA,
|
|
http_code=e.code)
|
|
except (urllib.error.URLError, TimeoutError, OSError) as e:
|
|
raise TolariaUnavailableError(
|
|
f"Tolaria nicht erreichbar ({endpoint}): {e}", RC_TOLARIA_UNAVAILABLE)
|
|
|
|
# -- Read ---------------------------------------------------------------
|
|
|
|
def read(self, vault_path: str) -> Optional[str]:
|
|
"""Read-Back eines Vault-Objekts. Gibt Inhalt oder None (nicht vorhanden)."""
|
|
try:
|
|
resp = self._post("content", {"path": vault_path})
|
|
except TolariaWriteError as e:
|
|
# Enger Not-Found-Fall: HTTP 400 + eindeutige 'Invalid or missing
|
|
# path'-Semantik -> Objekt existiert nicht -> None (kein Fehler).
|
|
# Alle anderen Fehler (andere 400, Auth 401/403, ...) bleiben fail-closed.
|
|
if e.http_code == 400 and TOLARIA_NOT_FOUND_MSG in (e.message or ""):
|
|
return None
|
|
raise
|
|
if "error" in resp:
|
|
return None # Invalid or missing path -> Objekt nicht vorhanden
|
|
return resp.get("content")
|
|
|
|
def list(self, vault_path: str = VAULT_PREFIX) -> List[Dict[str, Any]]:
|
|
"""Listet Vault-Metadaten (read-only)."""
|
|
resp = self._post("list", {"path": vault_path})
|
|
if isinstance(resp, list):
|
|
return resp
|
|
return []
|
|
|
|
# -- Write (NUR C5C-Writer) --------------------------------------------
|
|
|
|
def write(self, vault_path: str, content: str) -> Dict[str, Any]:
|
|
"""Schreibt ein Vault-Objekt (POST /save). Isolierte Write-Komponente."""
|
|
resp = self._post("save", {"path": vault_path, "content": content})
|
|
# HTTP 2xx + JSON null / leerer Payload -> erfolgreicher Transport.
|
|
# Nur hier (write /save) wird dieser Fall als Erfolg gewertet; der
|
|
# nachgelagerte Read-Back (verify) bleibt zwingend für einen C5C-Step.
|
|
if resp is None:
|
|
resp = {}
|
|
if "error" in resp:
|
|
raise TolariaWriteError(
|
|
f"Tolaria save fehlgeschlagen: {resp['error']}", RC_INVALID_SCHEMA)
|
|
return resp
|
|
|
|
# -- Verify (Read-Back + Hash) ------------------------------------------
|
|
|
|
def verify(self, vault_path: str, expected_content: str) -> Dict[str, Any]:
|
|
"""
|
|
Read-Back + Hash-/Metadata-Verify nach einem Write.
|
|
|
|
Vergleicht: path, content_hash, metadata_hash, representation, state,
|
|
derived_from, tags. Kein Write gilt ohne Read-Back als erfolgreich.
|
|
"""
|
|
rb = self.read(vault_path)
|
|
if rb is None:
|
|
raise ReadBackMismatchError(
|
|
f"Read-Back lieferte kein Objekt fuer {vault_path}",
|
|
RC_UNEXPECTED_TOLARIA_DRIFT)
|
|
fm_rb, body_rb = parse_frontmatter(rb)
|
|
fm_exp, body_exp = parse_frontmatter(expected_content)
|
|
|
|
content_ok = content_hash(body_rb) == content_hash(body_exp)
|
|
metadata_ok = metadata_hash(fm_rb) == metadata_hash(fm_exp)
|
|
|
|
# Feld-fuer-Feld-Vergleich (fuer Diagnose)
|
|
fields = {
|
|
"object_id": (fm_rb.get("id"), fm_exp.get("id")),
|
|
"knowledge_schema": (fm_rb.get("knowledge_schema"), fm_exp.get("knowledge_schema")),
|
|
"type": (fm_rb.get("type"), fm_exp.get("type")),
|
|
"role": (fm_rb.get("role"), fm_exp.get("role")),
|
|
"representation": (fm_rb.get("representation"), fm_exp.get("representation")),
|
|
"state": (fm_rb.get("state"), fm_exp.get("state")),
|
|
"derived_from": (fm_rb.get("derived_from"), fm_exp.get("derived_from")),
|
|
"tags": (fm_rb.get("tags"), fm_exp.get("tags")),
|
|
}
|
|
mismatches = {k: v for k, v in fields.items() if v[0] != v[1]}
|
|
|
|
ok = content_ok and metadata_ok and not mismatches
|
|
return {
|
|
"ok": ok,
|
|
"path": vault_path,
|
|
"content_hash_match": content_ok,
|
|
"metadata_hash_match": metadata_ok,
|
|
"field_mismatches": mismatches,
|
|
}
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Pre-Write Drift Check (harte Invariante, §5)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
def pre_write_drift_check(
|
|
client: TolariaClient,
|
|
vault_path: str,
|
|
before_content: Optional[str],
|
|
target_content: str,
|
|
) -> Tuple[str, Optional[str]]:
|
|
"""
|
|
Vor JEDEM Tolaria-Write: Forgejo TARGET gegen Tolaria CURRENT vergleichen.
|
|
|
|
* Tolaria CURRENT == erwarteter BEFORE-Zustand -> WRITE_ALLOWED
|
|
* Tolaria CURRENT == TARGET exakt -> ALREADY_AT_TARGET (kein Write)
|
|
* Tolaria CURRENT weder BEFORE noch TARGET -> UNEXPECTED_TOLARIA_DRIFT (FAIL CLOSED)
|
|
|
|
Rueckgabe: (DRIFT_*, current_content)
|
|
"""
|
|
current = client.read(vault_path)
|
|
|
|
# Tolaria CURRENT == TARGET -> bereits am Ziel, kein Write
|
|
if current is not None and _content_equal(current, target_content):
|
|
return DRIFT_ALREADY_AT_TARGET, current
|
|
|
|
# Tolaria CURRENT == erwarteter BEFORE -> Write erlaubt
|
|
if before_content is not None and current is not None and _content_equal(current, before_content):
|
|
return DRIFT_WRITE_ALLOWED, current
|
|
|
|
# CREATE-Fall: kein BEFORE, Tolaria hat das Objekt noch nicht -> Write erlaubt
|
|
if before_content is None and current is None:
|
|
return DRIFT_WRITE_ALLOWED, current
|
|
|
|
# Alles andere -> UNEXPECTED_TOLARIA_DRIFT (fail closed, kein Ueberschreiben)
|
|
return DRIFT_UNEXPECTED, current
|
|
|
|
|
|
def _content_equal(a: str, b: str) -> bool:
|
|
"""Byte-genauer Vergleich (keine Transformation, kein LLM Rewrite)."""
|
|
return a == b
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Propagation-Modelle (pro Operation)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
def _vault_path(rel_path: str) -> str:
|
|
"""Konvertiert einen Repo-relativen Pfad in einen Vault-Pfad."""
|
|
return f"{VAULT_PREFIX}/{rel_path.lstrip('/')}"
|
|
|
|
|
|
def _read_forgejo_after(reader: GitReader, sha: str, path_after: Optional[str]) -> Optional[str]:
|
|
"""Liest den exakten Forgejo-AFTER-Inhalt (keine Transformation)."""
|
|
if not path_after:
|
|
return None
|
|
return reader.file_content(sha, path_after)
|
|
|
|
|
|
def _read_forgejo_before(reader: GitReader, parent_sha: Optional[str], path_before: Optional[str]) -> Optional[str]:
|
|
"""Liest den Forgejo-BEFORE-Inhalt (fuer Pre-Write-Drift-Check)."""
|
|
if not path_before or not parent_sha:
|
|
return None
|
|
return reader.file_content(parent_sha, path_before)
|
|
|
|
|
|
def _validate_create(obj: Dict[str, Any], after_content: str) -> Optional[str]:
|
|
"""§6 CREATE-Validierung. Gibt reason_code zurueck oder None (gueltig)."""
|
|
oid = obj.get("object_id")
|
|
if not oid or not oid.startswith("object/"):
|
|
return RC_UNKNOWN_OBJECT_ID
|
|
fm, _ = parse_frontmatter(after_content)
|
|
if not fm.get("knowledge_schema"):
|
|
return RC_INVALID_SCHEMA
|
|
if detect_secret(after_content):
|
|
return RC_SECRET_DETECTED
|
|
derived = fm.get("derived_from")
|
|
if derived and not derived.startswith("object/"):
|
|
return RC_DANGLING_DERIVED_FROM
|
|
return None
|
|
|
|
|
|
def _validate_derived_from(fm: Dict[str, Any], existing_ids: set) -> Optional[str]:
|
|
"""§9: Keine dangling derived_from. Ziel muss produktiv existieren."""
|
|
derived = fm.get("derived_from")
|
|
if derived and derived not in existing_ids:
|
|
return RC_DANGLING_DERIVED_FROM
|
|
return None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# C5C Propagation Engine
|
|
# ---------------------------------------------------------------------------
|
|
|
|
class C5CPropagator:
|
|
"""
|
|
Propagiert validierte ObjectChanges eines Forgejo-Commits nach Tolaria.
|
|
|
|
Commit-Atomicity (§15): Ein Commit ist die Processing Unit. Alle Tolaria-Steps
|
|
muessen PASS sein, bevor der Commit zum naechsten C5-Schritt darf. Wenn ein
|
|
Objekt fehlschlaegt, wird der Commit NICHT als Tolaria-complete markiert.
|
|
|
|
C5C darf States bewegen: READY -> PROPAGATING_TOLARIA -> VERIFYING_TOLARIA
|
|
-> UPDATING_SEARCH (READY_FOR_SEARCH). C5C darf NICHT VERIFYING_SEARCH ->
|
|
APPLIED durchfuehren. last_applied_commit bleibt unveraendert.
|
|
"""
|
|
|
|
def __init__(
|
|
self,
|
|
store: C5AStore,
|
|
reader: GitReader,
|
|
client: Optional[TolariaClient] = None,
|
|
max_retries: int = DEFAULT_MAX_RETRIES,
|
|
backoff_seconds: Optional[List[int]] = None,
|
|
):
|
|
self.store = store
|
|
self.reader = reader
|
|
self.client = client or TolariaClient()
|
|
self.max_retries = max_retries
|
|
self.backoff_seconds = backoff_seconds or DEFAULT_BACKOFF_SECONDS
|
|
|
|
# -- Idempotenz ---------------------------------------------------------
|
|
|
|
def _commit_already_propagated(self, commit_sha: str) -> bool:
|
|
"""Commit wurde bereits vollstaendig propagiert (Tolaria-complete)."""
|
|
return self.store.commit_status(commit_sha) in (
|
|
ST_UPDATING_SEARCH, ST_APPLIED)
|
|
|
|
# -- Retry-Klassifikation (§17) -----------------------------------------
|
|
|
|
def _is_retryable(self, err: C5CError) -> bool:
|
|
"""Nur technische, retryable Fehler nutzen das Retry-Modell."""
|
|
if isinstance(err, TolariaUnavailableError):
|
|
return True
|
|
if err.reason_code in _NON_RETRYABLE_REASONS:
|
|
return False
|
|
# Default: nicht-retrybar (fail closed)
|
|
return False
|
|
|
|
def _handle_failure(self, commit_sha: str, err: C5CError) -> Dict[str, Any]:
|
|
"""Behandelt einen Fehler: Retry, DEAD oder Human Gate."""
|
|
if self._is_retryable(err):
|
|
retry = self.store.increment_retry(commit_sha)
|
|
if retry >= self.max_retries:
|
|
self.store.set_commit_error(commit_sha, err.reason_code or RC_TOLARIA_UNAVAILABLE, err.message)
|
|
self.store.transition_commit(commit_sha, ST_DEAD)
|
|
return {"commit_sha": commit_sha, "status": ST_DEAD,
|
|
"reason_code": err.reason_code, "retry_count": retry}
|
|
self.store.set_commit_error(commit_sha, err.reason_code or RC_TOLARIA_UNAVAILABLE, err.message)
|
|
self.store.transition_commit(commit_sha, ST_RETRY_PENDING)
|
|
return {"commit_sha": commit_sha, "status": ST_RETRY_PENDING,
|
|
"reason_code": err.reason_code, "retry_count": retry}
|
|
# Nicht-retrybar -> Human Gate / Fail Closed
|
|
rc = err.reason_code or RC_UNEXPECTED_TOLARIA_DRIFT
|
|
self.store.set_commit_error(commit_sha, rc, err.message)
|
|
self.store.transition_commit(commit_sha, ST_HUMAN_REVIEW_REQUIRED)
|
|
return {"commit_sha": commit_sha, "status": ST_HUMAN_REVIEW_REQUIRED,
|
|
"reason_code": rc}
|
|
|
|
# -- Einzelnes Objekt propagieren ---------------------------------------
|
|
|
|
def _propagate_object(self, commit_sha: str, parent_sha: Optional[str],
|
|
obj: Dict[str, Any], existing_ids: set) -> Dict[str, Any]:
|
|
"""
|
|
Propagiert EIN ObjectChange nach Tolaria. Gibt ein Ergebnis-Dict zurueck.
|
|
|
|
Rueckgabe-Felder: object_id, operation, status, reason_code (optional),
|
|
drift (optional), idempotency (optional).
|
|
"""
|
|
oid = obj.get("object_id")
|
|
op = obj.get("operation")
|
|
|
|
# Human-Gate-Operationen (DELETE/RENAME/MOVE) -> niemals automatisch
|
|
if op in _HUMAN_GATE_OPS:
|
|
rc = RC_AMBIGUOUS_DELETE if op == OP_DELETE_REQUEST else RC_UNKNOWN_OBJECT_ID
|
|
return {"object_id": oid, "operation": op, "status": ST_HUMAN_REVIEW_REQUIRED,
|
|
"reason_code": rc}
|
|
|
|
# Nicht-auto-propagierbare Operation -> Human Gate
|
|
if op not in _AUTO_PROPAGATE_OPS:
|
|
return {"object_id": oid, "operation": op, "status": ST_HUMAN_REVIEW_REQUIRED,
|
|
"reason_code": RC_INVALID_SCHEMA}
|
|
|
|
# object_id muss vorhanden sein
|
|
if not oid:
|
|
return {"object_id": oid, "operation": op, "status": ST_HUMAN_REVIEW_REQUIRED,
|
|
"reason_code": RC_UNKNOWN_OBJECT_ID}
|
|
|
|
path_after = obj.get("path_after")
|
|
path_before = obj.get("path_before")
|
|
if not path_after:
|
|
return {"object_id": oid, "operation": op, "status": ST_HUMAN_REVIEW_REQUIRED,
|
|
"reason_code": RC_UNKNOWN_OBJECT_ID}
|
|
|
|
# Forgejo-AFTER-Inhalt lesen (exakt, keine Transformation)
|
|
after_content = _read_forgejo_after(self.reader, commit_sha, path_after)
|
|
if after_content is None:
|
|
return {"object_id": oid, "operation": op, "status": ST_HUMAN_REVIEW_REQUIRED,
|
|
"reason_code": RC_UNKNOWN_OBJECT_ID}
|
|
|
|
# Secret-Safety (fail-closed VOR Verarbeitung)
|
|
if detect_secret(after_content):
|
|
return {"object_id": oid, "operation": op, "status": ST_HUMAN_REVIEW_REQUIRED,
|
|
"reason_code": RC_SECRET_DETECTED}
|
|
|
|
# CREATE-Validierung (§6)
|
|
if op == OP_CREATE:
|
|
rc = _validate_create(obj, after_content)
|
|
if rc:
|
|
return {"object_id": oid, "operation": op, "status": ST_HUMAN_REVIEW_REQUIRED,
|
|
"reason_code": rc}
|
|
|
|
# derived_from-Validierung (§9): keine dangling relation
|
|
fm_after, _ = parse_frontmatter(after_content)
|
|
rc = _validate_derived_from(fm_after, existing_ids)
|
|
if rc:
|
|
return {"object_id": oid, "operation": op, "status": ST_HUMAN_REVIEW_REQUIRED,
|
|
"reason_code": rc}
|
|
|
|
vault_path = _vault_path(path_after)
|
|
before_content = _read_forgejo_before(self.reader, parent_sha, path_before)
|
|
|
|
# Pre-Write Drift Check (§5) — harte Invariante.
|
|
# Idempotenz wird HIER bestimmt: Tolaria CURRENT == TARGET -> ALREADY_AT_TARGET
|
|
# (kein Write). Tolaria CURRENT == BEFORE -> WRITE_ALLOWED.
|
|
drift, current = pre_write_drift_check(
|
|
self.client, vault_path, before_content, after_content)
|
|
if drift == DRIFT_ALREADY_AT_TARGET:
|
|
return {"object_id": oid, "operation": op, "status": ST_VERIFYING_TOLARIA,
|
|
"idempotency": IDEM_ALREADY_AT_TARGET}
|
|
if drift == DRIFT_UNEXPECTED:
|
|
return {"object_id": oid, "operation": op, "status": ST_HUMAN_REVIEW_REQUIRED,
|
|
"reason_code": RC_UNEXPECTED_TOLARIA_DRIFT, "drift": True}
|
|
|
|
# WRITE_ALLOWED -> Write (nur C5C-Writer)
|
|
self.client.write(vault_path, after_content)
|
|
|
|
# Read-Back + Verify (§14) — kein Write gilt ohne Read-Back als erfolgreich
|
|
verify = self.client.verify(vault_path, after_content)
|
|
if not verify["ok"]:
|
|
raise ReadBackMismatchError(
|
|
f"Read-Back-Mismatch auf {vault_path}: {verify['field_mismatches']}",
|
|
RC_UNEXPECTED_TOLARIA_DRIFT)
|
|
|
|
return {"object_id": oid, "operation": op, "status": ST_VERIFYING_TOLARIA,
|
|
"idempotency": IDEM_RETRY_SAFE}
|
|
|
|
# -- Commit propagieren (Atomicity, §15) --------------------------------
|
|
|
|
def propagate_commit(self, commit_sha: str) -> Dict[str, Any]:
|
|
"""
|
|
Propagiert alle ObjectChanges eines Commits nach Tolaria.
|
|
|
|
Commit-Atomicity: Alle Objekte muessen PASS sein, bevor der Commit zum
|
|
naechsten C5-Schritt (UPDATING_SEARCH / READY_FOR_SEARCH) darf. Wenn ein
|
|
Objekt fehlschlaegt, wird der Commit NICHT als Tolaria-complete markiert.
|
|
|
|
C5C bewegt: READY -> PROPAGATING_TOLARIA -> VERIFYING_TOLARIA
|
|
-> UPDATING_SEARCH (READY_FOR_SEARCH). NICHT weiter zu APPLIED.
|
|
last_applied_commit bleibt unveraendert.
|
|
"""
|
|
commit = self.store.get_commit(commit_sha)
|
|
if commit is None:
|
|
return {"commit_sha": commit_sha, "status": ST_HUMAN_REVIEW_REQUIRED,
|
|
"reason_code": RC_UNKNOWN_OBJECT_ID}
|
|
|
|
cur = commit.get("status")
|
|
# Resume-Einstieg: READY / RETRY_PENDING (Erst-/Retry-Propagation) UND
|
|
# PROPAGATING_TOLARIA (Crash-Window-Resume: Commit steckte waehrend der
|
|
# Propagation fest, Objekte evtl. teilweise propagiert). Idempotenz wird
|
|
# pro Objekt via pre_write_drift_check (ALREADY_AT_TARGET) bestimmt,
|
|
# kein Doppel-Write. Read-Back bleibt zwingend (verify / content_equal).
|
|
if cur not in (ST_READY, ST_RETRY_PENDING, ST_PROPAGATING_TOLARIA):
|
|
return {"commit_sha": commit_sha, "status": cur,
|
|
"message": "Commit nicht im propagierbaren Zustand"}
|
|
|
|
# READY -> PROPAGATING_TOLARIA (nur beim Erst-/Retry-Einstieg; beim
|
|
# PROPAGATING_TOLARIA-Resume ist der State bereits gesetzt, kein
|
|
# Self-Transition-Entry noetig, der nicht in _ALLOWED_TRANSITIONS ist)
|
|
if cur != ST_PROPAGATING_TOLARIA:
|
|
self.store.transition_commit(commit_sha, ST_PROPAGATING_TOLARIA)
|
|
|
|
parent_sha = commit.get("parent_sha")
|
|
objs = self.store.list_object_changes(commit_sha)
|
|
|
|
# Existierende IDs sammeln (fuer derived_from-Validierung, §9)
|
|
existing_ids = self._collect_existing_ids(commit_sha)
|
|
|
|
results = []
|
|
all_pass = True
|
|
for obj in objs:
|
|
try:
|
|
r = self._propagate_object(commit_sha, parent_sha, obj, existing_ids)
|
|
except C5CError as e:
|
|
r = self._handle_failure(commit_sha, e)
|
|
all_pass = False
|
|
results.append(r)
|
|
break # Commit-Atomicity: bei Fehler Commit nicht als complete markieren
|
|
results.append(r)
|
|
if r.get("status") == ST_HUMAN_REVIEW_REQUIRED:
|
|
all_pass = False
|
|
# Human Gate -> Commit-Status setzen (fail closed)
|
|
rc = r.get("reason_code") or RC_UNEXPECTED_TOLARIA_DRIFT
|
|
self.store.set_commit_error(commit_sha, rc, f"ObjectChange {r.get('object_id')} {r.get('operation')}")
|
|
self.store.transition_commit(commit_sha, ST_HUMAN_REVIEW_REQUIRED)
|
|
break # Commit nicht als complete markieren
|
|
|
|
if not all_pass:
|
|
# Commit bleibt in PROPAGATING_TOLARIA oder wird Human Gate / Retry
|
|
# (transition_commit wurde bereits im Fehler-Handler gesetzt)
|
|
return {"commit_sha": commit_sha, "status": self.store.commit_status(commit_sha),
|
|
"results": results}
|
|
|
|
# Alle Objekte PASS -> VERIFYING_TOLARIA
|
|
self.store.transition_commit(commit_sha, ST_VERIFYING_TOLARIA)
|
|
|
|
# VERIFYING_TOLARIA -> UPDATING_SEARCH (READY_FOR_SEARCH)
|
|
# C5C darf NICHT weiter zu VERIFYING_SEARCH -> APPLIED.
|
|
self.store.transition_commit(commit_sha, ST_UPDATING_SEARCH)
|
|
|
|
return {"commit_sha": commit_sha, "status": ST_UPDATING_SEARCH,
|
|
"results": results, "ready_for_search": True}
|
|
|
|
def _collect_existing_ids(self, commit_sha: str) -> set:
|
|
"""Sammelt alle bereits produktiv existierenden object_ids (fuer §9)."""
|
|
ids = set()
|
|
for c in self.store.list_commits():
|
|
for oc in self.store.list_object_changes(c["commit_sha"]):
|
|
oid = oc.get("object_id")
|
|
if oid:
|
|
ids.add(oid)
|
|
return ids
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Live-Dry-Run (read-only, §20) — erzeugt Propagation Plan, KEINE Writes
|
|
# ---------------------------------------------------------------------------
|
|
|
|
class C5CDryRun:
|
|
"""
|
|
Read-only Live-Dry-Run gegen den realen Forgejo + Tolaria Stand.
|
|
|
|
Erzeugt einen Propagation Plan:
|
|
WOULD_WRITE | ALREADY_AT_TARGET | HUMAN_REVIEW | DRIFT | LEGACY_SPECIAL
|
|
|
|
Fuehrt KEINE produktive Propagation aus. KEINE Writes.
|
|
"""
|
|
|
|
def __init__(self, store: C5AStore, reader: GitReader, client: Optional[TolariaClient] = None):
|
|
self.store = store
|
|
self.reader = reader
|
|
self.client = client or TolariaClient()
|
|
|
|
def plan_commit(self, commit_sha: str) -> Dict[str, Any]:
|
|
"""Erzeugt den Propagation-Plan fuer einen Commit (read-only)."""
|
|
commit = self.store.get_commit(commit_sha)
|
|
if commit is None:
|
|
return {"commit_sha": commit_sha, "error": "Commit nicht gefunden"}
|
|
parent_sha = commit.get("parent_sha")
|
|
objs = self.store.list_object_changes(commit_sha)
|
|
existing_ids = self._collect_existing_ids()
|
|
|
|
plan = []
|
|
for obj in objs:
|
|
plan.append(self._plan_object(commit_sha, parent_sha, obj, existing_ids))
|
|
|
|
return {"commit_sha": commit_sha, "plan": plan}
|
|
|
|
def _plan_object(self, commit_sha: str, parent_sha: Optional[str],
|
|
obj: Dict[str, Any], existing_ids: set) -> Dict[str, Any]:
|
|
oid = obj.get("object_id")
|
|
op = obj.get("operation")
|
|
|
|
# Human-Gate-Operationen
|
|
if op in _HUMAN_GATE_OPS:
|
|
return {"object_id": oid, "operation": op, "plan": PLAN_HUMAN_REVIEW,
|
|
"reason_code": RC_AMBIGUOUS_DELETE if op == OP_DELETE_REQUEST else RC_UNKNOWN_OBJECT_ID}
|
|
if op not in _AUTO_PROPAGATE_OPS:
|
|
return {"object_id": oid, "operation": op, "plan": PLAN_HUMAN_REVIEW,
|
|
"reason_code": RC_INVALID_SCHEMA}
|
|
if not oid:
|
|
return {"object_id": oid, "operation": op, "plan": PLAN_HUMAN_REVIEW,
|
|
"reason_code": RC_UNKNOWN_OBJECT_ID}
|
|
|
|
path_after = obj.get("path_after")
|
|
if not path_after:
|
|
return {"object_id": oid, "operation": op, "plan": PLAN_HUMAN_REVIEW,
|
|
"reason_code": RC_UNKNOWN_OBJECT_ID}
|
|
|
|
after_content = _read_forgejo_after(self.reader, commit_sha, path_after)
|
|
if after_content is None:
|
|
return {"object_id": oid, "operation": op, "plan": PLAN_HUMAN_REVIEW,
|
|
"reason_code": RC_UNKNOWN_OBJECT_ID}
|
|
if detect_secret(after_content):
|
|
return {"object_id": oid, "operation": op, "plan": PLAN_HUMAN_REVIEW,
|
|
"reason_code": RC_SECRET_DETECTED}
|
|
if op == OP_CREATE:
|
|
rc = _validate_create(obj, after_content)
|
|
if rc:
|
|
return {"object_id": oid, "operation": op, "plan": PLAN_HUMAN_REVIEW,
|
|
"reason_code": rc}
|
|
fm_after, _ = parse_frontmatter(after_content)
|
|
rc = _validate_derived_from(fm_after, existing_ids)
|
|
if rc:
|
|
return {"object_id": oid, "operation": op, "plan": PLAN_HUMAN_REVIEW,
|
|
"reason_code": rc}
|
|
|
|
vault_path = _vault_path(path_after)
|
|
before_content = _read_forgejo_before(self.reader, parent_sha, path_before=obj.get("path_before"))
|
|
|
|
# Pre-Write Drift Check (read-only)
|
|
drift, _ = pre_write_drift_check(self.client, vault_path, before_content, after_content)
|
|
if drift == DRIFT_ALREADY_AT_TARGET:
|
|
return {"object_id": oid, "operation": op, "plan": PLAN_ALREADY_AT_TARGET}
|
|
if drift == DRIFT_UNEXPECTED:
|
|
return {"object_id": oid, "operation": op, "plan": PLAN_DRIFT,
|
|
"reason_code": RC_UNEXPECTED_TOLARIA_DRIFT}
|
|
return {"object_id": oid, "operation": op, "plan": PLAN_WOULD_WRITE}
|
|
|
|
def _collect_existing_ids(self) -> set:
|
|
ids = set()
|
|
for c in self.store.list_commits():
|
|
for oc in self.store.list_object_changes(c["commit_sha"]):
|
|
oid = oc.get("object_id")
|
|
if oid:
|
|
ids.add(oid)
|
|
return ids
|
|
|
|
def drift_report(self, commits: List[str]) -> Dict[str, Any]:
|
|
"""Erzeugt den Drift Report (§21) aus dem Live-Dry-Run."""
|
|
total = 0
|
|
counts = {PLAN_WOULD_WRITE: 0, PLAN_ALREADY_AT_TARGET: 0,
|
|
PLAN_HUMAN_REVIEW: 0, PLAN_DRIFT: 0, PLAN_LEGACY_SPECIAL: 0}
|
|
drift_objects = []
|
|
for sha in commits:
|
|
res = self.plan_commit(sha)
|
|
for item in res.get("plan", []):
|
|
total += 1
|
|
p = item.get("plan")
|
|
counts[p] = counts.get(p, 0) + 1
|
|
if p == PLAN_DRIFT:
|
|
drift_objects.append(item)
|
|
return {
|
|
"TOTAL_OBJECTS_CHECKED": total,
|
|
"ALREADY_AT_TARGET": counts[PLAN_ALREADY_AT_TARGET],
|
|
"WOULD_WRITE": counts[PLAN_WOULD_WRITE],
|
|
"UNEXPECTED_DRIFT": counts[PLAN_DRIFT],
|
|
"LEGACY_SPECIAL": counts[PLAN_LEGACY_SPECIAL],
|
|
"HUMAN_REVIEW": counts[PLAN_HUMAN_REVIEW],
|
|
"drift_objects": drift_objects,
|
|
}
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# No-Search / No-Master-Write Guarantee (statische Pruefung)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
def _collect_docstrings(tree: Any) -> set:
|
|
"""Sammelt alle Docstring-Strings (Modul-/Funktions-/Klassen-Docstrings)."""
|
|
import ast
|
|
docstrings = set()
|
|
for node in ast.walk(tree):
|
|
if isinstance(node, (ast.Module, ast.FunctionDef, ast.AsyncFunctionDef,
|
|
ast.ClassDef)):
|
|
body = node.body
|
|
if body and isinstance(body[0], ast.Expr):
|
|
val = body[0].value
|
|
if isinstance(val, ast.Constant) and isinstance(val.value, str):
|
|
docstrings.add(val.value)
|
|
return docstrings
|
|
|
|
|
|
def assert_no_search_calls() -> Dict[str, Any]:
|
|
"""
|
|
Beweist statisch, dass C5C NICHT /api/search/rebuild oder Search-Admin-API aufruft.
|
|
|
|
Prueft, dass keine Search-Endpoint-Strings in echten Code-Ausdruecken
|
|
(nicht Docstrings/Kommentaren) vorkommen.
|
|
"""
|
|
import ast
|
|
this_file = Path(__file__).resolve()
|
|
tree = ast.parse(this_file.read_text(encoding="utf-8"))
|
|
docstrings = _collect_docstrings(tree)
|
|
# Verbotene Endpoints aus Fragmenten zusammensetzen, damit sie nicht als
|
|
# zusammenhaengende wörtliche Strings im Quellcode stehen (Selbstreferenz vermeiden).
|
|
banned = ["/api/" + "search/" + "rebuild", "search/" + "rebuild", "api/" + "search"]
|
|
found = []
|
|
for node in ast.walk(tree):
|
|
if isinstance(node, ast.Constant) and isinstance(node.value, str):
|
|
if node.value in docstrings:
|
|
continue # Docstring, kein echter Code-Ausdruck
|
|
for b in banned:
|
|
if b in node.value:
|
|
found.append(node.value)
|
|
return {
|
|
"no_search_calls": len(found) == 0,
|
|
"search_endpoints_found": sorted(set(found)),
|
|
}
|
|
|
|
|
|
def assert_no_master_write() -> Dict[str, Any]:
|
|
"""
|
|
Beweist statisch, dass C5C KEIN git push / Forgejo-Write ausfuehrt.
|
|
|
|
Prueft, dass keine git-push-/Forgejo-Write-Befehle in echten Code-Ausdruecken
|
|
vorkommen und keine subprocess-Mutationsfunktionen importiert werden.
|
|
"""
|
|
import ast
|
|
this_file = Path(__file__).resolve()
|
|
tree = ast.parse(this_file.read_text(encoding="utf-8"))
|
|
docstrings = _collect_docstrings(tree)
|
|
banned_imports = {"subprocess", "os.system", "git"}
|
|
found_imports = 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_imports.add(root)
|
|
elif isinstance(node, ast.ImportFrom):
|
|
if node.module:
|
|
root = node.module.split(".")[0]
|
|
if root in banned_imports:
|
|
found_imports.add(root)
|
|
banned_calls = ["git " + "pus" + "h", "git " + "comm" + "it", "git " + "ad" + "d", "pus" + "h"]
|
|
found_calls = []
|
|
for node in ast.walk(tree):
|
|
if isinstance(node, ast.Constant) and isinstance(node.value, str):
|
|
if node.value in docstrings:
|
|
continue # Docstring, kein echter Code-Ausdruck
|
|
for b in banned_calls:
|
|
if b in node.value:
|
|
found_calls.append(node.value)
|
|
return {
|
|
"no_master_write": len(found_imports) == 0 and len(found_calls) == 0,
|
|
"banned_imports_found": sorted(found_imports),
|
|
"write_commands_found": sorted(set(found_calls)),
|
|
}
|