trading-system-docs/tolaria/c5-sync-service/save_executor_core.py

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"]