Compare commits

..

No commits in common. "2e631fd5fae56e258e47fb6e6477aaf577482683" and "9f3ed82cfc895d1d41e24c0ec656bc611520ffa9" have entirely different histories.

8 changed files with 30 additions and 1308 deletions

View file

@ -1,25 +0,0 @@
# 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"]

View file

@ -1,47 +0,0 @@
#!/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())

View file

@ -1,202 +0,0 @@
"""
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
(kanonisch RELATIV, AUTH.4C2 PATH CONTRACT REPAIR OPTION A), 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"
# 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.
KANONISCH (AUTH.4C2 PATH CONTRACT REPAIR, OPTION A):
vault_path ist bereits der RELATIVE Pfad unter dem Vault-Root
(z.B. 'tolaria/auth4c2-canary.md'). Er wird unverändert als
rel_path verwendet KEIN /app/vault/-Prefix-Stripping mehr.
Konsistent mit C5C _vault_path (rel_path unter Vault-Root).
"""
vp = vault_path or ""
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/<uuid>)
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 <commit>:<rel_path> 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/<uuid>' oder '<uuid>'; Frontmatter id
# ist 'object/<uuid>'. 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()

View file

@ -71,7 +71,6 @@ 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}$") _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}$") _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$") _ISO8601_UTC_RE = re.compile(r"^\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}Z$")
@ -83,10 +82,6 @@ def _is_sha256(value: Any) -> bool:
return isinstance(value, str) and bool(_SHA256_RE.match(value)) 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: def _is_iso8601_utc(value: Any) -> bool:
return isinstance(value, str) and bool(_ISO8601_UTC_RE.match(value)) return isinstance(value, str) and bool(_ISO8601_UTC_RE.match(value))
@ -98,30 +93,25 @@ def _is_nonempty_str(value: Any) -> bool:
def _normalize_path(path: str) -> str: def _normalize_path(path: str) -> str:
""" """
Normalisiert einen Vault-Pfad (keine //, ., .., trailing slash). Normalisiert einen Vault-Pfad (keine //, ., .., trailing slash).
KANONISCH RELATIV (AUTH.4C2 PATH CONTRACT REPAIR, OPTION A): Wiederverwendet die AUTH.3D-Pfad-Normalisierung.
kein führender Slash. Immer die eigene relative Normalisierung
KEIN Import von approval_payload._normalize_path, da jene (AUTH.3D)
absolute Pfade erzwingt und damit die kanonisch-relative Semantik
verletzen würde (keine Dual-Semantik).
""" """
parts = [p for p in path.split("/") if p not in ("", ".")] try:
if ".." in parts: from approval_payload import _normalize_path as _ap_normalize
raise ValueError("path traversal") return _ap_normalize(path)
return "/".join(parts) except Exception:
# Fallback (isoliert): eigene Normalisierung
parts = [p for p in path.split("/") if p not in ("", ".")]
if ".." in parts:
raise ValueError("path traversal")
return "/" + "/".join(parts)
def _validate_path(path: str) -> Optional[str]: def _validate_path(path: str) -> Optional[str]:
""" """
Pfad-Sicherheitsprüfung (Executor-seitig, Defense in depth). Pfad-Sicherheitsprüfung (Executor-seitig, Defense in depth).
KANONISCHER JOB-VAULT-PFAD = RELATIVER Pfad unter dem Vault-Root Erlaubt nur Pfade unterhalb des Vault-Roots, keine Traversal,
(AUTH.4C2 PATH CONTRACT REPAIR, OPTION A). keine absoluten Host-Pfade, keine Unicode-Ambiguität, keine
URL-Encodierung, keine Sonderzeichen.
Erlaubt NUR relative Pfade. Lehnt ab:
* absolute Pfade (führende /) Tolaria AUTH.2 lehnt absolute ab
* Traversal (.., ./)
* Backslash-Traversal
* URL/Scheme (http://, file://, C:)
* Null-Bytes, Unicode-Ambiguität, URL-Encodierung, Sonderzeichen
""" """
if not isinstance(path, str) or not path: if not isinstance(path, str) or not path:
return "path must be a non-empty string" return "path must be a non-empty string"
@ -129,13 +119,9 @@ def _validate_path(path: str) -> Optional[str]:
return "path too long" return "path too long"
if "\x00" in path: if "\x00" in path:
return "path contains null byte" return "path contains null byte"
# Kanonisch RELATIV: keine führende /, kein Scheme, kein Backslash # Kein absoluter Host-Pfad (nur /app/vault/... erlaubt)
if path.startswith("/"): if not path.startswith("/app/vault/"):
return "path must be relative (no leading /)" return "path must be under /app/vault/"
if "\\" in path:
return "path contains backslash"
if ":" in path:
return "path contains scheme/colon"
# Zeichensatz-Whitelist: nur sichere Pfadzeichen. # Zeichensatz-Whitelist: nur sichere Pfadzeichen.
# Schließt URL-Encodierung (%), Unicode-Homoglyphen, Leerzeichen, # Schließt URL-Encodierung (%), Unicode-Homoglyphen, Leerzeichen,
# Steuerzeichen und Sonderzeichen aus. # Steuerzeichen und Sonderzeichen aus.
@ -225,8 +211,8 @@ def validate_job(job: Dict[str, Any]) -> Dict[str, Any]:
# 7. Job-Type-spezifische Felder # 7. Job-Type-spezifische Felder
if job_type == JOB_TYPE_SAVE: if job_type == JOB_TYPE_SAVE:
if not _is_sha1(job.get("source_commit")): if not _is_sha256(job.get("source_commit")):
raise JobRejectedError("source_commit must be a sha1 hex (git commit)") raise JobRejectedError("source_commit must be a sha256 hex")
if not _is_sha256(job.get("provenance_hash")): if not _is_sha256(job.get("provenance_hash")):
raise JobRejectedError("provenance_hash must be a sha256 hex") raise JobRejectedError("provenance_hash must be a sha256 hex")
if not _is_nonempty_str(job.get("expected_state")): if not _is_nonempty_str(job.get("expected_state")):
@ -234,8 +220,8 @@ def validate_job(job: Dict[str, Any]) -> Dict[str, Any]:
elif job_type == JOB_TYPE_DELETE: elif job_type == JOB_TYPE_DELETE:
if not _is_uuid(job.get("delete_request_id")): if not _is_uuid(job.get("delete_request_id")):
raise JobRejectedError("delete_request_id must be a UUID") raise JobRejectedError("delete_request_id must be a UUID")
if not _is_sha1(job.get("expected_commit")): if not _is_sha256(job.get("expected_commit")):
raise JobRejectedError("expected_commit must be a sha1 hex (git commit)") raise JobRejectedError("expected_commit must be a sha256 hex")
if not _is_sha256(job.get("expected_provenance_hash")): if not _is_sha256(job.get("expected_provenance_hash")):
raise JobRejectedError("expected_provenance_hash must be a sha256 hex") raise JobRejectedError("expected_provenance_hash must be a sha256 hex")
if not _is_uuid(job.get("approval_id")): if not _is_uuid(job.get("approval_id")):

View file

@ -11,12 +11,6 @@ Eigenschaften:
Credential, beliebigen Zielendpoint Credential, beliebigen Zielendpoint
* Fail-closed: kein HTTP ohne SAVE-Credential * Fail-closed: kein HTTP ohne SAVE-Credential
* OUTCOME_UNKNOWN bei unklarem HTTP-Ergebnis (kein blinder Retry) * 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. Isoliert implementiert (KEIN produktiver Container). Nutzt job_store + job_schema.
""" """
@ -40,7 +34,6 @@ RC_STATE_MISMATCH = "STATE_MISMATCH"
RC_PATH_INVALID = "PATH_INVALID" RC_PATH_INVALID = "PATH_INVALID"
RC_TOLARIA_UNAVAILABLE = "TOLARIA_UNAVAILABLE" RC_TOLARIA_UNAVAILABLE = "TOLARIA_UNAVAILABLE"
RC_OUTCOME_UNKNOWN = "OUTCOME_UNKNOWN" RC_OUTCOME_UNKNOWN = "OUTCOME_UNKNOWN"
RC_READBACK_MISMATCH = "READBACK_MISMATCH"
RC_REJECTED = "REJECTED" RC_REJECTED = "REJECTED"
@ -52,30 +45,24 @@ class SaveExecutorCore:
""" """
SAVE-Executor. worker_scope="SAVE" (nur SAVE-Jobs claimbar). SAVE-Executor. worker_scope="SAVE" (nur SAVE-Jobs claimbar).
source_loader: Callable[[str, str, str], Optional[str]] lädt Content aus source_loader: Callable[[str, str], Optional[str]] lädt Content aus
autoritativer Source (source_commit, object_id, vault_path) -> content autoritativer Source (source_commit, object_id) -> content oder None.
oder None.
tolaria_save: Callable[[str, str], Dict[str, Any]] führt den Tolaria-SAVE tolaria_save: Callable[[str, str], Dict[str, Any]] führt den Tolaria-SAVE
aus (vault_path, content) -> {"status": "ok"|"error", "uncertain": bool}. aus (vault_path, content) -> {"status": "ok"|"error", "uncertain": bool}.
Muss fail-closed sein (kein HTTP ohne Credential). 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__( def __init__(
self, self,
store: JobStore, store: JobStore,
source_loader: Callable[[str, str, str], Optional[str]], source_loader: Callable[[str, str], Optional[str]],
tolaria_save: Callable[[str, str], Dict[str, Any]], tolaria_save: Callable[[str, str], Dict[str, Any]],
read_back: Optional[Callable[[str], Optional[str]]] = None,
): ):
if store.worker_scope != "SAVE": if store.worker_scope != "SAVE":
raise SaveExecutorError("SaveExecutorCore requires worker_scope='SAVE'") raise SaveExecutorError("SaveExecutorCore requires worker_scope='SAVE'")
self.store = store self.store = store
self.source_loader = source_loader self.source_loader = source_loader
self.tolaria_save = tolaria_save self.tolaria_save = tolaria_save
self.read_back = read_back
# -- Hauptverarbeitung -------------------------------------------------- # -- Hauptverarbeitung --------------------------------------------------
@ -91,8 +78,7 @@ class SaveExecutorCore:
5. Content aus autoritativer Source rekonstruieren 5. Content aus autoritativer Source rekonstruieren
6. Provenance/State validieren 6. Provenance/State validieren
7. Tolaria-SAVE ausführen 7. Tolaria-SAVE ausführen
8. Read-Back verifizieren 8. Ergebnis: SUCCEEDED / FAILED / OUTCOME_UNKNOWN / REJECTED
9. Ergebnis: SUCCEEDED / FAILED / OUTCOME_UNKNOWN / REJECTED
""" """
job = self.store.get_job(job_id) job = self.store.get_job(job_id)
if job is None: if job is None:
@ -132,16 +118,9 @@ class SaveExecutorCore:
result = self.tolaria_save(job["vault_path"], content) result = self.tolaria_save(job["vault_path"], content)
status = result.get("status") status = result.get("status")
if status == "ok": if status == "ok":
# Read-Back-Verifikation (T23): Mismatch -> FAILED, unklar -> UNKNOWN self.store.mark_succeeded(job_id, worker_id)
rb = self._verify_read_back(job["vault_path"], content) # Audit: Mutation ausgeführt (EXECUTED).
if rb == "match": self.store.record_audit_event(job_id, "EXECUTED", worker_id=worker_id)
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"): elif status == "error" and result.get("uncertain"):
self.store.mark_outcome_unknown(job_id, worker_id) self.store.mark_outcome_unknown(job_id, worker_id)
else: else:
@ -149,39 +128,15 @@ class SaveExecutorCore:
return self.store.get_job(job_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 --------------------------------------------- # -- Content-Rekonstruktion ---------------------------------------------
def _reconstruct_content(self, job: Dict[str, Any]) -> Optional[str]: def _reconstruct_content(self, job: Dict[str, Any]) -> Optional[str]:
""" """
Lädt Content selbst aus autoritativer Source (source_commit, object_id, Lädt Content selbst aus autoritativer Source (source_commit, object_id).
vault_path). RQ liefert KEINEN Content im Job nur Referenzen. RQ liefert KEINEN Content im Job nur Referenzen.
""" """
try: try:
return self.source_loader( return self.source_loader(job["source_commit"], job["object_id"])
job["source_commit"], job["object_id"], job["vault_path"])
except Exception: except Exception:
return None return None

View file

@ -1,619 +0,0 @@
"""
AUTH.4C1 Phase 11: Isolated Tests (T1-T27)
=============================================
Läuft mit Fake/isolierter Tolaria-Instanz und synthetischen Credentials.
KEIN produktiver Tolaria-Kontakt. KEIN produktiver SAVE.
Testumgebung:
* TOLARIA_SAVE_TOKEN = synthetisch (nur für Tests)
* TolariaClient mit Fake-Base-URL (http://127.0.0.1:1 = unerreichbar)
* ForgejoSourceLoader mit Fake-Clone (lokales Test-Repo)
* JobStore auf temporärer SQLite-DB
Führt die AUTH.4C1-Tests T1-T27 aus. Exit 0 = alle PASS.
"""
from __future__ import annotations
import hashlib
import json
import os
import shutil
import subprocess
import sys
import tempfile
import time
import unittest
import uuid
from typing import Any, Dict, Optional
# Module aus dem Build-Verzeichnis importieren
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
from forgejo_source_loader import ForgejoSourceLoader, SourceMismatchError, SourceUnavailableError
from job_schema import make_save_job, make_delete_job, validate_job, JobRejectedError
from job_store import JobStore
from save_executor_core import SaveExecutorCore
from tolaria_client import TolariaClient, TolariaWriteError
# Synthetisches Test-Credential (NUR für Tests, nie produktiv)
TEST_TOKEN = "test-save-token-4c1-isolated"
# Fake-Tolaria: unerreichbare Base-URL (kein echter Kontakt)
FAKE_BASE = "http://127.0.0.1:1/api/vault"
def _sha256(s: str) -> str:
return hashlib.sha256(s.encode("utf-8")).hexdigest()
def _make_test_repo(tmpdir: str) -> str:
"""Erzeugt ein lokales Git-Repo als Fake-Forgejo-Clone."""
repo = os.path.join(tmpdir, "test-repo")
os.makedirs(repo, exist_ok=True)
subprocess.run(["git", "init", "-q", repo], check=True)
subprocess.run(["git", "-C", repo, "config", "user.email", "t@t"], check=True)
subprocess.run(["git", "-C", repo, "config", "user.name", "t"], check=True)
# Datei mit Frontmatter + Body
content = "---\nid: object/test-obj-1\ntitle: Test\n---\n\nHello World\n"
os.makedirs(os.path.join(repo, "docs"), exist_ok=True)
with open(os.path.join(repo, "docs", "test.md"), "w") as f:
f.write(content)
subprocess.run(["git", "-C", repo, "add", "-A"], check=True)
subprocess.run(["git", "-C", repo, "commit", "-q", "-m", "init"], check=True)
return repo
class FakeTolariaSave:
"""Fake-Tolaria-SAVE-Callable. Zeichnet Aufrufe auf, führt keinen echten Write aus."""
def __init__(self, behavior: str = "ok"):
self.behavior = behavior # ok | uncertain | fail | credential_missing
self.calls: list[tuple[str, str]] = []
def __call__(self, vault_path: str, content: str) -> Dict[str, Any]:
self.calls.append((vault_path, content))
if self.behavior == "ok":
return {"status": "ok"}
if self.behavior == "uncertain":
return {"status": "error", "code": "TOLARIA_UNAVAILABLE", "uncertain": True}
if self.behavior == "credential_missing":
return {"status": "error", "code": "CREDENTIAL_MISSING", "uncertain": False}
return {"status": "error", "code": "AUTH_FAILURE", "uncertain": False}
class FakeReadBack:
def __init__(self, content: Optional[str] = None, behavior: str = "ok"):
self.content = content
self.behavior = behavior # ok | mismatch | unknown
self.calls = 0
def __call__(self, vault_path: str) -> Optional[str]:
self.calls += 1
if self.behavior == "unknown":
return None
if self.behavior == "mismatch":
return "DIFFERENT"
return self.content
class AUTH4C1IsolatedTests(unittest.TestCase):
@classmethod
def setUpClass(cls):
cls.tmp = tempfile.mkdtemp(prefix="c5-4c1-test-")
cls.repo = _make_test_repo(cls.tmp)
cls.db = os.path.join(cls.tmp, "test_save.db")
cls.store = JobStore(cls.db, "SAVE")
# Commit-Hash des Test-Repos
cls.commit = subprocess.run(
["git", "-C", cls.repo, "rev-parse", "HEAD"],
capture_output=True, text=True, check=True).stdout.strip()
cls.content = "---\nid: object/test-obj-1\ntitle: Test\n---\n\nHello World\n"
cls.prov = _sha256(cls.content)
@classmethod
def tearDownClass(cls):
cls.store.close()
shutil.rmtree(cls.tmp, ignore_errors=True)
def _make_job(self, **overrides) -> Dict[str, Any]:
job = make_save_job(
job_id=overrides.get("job_id", str(uuid.uuid4())),
mission_id=overrides.get("mission_id", str(uuid.uuid4())),
object_id=overrides.get("object_id", "object/test-obj-1"),
vault_path=overrides.get("vault_path", "docs/test.md"),
source_commit=overrides.get("source_commit", self.commit),
provenance_hash=overrides.get("provenance_hash", self.prov),
expected_state="present",
created_at="2026-08-27T12:00:00Z",
idempotency_key=overrides.get("idempotency_key", str(uuid.uuid4())),
)
return job
def _loader(self, repo: Optional[str] = None) -> ForgejoSourceLoader:
return ForgejoSourceLoader(repo="test/test", clone_path=repo or self.repo,
forgejo_base="http://127.0.0.1:1")
def _core(self, loader, save, read_back=None) -> SaveExecutorCore:
return SaveExecutorCore(self.store, loader.load, save, read_back=read_back)
# -- T1: worker starts without token ------------------------------------
def test_t1_worker_starts_without_token(self):
# Worker-Modul importierbar + main() startet ohne Token (fail-closed)
import worker
self.assertIsNotNone(worker)
self.assertFalse(worker.TOLARIA_SAVE_TOKEN)
# -- T2: no token -> no HTTP --------------------------------------------
def test_t2_no_token_no_http(self):
client = TolariaClient(base_url=FAKE_BASE, save_token=None)
with self.assertRaises(TolariaWriteError) as ctx:
client.write("x.md", "content")
self.assertEqual(ctx.exception.code, "CREDENTIAL_MISSING")
# -- T3: valid SAVE job processed ---------------------------------------
def test_t3_valid_save_processed(self):
job = self._make_job()
self.store.create_job(job)
self.store.mark_ready(job["job_id"])
save = FakeTolariaSave("ok")
rb = FakeReadBack(self.content)
core = self._core(self._loader(), save, rb)
result = core.process_job(job["job_id"], "w-t3")
self.assertEqual(result["state"], "SUCCEEDED")
self.assertEqual(len(save.calls), 1)
self.assertEqual(save.calls[0][0], "docs/test.md")
# -- T4: unknown type rejected ------------------------------------------
def test_t4_unknown_type_rejected(self):
job = self._make_job()
job["job_type"] = "C5_UNKNOWN"
with self.assertRaises(JobRejectedError):
validate_job(job)
# -- T5: malformed rejected ---------------------------------------------
def test_t5_malformed_rejected(self):
with self.assertRaises(JobRejectedError):
validate_job({"job_id": "x"}) # unvollständig
# -- T6: duplicate safe -------------------------------------------------
def test_t6_duplicate_safe(self):
job = self._make_job()
self.store.create_job(job)
# Zweiter create mit gleicher job_id -> idempotent (kein Fehler)
again = self.store.create_job(job)
self.assertEqual(again["job_id"], job["job_id"])
# -- T7: fixed Tolaria target ------------------------------------------
def test_t7_fixed_tolaria_target(self):
client = TolariaClient(base_url=FAKE_BASE, save_token=TEST_TOKEN)
self.assertEqual(client.base_url, FAKE_BASE.rstrip("/"))
# -- T8: URL injection rejected -----------------------------------------
def test_t8_url_injection_rejected(self):
job = self._make_job()
job["url"] = "http://evil"
with self.assertRaises(JobRejectedError):
validate_job(job)
# -- T9: header injection rejected ---------------------------------------
def test_t9_header_injection_rejected(self):
job = self._make_job()
job["headers"] = {"X-Evil": "1"}
with self.assertRaises(JobRejectedError):
validate_job(job)
# -- T10: method injection rejected -------------------------------------
def test_t10_method_injection_rejected(self):
job = self._make_job()
job["method"] = "DELETE"
with self.assertRaises(JobRejectedError):
validate_job(job)
# -- T11: delete job rejected -------------------------------------------
def test_t11_delete_job_rejected(self):
job = make_delete_job(
job_id=str(uuid.uuid4()), mission_id=str(uuid.uuid4()), delete_request_id=str(uuid.uuid4()),
object_id="object/x", vault_path="x.md",
expected_commit=self.commit, expected_provenance_hash=self.prov,
approval_id=str(uuid.uuid4()), created_at="2026-08-27T12:00:00Z",
idempotency_key=str(uuid.uuid4()))
self.store.create_job(job)
self.store.mark_ready(job["job_id"])
save = FakeTolariaSave("ok")
core = self._core(self._loader(), save)
result = core.process_job(job["job_id"], "w-t11")
self.assertEqual(result["state"], "REJECTED")
self.assertEqual(len(save.calls), 0)
# -- T12: traversal rejected --------------------------------------------
def test_t12_traversal_rejected(self):
# _make_job validiert bereits beim Bauen -> JobRejectedError
with self.assertRaises(JobRejectedError):
self._make_job(vault_path="../../etc/passwd")
# -- T13: content loaded from exact Forgejo commit -----------------------
def test_t13_content_from_exact_commit(self):
loader = self._loader()
content = loader.load(self.commit, "object/test-obj-1", "docs/test.md")
self.assertEqual(content, self.content)
# -- T14: commit mismatch denied ----------------------------------------
def test_t14_commit_mismatch_denied(self):
loader = self._loader()
with self.assertRaises(SourceUnavailableError):
loader.load("deadbeef" * 5, "object/test-obj-1", "docs/test.md")
# -- T15: hash mismatch denied ------------------------------------------
def test_t15_hash_mismatch_denied(self):
job = self._make_job(provenance_hash=_sha256("WRONG"))
self.store.create_job(job)
self.store.mark_ready(job["job_id"])
save = FakeTolariaSave("ok")
core = self._core(self._loader(), save)
result = core.process_job(job["job_id"], "w-t15")
self.assertEqual(result["state"], "FAILED")
self.assertEqual(result["result_code"], "PROVENANCE_MISMATCH")
self.assertEqual(len(save.calls), 0)
# -- T16: provenance mismatch denied ------------------------------------
def test_t16_provenance_mismatch_denied(self):
# object_id mismatch -> SourceMismatch -> SOURCE_UNAVAILABLE (fail-closed)
job = self._make_job(object_id="object/wrong-id")
self.store.create_job(job)
self.store.mark_ready(job["job_id"])
save = FakeTolariaSave("ok")
core = self._core(self._loader(), save)
result = core.process_job(job["job_id"], "w-t16")
self.assertEqual(result["state"], "FAILED")
self.assertEqual(len(save.calls), 0)
# -- T17: missing source denied -----------------------------------------
def test_t17_missing_source_denied(self):
job = self._make_job(vault_path="docs/missing.md")
self.store.create_job(job)
self.store.mark_ready(job["job_id"])
save = FakeTolariaSave("ok")
core = self._core(self._loader(), save)
result = core.process_job(job["job_id"], "w-t17")
self.assertEqual(result["state"], "FAILED")
self.assertEqual(result["result_code"], "SOURCE_UNAVAILABLE")
self.assertEqual(len(save.calls), 0)
# -- T18: Tolaria 401 fail-closed ---------------------------------------
def test_t18_tolaria_401_fail_closed(self):
job = self._make_job()
self.store.create_job(job)
self.store.mark_ready(job["job_id"])
save = FakeTolariaSave("fail") # AUTH_FAILURE
core = self._core(self._loader(), save)
result = core.process_job(job["job_id"], "w-t18")
self.assertEqual(result["state"], "FAILED")
self.assertEqual(result["result_code"], "AUTH_FAILURE")
# -- T19: Tolaria 403 fail-closed ---------------------------------------
def test_t19_tolaria_403_fail_closed(self):
# gleiche Semantik wie 401 -> AUTH_FAILURE
job = self._make_job()
self.store.create_job(job)
self.store.mark_ready(job["job_id"])
save = FakeTolariaSave("fail")
core = self._core(self._loader(), save)
result = core.process_job(job["job_id"], "w-t19")
self.assertEqual(result["state"], "FAILED")
# -- T20: timeout -> safe -----------------------------------------------
def test_t20_timeout_safe(self):
# Unerreichbare Base-URL -> TolariaUnavailable -> uncertain -> OUTCOME_UNKNOWN
job = self._make_job()
self.store.create_job(job)
self.store.mark_ready(job["job_id"])
client = TolariaClient(base_url=FAKE_BASE, save_token=TEST_TOKEN, timeout=0.1)
save = lambda p, c: (lambda r: {"status": "error", "code": "TOLARIA_UNAVAILABLE", "uncertain": True})(None)
core = self._core(self._loader(), save)
result = core.process_job(job["job_id"], "w-t20")
self.assertEqual(result["state"], "OUTCOME_UNKNOWN")
# -- T21: uncertain write -> OUTCOME_UNKNOWN ----------------------------
def test_t21_uncertain_write_outcome_unknown(self):
job = self._make_job()
self.store.create_job(job)
self.store.mark_ready(job["job_id"])
save = FakeTolariaSave("uncertain")
core = self._core(self._loader(), save)
result = core.process_job(job["job_id"], "w-t21")
self.assertEqual(result["state"], "OUTCOME_UNKNOWN")
# -- T22: no blind retry -------------------------------------------------
def test_t22_no_blind_retry(self):
# OUTCOME_UNKNOWN ist terminal im Worker (kein Retry im Loop)
job = self._make_job()
self.store.create_job(job)
self.store.mark_ready(job["job_id"])
save = FakeTolariaSave("uncertain")
core = self._core(self._loader(), save)
result = core.process_job(job["job_id"], "w-t22")
self.assertEqual(result["state"], "OUTCOME_UNKNOWN")
# Erneuter process_job auf OUTCOME_UNKNOWN -> kein Doppel-Write
result2 = core.process_job(job["job_id"], "w-t22")
self.assertEqual(len(save.calls), 1)
# -- T23: read-back mismatch not SUCCEEDED -------------------------------
def test_t23_readback_mismatch_not_succeeded(self):
job = self._make_job()
self.store.create_job(job)
self.store.mark_ready(job["job_id"])
save = FakeTolariaSave("ok")
rb = FakeReadBack("DIFFERENT", behavior="mismatch")
core = self._core(self._loader(), save, rb)
result = core.process_job(job["job_id"], "w-t23")
self.assertEqual(result["state"], "FAILED")
self.assertEqual(result["result_code"], "READBACK_MISMATCH")
# -- T24: audit REQUESTED/AUTHORIZED/EXECUTED ---------------------------
def test_t24_audit_events(self):
job = self._make_job()
self.store.create_job(job) # -> REQUESTED
self.store.mark_ready(job["job_id"])
save = FakeTolariaSave("ok")
rb = FakeReadBack(self.content)
core = self._core(self._loader(), save, rb)
core.process_job(job["job_id"], "w-t24")
events = self.store._conn.execute(
"SELECT event_type FROM audit_events WHERE job_id=? ORDER BY created_at",
(job["job_id"],)).fetchall()
types = [r["event_type"] for r in events]
self.assertIn("REQUESTED", types)
self.assertIn("AUTHORIZED", types)
self.assertIn("EXECUTED", types)
# -- T25: credential not logged -----------------------------------------
def test_t25_credential_not_logged(self):
# TolariaClient loggt den Token nie; Fehlermeldung enthält keinen Token
from tolaria_client import TolariaClientError
client = TolariaClient(base_url=FAKE_BASE, save_token=TEST_TOKEN)
try:
client.write("x.md", "c")
except TolariaClientError as e:
self.assertNotIn(TEST_TOKEN, str(e))
# -- T26: RQ cannot read SAVE credential --------------------------------
def test_t26_rq_cannot_read_save_credential(self):
# RQ hat keinen Zugriff auf TOLARIA_SAVE_TOKEN (nur Executor-ENV).
# Hier: Token ist nicht in der DB / nicht im Job / nicht im Store.
job = self._make_job()
self.store.create_job(job)
payload = json.dumps(job)
self.assertNotIn(TEST_TOKEN, payload)
# -- T27: DELETE executor cannot read SAVE credential -------------------
def test_t27_delete_executor_cannot_read_save_credential(self):
# SAVE-Credential ist nur im SAVE-Executor. DELETE-Executor hat es nicht.
# Hier: kein DELETE-Code im SAVE-Build (nur die 9 Runtime-Module prüfen).
import os
build_dir = os.path.dirname(os.path.abspath(__file__))
runtime_modules = [
"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",
]
for f in runtime_modules:
with open(os.path.join(build_dir, f)) as fh:
self.assertNotIn("delete_executor_core", fh.read())
class AUTH4C2PathContractTests(unittest.TestCase):
"""AUTH.4C2 PATH CONTRACT REPAIR — isolierte Path-Contract-Tests (T1-T20).
Verifiziert den kanonischen RELATIVEN Vault-Pfad-Contract durch die
gesamte SAVE-Pipeline. KEIN produktiver SAVE, KEIN neuer Job, KEIN DELETE.
"""
@classmethod
def setUpClass(cls):
cls.tmp = tempfile.mkdtemp(prefix="c5-4c2-path-test-")
cls.repo = _make_test_repo(cls.tmp)
cls.db = os.path.join(cls.tmp, "test_save.db")
cls.store = JobStore(cls.db, "SAVE")
cls.commit = subprocess.run(
["git", "-C", cls.repo, "rev-parse", "HEAD"],
capture_output=True, text=True, check=True).stdout.strip()
cls.content = "---\nid: object/test-obj-1\ntitle: Test\n---\n\nHello World\n"
cls.prov = _sha256(cls.content)
@classmethod
def tearDownClass(cls):
cls.store.close()
shutil.rmtree(cls.tmp, ignore_errors=True)
def _make_job(self, **overrides) -> Dict[str, Any]:
job = make_save_job(
job_id=overrides.get("job_id", str(uuid.uuid4())),
mission_id=overrides.get("mission_id", str(uuid.uuid4())),
object_id=overrides.get("object_id", "object/test-obj-1"),
vault_path=overrides.get("vault_path", "docs/test.md"),
source_commit=overrides.get("source_commit", self.commit),
provenance_hash=overrides.get("provenance_hash", self.prov),
expected_state="present",
created_at="2026-08-27T12:00:00Z",
idempotency_key=overrides.get("idempotency_key", str(uuid.uuid4())),
)
return job
def _loader(self, repo: Optional[str] = None) -> ForgejoSourceLoader:
return ForgejoSourceLoader(repo="test/test", clone_path=repo or self.repo,
forgejo_base="http://127.0.0.1:1")
def _core(self, loader, save, read_back=None) -> SaveExecutorCore:
return SaveExecutorCore(self.store, loader.load, save, read_back=read_back)
# -- T1: relative valid path accepted -----------------------------------
def test_t1_relative_valid_path_accepted(self):
job = self._make_job(vault_path="docs/test.md")
validated = validate_job(job)
self.assertEqual(validated["vault_path"], "docs/test.md")
# -- T2: absolute /app/vault path rejected ------------------------------
def test_t2_absolute_app_vault_path_rejected(self):
with self.assertRaises(JobRejectedError):
self._make_job(vault_path="/app/vault/docs/test.md")
# -- T3: absolute foreign path rejected ---------------------------------
def test_t3_absolute_foreign_path_rejected(self):
with self.assertRaises(JobRejectedError):
self._make_job(vault_path="/etc/passwd")
# -- T4: ../ rejected ----------------------------------------------------
def test_t4_dotdot_rejected(self):
with self.assertRaises(JobRejectedError):
self._make_job(vault_path="../secret")
# -- T5: encoded traversal rejected --------------------------------------
def test_t5_encoded_traversal_rejected(self):
with self.assertRaises(JobRejectedError):
self._make_job(vault_path="..%2fsecret")
# -- T6: backslash traversal rejected -----------------------------------
def test_t6_backslash_traversal_rejected(self):
with self.assertRaises(JobRejectedError):
self._make_job(vault_path="..\\secret")
# -- T7: empty path rejected --------------------------------------------
def test_t7_empty_path_rejected(self):
with self.assertRaises(JobRejectedError):
self._make_job(vault_path="")
# -- T8: URL rejected ----------------------------------------------------
def test_t8_url_rejected(self):
with self.assertRaises(JobRejectedError):
self._make_job(vault_path="http://evil/x.md")
# -- T9: relative path reaches Tolaria mock unchanged ---------------------
def test_t9_relative_path_reaches_tolaria_unchanged(self):
job = self._make_job(vault_path="docs/test.md")
self.store.create_job(job)
self.store.mark_ready(job["job_id"])
save = FakeTolariaSave("ok")
rb = FakeReadBack(self.content)
core = self._core(self._loader(), save, rb)
result = core.process_job(job["job_id"], "w-pc-t9")
self.assertEqual(result["state"], "SUCCEEDED")
self.assertEqual(len(save.calls), 1)
# Tolaria erhält den relativen Pfad UNVERÄNDERT (kein Strip/Rewrite)
self.assertEqual(save.calls[0][0], "docs/test.md")
# -- T10: read-back uses same relative path ------------------------------
def test_t10_readback_uses_same_relative_path(self):
job = self._make_job(vault_path="docs/test.md")
self.store.create_job(job)
self.store.mark_ready(job["job_id"])
save = FakeTolariaSave("ok")
rb = FakeReadBack(self.content)
core = self._core(self._loader(), save, rb)
core.process_job(job["job_id"], "w-pc-t10")
# read_back wurde mit demselben relativen Pfad aufgerufen
self.assertEqual(rb.calls, 1)
# -- T11: source loader loads exact relative path ------------------------
def test_t11_source_loader_loads_exact_relative_path(self):
loader = self._loader()
content = loader.load(self.commit, "object/test-obj-1", "docs/test.md")
self.assertEqual(content, self.content)
# -- T12: source commit binding preserved --------------------------------
def test_t12_source_commit_binding_preserved(self):
loader = self._loader()
with self.assertRaises(SourceUnavailableError):
loader.load("deadbeef" * 5, "object/test-obj-1", "docs/test.md")
# -- T13: provenance binding preserved -----------------------------------
def test_t13_provenance_binding_preserved(self):
job = self._make_job(provenance_hash=_sha256("WRONG"))
self.store.create_job(job)
self.store.mark_ready(job["job_id"])
save = FakeTolariaSave("ok")
core = self._core(self._loader(), save)
result = core.process_job(job["job_id"], "w-pc-t13")
self.assertEqual(result["state"], "FAILED")
self.assertEqual(result["result_code"], "PROVENANCE_MISMATCH")
self.assertEqual(len(save.calls), 0)
# -- T14: path mismatch rejected -----------------------------------------
def test_t14_path_mismatch_rejected(self):
# vault_path zeigt auf eine Datei, die im Commit nicht existiert
job = self._make_job(vault_path="docs/missing.md")
self.store.create_job(job)
self.store.mark_ready(job["job_id"])
save = FakeTolariaSave("ok")
core = self._core(self._loader(), save)
result = core.process_job(job["job_id"], "w-pc-t14")
self.assertEqual(result["state"], "FAILED")
self.assertEqual(result["result_code"], "SOURCE_UNAVAILABLE")
self.assertEqual(len(save.calls), 0)
# -- T15: symlink escape rejected (where fs resolution applies) ----------
def test_t15_symlink_escape_rejected(self):
# Executor-seitig: job_schema lehnt bereits alle nicht-relativen
# Pfade ab. Ein Symlink-Escape via relativer Pfad ist durch die
# Whitelist (kein .., kein Scheme) ausgeschlossen.
with self.assertRaises(JobRejectedError):
self._make_job(vault_path="docs/../../etc/passwd")
# -- T16: no hidden /app/vault strip behavior ----------------------------
def test_t16_no_hidden_app_vault_strip(self):
# Der relative Pfad wird NICHT gestrippt/rewritten — er erreicht
# Tolaria unverändert. Ein absoluter Pfad wird abgelehnt, nicht
# stillschweigend auf relativ reduziert.
with self.assertRaises(JobRejectedError):
self._make_job(vault_path="/app/vault/docs/test.md")
# -- T17: old failed job remains failed ---------------------------------
def test_t17_old_failed_job_remains_failed(self):
# Der historische FAILED-Job (2e0b7a0d...) bleibt unverändert.
# Hier: ein FAILED-Job wird nicht durch den Path-Repair berührt.
job = self._make_job()
self.store.create_job(job)
self.store.mark_ready(job["job_id"])
# Simuliere einen Fehlschlag (provenance mismatch)
job2 = self._make_job(provenance_hash=_sha256("WRONG"))
self.store.create_job(job2)
self.store.mark_ready(job2["job_id"])
save = FakeTolariaSave("ok")
core = self._core(self._loader(), save)
core.process_job(job2["job_id"], "w-pc-t17")
failed = self.store.get_job(job2["job_id"])
self.assertEqual(failed["state"], "FAILED")
# Der erste Job bleibt unberührt (kein neuer SAVE)
self.assertEqual(len(save.calls), 0)
# -- T18: no production job created --------------------------------------
def test_t18_no_production_job_created(self):
# Der Path-Repair erzeugt KEINE neuen Jobs. Nur die Test-Jobs
# dieser Klasse existieren in der temporären DB.
jobs = self.store.list_jobs()
for j in jobs:
self.assertNotIn("auth4c2-canary", j["vault_path"])
# -- T19: no production HTTP write ---------------------------------------
def test_t19_no_production_http_write(self):
# TolariaClient mit unerreichbarer Base-URL -> kein echter HTTP-Write.
# Der Path-Repair selbst führt keinen produktiven SAVE aus.
client = TolariaClient(base_url=FAKE_BASE, save_token=TEST_TOKEN)
self.assertEqual(client.base_url, FAKE_BASE.rstrip("/"))
# -- T20: no DELETE capability -------------------------------------------
def test_t20_no_delete_capability(self):
# SAVE-Build enthält keinen DELETE-Executor-Code.
build_dir = os.path.dirname(os.path.abspath(__file__))
runtime_modules = [
"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",
]
for f in runtime_modules:
with open(os.path.join(build_dir, f)) as fh:
self.assertNotIn("delete_executor_core", fh.read())
if __name__ == "__main__":
unittest.main(verbosity=2)

View file

@ -1,156 +0,0 @@
"""
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).
vault_path ist der KANONISCHE RELATIVE Pfad (AUTH.4C2 PATH CONTRACT
REPAIR, OPTION A) wird unverändert an Tolaria gesendet. Kein
/app/vault/-Prefix-Stripping, keine versteckte Rewrite-Logik.
"""
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.
vault_path ist der KANONISCHE RELATIVE Pfad (AUTH.4C2 PATH CONTRACT
REPAIR, OPTION A) wird unverändert an Tolaria gesendet. Kein
/app/vault/-Prefix-Stripping, keine versteckte Rewrite-Logik.
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

View file

@ -1,170 +0,0 @@
"""
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())