196 lines
7.8 KiB
Python
196 lines
7.8 KiB
Python
"""
|
|
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"]
|