""" 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) 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_REJECTED = "REJECTED" class SaveExecutorError(Exception): pass 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. 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). """ def __init__( self, store: JobStore, source_loader: Callable[[str, str], Optional[str]], tolaria_save: Callable[[str, str], Dict[str, Any]], ): 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 # -- 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. 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": 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 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) # -- 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. """ try: return self.source_loader(job["source_commit"], job["object_id"]) 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"]