""" AUTH.3E — save_executor_core.py ================================ SAVE-Executor-Logik (SAVE-only). Eigenschaften: * SAVE-only: claimt NUR C5_SAVE_OBJECT-Jobs * Content-Rekonstruktion: lädt Content selbst aus autoritativer Source (Forgejo), statt RQ blind zu vertrauen (DATA FROM RQ != AUTHORITY) * RQ darf NICHT bestimmen: Tolaria Base URL, HTTP Method, Authorization Header, 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. """ from __future__ import annotations import hashlib from typing import Any, Callable, Dict, Optional from job_schema import JOB_TYPE_SAVE from job_store import JobStore # --------------------------------------------------------------------------- # Result-Codes # --------------------------------------------------------------------------- RC_OK = "OK" RC_CREDENTIAL_MISSING = "CREDENTIAL_MISSING" RC_SOURCE_UNAVAILABLE = "SOURCE_UNAVAILABLE" RC_PROVENANCE_MISMATCH = "PROVENANCE_MISMATCH" 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" class SaveExecutorError(Exception): pass class SaveExecutorCore: """ SAVE-Executor. worker_scope="SAVE" (nur SAVE-Jobs claimbar). 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, 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 -------------------------------------------------- def process_job(self, job_id: str, worker_id: str) -> Dict[str, Any]: """ Verarbeitet einen SAVE-Job durch die State Machine. Ablauf: 1. Job laden (muss existieren) 2. Job-Type-Scope prüfen (nur SAVE) 3. READY -> CLAIMED (atomarer Claim) 4. CLAIMED -> EXECUTING 5. Content aus autoritativer Source rekonstruieren 6. Provenance/State validieren 7. Tolaria-SAVE ausführen 8. Read-Back verifizieren 9. Ergebnis: SUCCEEDED / FAILED / OUTCOME_UNKNOWN / REJECTED """ job = self.store.get_job(job_id) if job is None: raise SaveExecutorError(f"job not found: {job_id}") # Scope: nur SAVE-Jobs if job["job_type"] != JOB_TYPE_SAVE: self.store.mark_rejected(job_id, RC_REJECTED) return self.store.get_job(job_id) # Atomarer Claim try: self.store.claim_job(job_id, worker_id) except Exception: # Nicht claimbar (bereits geclaimt) -> kein Doppel-Write return self.store.get_job(job_id) self.store.begin_execution(job_id, worker_id) # Content-Rekonstruktion aus autoritativer Source content = self._reconstruct_content(job) if content is None: self.store.mark_failed(job_id, RC_SOURCE_UNAVAILABLE, worker_id) return self.store.get_job(job_id) # Provenance validieren (RECOMPUTE) if not self._validate_provenance(job, content): self.store.mark_failed(job_id, RC_PROVENANCE_MISMATCH, worker_id) return self.store.get_job(job_id) # Audit: Executor hat Job + Provenance validiert (AUTHORIZED). # Ein REQUESTED-Eintrag ist KEINE Autorisierung — nur AUTHORIZED # autorisiert die Mutation. self.store.record_audit_event(job_id, "AUTHORIZED", worker_id=worker_id) # Tolaria-SAVE ausführen (fail-closed) result = self.tolaria_save(job["vault_path"], content) status = result.get("status") if status == "ok": # 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: self.store.mark_failed(job_id, result.get("code", RC_TOLARIA_UNAVAILABLE), worker_id) 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, vault_path). RQ liefert KEINEN Content im Job — nur Referenzen. """ try: return self.source_loader( job["source_commit"], job["object_id"], job["vault_path"]) except Exception: return None # -- Provenance-Validierung --------------------------------------------- def _validate_provenance(self, job: Dict[str, Any], content: str) -> bool: """ RECOMPUTE: berechnet den Provenance-Hash aus dem rekonstruierten Content und vergleicht mit dem Job-Feld. RQ darf den Hash nicht blind bestimmen. """ computed = hashlib.sha256(content.encode("utf-8")).hexdigest() return computed == job["provenance_hash"]