trading-system-docs/red-queen-architecture/control-plane/cp2a1/gate_evaluator.py

314 lines
11 KiB
Python

#!/usr/bin/env python3
"""
PRE_HERMES CP2A1 — Trusted Gate Evaluator (External Gate Enforcement Foundation).
STRICTLY NON-PRODUCTIVE / NO AUTONOMY. Reine Gate-Entscheidungs-Boundary.
Security-Prinzip:
RED_QUEEN_DECISION != AUTHORIZATION
RED_QUEEN_GATE_CHECK != TRUSTED_ENFORCEMENT
Die Security Boundary liegt AUSSERHALB des hermes-red-queen Containers.
Dieser Evaluator:
- liest ausschliesslich den autoritativen CP1-State (/opt/control-plane/state/*)
- liest current_boot_id frisch (/proc/sys/kernel/random/boot_id)
- verwendet die bestehende CP1-Evaluator-Logik (control_reader.py aus der SoT)
- berechnet den Effective-State
- liefert eine kleine deterministische Gate-Entscheidung (ALLOW/DENY)
- erzeugt Audit
- fuehrt KEINE produktive Aktion aus
- haelt KEINE mutierenden Credentials
- hat KEINE Executor-Integration
CP1-Logik-Reuse: importiert control_reader (SoT-Modul). KEINE Copy/Paste-Drift.
Die ENV-basierten Pfade (C5_CONTROL_STATE_DIR / C5_BOOT_ID_FILE) werden HART
UEBERSCHRIEBEN, damit RQ den Evaluator nicht dazu bringen kann, alternative
State-Pfade zu lesen (user-supplied state path = verboten).
Fail-closed: Jede Unklarheit -> DENY. Emergency unklar = ON.
"""
import json
import os
import sys
import time
import uuid
from datetime import datetime, timezone
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
# ---------------------------------------------------------------------------
# Hart verdrahtete autoritative Pfade (NICHT ENV-ueberschreibbar im Evaluator)
# ---------------------------------------------------------------------------
_DEFAULT_STATE_DIR = "/opt/control-plane/state"
_DEFAULT_BOOT_ID_FILE = "/proc/sys/kernel/random/boot_id"
AUTHORITATIVE_STATE_DIR = _DEFAULT_STATE_DIR
AUTHORITATIVE_BOOT_ID_FILE = _DEFAULT_BOOT_ID_FILE
AUDIT_DIR = "/audit"
# ---------------------------------------------------------------------------
# CP1-Logik-Reuse: SoT-Modul importieren und Pfade hart setzen
# ---------------------------------------------------------------------------
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
import control_reader # noqa: E402 (SoT-Modul, eine semantische Wahrheit)
# ENV-Ueberschreibbarkeit eliminieren: Pfade hart verdrahten.
control_reader.STATE_DIR = AUTHORITATIVE_STATE_DIR
control_reader.BOOT_ID_FILE = AUTHORITATIVE_BOOT_ID_FILE
# ---------------------------------------------------------------------------
# Statische Action-Class-Whitelist (NUR bekannte non-mutating Decision Classes)
# ---------------------------------------------------------------------------
ALLOWED_ACTION_CLASSES = frozenset(
{
"AUTONOMOUS_TICK_START",
"AUTONOMOUS_TICK_CONTINUE",
}
)
# Noch NICHT erlaubt (spaetere Phasen):
# SAVE, DELETE, FORGEJO_WRITE, NOTION_WRITE, TELEGRAM_SEND,
# HOST_MUTATION, EXTERNAL_MUTATION
# ---------------------------------------------------------------------------
# Reason Codes (deterministisch, KEINE freien LLM-Texte als Security Decision)
# ---------------------------------------------------------------------------
RC_ALLOW = "ALLOW"
RC_GLOBAL_AUTONOMY_OFF = "GLOBAL_AUTONOMY_OFF"
RC_EMERGENCY_ON = "EMERGENCY_ON"
RC_BOOT_ID_MISMATCH = "BOOT_ID_MISMATCH"
RC_STATE_MISSING = "STATE_MISSING"
RC_STATE_MALFORMED = "STATE_MALFORMED"
RC_BOOT_ID_ERROR = "BOOT_ID_ERROR"
RC_UNKNOWN_ACTION_CLASS = "UNKNOWN_ACTION_CLASS"
RC_INVALID_REQUEST = "INVALID_REQUEST"
RC_INTERNAL_ERROR = "INTERNAL_ERROR"
# Limits
MAX_REQUEST_BYTES = 4096
RATE_LIMIT_WINDOW_SEC = 60
RATE_LIMIT_MAX_REQUESTS = 120 # pro Window (RQ-Zugriff)
# ---------------------------------------------------------------------------
# Audit (append-only, gegen RQ-Manipulation geschuetzt: eigenes Volume)
# ---------------------------------------------------------------------------
def _now_iso():
return datetime.now(timezone.utc).isoformat()
def _audit(entry: dict):
"""Schreibt einen Audit-Eintrag append-only. Fail-open fuer Audit (kein Security-Decision)."""
try:
os.makedirs(AUDIT_DIR, exist_ok=True)
path = os.path.join(AUDIT_DIR, "gate_audit.log")
line = json.dumps(entry, sort_keys=True) + "\n"
with open(path, "a") as f:
f.write(line)
except Exception:
# Audit-Fehler darf die Gate-Entscheidung nicht beeinflussen (fail-open fuer Audit).
pass
# ---------------------------------------------------------------------------
# Rate Limiter (einfach, in-memory, pro Client-IP)
# ---------------------------------------------------------------------------
class _RateLimiter:
def __init__(self, window_sec, max_requests):
self.window_sec = window_sec
self.max_requests = max_requests
self._hits = {} # ip -> list[timestamp]
def allow(self, ip: str) -> bool:
now = time.time()
cutoff = now - self.window_sec
hits = [t for t in self._hits.get(ip, []) if t > cutoff]
if len(hits) >= self.max_requests:
self._hits[ip] = hits
return False
hits.append(now)
self._hits[ip] = hits
return True
_rate_limiter = _RateLimiter(RATE_LIMIT_WINDOW_SEC, RATE_LIMIT_MAX_REQUESTS)
# ---------------------------------------------------------------------------
# Gate-Entscheidung (deterministisch, fail-closed)
# ---------------------------------------------------------------------------
def _effective_state():
"""Berechnet den Effective-State via CP1-Logik. Fail-closed: Exception -> None."""
try:
return control_reader.read_control_state()
except Exception:
return None
def _tick_decision(state):
"""
AUTONOMOUS_TICK_START / CONTINUE = ALLOW nur wenn:
GLOBAL_AUTONOMY_EFFECTIVE == ON AND EMERGENCY_EFFECTIVE == OFF
Alles andere -> DENY mit deterministischem Reason Code.
"""
if state is None:
return {"result": "DENY", "reason_code": RC_INTERNAL_ERROR}
# Emergency: fail-closed, unklar = ON
emergency_eff = state.get("emergency_effective")
if emergency_eff != "OFF":
return {"result": "DENY", "reason_code": RC_EMERGENCY_ON}
# Boot-ID: muss vorhanden und plausibel sein (explizit, praeziser Reason)
current_boot_id = state.get("current_boot_id")
if not current_boot_id:
return {"result": "DENY", "reason_code": RC_BOOT_ID_ERROR}
# Global Autonomy: muss ON sein
global_eff = state.get("global_autonomy_effective")
if global_eff != "ON":
return {"result": "DENY", "reason_code": RC_GLOBAL_AUTONOMY_OFF}
return {"result": "ALLOW", "reason_code": RC_ALLOW}
def _check_action(action_class: str):
"""Validiert Action Class und liefert die Gate-Entscheidung."""
if action_class not in ALLOWED_ACTION_CLASSES:
return {"result": "DENY", "reason_code": RC_UNKNOWN_ACTION_CLASS}
state = _effective_state()
decision = _tick_decision(state)
# Audit-Eintrag
_audit(
{
"request_id": str(uuid.uuid4()),
"timestamp": _now_iso(),
"action_class": action_class,
"result": decision["result"],
"reason_code": decision["reason_code"],
"current_boot_id": (state or {}).get("current_boot_id"),
"effective_global": (state or {}).get("global_autonomy_effective"),
"effective_emergency": (state or {}).get("emergency_effective"),
}
)
return decision
# ---------------------------------------------------------------------------
# Strict JSON-Schema-Validierung (keine unbekannten Felder, keine Duplikate)
# ---------------------------------------------------------------------------
def _validate_check_request(raw_body: bytes):
"""
Validiert den /check-Request streng.
Erlaubt NUR: {"action_class": "<string>"}
Unbekannte Felder, Duplikate, null, falsche Typen -> INVALID_REQUEST.
Rueckgabe: (action_class | None, error_reason | None)
"""
if not raw_body or len(raw_body) > MAX_REQUEST_BYTES:
return None, RC_INVALID_REQUEST
try:
text = raw_body.decode("utf-8")
except UnicodeDecodeError:
return None, RC_INVALID_REQUEST
try:
obj = json.loads(text)
except json.JSONDecodeError:
return None, RC_INVALID_REQUEST
if not isinstance(obj, dict):
return None, RC_INVALID_REQUEST
# Nur das Feld "action_class" erlaubt
if set(obj.keys()) != {"action_class"}:
return None, RC_INVALID_REQUEST
action_class = obj.get("action_class")
if not isinstance(action_class, str):
return None, RC_INVALID_REQUEST
if not action_class.strip():
return None, RC_INVALID_REQUEST
return action_class, None
# ---------------------------------------------------------------------------
# HTTP Handler
# ---------------------------------------------------------------------------
class _Handler(BaseHTTPRequestHandler):
protocol_version = "HTTP/1.1"
def _send_json(self, status, payload):
body = json.dumps(payload, sort_keys=True).encode("utf-8")
self.send_response(status)
self.send_header("Content-Type", "application/json")
self.send_header("Content-Length", str(len(body)))
self.send_header("Cache-Control", "no-store")
self.end_headers()
self.wfile.write(body)
def _client_ip(self):
return self.client_address[0] if self.client_address else "unknown"
def do_GET(self):
ip = self._client_ip()
if not _rate_limiter.allow(ip):
self._send_json(429, {"error": "rate_limited"})
return
if self.path == "/health":
self._send_json(200, {"status": "ok", "service": "gate-evaluator"})
return
if self.path == "/effective":
state = _effective_state()
if state is None:
self._send_json(500, {"error": "internal_error"})
return
# Observability-Projektion (KEINE Autoritaet)
self._send_json(200, state)
return
self._send_json(404, {"error": "not_found"})
def do_POST(self):
ip = self._client_ip()
if not _rate_limiter.allow(ip):
self._send_json(429, {"error": "rate_limited"})
return
if self.path != "/check":
self._send_json(404, {"error": "not_found"})
return
# Request-Groesse begrenzen
try:
length = int(self.headers.get("Content-Length", "0"))
except ValueError:
self._send_json(400, {"result": "DENY", "reason_code": RC_INVALID_REQUEST})
return
if length <= 0 or length > MAX_REQUEST_BYTES:
self._send_json(400, {"result": "DENY", "reason_code": RC_INVALID_REQUEST})
return
raw_body = self.rfile.read(length)
action_class, err = _validate_check_request(raw_body)
if err is not None:
self._send_json(400, {"result": "DENY", "reason_code": err})
return
decision = _check_action(action_class)
status = 200 if decision["result"] == "ALLOW" else 403
self._send_json(status, decision)
def log_message(self, fmt, *args):
# Stille Logs (kein Request-Dump auf stdout)
pass
def main():
port = int(os.environ.get("GATE_EVALUATOR_PORT", "8080"))
server = ThreadingHTTPServer(("0.0.0.0", port), _Handler)
server.daemon_threads = True
server.serve_forever()
if __name__ == "__main__":
main()