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

431 lines
18 KiB
Python

"""
AUTH.3E — job_store.py
======================
SQLite Inbox für den Executor Command Channel.
Eigenschaften:
* Getrennte DB-Dateien pro Executor (c5a_save.db / c5a_delete.db)
* Atomarer Claim (READY -> CLAIMED) via SQLite-Transaktion (kein TOCTOU)
* Lease (lease_until) für Crash-Recovery
* Idempotenz (job_id / idempotency_key UNIQUE)
* Immutable fields after claim (UPDATE verboten für sicherheitskritische Felder)
* FAIL CLOSED bei DB-Fehler
Isoliert implementiert (KEIN produktiver Container). Nutzt job_schema + job_state_machine.
"""
from __future__ import annotations
import json
import sqlite3
import time
import uuid
from typing import Any, Dict, List, Optional
from job_schema import (
JOB_TYPE_DELETE,
JOB_TYPE_SAVE,
JobRejectedError,
validate_job,
)
from job_state_machine import (
ST_CLAIMED,
ST_CREATED,
ST_EXECUTING,
ST_FAILED,
ST_OUTCOME_UNKNOWN,
ST_READY,
ST_RECONCILED,
ST_REJECTED,
ST_SUCCEEDED,
InvalidTransitionError,
transition,
)
from job_claim import (
JobAlreadyClaimedError,
JobNotFoundError,
atomic_claim,
recover_stale_claims as _recover_stale_claims,
renew_lease as _renew_lease,
)
# ---------------------------------------------------------------------------
# Immutable fields after claim (UPDATE verboten)
# ---------------------------------------------------------------------------
IMMUTABLE_FIELDS = frozenset({
"job_type", "object_id", "vault_path", "approval_id",
"expected_commit", "expected_provenance_hash", "delete_request_id",
"source_commit", "provenance_hash", "mission_id",
})
class JobStoreError(Exception):
"""Basis-Fehler für JobStore."""
class JobImmutableFieldError(JobStoreError):
pass
class JobStore:
"""
SQLite Inbox. Ein Store pro Executor (SAVE oder DELETE).
worker_scope: "SAVE" oder "DELETE" — bestimmt, welche Job-Types claimbar sind.
"""
def __init__(self, db_path: str, worker_scope: str):
self.db_path = db_path
self.worker_scope = worker_scope.upper()
if self.worker_scope not in ("SAVE", "DELETE"):
raise JobStoreError(f"invalid worker_scope: {worker_scope!r}")
self._conn = sqlite3.connect(db_path)
self._conn.row_factory = sqlite3.Row
self._init_schema()
def _init_schema(self) -> None:
self._conn.execute("""
CREATE TABLE IF NOT EXISTS jobs (
job_id TEXT PRIMARY KEY,
job_version INTEGER NOT NULL,
mission_id TEXT NOT NULL,
job_type TEXT NOT NULL,
object_id TEXT NOT NULL,
vault_path TEXT NOT NULL,
payload TEXT NOT NULL, -- vollständiger Job (JSON)
state TEXT NOT NULL,
idempotency_key TEXT NOT NULL UNIQUE,
worker_id TEXT,
claim_id TEXT,
claimed_at INTEGER,
lease_until INTEGER,
attempt_count INTEGER NOT NULL DEFAULT 0,
result_code TEXT,
created_at TEXT NOT NULL,
updated_at INTEGER NOT NULL
)
""")
self._conn.execute("""
CREATE INDEX IF NOT EXISTS idx_jobs_state ON jobs(state)
""")
self._conn.execute("""
CREATE INDEX IF NOT EXISTS idx_jobs_type ON jobs(job_type)
""")
# Audit-Ereignis-Log (REQUESTED / AUTHORIZED / EXECUTED).
# Trennt die Phasen: RQ hat Job erzeugt (REQUESTED), Executor hat
# Job + Approval validiert (AUTHORIZED), Mutation ausgeführt (EXECUTED).
self._conn.execute("""
CREATE TABLE IF NOT EXISTS audit_events (
event_id TEXT PRIMARY KEY,
job_id TEXT NOT NULL,
event_type TEXT NOT NULL,
worker_id TEXT,
result_code TEXT,
created_at INTEGER NOT NULL
)
""")
self._conn.execute("""
CREATE INDEX IF NOT EXISTS idx_audit_job ON audit_events(job_id)
""")
self._conn.commit()
# -- Job-Type-Scope -----------------------------------------------------
def _job_type_allowed(self, job_type: str) -> bool:
"""SAVE-Executor claimt nur SAVE-Jobs; DELETE-Executor nur DELETE-Jobs."""
if self.worker_scope == "SAVE":
return job_type == JOB_TYPE_SAVE
return job_type == JOB_TYPE_DELETE
# -- Erzeugen (RQ-Seite) ------------------------------------------------
def create_job(self, job: Dict[str, Any]) -> Dict[str, Any]:
"""
Erzeugt einen Job (Status CREATED). Idempotent per job_id.
Wirft JobRejectedError bei Schema-Verletzung.
"""
validated = validate_job(job)
job_id = validated["job_id"]
now = int(time.time() * 1000)
# Idempotenz: gleiche job_id -> bestehenden Job zurückgeben (kein Fehler).
existing = self.get_job(job_id)
if existing is not None:
return existing
try:
self._conn.execute(
"""
INSERT INTO jobs
(job_id, job_version, mission_id, job_type, object_id,
vault_path, payload, state, idempotency_key, created_at, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""",
(
job_id, validated["job_version"], validated["mission_id"],
validated["job_type"], validated["object_id"],
validated["vault_path"], json.dumps(validated),
ST_CREATED, validated["idempotency_key"],
validated["created_at"], now,
),
)
self._conn.commit()
except sqlite3.IntegrityError as e:
# UNIQUE-Verletzung (idempotency_key bereits vergeben) -> Anomalie.
raise JobStoreError(f"create_job failed (idempotency_key collision): {e}") from e
# Audit: RQ hat den Job erzeugt (REQUESTED). Kein REQUESTED-Eintrag ist
# eine Autorisierung — nur AUTHORIZED (nach Validierung) autorisiert.
self._record_audit(job_id, "REQUESTED")
return self.get_job(job_id)
def get_job(self, job_id: str) -> Optional[Dict[str, Any]]:
row = self._conn.execute(
"SELECT * FROM jobs WHERE job_id = ?", (job_id,)
).fetchone()
return self._row_to_dict(row) if row else None
def get_job_by_idempotency_key(self, idempotency_key: str) -> Optional[Dict[str, Any]]:
row = self._conn.execute(
"SELECT * FROM jobs WHERE idempotency_key = ?", (idempotency_key,)
).fetchone()
return self._row_to_dict(row) if row else None
def list_jobs(self, state: Optional[str] = None) -> List[Dict[str, Any]]:
if state:
rows = self._conn.execute(
"SELECT * FROM jobs WHERE state = ? ORDER BY created_at", (state,)
).fetchall()
else:
rows = self._conn.execute(
"SELECT * FROM jobs ORDER BY created_at"
).fetchall()
return [self._row_to_dict(r) for r in rows]
# -- State-Transition (mit Immutable-Fields-Schutz) ---------------------
def _record_audit(self, job_id: str, event_type: str, *,
worker_id: Optional[str] = None,
result_code: Optional[str] = None) -> None:
"""Schreibt ein Audit-Ereignis (REQUESTED/AUTHORIZED/EXECUTED)."""
now = int(time.time() * 1000)
self._conn.execute(
"""
INSERT INTO audit_events
(event_id, job_id, event_type, worker_id, result_code, created_at)
VALUES (?, ?, ?, ?, ?, ?)
""",
(str(uuid.uuid4()), job_id, event_type, worker_id, result_code, now),
)
self._conn.commit()
def record_audit_event(self, job_id: str, event_type: str, *,
worker_id: Optional[str] = None,
result_code: Optional[str] = None) -> None:
"""Öffentlicher Audit-Eintrag (AUTHORIZED/EXECUTED) für Executor-Cores."""
self._record_audit(job_id, event_type, worker_id=worker_id,
result_code=result_code)
def _transition(self, job_id: str, to_state: str, *, result_code: Optional[str] = None,
worker_id: Optional[str] = None) -> Dict[str, Any]:
"""
Führt eine State-Transition aus. Wirft InvalidTransitionError bei
ungültiger Transition (FAIL CLOSED).
"""
job = self.get_job(job_id)
if job is None:
raise JobNotFoundError(f"job not found: {job_id}")
from_state = job["state"]
transition(from_state, to_state) # wirft bei ungültiger Transition
# AUTH.4C3: Immutable-Fields-Guard. Sicherheitsrelevante Felder dürfen
# nach OUTCOME_UNKNOWN nicht mehr verändert werden. Der Guard verifiziert
# die immutable Felder im persistierten Job gegen den Payload (der beim
# Claim/Execution nicht verändert wird). Falls ein Feld fehlt oder
# inkonsistent ist -> FAIL CLOSED (kein State-Change).
self._assert_immutable_fields(job)
now = int(time.time() * 1000)
self._conn.execute(
"""
UPDATE jobs SET state = ?, result_code = ?, updated_at = ?,
worker_id = COALESCE(?, worker_id)
WHERE job_id = ?
""",
(to_state, result_code, now, worker_id, job_id),
)
self._conn.commit()
return self.get_job(job_id)
def mark_ready(self, job_id: str) -> Dict[str, Any]:
return self._transition(job_id, ST_READY)
def mark_rejected(self, job_id: str, result_code: str) -> Dict[str, Any]:
return self._transition(job_id, ST_REJECTED, result_code=result_code)
def mark_succeeded(self, job_id: str, worker_id: str) -> Dict[str, Any]:
return self._transition(job_id, ST_SUCCEEDED, worker_id=worker_id)
def mark_failed(self, job_id: str, result_code: str, worker_id: str) -> Dict[str, Any]:
return self._transition(job_id, ST_FAILED, result_code=result_code, worker_id=worker_id)
def mark_outcome_unknown(self, job_id: str, worker_id: str) -> Dict[str, Any]:
return self._transition(job_id, ST_OUTCOME_UNKNOWN, worker_id=worker_id)
def mark_reconciled(self, job_id: str, worker_id: str) -> Dict[str, Any]:
"""AUTH.4C3: OUTCOME_UNKNOWN -> RECONCILED (read-only Reconciliation)."""
return self._transition(job_id, ST_RECONCILED, worker_id=worker_id)
def set_result_code(self, job_id: str, result_code: str) -> Dict[str, Any]:
"""
AUTH.4C3: Setzt NUR den result_code eines Jobs (Diagnose-Information),
OHNE State-Transition. Der State bleibt unverändert (z.B. OUTCOME_UNKNOWN).
Kein Mutation-Retry, kein State-Change. Verbessert Operator-Visibility:
der Operator sieht, warum die Reconciliation nicht erfolgreich war.
"""
job = self.get_job(job_id)
if job is None:
raise JobNotFoundError(f"job not found: {job_id}")
now = int(time.time() * 1000)
self._conn.execute(
"UPDATE jobs SET result_code = ?, updated_at = ? WHERE job_id = ?",
(result_code, now, job_id),
)
self._conn.commit()
return self.get_job(job_id)
# -- Immutable-Fields-Guard (AUTH.4C3) ----------------------------------
def _assert_immutable_fields(self, job: Dict[str, Any]) -> None:
"""
Verifiziert, dass die sicherheitsrelevanten immutable Felder im
persistierten Job konsistent sind. Prüft die rohen DB-Spalten
(object_id, vault_path) gegen das Payload-JSON. Falls ein immutable Feld
im Payload von seinem DB-Spalten-Wert abweicht -> FAIL CLOSED
(JobImmutableFieldError, kein State-Change). Verhindert, dass ein
manipulierter Job (z.B. geänderte object_id/vault_path) eine Transition
durchläuft.
Hinweis: _row_to_dict merged das Payload in das Top-Level-Dict, daher
ist ein Payload-vs-Top-Level-Vergleich wirkungslos. Stattdessen wird
die rohe DB-Spalte direkt gelesen und gegen das Payload geprüft.
"""
payload = job.get("payload")
if not isinstance(payload, dict):
return
job_id = job.get("job_id")
if not job_id:
return
row = self._conn.execute(
"SELECT object_id, vault_path FROM jobs WHERE job_id = ?", (job_id,)
).fetchone()
if row is None:
return
db_object_id, db_vault_path = row
# object_id: DB-Spalte muss mit Payload übereinstimmen
if "object_id" in payload and payload["object_id"] != db_object_id:
raise JobImmutableFieldError(
f"job {job_id} immutable field drift: object_id")
# vault_path: DB-Spalte muss mit Payload übereinstimmen
if "vault_path" in payload and payload["vault_path"] != db_vault_path:
raise JobImmutableFieldError(
f"job {job_id} immutable field drift: vault_path")
# -- Claim / Lease ------------------------------------------------------
def claim_job(self, job_id: str, worker_id: str, lease_seconds: int = 60) -> Dict[str, Any]:
"""
Atomarer Claim: READY -> CLAIMED, nur wenn Lease abgelaufen oder nie gesetzt.
Kein TOCTOU (SQLite-Transaktion mit Status-Bedingung). Delegiert an job_claim.
"""
claimed = atomic_claim(self._conn, job_id, worker_id, lease_seconds)
return claimed
def begin_execution(self, job_id: str, worker_id: str) -> Dict[str, Any]:
"""CLAIMED -> EXECUTING (Claim bestätigt)."""
return self._transition(job_id, ST_EXECUTING, worker_id=worker_id)
def renew_lease(self, job_id: str, lease_seconds: int = 60) -> Dict[str, Any]:
"""Verlängert die Lease eines CLAIMED/EXECUTING-Jobs."""
_renew_lease(self._conn, job_id, lease_seconds)
return self.get_job(job_id)
# -- Crash-Recovery -----------------------------------------------------
def recover_stale_claims(self, worker_id: str, lease_seconds: int = 60) -> List[Dict[str, Any]]:
"""
Findet stale CLAIMED-Jobs (Lease abgelaufen) und überführt sie:
* SAVE: zurück zu READY (wieder claimbar, RETRYABLE)
* DELETE: zu OUTCOME_UNKNOWN (NICHT blind wiederholen)
"""
# job_claim.recover_stale_claims überführt SAVE->READY, DELETE->OUTCOME_UNKNOWN
# anhand des job_type. Wir rufen es pro Job-Type auf (Scope-getrennt).
recovered = []
for job_type in (JOB_TYPE_SAVE, JOB_TYPE_DELETE):
if self._job_type_allowed(job_type):
_recover_stale_claims(self._conn, job_type)
# Re-read recovered jobs
for row in self._conn.execute(
"SELECT job_id FROM jobs WHERE state IN (?, ?)",
(ST_READY, ST_OUTCOME_UNKNOWN),
).fetchall():
job = self.get_job(row["job_id"])
if job is not None:
recovered.append(job)
return recovered
# -- Idempotenz ---------------------------------------------------------
def is_duplicate(self, job_id: str, idempotency_key: str) -> bool:
"""True, wenn job_id ODER idempotency_key bereits existiert."""
row = self._conn.execute(
"SELECT 1 FROM jobs WHERE job_id = ? OR idempotency_key = ? LIMIT 1",
(job_id, idempotency_key),
).fetchone()
return row is not None
# -- Audit --------------------------------------------------------------
def audit_trail(self, job_id: str) -> List[Dict[str, Any]]:
"""
Gibt den Audit-Trail eines Jobs zurück: die chronologische Liste der
Audit-Ereignisse (REQUESTED / AUTHORIZED / EXECUTED) plus den aktuellen
Job-Snapshot. Ein REQUESTED-Eintrag ist KEINE Autorisierung — nur
AUTHORIZED (nach AUTH.3D-Validierung) autorisiert eine Mutation.
"""
events = self._conn.execute(
"SELECT * FROM audit_events WHERE job_id = ? ORDER BY created_at",
(job_id,),
).fetchall()
trail = [{
"event_id": e["event_id"],
"job_id": e["job_id"],
"event_type": e["event_type"],
"worker_id": e["worker_id"],
"result_code": e["result_code"],
"created_at": e["created_at"],
} for e in events]
job = self.get_job(job_id)
if job is not None:
trail.append({
"job_id": job["job_id"],
"mission_id": job["mission_id"],
"job_type": job["job_type"],
"object_id": job["object_id"],
"state": job["state"],
"worker_id": job["worker_id"],
"attempt_count": job["attempt_count"],
"result_code": job["result_code"],
"created_at": job["created_at"],
"updated_at": job["updated_at"],
})
return trail
# -- Helpers ------------------------------------------------------------
def _row_to_dict(self, row: sqlite3.Row) -> Dict[str, Any]:
d = dict(row)
payload = json.loads(d["payload"])
# Merge Payload-Felder in das Top-Level-Dict, damit Executor-Cores
# auf approval_id/expected_commit/source_commit etc. zugreifen können.
d.update(payload)
d["payload"] = payload
return d
def close(self) -> None:
self._conn.close()