From 2ba503df0eb5878a0883a816684350d94f7498b6 Mon Sep 17 00:00:00 2001 From: Red Queen Date: Thu, 27 Aug 2026 13:23:43 +0000 Subject: [PATCH] =?UTF-8?q?AUTH.4C1:=20SoT=20reconciliation=20=E2=80=94=20?= =?UTF-8?q?runtime=20wiring=20(entrypoint,=20worker,=20tolaria=5Fclient,?= =?UTF-8?q?=20forgejo=5Fsource=5Floader,=20Dockerfile)=20+=20Read-Back-Ver?= =?UTF-8?q?ifikation=20+=20SHA1-Korrektur?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- tolaria/c5-sync-service/Dockerfile | 25 +++ tolaria/c5-sync-service/entrypoint.py | 47 ++++ .../c5-sync-service/forgejo_source_loader.py | 203 ++++++++++++++++++ tolaria/c5-sync-service/job_schema.py | 13 +- tolaria/c5-sync-service/save_executor_core.py | 65 +++++- tolaria/c5-sync-service/tolaria_client.py | 147 +++++++++++++ tolaria/c5-sync-service/worker.py | 170 +++++++++++++++ 7 files changed, 656 insertions(+), 14 deletions(-) create mode 100644 tolaria/c5-sync-service/Dockerfile create mode 100644 tolaria/c5-sync-service/entrypoint.py create mode 100644 tolaria/c5-sync-service/forgejo_source_loader.py create mode 100644 tolaria/c5-sync-service/tolaria_client.py create mode 100644 tolaria/c5-sync-service/worker.py diff --git a/tolaria/c5-sync-service/Dockerfile b/tolaria/c5-sync-service/Dockerfile new file mode 100644 index 0000000..d509702 --- /dev/null +++ b/tolaria/c5-sync-service/Dockerfile @@ -0,0 +1,25 @@ +# AUTH.4C1 — c5-save-executor (SAVE-only, fail-closed, verdrahtet) +FROM python:3.11-slim +# Kein echtes Credential im Image. Kein produktiver Write ohne injiziertes +# TOLARIA_SAVE_TOKEN (env_file, 0600). SAVE-only: kein DELETE-Code. +# git wird für den Forgejo-Source-Loader (git show {commit}:{path}) benötigt. +RUN apt-get update && apt-get install -y --no-install-recommends git \ + && rm -rf /var/lib/apt/lists/* +RUN useradd --uid 10011 --create-home --shell /usr/sbin/nologin c5save +WORKDIR /app +COPY --chown=10011:10011 \ + save_executor_core.py \ + job_store.py \ + job_schema.py \ + job_state_machine.py \ + job_claim.py \ + tolaria_client.py \ + forgejo_source_loader.py \ + worker.py \ + entrypoint.py \ + /app/ +USER 10011:10011 +ENV C5_SAVE_DB=/data/c5a_save.db +# FIXED Tolaria Base URL (interne Docker-DNS-Adresse, NIE aus dem Job) +ENV C5_TOLARIA_BASE=http://tolaria:5173/api/vault +ENTRYPOINT ["python3", "/app/entrypoint.py"] diff --git a/tolaria/c5-sync-service/entrypoint.py b/tolaria/c5-sync-service/entrypoint.py new file mode 100644 index 0000000..4031c07 --- /dev/null +++ b/tolaria/c5-sync-service/entrypoint.py @@ -0,0 +1,47 @@ +#!/usr/bin/env python3 +""" +AUTH.4C1 — c5-save-executor Runtime Entrypoint (SAVE-only, fail-closed). +Startet den verdrahteten SAVE-Worker-Loop. Health/Status-Modus. +""" +import json +import os +import sys + +from job_store import JobStore + +SCOPE = "SAVE" +DB_PATH = os.environ.get("C5_SAVE_DB", "/data/c5a_save.db") +VERSION = "0.1.0-auth4c1-wired" + + +def health() -> dict: + credential_present = bool(os.environ.get("TOLARIA_SAVE_TOKEN")) + db_ready = False + try: + store = JobStore(DB_PATH, SCOPE) + store.close() + db_ready = True + except Exception: + db_ready = False + return { + "service": "c5-save-executor", + "version": VERSION, + "mode": "fail_closed" if not credential_present else "credential_present", + "credential_present": credential_present, + "job_db_ready": db_ready, + "state": "idle" if not credential_present else "blocked", + "last_error": None, + } + + +def main() -> int: + if "--health" in sys.argv or "--status" in sys.argv: + print(json.dumps(health(), indent=2)) + return 0 + # Worker-Loop starten (verdrahtet, fail-closed) + from worker import main as worker_main + return worker_main() + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/tolaria/c5-sync-service/forgejo_source_loader.py b/tolaria/c5-sync-service/forgejo_source_loader.py new file mode 100644 index 0000000..7966003 --- /dev/null +++ b/tolaria/c5-sync-service/forgejo_source_loader.py @@ -0,0 +1,203 @@ +""" +AUTH.3E — forgejo_source_loader.py +=================================== +Forgejo Source Loader für den SAVE-Executor. + +Lädt Content aus der Forgejo Source of Truth, exakt commit-bound. + +Eigenschaften: + * Commit-bound: liest Content aus genau dem source_commit (kein Branch-Latest). + * object_id -> Pfad: leitet den Repo-relativen Pfad aus dem vault_path ab + (VAULT_PREFIX + rel_path), konsistent mit C5C (_vault_path). + * Keine Write-Fähigkeit. Keine automatische Upstream-Integration. + * FAIL CLOSED bei: commit missing, object missing, hash mismatch, + provenance mismatch, path mismatch, Forgejo unavailable. + * Read-only git-Befehle (git show / git cat-file) gegen den lokalen Clone + ODER HTTP-API (public Repo, kein Credential nötig). + +Zugriff: Repo nexo312/trading-system-docs ist PUBLIC -> kein Credential nötig. +""" + +from __future__ import annotations + +import hashlib +import os +import subprocess +from typing import Any, Dict, Optional + +# Forgejo-Repo (autoritative Source of Truth) +DEFAULT_FORGEJO_REPO = "nexo312/trading-system-docs" +ENV_FORGEJO_REPO = "C5_FORGEJO_REPO" + +# Vault-Pfad-Praefix (konsistent mit C5C VAULT_PREFIX) +VAULT_PREFIX = "/app/vault" + +# Lokaler Clone-Pfad (falls vorhanden) — sonst HTTP-API +DEFAULT_CLONE_PATH = "/data/forgejo-clone" +ENV_CLONE_PATH = "C5_FORGEJO_CLONE" + +# Forgejo HTTP-API Basis (public Repo, read-only) +DEFAULT_FORGEJO_BASE = "http://forgejo-c4u8yyi1eaz1gepn3pqmr5fb:3000" +ENV_FORGEJO_BASE = "C5_FORGEJO_BASE" + + +class SourceLoaderError(Exception): + pass + + +class SourceUnavailableError(SourceLoaderError): + pass + + +class SourceMismatchError(SourceLoaderError): + pass + + +class ForgejoSourceLoader: + """ + Lädt Content aus Forgejo, exakt commit-bound. + + load(source_commit, object_id, vault_path) -> content (str) + """ + + def __init__(self, repo: Optional[str] = None, + clone_path: Optional[str] = None, + forgejo_base: Optional[str] = None): + self.repo = (repo or os.environ.get(ENV_FORGEJO_REPO) + or DEFAULT_FORGEJO_REPO) + self.clone_path = (clone_path or os.environ.get(ENV_CLONE_PATH) + or DEFAULT_CLONE_PATH) + self.forgejo_base = (forgejo_base or os.environ.get(ENV_FORGEJO_BASE) + or DEFAULT_FORGEJO_BASE).rstrip("/") + + # -- Pfad-Ableitung ----------------------------------------------------- + + def _rel_path(self, vault_path: str) -> str: + """Leitet den Repo-relativen Pfad aus dem vault_path ab. + + vault_path = /app/vault/ -> rel_path. + Konsistent mit C5C _vault_path (VAULT_PREFIX + rel_path). + """ + vp = vault_path or "" + if vp.startswith(VAULT_PREFIX): + rel = vp[len(VAULT_PREFIX):].lstrip("/") + else: + rel = vp.lstrip("/") + if not rel: + raise SourceMismatchError("vault_path ergibt keinen rel_path") + return rel + + # -- Content-Load (commit-bound) ---------------------------------------- + + def load(self, source_commit: str, object_id: str, + vault_path: str) -> str: + """Lädt Content aus dem exakten Commit. + + Returns: autoritativer Content (str). + Raises: SourceUnavailableError / SourceMismatchError (FAIL CLOSED). + """ + rel_path = self._rel_path(vault_path) + + # 1. Commit-bound lesen (kein Branch-Latest) + content = self._read_from_commit(source_commit, rel_path) + if content is None: + raise SourceUnavailableError( + f"object {rel_path} nicht in commit {source_commit}") + + # 2. object_id-Konsistenz prüfen (Frontmatter id: object/) + self._validate_object_id(content, object_id, rel_path) + + return content + + def _read_from_commit(self, commit: str, rel_path: str) -> Optional[str]: + """Liest Datei aus exakt dem Commit. Bevorzugt lokalen Clone, sonst HTTP.""" + # Versuche lokalen Clone (falls gemountet) + if os.path.isdir(self.clone_path): + try: + return self._read_from_clone(commit, rel_path) + except SourceUnavailableError: + # Clone nicht verfügbar -> Fallback auf HTTP-API + pass + # HTTP-API (public Repo, read-only) + return self._read_from_http(commit, rel_path) + + def _read_from_clone(self, commit: str, rel_path: str) -> Optional[str]: + """git show : gegen lokalen Clone.""" + try: + proc = subprocess.run( + ["git", "-C", self.clone_path, "show", f"{commit}:{rel_path}"], + capture_output=True, text=True, timeout=30) + except (subprocess.SubprocessError, OSError) as e: + raise SourceUnavailableError(f"git show fehlgeschlagen: {e}") + if proc.returncode != 0: + return None # object nicht in commit + return proc.stdout + + def _read_from_http(self, commit: str, rel_path: str) -> Optional[str]: + """HTTP-API: GET /api/v1/repos/{repo}/raw/{commit}/{rel_path} (public).""" + import urllib.error + import urllib.request + url = (f"{self.forgejo_base}/api/v1/repos/{self.repo}" + f"/raw/{commit}/{rel_path}") + req = urllib.request.Request(url, method="GET") + try: + with urllib.request.urlopen(req, timeout=30) as resp: + return resp.read().decode("utf-8") + except urllib.error.HTTPError as e: + if e.code in (404, 410): + return None # object nicht in commit + raise SourceUnavailableError(f"Forgejo HTTP {e.code} auf {rel_path}") + except (urllib.error.URLError, TimeoutError, OSError) as e: + raise SourceUnavailableError(f"Forgejo nicht erreichbar: {e}") + + # -- object_id-Konsistenz ---------------------------------------------- + + def _validate_object_id(self, content: str, object_id: str, + rel_path: str) -> None: + """Prüft, dass das Frontmatter die erwartete object_id enthält.""" + fm = self._parse_frontmatter(content) + fm_id = fm.get("id") + if not fm_id: + raise SourceMismatchError( + f"object {rel_path} hat keine id im Frontmatter") + # object_id im Job ist 'object/' oder ''; Frontmatter id + # ist 'object/'. Normalisiere. + expected = object_id + if not expected.startswith("object/"): + expected = f"object/{expected}" + if fm_id != expected: + raise SourceMismatchError( + f"object_id mismatch: Frontmatter={fm_id} Job={expected}") + + @staticmethod + def _parse_frontmatter(content: str) -> Dict[str, Any]: + """Minimaler Frontmatter-Parser (--- ... ---).""" + if not content.startswith("---"): + return {} + end = content.find("\n---", 3) + if end == -1: + return {} + fm_text = content[3:end] + fm: Dict[str, Any] = {} + for line in fm_text.splitlines(): + if ":" in line: + k, _, v = line.partition(":") + fm[k.strip()] = v.strip().strip('"').strip("'") + return fm + + # -- Hash/Provenance (RECOMPUTE) --------------------------------------- + + @staticmethod + def content_hash(content: str) -> str: + """SHA-256 des fachlichen Bodies (nach Frontmatter).""" + body = content + if content.startswith("---"): + end = content.find("\n---", 3) + if end != -1: + body = content[end + 4:] + return hashlib.sha256(body.encode("utf-8")).hexdigest() + + @staticmethod + def provenance_hash(content: str) -> str: + """SHA-256 des gesamten Contents (konsistent mit save_executor_core).""" + return hashlib.sha256(content.encode("utf-8")).hexdigest() diff --git a/tolaria/c5-sync-service/job_schema.py b/tolaria/c5-sync-service/job_schema.py index dd9c5b1..0741cdc 100644 --- a/tolaria/c5-sync-service/job_schema.py +++ b/tolaria/c5-sync-service/job_schema.py @@ -71,6 +71,7 @@ ALLOWED_FIELDS = { # --------------------------------------------------------------------------- _UUID_RE = re.compile(r"^[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{12}$") _SHA256_RE = re.compile(r"^[0-9a-fA-F]{64}$") +_SHA1_RE = re.compile(r"^[0-9a-fA-F]{40}$") _ISO8601_UTC_RE = re.compile(r"^\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}Z$") @@ -82,6 +83,10 @@ def _is_sha256(value: Any) -> bool: return isinstance(value, str) and bool(_SHA256_RE.match(value)) +def _is_sha1(value: Any) -> bool: + return isinstance(value, str) and bool(_SHA1_RE.match(value)) + + def _is_iso8601_utc(value: Any) -> bool: return isinstance(value, str) and bool(_ISO8601_UTC_RE.match(value)) @@ -211,8 +216,8 @@ def validate_job(job: Dict[str, Any]) -> Dict[str, Any]: # 7. Job-Type-spezifische Felder if job_type == JOB_TYPE_SAVE: - if not _is_sha256(job.get("source_commit")): - raise JobRejectedError("source_commit must be a sha256 hex") + if not _is_sha1(job.get("source_commit")): + raise JobRejectedError("source_commit must be a sha1 hex (git commit)") if not _is_sha256(job.get("provenance_hash")): raise JobRejectedError("provenance_hash must be a sha256 hex") if not _is_nonempty_str(job.get("expected_state")): @@ -220,8 +225,8 @@ def validate_job(job: Dict[str, Any]) -> Dict[str, Any]: elif job_type == JOB_TYPE_DELETE: if not _is_uuid(job.get("delete_request_id")): raise JobRejectedError("delete_request_id must be a UUID") - if not _is_sha256(job.get("expected_commit")): - raise JobRejectedError("expected_commit must be a sha256 hex") + if not _is_sha1(job.get("expected_commit")): + raise JobRejectedError("expected_commit must be a sha1 hex (git commit)") if not _is_sha256(job.get("expected_provenance_hash")): raise JobRejectedError("expected_provenance_hash must be a sha256 hex") if not _is_uuid(job.get("approval_id")): diff --git a/tolaria/c5-sync-service/save_executor_core.py b/tolaria/c5-sync-service/save_executor_core.py index 1447c44..f0d946d 100644 --- a/tolaria/c5-sync-service/save_executor_core.py +++ b/tolaria/c5-sync-service/save_executor_core.py @@ -11,6 +11,12 @@ Eigenschaften: Credential, beliebigen Zielendpoint * Fail-closed: kein HTTP ohne SAVE-Credential * OUTCOME_UNKNOWN bei unklarem HTTP-Ergebnis (kein blinder Retry) + * Read-Back-Verifikation: nach SAVE wird der geschriebene Content gelesen; + Mismatch -> FAILED (nicht SUCCEEDED), unklares Read-Back -> OUTCOME_UNKNOWN + +AUTH.4C1-Verdrahtung (gegenüber AUTH.3E-Skeleton): + * source_loader-Signatur um vault_path erweitert (Pfad-Ableitung im Loader) + * optionales read_back-Callback für Read-Back-Verifikation (T23) Isoliert implementiert (KEIN produktiver Container). Nutzt job_store + job_schema. """ @@ -34,6 +40,7 @@ RC_STATE_MISMATCH = "STATE_MISMATCH" RC_PATH_INVALID = "PATH_INVALID" RC_TOLARIA_UNAVAILABLE = "TOLARIA_UNAVAILABLE" RC_OUTCOME_UNKNOWN = "OUTCOME_UNKNOWN" +RC_READBACK_MISMATCH = "READBACK_MISMATCH" RC_REJECTED = "REJECTED" @@ -45,24 +52,30 @@ class SaveExecutorCore: """ SAVE-Executor. worker_scope="SAVE" (nur SAVE-Jobs claimbar). - source_loader: Callable[[str, str], Optional[str]] — lädt Content aus - autoritativer Source (source_commit, object_id) -> content oder None. + source_loader: Callable[[str, str, str], Optional[str]] — lädt Content aus + autoritativer Source (source_commit, object_id, vault_path) -> content + oder None. tolaria_save: Callable[[str, str], Dict[str, Any]] — führt den Tolaria-SAVE aus (vault_path, content) -> {"status": "ok"|"error", "uncertain": bool}. Muss fail-closed sein (kein HTTP ohne Credential). + read_back: Optional[Callable[[str], Optional[str]]] — liest den geschriebenen + Content zurück (vault_path) -> content oder None. Wird NACH einem + bestätigten SAVE zur Verifikation aufgerufen. """ def __init__( self, store: JobStore, - source_loader: Callable[[str, str], Optional[str]], + source_loader: Callable[[str, str, str], Optional[str]], tolaria_save: Callable[[str, str], Dict[str, Any]], + read_back: Optional[Callable[[str], Optional[str]]] = None, ): if store.worker_scope != "SAVE": raise SaveExecutorError("SaveExecutorCore requires worker_scope='SAVE'") self.store = store self.source_loader = source_loader self.tolaria_save = tolaria_save + self.read_back = read_back # -- Hauptverarbeitung -------------------------------------------------- @@ -78,7 +91,8 @@ class SaveExecutorCore: 5. Content aus autoritativer Source rekonstruieren 6. Provenance/State validieren 7. Tolaria-SAVE ausführen - 8. Ergebnis: SUCCEEDED / FAILED / OUTCOME_UNKNOWN / REJECTED + 8. Read-Back verifizieren + 9. Ergebnis: SUCCEEDED / FAILED / OUTCOME_UNKNOWN / REJECTED """ job = self.store.get_job(job_id) if job is None: @@ -118,9 +132,16 @@ class SaveExecutorCore: result = self.tolaria_save(job["vault_path"], content) status = result.get("status") if status == "ok": - self.store.mark_succeeded(job_id, worker_id) - # Audit: Mutation ausgeführt (EXECUTED). - self.store.record_audit_event(job_id, "EXECUTED", worker_id=worker_id) + # Read-Back-Verifikation (T23): Mismatch -> FAILED, unklar -> UNKNOWN + rb = self._verify_read_back(job["vault_path"], content) + if rb == "match": + self.store.mark_succeeded(job_id, worker_id) + # Audit: Mutation ausgeführt (EXECUTED). + self.store.record_audit_event(job_id, "EXECUTED", worker_id=worker_id) + elif rb == "mismatch": + self.store.mark_failed(job_id, RC_READBACK_MISMATCH, worker_id) + else: # "unknown" + self.store.mark_outcome_unknown(job_id, worker_id) elif status == "error" and result.get("uncertain"): self.store.mark_outcome_unknown(job_id, worker_id) else: @@ -128,15 +149,39 @@ class SaveExecutorCore: return self.store.get_job(job_id) + # -- Read-Back-Verifikation --------------------------------------------- + + def _verify_read_back(self, vault_path: str, written_content: str) -> str: + """ + Verifiziert den geschriebenen Content per Read-Back. + + Returns: + "match" — Read-Back == geschriebener Content + "mismatch" — Read-Back vorhanden, aber != geschriebener Content + "unknown" — Read-Back nicht verfügbar / unklar (kein SUCCEEDED) + """ + if self.read_back is None: + # Kein Read-Back konfiguriert -> konservativ: nicht als SUCCEEDED + # bestätigen, wenn Verifikation gefordert ist. Hier: unknown. + return "unknown" + try: + rb_content = self.read_back(vault_path) + except Exception: + return "unknown" + if rb_content is None: + return "unknown" + return "match" if rb_content == written_content else "mismatch" + # -- Content-Rekonstruktion --------------------------------------------- def _reconstruct_content(self, job: Dict[str, Any]) -> Optional[str]: """ - Lädt Content selbst aus autoritativer Source (source_commit, object_id). - RQ liefert KEINEN Content im Job — nur Referenzen. + Lädt Content selbst aus autoritativer Source (source_commit, object_id, + vault_path). RQ liefert KEINEN Content im Job — nur Referenzen. """ try: - return self.source_loader(job["source_commit"], job["object_id"]) + return self.source_loader( + job["source_commit"], job["object_id"], job["vault_path"]) except Exception: return None diff --git a/tolaria/c5-sync-service/tolaria_client.py b/tolaria/c5-sync-service/tolaria_client.py new file mode 100644 index 0000000..d8b5c03 --- /dev/null +++ b/tolaria/c5-sync-service/tolaria_client.py @@ -0,0 +1,147 @@ +""" +AUTH.3A — tolaria_client.py +============================ +Isolierter Tolaria-Client für den SAVE-Executor (SAVE-only). + +Eigenschaften: + * FIXED Base URL (aus ENV C5_TOLARIA_BASE oder Default http://tolaria:5173/api/vault) + — NIE aus dem Job. Keine arbitrary URL. + * SAVE Endpoint fest: /save. Read-Back fest: /content. + * Keine arbitrary headers/method. + * Fail-closed: kein HTTP ohne SAVE-Credential (TOLARIA_SAVE_TOKEN). + * OUTCOME_UNKNOWN bei unklarem HTTP-Ergebnis (kein blinder Retry). + * Token-Wert wird NIE geloggt. + +Nur SAVE-Scope. Kein DELETE. Kein Master-Token. +""" + +from __future__ import annotations + +import json +import os +import urllib.error +import urllib.request +from typing import Any, Dict, Optional + +# FIXED Tolaria Base URL (produktiv: interne Docker-DNS-Adresse) +DEFAULT_TOLARIA_BASE = "http://tolaria:5173/api/vault" +ENV_TOLARIA_BASE = "C5_TOLARIA_BASE" +ENV_TOLARIA_SAVE_TOKEN = "TOLARIA_SAVE_TOKEN" + +# Feste Endpoints (keine arbitrary URL/method) +ENDPOINT_SAVE = "save" +ENDPOINT_CONTENT = "content" + +# Eindeutige, enge Not-Found-Semantik der Tolaria-Vault-API +TOLARIA_NOT_FOUND_MSG = "Invalid or missing path" + + +class TolariaClientError(Exception): + pass + + +class TolariaUnavailableError(TolariaClientError): + pass + + +class TolariaWriteError(TolariaClientError): + def __init__(self, message: str, code: str, http_code: Optional[int] = None): + super().__init__(message) + self.code = code + self.http_code = http_code + + +class TolariaClient: + """ + Isolierter Tolaria-Client (SAVE-only). + + read() — POST /content (Read-Back) + write() — POST /save (NUR SAVE; fail-closed ohne Credential) + """ + + def __init__(self, base_url: Optional[str] = None, timeout: float = 15.0, + save_token: Optional[str] = None): + # FIXED Base URL: aus ENV oder Default. NIE aus dem Job. + self.base_url = (base_url or os.environ.get(ENV_TOLARIA_BASE) + or DEFAULT_TOLARIA_BASE).rstrip("/") + self.timeout = timeout + # SAVE-Credential explizit injiziert (kein verstecktes globales). + # Fehlend/leer -> fail-closed beim mutierenden Aufruf. + self.save_token = save_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"} + # Authorization-Header NUR wenn ein Token explizit übergeben wird + # (mutierender SAVE). READ sendet KEIN Credential (Least Privilege). + 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: + if e.code >= 500: + raise TolariaUnavailableError( + f"Tolaria HTTP {e.code} auf {endpoint}", + "TOLARIA_UNAVAILABLE") + raise TolariaWriteError( + f"Tolaria HTTP {e.code} auf {endpoint}: " + f"{e.read().decode('utf-8', 'replace')[:200]}", + "AUTH_FAILURE" if e.code in (401, 403) else "INVALID_SCHEMA", + http_code=e.code) + except (urllib.error.URLError, TimeoutError, OSError) as e: + raise TolariaUnavailableError( + f"Tolaria nicht erreichbar ({endpoint}): {e}", + "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(ENDPOINT_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). + if e.http_code == 400 and TOLARIA_NOT_FOUND_MSG in (e.args[0] or ""): + return None + raise + if "error" in resp: + return None + return resp.get("content") + + # -- Write (NUR SAVE) --------------------------------------------------- + + def _require_token(self) -> str: + """Fail-closed: fehlendes/leeres SAVE-Credential -> kein HTTP.""" + token = self.save_token + if not token or not isinstance(token, str) or not token.strip(): + raise TolariaWriteError( + "Tolaria SAVE-Credential fehlt oder ist leer " + "(fail-closed, kein Request gesendet)", + "CREDENTIAL_MISSING") + return token + + def write(self, vault_path: str, content: str) -> Dict[str, Any]: + """Schreibt ein Vault-Objekt (POST /save). SAVE-Scope. + + Fail-closed: fehlendes/leeres SAVE-Credential -> lokaler Abbruch, + HTTP wird NICHT aufgerufen. + """ + token = self._require_token() + resp = self._post(ENDPOINT_SAVE, {"path": vault_path, "content": content}, + auth_token=token) + if resp is None: + resp = {} + if "error" in resp: + raise TolariaWriteError( + f"Tolaria save fehlgeschlagen: {resp['error']}", + "INVALID_SCHEMA") + return resp diff --git a/tolaria/c5-sync-service/worker.py b/tolaria/c5-sync-service/worker.py new file mode 100644 index 0000000..82dce2c --- /dev/null +++ b/tolaria/c5-sync-service/worker.py @@ -0,0 +1,170 @@ +""" +AUTH.4C1 — worker.py +==================== +SAVE-Worker-Loop für den c5-save-executor. + +Ablauf (pro Job): + poll -> atomic claim -> validate -> load source -> authorize local + prerequisites -> execute -> read-back -> audit -> complete + +Eigenschaften: + * Kein Planner. Kein Agent. Keine LLM-Entscheidung. + * POLL_INTERVAL, LEASE, MAX_ATTEMPTS, structured backoff, graceful SIGTERM. + * OUTCOME_UNKNOWN bei unklarem Write-Ergebnis. KEIN blinder Retry. + * Fail-closed: ohne SAVE-Credential wird kein HTTP-SAVE ausgeführt. +""" + +from __future__ import annotations + +import json +import logging +import os +import signal +import sys +import time +from typing import Any, Dict, Optional + +from forgejo_source_loader import ForgejoSourceLoader +from job_store import JobStore +from save_executor_core import SaveExecutorCore +from tolaria_client import TolariaClient + +SCOPE = "SAVE" +DB_PATH = os.environ.get("C5_SAVE_DB", "/data/c5a_save.db") +VERSION = "0.1.0-auth4c1-wired" + +# Konfiguration (ENV mit konservativen Defaults) +POLL_INTERVAL = float(os.environ.get("C5_POLL_INTERVAL", "5.0")) +LEASE_SECONDS = int(os.environ.get("C5_LEASE_SECONDS", "60")) +MAX_ATTEMPTS = int(os.environ.get("C5_MAX_ATTEMPTS", "3")) +BACKOFF_BASE = float(os.environ.get("C5_BACKOFF_BASE", "2.0")) +BACKOFF_MAX = float(os.environ.get("C5_BACKOFF_MAX", "60.0")) + +# Tolaria-Client (fixed base URL, SAVE-only) +TOLARIA_BASE = os.environ.get("C5_TOLARIA_BASE", "http://tolaria:5173/api/vault") +TOLARIA_SAVE_TOKEN = os.environ.get("TOLARIA_SAVE_TOKEN") + +logging.basicConfig( + level=logging.INFO, + format="%(asctime)s %(levelname)s %(name)s: %(message)s", +) +log = logging.getLogger("c5-save-worker") + +# Graceful Shutdown +_shutdown = False + + +def _handle_sigterm(signum, frame): + global _shutdown + log.info("SIGTERM empfangen — Worker fährt sauber herunter.") + _shutdown = True + + +def _handle_sigint(signum, frame): + global _shutdown + log.info("SIGINT empfangen — Worker fährt sauber herunter.") + _shutdown = True + + +signal.signal(signal.SIGTERM, _handle_sigterm) +signal.signal(signal.SIGINT, _handle_sigint) + + +def _backoff_delay(attempt: int) -> float: + """Structured exponential backoff (kein blinder Retry, nur transient).""" + delay = min(BACKOFF_MAX, BACKOFF_BASE * (2 ** (attempt - 1))) + return delay + + +def _build_worker() -> tuple[JobStore, SaveExecutorCore]: + """Baut Store + Core mit injizierten Callables (Verdrahtung).""" + store = JobStore(DB_PATH, SCOPE) + loader = ForgejoSourceLoader() + client = TolariaClient(base_url=TOLARIA_BASE, save_token=TOLARIA_SAVE_TOKEN) + + def source_loader(source_commit: str, object_id: str, vault_path: str) -> Optional[str]: + return loader.load(source_commit, object_id, vault_path) + + def tolaria_save(vault_path: str, content: str) -> Dict[str, Any]: + try: + client.write(vault_path, content) + return {"status": "ok"} + except Exception as e: + # Unklar vs. bestätigt unterscheiden + code = getattr(e, "code", "TOLARIA_UNAVAILABLE") + if code == "CREDENTIAL_MISSING": + return {"status": "error", "code": "CREDENTIAL_MISSING", "uncertain": False} + if code in ("AUTH_FAILURE", "INVALID_SCHEMA"): + return {"status": "error", "code": code, "uncertain": False} + # Netzwerk/Timeout/5xx -> unklar (kein blinder Retry) + return {"status": "error", "code": "TOLARIA_UNAVAILABLE", "uncertain": True} + + def read_back(vault_path: str) -> Optional[str]: + try: + return client.read(vault_path) + except Exception: + return None + + core = SaveExecutorCore(store, source_loader, tolaria_save, read_back=read_back) + return store, core + + +def _process_ready_job(store: JobStore, core: SaveExecutorCore, job: Dict[str, Any]) -> None: + """Verarbeitet einen READY-Job (mit MAX_ATTEMPTS + backoff).""" + job_id = job["job_id"] + attempt = 0 + while attempt < MAX_ATTEMPTS: + attempt += 1 + log.info("Verarbeite Job %s (Versuch %d/%d)", job_id, attempt, MAX_ATTEMPTS) + try: + result = core.process_job(job_id, f"worker-{os.getpid()}") + state = result.get("state") + log.info("Job %s -> %s (result_code=%s)", + job_id, state, result.get("result_code")) + # Terminal-States: fertig. OUTCOME_UNKNOWN: kein blinder Retry. + if state in ("SUCCEEDED", "FAILED", "REJECTED", "OUTCOME_UNKNOWN"): + return + # READY/CLAIMED/EXECUTING (transient) -> backoff + erneut versuchen + if attempt < MAX_ATTEMPTS: + time.sleep(_backoff_delay(attempt)) + except Exception as e: + log.error("Job %s Fehler: %s", job_id, e) + if attempt < MAX_ATTEMPTS: + time.sleep(_backoff_delay(attempt)) + # MAX_ATTEMPTS erschöpft -> OUTCOME_UNKNOWN (kein blinder Retry) + try: + store.mark_outcome_unknown(job_id, f"worker-{os.getpid()}") + except Exception: + pass + + +def main() -> int: + log.info("c5-save-worker v%s startet (scope=%s, poll=%.1fs, lease=%ds, max_attempts=%d)", + VERSION, SCOPE, POLL_INTERVAL, LEASE_SECONDS, MAX_ATTEMPTS) + if not TOLARIA_SAVE_TOKEN: + log.warning("TOLARIA_SAVE_TOKEN fehlt — fail-closed: kein HTTP-SAVE wird ausgeführt.") + store, core = _build_worker() + try: + while not _shutdown: + try: + # READY-Jobs claimen und verarbeiten + ready_jobs = store.list_jobs(state="READY") + for job in ready_jobs: + if _shutdown: + break + _process_ready_job(store, core, job) + except Exception as e: + log.error("Worker-Loop-Fehler: %s", e) + # Poll-Intervall (graceful: bei SIGTERM sofort raus) + for _ in range(int(POLL_INTERVAL * 10)): + if _shutdown: + break + time.sleep(0.1) + finally: + store.close() + log.info("Worker beendet.") + return 0 + + +if __name__ == "__main__": + sys.exit(main())