trading-system-docs/tolaria/c5-sync-service/rq_c5c.py
Red Queen a11c1bbe53 AUTH.3A: C5 Caller Auth Integration (SAVE/DELETE Credential, fail-closed, Tests, Sensitivity, Gap-Doku)
- rq_c5c.py: TolariaClient save_token/delete_token DI, _require_token fail-closed,
  write()/delete() senden Bearer (SAVE/DELETE), read()/list() ohne Credential
- rq_c5_cli.py: _delete_executor liest nur DELETE-Credential (Least Privilege)
- test_c5c.py: write()-Tests injizieren synthetisches SAVE-Token
- test_c5_auth3a.py: isolierte AUTH.3A-Testsuite (19 Tests, Fake/Mock Tolaria)
- auth3a-sensitivity.sh: 8 Sensitivitaets-Mutationen (A-H) -> ROT
- AUTH3A_HUMAN_APPROVAL_AUTHENTICITY_GAP.md: Gap dokumentiert (OPEN, nicht repariert)

Keine echten Tokens. Keine ENV-Mutation. Kein Deployment. Keine produktive Auth-Aktivierung.
2026-08-27 08:42:26 +00:00

940 lines
39 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,
ST_DELETE_APPROVED,
ST_DELETING,
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"
# AUTH.3A ENV-Contract (PHASE 3): Eindeutige Namen, festgelegt in AUTH.2
# (vite.config.ts: process.env.TOLARIA_SAVE_TOKEN / TOLARIA_DELETE_TOKEN).
# AUTH.3A implementiert NUR die code-seitige Faehigkeit, diese aus der
# Runtime-Konfiguration zu lesen. Es werden KEINE echten ENV-Werte gesetzt.
# Least Privilege: SAVE-Caller liest NUR TOLARIA_SAVE_TOKEN, DELETE-Caller
# liest NUR TOLARIA_DELETE_TOKEN. Kein Prozess laedt beide automatisch.
ENV_TOLARIA_SAVE_TOKEN = "TOLARIA_SAVE_TOKEN"
ENV_TOLARIA_DELETE_TOKEN = "TOLARIA_DELETE_TOKEN"
# 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,
save_token: Optional[str] = None,
delete_token: Optional[str] = None,
):
self.base_url = (base_url or os.environ.get(ENV_TOLARIA_BASE)
or DEFAULT_TOLARIA_BASE).rstrip("/")
self.timeout = timeout
# AUTH.3A: Explizite Credential-Injection (kein verstecktes globales
# Credential, kein Default-Token). SAVE- und DELETE-Credential sind
# strikt getrennt (AUTH.1 §3). Fehlende/leere Werte -> fail-closed
# beim jeweiligen mutierenden Aufruf (AUTH.1 §4).
self.save_token = save_token
self.delete_token = delete_token
# -- HTTP-Helfer --------------------------------------------------------
def _post(
self,
endpoint: str,
payload: Dict[str, Any],
auth_token: Optional[str] = None,
) -> Dict[str, Any]:
url = f"{self.base_url}/{endpoint.lstrip('/')}"
data = json.dumps(payload).encode("utf-8")
headers: Dict[str, str] = {"Content-Type": "application/json"}
# AUTH.3A: Authorization-Header NUR wenn ein Token explizit uebergeben
# wird (mutierende SAVE/DELETE-Aufrufe). READ-Aufrufe senden KEIN
# Credential (Least Privilege, AUTH.1 §2). Der Token-Wert wird nie
# geloggt (AUTH.1 §6).
if auth_token is not None:
headers["Authorization"] = f"Bearer {auth_token}"
req = urllib.request.Request(
url, data=data, headers=headers,
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 _require_token(self, token: Optional[str], scope: str) -> str:
"""Fail-closed: fehlendes/leeres/malformed Credential -> kein HTTP.
AUTH.1 §4: Ein mutierender Request darf NIE ohne gültiges Credential
gesendet werden. Fehlendes oder leeres Token -> lokaler Abbruch
(TolariaWriteError, RC_AUTH_FAILURE), HTTP wird NICHT aufgerufen.
"""
if not token or not isinstance(token, str) or not token.strip():
raise TolariaWriteError(
f"Tolaria {scope}-Credential fehlt oder ist leer "
f"(fail-closed, kein Request gesendet)",
RC_AUTH_FAILURE,
)
return token
def write(self, vault_path: str, content: str) -> Dict[str, Any]:
"""Schreibt ein Vault-Objekt (POST /save). Isolierte Write-Komponente.
AUTH.3A: SAVE-Scope. Sendet Authorization: Bearer <save_token>.
Fail-closed: fehlendes/leeres SAVE-Credential -> lokaler Abbruch,
HTTP wird NICHT aufgerufen (AUTH.1 §4).
"""
token = self._require_token(self.save_token, "SAVE")
resp = self._post("save", {"path": vault_path, "content": content},
auth_token=token)
# 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
def delete(self, vault_path: str) -> Dict[str, Any]:
"""
Loescht ein Vault-Objekt (POST /delete).
NUR durch den DeleteExecutor (C5 DELETE-Execution) aufrufbar, der zuvor
eine persistierte, exakt passende Human-Approval validiert hat. Kein
freier Pfad-Parameter; vault_path stammt ausschliesslich aus dem
ObjectChange, das durch die Approval gebunden ist (Invariante I/J).
Endpoint-/Payload-Contract (POST /delete mit {"path": ...}) ist aus dem
vorhandenen Tolaria-Code/Router und der dokumentierten Incident-Evidence
bestimmt. KEINE produktive Endpoint-Probe (Incident-Regel).
AUTH.3A: DELETE-Scope. Sendet Authorization: Bearer <delete_token>.
Fail-closed: fehlendes/leeres DELETE-Credential -> lokaler Abbruch,
HTTP wird NICHT aufgerufen (AUTH.1 §4).
"""
token = self._require_token(self.delete_token, "DELETE")
resp = self._post("delete", {"path": vault_path}, auth_token=token)
if resp is None:
resp = {}
if "error" in resp:
raise TolariaWriteError(
f"Tolaria delete 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)),
}