170 lines
6 KiB
Python
170 lines
6 KiB
Python
"""
|
|
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())
|