trading-system-docs/a3/rq_safety.py

1055 lines
44 KiB
Python

#!/usr/bin/env python3
"""
Red Queen — A3: Deterministic Safety Layer (V1).
Zwischen Mission/WP-State (A2 `missions.db`) und dem spaeteren Orchestrator /
Action/Retry/Delegation sitzt dieser deterministische Safety-Layer. Er entscheidet
deterministisch (Zaehler, Limits, Tabellen) und NICHT ueber 'LLM-Lust':
CONTINUE / RETRY / DEBUG / SECOND_OPINION / BLOCK / ESCALATE / CIRCUIT_BREAK
Ziel: Red Queen erhaelt die Safety-Logik, BEVOR spaetere autonome Mission-Loops
existieren. Nach A3 gibt es KEINE autonome Orchestrierung — kein dispatcher,
kein loop, kein heartbeat, kein cron, kein self-improvement. Dieses Modul ist
eine reine Library / Safety-Capability.
Grundsaetze (A1 SAFETY_CONTRACT + A3-Spezifikation):
* DETERMINISTISCH : harte Entscheidung aus Zaehlern/Limits/Tabellen in `safety.db`.
* IDEMPOTENT : derselbe Attempt/Event/Trigger wird nie doppelt gezaehlt.
* RESTART-STABIL : Zustand aus DB wiederherstellbar; nichts wird beim Restart vergessen.
* FAIL-CLOSED : unbekannter/inkonsistenter Safety-State -> SAFETY_STATE_ERROR,
keine Mutation, Evidence wird gesichert.
* APPEND-ONLY : Attempt-Ledger ist unverstaendlich, kein Loeschen/Ueberschreiben.
* KEINE SECRETS : im Ledger/Events/Evidence/Telegram keine Tokens/Passwoerter.
DB-Entscheidung (Maker-Entscheidung, siehe README): EIGENE `safety.db`.
- A2 `missions.db` bleibt voellig unangetastet -> maximale Rollback-Faehigkeit fuer A2.
- Safety-State (Circuit, Attempts, Events, Evidence) ist unabhängig vom Mission-State.
- Kein ALTER auf bestehenden A2-Tabellen; nur neue Tabellen in eigener DB.
- Beide DBs werden immer ueber waehlbare Pfade getestet (temp), nie produktiv beruehrt.
"""
from __future__ import annotations
import datetime
import hashlib
import json
import re
import sqlite3
from typing import Any, Callable, Dict, List, Optional, Tuple
# --------------------------------------------------------------------------- #
# Fehler-Codes (maschinenlesbar, stabil) — spiegelt A2 `RqError`-Kontrakt
# --------------------------------------------------------------------------- #
class SafetyError(Exception):
"""Geworfener, maschinenlesbarer Safety-Fehler mit stabilem `code`."""
def __init__(self, code: str, message: str, detail: Optional[Dict[str, Any]] = None):
super().__init__(message)
self.code = code
self.message = message
self.detail = detail or {}
def to_dict(self) -> Dict[str, Any]:
return {"code": self.code, "message": self.message, "detail": self.detail}
def _err(code: str, msg: str, **detail) -> SafetyError:
return SafetyError(code, msg, detail)
# --------------------------------------------------------------------------- #
# Konstanten (Verbindliche, konservative Limits aus A1 §2)
# --------------------------------------------------------------------------- #
MAX_MAKER_CHECKER_REPAIRS = 3 # Maker->Checker Repair MAX 3
MAX_SAME_ERROR_SIGNATURE = 2 # Gleiche Error-Signatur MAX 2
MAX_ITERATIONS = 50 # AEUSSERSTE Runtime-Notbremse (nicht operativ)
OSCILLATION_SIGNATURE_REPEATS = 3 # gleiche Signatur >=3 -> Oscillation-Verdacht (A1 §4.1)
OSCILLATION_TARGET_CHANGES = 3 # gleiche Datei/Change-Ziel >=3 ohne Progress
ABAB_PATTERN_LEN = 4 # A-B-A-B
NO_PROGRESS_FAILS_LIMIT = 2 # >=2 FAIL ohne messbaren Fortschritt -> Debug
# Zustaende & Entscheidungen
CIRCUIT_CLOSED = "CLOSED"
CIRCUIT_OPEN = "OPEN"
CIRCUIT_STATES = frozenset({CIRCUIT_CLOSED, CIRCUIT_OPEN})
SEVERITIES = frozenset({"INFO", "WARNING", "HIGH", "CRITICAL"})
SCOPE_GLOBAL = "GLOBAL"
SCOPE_MISSION = "MISSION"
SCOPE_WORK_PACKAGE = "WORK_PACKAGE"
SCOPE_COMPONENT = "COMPONENT"
SCOPE_TYPES = frozenset({SCOPE_GLOBAL, SCOPE_MISSION, SCOPE_WORK_PACKAGE, SCOPE_COMPONENT})
DECISIONS = frozenset(
{
"CONTINUE",
"RETRY",
"DEBUG",
"SECOND_OPINION",
"BLOCK",
"ESCALATE",
"CIRCUIT_BREAK",
}
)
ALLOWED_CONTINUE = "CONTINUE"
ALLOWED_RETRY = "RETRY"
ALLOWED_STRATEGY_CHANGE = "STRATEGY_CHANGE"
ALLOWED_DEBUG = "DEBUG"
ALLOWED_SECOND_OPINION = "SECOND_OPINION"
ALLOWED_READ_ONLY_DIAGNOSIS = "READ_ONLY_DIAGNOSIS"
ALLOWED_HUMAN_GATE = "HUMAN_GATE"
ALLOWED_NONE = "NONE"
# Reason Codes (§22) — stabile, maschinenlesbare Codes
REASON_RETRY_AVAILABLE = "RETRY_AVAILABLE"
REASON_RETRY_LIMIT = "RETRY_LIMIT"
REASON_SAME_ERROR_LIMIT = "SAME_ERROR_LIMIT"
REASON_FAILED_STRATEGY_REPEAT = "FAILED_STRATEGY_REPEAT"
REASON_NO_MEASURABLE_PROGRESS = "NO_MEASURABLE_PROGRESS"
REASON_OSCILLATION_ABAB = "OSCILLATION_ABAB"
REASON_CIRCUIT_ALREADY_OPEN = "CIRCUIT_ALREADY_OPEN"
REASON_STATE_INCONSISTENT = "STATE_INCONSISTENT"
REASON_CRITICAL_TRIGGER = "CRITICAL_TRIGGER"
REASON_HUMAN_GATE_REQUIRED = "HUMAN_GATE_REQUIRED"
REASON_ITERATION_LIMIT = "ITERATION_LIMIT"
# Event-Typen (§17)
EVENT_RETRY_ALLOWED = "RETRY_ALLOWED"
EVENT_RETRY_DENIED = "RETRY_DENIED"
EVENT_DEBUG_REQUIRED = "DEBUG_REQUIRED"
EVENT_SECOND_OPINION_REQUIRED = "SECOND_OPINION_REQUIRED"
EVENT_RETRY_LIMIT_REACHED = "RETRY_LIMIT_REACHED"
EVENT_ERROR_SIGNATURE_REPEAT = "ERROR_SIGNATURE_REPEAT"
EVENT_STRATEGY_REPEAT = "STRATEGY_REPEAT"
EVENT_NO_PROGRESS = "NO_PROGRESS"
EVENT_OSCILLATION_DETECTED = "OSCILLATION_DETECTED"
EVENT_CIRCUIT_OPENED = "CIRCUIT_OPENED"
EVENT_CIRCUIT_RESET_REQUESTED = "CIRCUIT_RESET_REQUESTED"
EVENT_CIRCUIT_CLOSED = "CIRCUIT_CLOSED"
EVENT_SAFETY_STATE_ERROR = "SAFETY_STATE_ERROR"
EVENT_ESCALATION_REQUIRED = "ESCALATION_REQUIRED"
EVENT_HUMAN_DECISION_REQUIRED = "HUMAN_DECISION_REQUIRED"
# Critical-Trigger (§12 / A1 §5.1) die einen GLOBAL-Scope-Circuit oeffnen
GLOBAL_SCOPE_TRIGGERS = frozenset(
{
"SSH_RECOVERY_JEOPARDIZED",
"SECRET_EXPOSURE_UNKNOWN_SCOPE",
"PERSISTENCE_CORRUPTED",
"IDENTITY_AUTH_MISMATCH_CRITICAL",
"RUNTIME_INTEGRITY_JEOPARDIZED",
}
)
# Ergebnis-/Fortschritt-Werte
RESULT_PASS = "PASS"
RESULT_FAIL = "FAIL"
RESULT_UNKNOWN = "UNKNOWN"
RESULT_VALUES = frozenset({RESULT_PASS, RESULT_FAIL, RESULT_UNKNOWN})
PROGRESS_YES = "YES"
PROGRESS_NO = "NO"
PROGRESS_UNKNOWN = "UNKNOWN"
PROGRESS_VALUES = frozenset({PROGRESS_YES, PROGRESS_NO, PROGRESS_UNKNOWN})
# --------------------------------------------------------------------------- #
# DB-Helfer (A2-Konvention)
# --------------------------------------------------------------------------- #
def _connect(path: str) -> sqlite3.Connection:
conn = sqlite3.connect(path)
conn.row_factory = sqlite3.Row
conn.execute("PRAGMA journal_mode=WAL")
conn.execute("PRAGMA foreign_keys=ON")
return conn
def _utcnow() -> str:
return datetime.datetime.now(datetime.timezone.utc).isoformat(timespec="seconds")
# --------------------------------------------------------------------------- #
# Secret-Redaction (§28) — bewusst einfach & konservativ
# --------------------------------------------------------------------------- #
_SECRET_PATTERNS = [
re.compile(r"(?i)(api[_-]?key|secret|token|password|passwd|pwd|pw|authorization|bearer|access[_-]?key)\s*[=:]\s*[^\s,;]+"),
re.compile(r"(?i)https?://[^\s/@]+:[^\s/@]+@"),
re.compile(r"(?i)(-----BEGIN[ A-Z]*PRIVATE KEY-----.*?-----END[ A-Z]*PRIVATE KEY-----)", re.DOTALL),
re.compile(r"(?i)\b(?:ghp|gho|ghu|ghs|sk-|xox[baprs]-)[a-z0-9_-]{10,}\b"),
re.compile(r"(?i)((?:password|secret|token|key)\s*[=:]\s*)(?:['\"]?)[^\s'\",;]+"),
]
def redact_secret(text: Optional[str]) -> Optional[str]:
"""Maskiert credential-artige Werte. Leere/None-Eingabe bleibt unveraendert."""
if not text:
return text
out = str(text)
for pat in _SECRET_PATTERNS:
out = pat.sub("REDACTED", out)
return out
# --------------------------------------------------------------------------- #
# Error-Signature Detection (§7)
# --------------------------------------------------------------------------- #
_VOLATILE_RE = [
(re.compile(r"\b0x[0-9a-fA-F]{4,}\b"), "<HEX>"),
(re.compile(r"\b(pid|ppid)\s*[=:]\s*\d+\b"), r"\1=<PID>"),
(re.compile(r"\b\d{4}-\d{2}-\d{2}[T ]\d{2}:\d{2}:\d{2}"), "<TS>"),
(re.compile(r"\b[A-Fa-f0-9]{32}\b"), "<H32>"),
(re.compile(r"\b[A-Fa-f0-9]{64}\b"), "<H64>"),
(re.compile(r"\b[A-Fa-f0-9]{8}-[A-Fa-f0-9]{4}-[A-Fa-f0-9]{4}-[A-Fa-f0-9]{4}-[A-Fa-f0-9]{12}\b"), "<UUID>"),
(re.compile(r"(?i)(port\s*[=:]\s*)\d{2,5}"), r"\1<PORT>"),
(re.compile(r"(?i)(session|request|trace|run)[_-]?id[^:,;]*\s*[=:]\s*[^\s,;]+"), r"\1=<ID>"),
]
def normalize_error(raw: Optional[str]) -> str:
"""Normalisiert volatile Bestandteile zu stabilen Platzhaltern.
NICHT so aggressiv, dass unterschiedliche Fehler faelschlich gleich werden.
Es werden nur klar volatile, generische Kategorien ersetzt; der Rest bleibt.
"""
if not raw:
return ""
s = str(raw)
for pat, repl in _VOLATILE_RE:
s = pat.sub(repl, s)
return s.strip()
def signature_hash(normalized: str) -> str:
"""Deterministischer SHA-256-Hash (gekuerzt auf 16 hex) einer normalisierten Signatur."""
return hashlib.sha256((normalized or "").encode("utf-8")).hexdigest()[:16]
def error_signature(raw: Optional[str]) -> Tuple[str, str]:
"""Liefert (NORMALIZED_SIGNATURE, SIGNATURE_HASH) fuer einen rohen Fehler."""
norm = normalize_error(raw)
return norm, signature_hash(norm)
# --------------------------------------------------------------------------- #
# Strategy Fingerprinting (§8) — deterministische Merkmale
# --------------------------------------------------------------------------- #
def strategy_fingerprint(
operation_type: str,
target_component: str,
target_files: Optional[List[str]] = None,
actions_class: Optional[str] = None,
normalized_strategy: str = "",
intended_change: Optional[str] = None,
) -> Tuple[str, str]:
"""Deterministisches Merkmals-Tupel (LABEL, HASH) einer Strategie.
Erkennt, ob eine bereits gescheiterte Strategie im Wesentlichen erneut
angewendet wird. Files werden sortiert, damit die Reihenfolge egal ist.
"""
files = sorted(set(target_files or []))
parts = [
"op=" + (operation_type or ""),
"tgt=" + (target_component or ""),
"files=" + ",".join(files),
"act=" + (actions_class or ""),
"strat=" + (normalized_strategy or ""),
"change=" + (intended_change or ""),
]
label = "|".join(parts)
h = hashlib.sha256(label.encode("utf-8")).hexdigest()[:16]
return label, h
# --------------------------------------------------------------------------- #
# SafetyStore
# --------------------------------------------------------------------------- #
class SafetyStore:
"""Persistenter, deterministischer Safety-Layer auf einer eigenen `safety.db`."""
def __init__(self, db_path: str, now_fn: Callable[[], str] = _utcnow):
self.db_path = db_path
self._now = now_fn
self._ensure_schema()
# -- Schema ---------------------------------------------------------------
def _ensure_schema(self) -> None:
conn = _connect(self.db_path)
try:
conn.executescript(
"""
CREATE TABLE IF NOT EXISTS attempts (
id TEXT PRIMARY KEY,
idempotency_key TEXT,
timestamp TEXT NOT NULL,
mission_id TEXT,
wp_id TEXT,
actor TEXT,
task TEXT,
hypothesis TEXT,
strategy TEXT,
strategy_fingerprint TEXT,
strategy_fingerprint_hash TEXT,
change TEXT,
result TEXT,
error_raw TEXT,
error_signature TEXT,
error_signature_hash TEXT,
progress_metric TEXT,
progress_before TEXT,
progress_after TEXT,
progress_delta TEXT,
progress TEXT
);
CREATE INDEX IF NOT EXISTS idx_att_miss ON attempts(mission_id);
CREATE INDEX IF NOT EXISTS idx_att_wp ON attempts(wp_id);
CREATE INDEX IF NOT EXISTS idx_att_sig ON attempts(error_signature_hash);
CREATE INDEX IF NOT EXISTS idx_att_fp ON attempts(strategy_fingerprint_hash);
CREATE TABLE IF NOT EXISTS circuit_state (
scope_type TEXT NOT NULL,
scope_id TEXT NOT NULL,
state TEXT NOT NULL,
severity TEXT NOT NULL,
trigger TEXT,
reason TEXT,
evidence_ref TEXT,
opened_at TEXT,
closed_at TEXT,
updated_at TEXT NOT NULL,
PRIMARY KEY (scope_type, scope_id)
);
CREATE TABLE IF NOT EXISTS safety_events (
id INTEGER PRIMARY KEY AUTOINCREMENT,
event_id TEXT NOT NULL UNIQUE,
timestamp TEXT NOT NULL,
type TEXT NOT NULL,
severity TEXT NOT NULL,
scope_type TEXT,
scope_id TEXT,
mission_id TEXT,
wp_id TEXT,
reason TEXT,
reason_code TEXT,
evidence_ref TEXT
);
CREATE INDEX IF NOT EXISTS idx_sev_type ON safety_events(type);
CREATE INDEX IF NOT EXISTS idx_sev_miss ON safety_events(mission_id);
CREATE TABLE IF NOT EXISTS safety_evidence (
id INTEGER PRIMARY KEY AUTOINCREMENT,
ref TEXT NOT NULL UNIQUE,
kind TEXT,
data TEXT,
created_at TEXT NOT NULL
);
"""
)
conn.commit()
finally:
conn.close()
# -- Restart-stabile ID-Helpers (Max+1) ------------------------------------
def _next_event_id(self, conn) -> str:
rows = conn.execute("SELECT event_id FROM safety_events").fetchall()
n = 0
for r in rows:
s = r["event_id"]
if s.startswith("SE-") and s[3:].isdigit():
n = max(n, int(s[3:]))
return f"SE-{n + 1:04d}"
def _next_evidence_ref(self, conn) -> str:
row = conn.execute("SELECT COUNT(*) AS c FROM safety_evidence").fetchone()
return f"E-{int(row['c']) + 1:04d}"
def _next_attempt_id(self, conn, mission_id: str, wp_id: Optional[str]) -> str:
prefix = f"att-{mission_id}" + (f"-{wp_id}" if wp_id else "")
rows = conn.execute(
"SELECT id FROM attempts WHERE id LIKE ?", (prefix + "-%",)
).fetchall()
n = 0
for r in rows:
suffix = r["id"][len(prefix) + 1:]
if suffix.isdigit():
n = max(n, int(suffix))
return f"{prefix}-{n + 1:04d}"
# -- Attempt Ledger --------------------------------------------------------
def record_attempt(
self,
mission_id: str,
wp_id: Optional[str] = None,
*,
actor: str = "red-queen",
task: Optional[str] = None,
hypothesis: Optional[str] = None,
strategy: Optional[str] = None,
target_component: Optional[str] = None,
target_files: Optional[List[str]] = None,
actions_class: Optional[str] = None,
intended_change: Optional[str] = None,
change: Optional[str] = None,
result: str = RESULT_UNKNOWN,
error: Optional[str] = None,
progress_metric: Optional[str] = None,
progress_before: Optional[str] = None,
progress_after: Optional[str] = None,
progress_delta: Optional[str] = None,
progress: str = PROGRESS_UNKNOWN,
idempotency_key: Optional[str] = None,
timestamp: Optional[str] = None,
) -> Dict[str, Any]:
"""Append-only Attempt-Eintrag. Idempotent via `idempotency_key`.
Unbekannte Werte -> fail-closed (INVALID_ARGS), KEINE Mutation.
Secrets werden grundsaetzlich redacted, bevor sie die DB erreichen.
"""
if progress not in PROGRESS_VALUES:
raise _err("INVALID_ARGS", f"invalid progress value {progress!r}", progress=progress)
if result not in RESULT_VALUES:
raise _err("INVALID_ARGS", f"invalid result value {result!r}", result=result)
# Secret-Safety (§28): Credential-artige Werte werden VOR jeder
# Ableitung (Fingerprint/Signatur) und vor der Persistenz redacted, damit
# weder das Ledger noch abgeleitete Hashes/Fingerprints Secrets enthalten.
error_red = redact_secret(error)
strategy_red = redact_secret(strategy)
change_red = redact_secret(change)
intended_change_red = redact_secret(intended_change)
fp_label, fp_hash = strategy_fingerprint(
operation_type=actions_class or "",
target_component=target_component or "",
target_files=target_files,
actions_class=actions_class,
normalized_strategy=strategy_red or "",
intended_change=intended_change_red,
)
sig_norm, sig_hash = error_signature(error_red)
conn = _connect(self.db_path)
try:
if idempotency_key:
existing = conn.execute(
"SELECT id FROM attempts WHERE idempotency_key = ?", (idempotency_key,)
).fetchone()
if existing is not None:
return {"attempt_id": existing["id"], "idempotent": True, "deduplicated": True}
now = timestamp or self._now()
attempt_id = self._next_attempt_id(conn, mission_id, wp_id)
conn.execute(
"INSERT INTO attempts "
"(id,idempotency_key,timestamp,mission_id,wp_id,actor,task,"
" hypothesis,strategy,strategy_fingerprint,strategy_fingerprint_hash,"
" change,result,error_raw,error_signature,error_signature_hash,"
" progress_metric,progress_before,progress_after,progress_delta,progress) "
"VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)",
(
attempt_id, idempotency_key, now, mission_id, wp_id, actor, task,
redact_secret(hypothesis), strategy_red,
fp_label, fp_hash,
change_red, result,
error_red, sig_norm, sig_hash,
progress_metric, progress_before, progress_after, progress_delta, progress,
),
)
conn.commit()
return {"attempt_id": attempt_id, "idempotent": False, "deduplicated": False}
finally:
conn.close()
def attempts(
self, mission_id: Optional[str] = None, wp_id: Optional[str] = None
) -> List[Dict[str, Any]]:
conn = _connect(self.db_path)
try:
q = "SELECT * FROM attempts"
params: List[Any] = []
conds = []
if mission_id:
conds.append("mission_id = ?")
params.append(mission_id)
if wp_id:
conds.append("wp_id = ?")
params.append(wp_id)
if conds:
q += " WHERE " + " AND ".join(conds)
q += " ORDER BY id ASC"
rows = conn.execute(q, params).fetchall()
return [dict(r) for r in rows]
finally:
conn.close()
def attempts_count(self, mission_id: str, wp_id: Optional[str] = None) -> int:
return len(self.attempts(mission_id, wp_id))
# -- Safety Events ----------------------------------------------------------
def _record_event(
self,
conn,
*,
event_type: str,
severity: str,
mission_id: Optional[str],
wp_id: Optional[str],
scope_type: Optional[str],
scope_id: Optional[str],
reason: Optional[str],
reason_code: Optional[str],
evidence_ref: Optional[str],
timestamp: Optional[str] = None,
) -> str:
if severity not in SEVERITIES:
raise _err("SAFETY_STATE_ERROR", f"invalid severity {severity!r}", severity=severity)
now = timestamp or self._now()
event_id = self._next_event_id(conn)
conn.execute(
"INSERT INTO safety_events "
"(event_id,timestamp,type,severity,scope_type,scope_id,mission_id,wp_id,reason,reason_code,evidence_ref) "
"VALUES (?,?,?,?,?,?,?,?,?,?,?)",
(event_id, now, event_type, severity, scope_type, scope_id,
mission_id, wp_id, redact_secret(reason), reason_code, evidence_ref),
)
return event_id
def safety_event(
self,
event_type: str,
severity: str = "WARNING",
*,
mission_id: Optional[str] = None,
wp_id: Optional[str] = None,
scope_type: Optional[str] = None,
scope_id: Optional[str] = None,
reason: Optional[str] = None,
reason_code: Optional[str] = None,
evidence_ref: Optional[str] = None,
timestamp: Optional[str] = None,
) -> Dict[str, Any]:
conn = _connect(self.db_path)
try:
event_id = self._record_event(
conn, event_type=event_type, severity=severity, mission_id=mission_id,
wp_id=wp_id, scope_type=scope_type, scope_id=scope_id,
reason=reason, reason_code=reason_code, evidence_ref=evidence_ref,
timestamp=timestamp,
)
conn.commit()
return {"event_id": event_id, "type": event_type, "severity": severity}
finally:
conn.close()
def safety_events(
self, mission_id: Optional[str] = None, event_type: Optional[str] = None
) -> List[Dict[str, Any]]:
conn = _connect(self.db_path)
try:
q = "SELECT * FROM safety_events"
params: List[Any] = []
conds = []
if mission_id:
conds.append("mission_id = ?")
params.append(mission_id)
if event_type:
conds.append("type = ?")
params.append(event_type)
if conds:
q += " WHERE " + " AND ".join(conds)
q += " ORDER BY id ASC"
rows = conn.execute(q, params).fetchall()
return [dict(r) for r in rows]
finally:
conn.close()
# -- Safety Evidence ---------------------------------------------------------
def safety_evidence(self, kind: str, data: Dict[str, Any], timestamp: Optional[str] = None) -> Dict[str, Any]:
conn = _connect(self.db_path)
try:
ref = self._next_evidence_ref(conn)
now = timestamp or self._now()
serialized = redact_secret(json.dumps(data, ensure_ascii=False, default=str))
conn.execute(
"INSERT INTO safety_evidence (ref,data,created_at) VALUES (?,?,?)",
(ref, serialized, now),
)
conn.commit()
return {"ref": ref, "kind": kind}
finally:
conn.close()
def evidence(self, ref: Optional[str] = None) -> List[Dict[str, Any]]:
conn = _connect(self.db_path)
try:
if ref:
rows = conn.execute(
"SELECT * FROM safety_evidence WHERE ref = ?", (ref,)
).fetchall()
else:
rows = conn.execute(
"SELECT * FROM safety_evidence ORDER BY id ASC"
).fetchall()
out = []
for r in rows:
d = dict(r)
try:
d["data"] = json.loads(d["data"])
except Exception:
pass
out.append(d)
return out
finally:
conn.close()
# -- Circuit Breaker (§12-15) -------------------------------------------------
def circuit_state(self, scope_type: str, scope_id: str) -> Dict[str, Any]:
if scope_type not in SCOPE_TYPES:
raise _err("SAFETY_STATE_ERROR", f"invalid scope_type {scope_type!r}", scope_type=scope_type)
conn = _connect(self.db_path)
try:
row = conn.execute(
"SELECT * FROM circuit_state WHERE scope_type=? AND scope_id=?",
(scope_type, scope_id),
).fetchone()
if row is None:
return {"scope_type": scope_type, "scope_id": scope_id, "state": CIRCUIT_CLOSED}
return dict(row)
finally:
conn.close()
def _circuit_is_open(self, conn, scope_type: str, scope_id: str) -> bool:
row = conn.execute(
"SELECT state FROM circuit_state WHERE scope_type=? AND scope_id=?",
(scope_type, scope_id),
).fetchone()
return row is not None and row["state"] == CIRCUIT_OPEN
def circuit_open_for_scope(self, scope_type: str, scope_id: str) -> bool:
conn = _connect(self.db_path)
try:
return self._circuit_is_open(conn, scope_type, scope_id)
finally:
conn.close()
def open_circuit(
self,
scope_type: str,
scope_id: str,
*,
trigger: str,
severity: str = "HIGH",
reason: Optional[str] = None,
mission_id: Optional[str] = None,
wp_id: Optional[str] = None,
evidence_ref: Optional[str] = None,
timestamp: Optional[str] = None,
) -> Dict[str, Any]:
"""Oeffnet einen Circuit. Idempotent: bereits offen + gleicher Trigger
-> KEIN Event-Sturm. FAIL-CLOSED bei unbekannter severity -> keine Mutation.
GLOBAL-Scope nur mit kritischem GLOBAL-Trigger (§15 Scope-Policy).
"""
if scope_type not in SCOPE_TYPES:
raise _err("SAFETY_STATE_ERROR", f"invalid scope_type {scope_type!r}", scope_type=scope_type)
if severity not in SEVERITIES:
raise _err("SAFETY_STATE_ERROR", f"invalid severity {severity!r}", severity=severity)
if scope_type == SCOPE_GLOBAL and trigger not in GLOBAL_SCOPE_TRIGGERS:
raise _err(
"INVALID_GLOBAL_TRIGGER",
f"GLOBAL scope requires a critical global trigger, got {trigger!r}",
trigger=trigger, scope_id=scope_id,
)
now = timestamp or self._now()
conn = _connect(self.db_path)
try:
existing = self.circuit_state(scope_type, scope_id)
if existing["state"] == CIRCUIT_OPEN:
if existing.get("trigger") == trigger:
return {
"scope_type": scope_type, "scope_id": scope_id,
"state": CIRCUIT_OPEN, "idempotent": True, "event_emitted": False,
}
conn.execute(
"UPDATE circuit_state SET trigger=?, reason=?, updated_at=? "
"WHERE scope_type=? AND scope_id=?",
(trigger, reason, now, scope_type, scope_id),
)
ev = self._record_event(
conn, event_type=EVENT_CIRCUIT_OPENED, severity=severity,
mission_id=mission_id, wp_id=wp_id, scope_type=scope_type,
scope_id=scope_id, reason=f"re-trigger {trigger}: {reason}",
reason_code=REASON_CRITICAL_TRIGGER, evidence_ref=evidence_ref, timestamp=now,
)
conn.commit()
return {
"scope_type": scope_type, "scope_id": scope_id, "state": CIRCUIT_OPEN,
"idempotent": False, "event_emitted": True, "event_id": ev,
}
conn.execute(
"INSERT INTO circuit_state (scope_type,scope_id,state,severity,trigger,reason,"
"evidence_ref,opened_at,closed_at,updated_at) VALUES (?,?,?,?,?,?,?,?,NULL,?)",
(scope_type, scope_id, CIRCUIT_OPEN, severity, trigger, reason,
evidence_ref, now, now),
)
ev = self._record_event(
conn, event_type=EVENT_CIRCUIT_OPENED, severity=severity,
mission_id=mission_id, wp_id=wp_id, scope_type=scope_type, scope_id=scope_id,
reason=reason or trigger, reason_code=REASON_CRITICAL_TRIGGER,
evidence_ref=evidence_ref, timestamp=now,
)
conn.commit()
return {
"scope_type": scope_type, "scope_id": scope_id, "state": CIRCUIT_OPEN,
"idempotent": False, "event_emitted": True, "event_id": ev,
}
finally:
conn.close()
def request_circuit_reset(
self, scope_type: str, scope_id: str, *,
requested_by: str = "red-queen", cause: str, recovery_evidence: str,
mission_id: Optional[str] = None, wp_id: Optional[str] = None,
timestamp: Optional[str] = None,
) -> Dict[str, Any]:
"""Registriert einen Reset-Wunsch. Schliesst den Circuit NICHT automatisch."""
cur = self.circuit_state(scope_type, scope_id)
if cur["state"] != CIRCUIT_OPEN:
raise _err("CIRCUIT_NOT_OPEN", f"circuit not open for {scope_type}/{scope_id}")
conn = _connect(self.db_path)
try:
ev = self._record_event(
conn, event_type=EVENT_CIRCUIT_RESET_REQUESTED, severity="WARNING",
mission_id=mission_id, wp_id=wp_id, scope_type=scope_type, scope_id=scope_id,
reason=f"reset requested by {requested_by}; cause={cause}; evidence={recovery_evidence}",
reason_code=REASON_HUMAN_GATE_REQUIRED, evidence_ref=None, timestamp=timestamp,
)
conn.commit()
return {"event_id": ev, "status": "requested", "scope_type": scope_type, "scope_id": scope_id}
finally:
conn.close()
def close_circuit(
self, scope_type: str, scope_id: str, *,
approved_by: str, gate: str = "documented_recovery",
cause: str, recovery_evidence: str,
mission_id: Optional[str] = None, wp_id: Optional[str] = None,
timestamp: Optional[str] = None,
) -> Dict[str, Any]:
"""Schliesst einen offenen Circuit mit Gate-Policy (§15).
LOW/MEDIUM -> dokumentierte Ursache + Recovery-Evidence.
HIGH/CRITICAL/GLOBAL -> Human Gate (gate='human_gate'/'external_review').
Ohne ausreichendes Gate -> FAIL-CLOSED (HUMAN_GATE_REQUIRED), keine Mutation.
"""
cur = self.circuit_state(scope_type, scope_id)
if cur["state"] != CIRCUIT_OPEN:
return {
"scope_type": scope_type, "scope_id": scope_id, "state": CIRCUIT_CLOSED,
"idempotent": True,
}
severity = cur.get("severity", "HIGH")
needs_human = severity in ("HIGH", "CRITICAL") or scope_type == SCOPE_GLOBAL
if needs_human:
if gate not in ("human_gate", "external_review"):
raise _err(
"HUMAN_GATE_REQUIRED",
f"closing {severity}/{scope_type} circuit requires human gate or external review",
scope_type=scope_type, scope_id=scope_id, severity=severity, gate=gate,
)
elif not cause or not recovery_evidence:
raise _err(
"INVALID_ARGS",
"LOW/MEDIUM circuit close requires documented cause + recovery evidence",
scope_type=scope_type, scope_id=scope_id,
)
now = timestamp or self._now()
conn = _connect(self.db_path)
try:
conn.execute(
"UPDATE circuit_state SET state=?, closed_at=?, updated_at=? "
"WHERE scope_type=? AND scope_id=?",
(CIRCUIT_CLOSED, now, now, scope_type, scope_id),
)
ev = self._record_event(
conn, event_type=EVENT_CIRCUIT_CLOSED, severity="INFO",
mission_id=mission_id, wp_id=wp_id, scope_type=scope_type, scope_id=scope_id,
reason=f"closed by {approved_by} via gate={gate}; cause={cause}",
reason_code=None, evidence_ref=None, timestamp=now,
)
conn.commit()
return {"scope_type": scope_type, "scope_id": scope_id, "state": CIRCUIT_CLOSED,
"idempotent": False, "event_id": ev}
finally:
conn.close()
# -- Fail-Closed Safety State ---------------------------------------------
def check_safety_state(
self, mission_id: Optional[str] = None, wp_id: Optional[str] = None
) -> Dict[str, Any]:
"""Konsistenzpruefung. Bei unbekanntem/inkonsistentem Circuit-State
-> FAIL-CLOSED: SafetyError(STATE_INCONSISTENT), Evidence + Safety-Event,
KEINE weitere Mutation."""
conn = _connect(self.db_path)
try:
rows = conn.execute("SELECT scope_type, scope_id, state FROM circuit_state").fetchall()
for r in rows:
if r["state"] not in CIRCUIT_STATES:
ref = self._next_evidence_ref(conn)
conn.execute(
"INSERT INTO safety_evidence (ref,data,created_at) VALUES (?,?,?)",
(ref, redact_secret(json.dumps({
"scope_type": r["scope_type"], "scope_id": r["scope_id"],
"corrupt_state": r["state"],
}, default=str)), self._now()),
)
self._record_event(
conn, event_type=EVENT_SAFETY_STATE_ERROR, severity="CRITICAL",
mission_id=mission_id, wp_id=wp_id, scope_type=r["scope_type"],
scope_id=r["scope_id"],
reason=f"corrupt circuit state {r['state']!r}",
reason_code=REASON_STATE_INCONSISTENT, evidence_ref=ref,
)
conn.commit()
raise _err(
"STATE_INCONSISTENT",
f"corrupt circuit state {r['state']!r} for {r['scope_type']}/{r['scope_id']}",
scope_type=r["scope_type"], scope_id=r["scope_id"], state=r["state"],
)
return {"ok": True}
finally:
conn.close()
# -- Oscillation Detection (deterministisch) -------------------------------
def _strategy_seq(self, attempts: List[Dict[str, Any]]) -> List[str]:
return [a.get("strategy_fingerprint_hash") or a.get("strategy") or "" for a in attempts]
def _detect_abab(self, seq: List[str]) -> bool:
"""A-B-A-B im letzten 4er-Fenster. Leere/identische Eintraege sind kein ABAB."""
if len(seq) < ABAB_PATTERN_LEN:
return False
win = seq[-ABAB_PATTERN_LEN:]
a, b, a2, b2 = win
return a and b and a != b and a == a2 and b == b2
def _detect_oscillation(
self, attempts: List[Dict[str, Any]], new_sig_hash: str
) -> Dict[str, Any]:
"""Konservativ V1: A-B-A-B der Strategie-Fingerprints (Oscillation-Kandidat).
Wiederholte Fehler-Signatur wird NICHT hier gezaehlt — sie wird vorher durch
den SAME_ERROR_LIMIT-Retry-Controller behandelt (§11: gleiche Fehler-Signatur
MAX 2; A-B-A-B = Circuit-Breaker-Kandidat). Die Funktion meldet also nur
echte oszillierende Strategie-Muster (mehrere DISTINCT Strategien im Wechsel).
"""
reasons = []
if self._detect_abab(self._strategy_seq(attempts)):
reasons.append("abab_strategy_pattern")
if reasons:
return {"oscillation": True, "reasons": reasons}
return {"oscillation": False, "reasons": []}
# -- Retry Controller / Safety Decision API (§21) -----------------------------
def evaluate_next_action(
self,
mission_id: str,
wp_id: Optional[str] = None,
*,
error: Optional[str] = None,
strategy_label: Optional[str] = None,
target_component: Optional[str] = None,
target_files: Optional[List[str]] = None,
is_mutating: bool = True,
persist: bool = True,
) -> Dict[str, Any]:
"""Deterministische Antwort: 'Darf ich weitermachen und welche Aktionsklasse?'
Reihenfolge (hart, deterministisch):
1) Fail-closed Safety-State-Pruefung
2) Circuit fuer relevante Scopes -> CIRCUIT_ALREADY_OPEN
3) Max-Iteration-Notbremse
4) Oscillation (A-B-A-B / repeated signature) -> OSCILLATION_ABAB
5) Retry-Limits (Signature / Strategy / Maker-Checker-Repair)
6) kein messbarer Fortschritt bei wiederholtem FAIL
7) sonst RETRY/CONTINUE
"""
out = {
"DECISION": None,
"REASON_CODE": None,
"SEVERITY": "INFO",
"CIRCUIT_STATE": CIRCUIT_CLOSED,
"ALLOWED_ACTION": ALLOWED_NONE,
"ACTION_TEXT": "",
"EVIDENCE": [],
}
# 1) Fail-closed Safety State
try:
self.check_safety_state(mission_id, wp_id)
except SafetyError as e:
out.update(
DECISION="BLOCK", REASON_CODE=e.code, SEVERITY="CRITICAL",
CIRCUIT_STATE=CIRCUIT_OPEN, ALLOWED_ACTION=ALLOWED_NONE,
ACTION_TEXT=f"fail-closed: {e.message}",
)
return self._persist(out, persist, EVENT_SAFETY_STATE_ERROR, mission_id, wp_id)
# 2) Circuit-Check scoped
circuit = self._find_open_circuit(mission_id, wp_id, target_component)
if circuit is not None:
st, sid, cstate = circuit
if is_mutating:
out.update(
DECISION="BLOCK", REASON_CODE=REASON_CIRCUIT_ALREADY_OPEN,
SEVERITY="CRITICAL", CIRCUIT_STATE=cstate,
ALLOWED_ACTION=ALLOWED_READ_ONLY_DIAGNOSIS,
ACTION_TEXT=f"circuit OPEN for {st}/{sid}; mutating ops blocked",
)
else:
out.update(
DECISION="CONTINUE", REASON_CODE=REASON_RETRY_AVAILABLE,
SEVERITY="INFO", CIRCUIT_STATE=cstate,
ALLOWED_ACTION=ALLOWED_READ_ONLY_DIAGNOSIS,
ACTION_TEXT="read-only diagnosis allowed while circuit open",
)
return self._persist(out, persist, None, mission_id, wp_id)
attempts = self.attempts(mission_id, wp_id)
count = len(attempts)
# 3) Max-Iteration-Notbremse
if count >= MAX_ITERATIONS:
out.update(
DECISION="CIRCUIT_BREAK", REASON_CODE=REASON_ITERATION_LIMIT,
SEVERITY="CRITICAL", ALLOWED_ACTION=ALLOWED_HUMAN_GATE,
ACTION_TEXT=f"global iteration limit {MAX_ITERATIONS} reached; forced stop",
)
return self._persist(out, persist, EVENT_OSCILLATION_DETECTED, mission_id, wp_id)
sig_norm, sig_hash = error_signature(error)
# 4) Oscillation
osc = self._detect_oscillation(attempts, sig_hash)
if osc["oscillation"]:
out.update(
DECISION="CIRCUIT_BREAK", REASON_CODE=REASON_OSCILLATION_ABAB,
SEVERITY="CRITICAL", ALLOWED_ACTION=ALLOWED_HUMAN_GATE,
ACTION_TEXT="oscillation detected: " + "; ".join(osc["reasons"]),
)
return self._persist(out, persist, EVENT_OSCILLATION_DETECTED, mission_id, wp_id)
# 5) Retry-Limits
fail_attempts = [a for a in attempts if a.get("result") == RESULT_FAIL]
# SAFETY_CONTRACT §2.4: Zaehler werden bei tatsaechlichem Fortschritt
# zurueckgesetzt. Messbarer Progress (YES) hebt die Error-Signatur-Sperre auf,
# damit eine Fehler+Messbare-Fortschritt-Situation NICHT vorschnell blockt (§26 TEST D).
has_measurable_progress = any(
a.get("progress") == PROGRESS_YES for a in attempts
)
sig_count_total = sum(1 for a in fail_attempts if a.get("error_signature_hash") == sig_hash)
if sig_hash and not has_measurable_progress and sig_count_total >= MAX_SAME_ERROR_SIGNATURE:
out.update(
DECISION="DEBUG", REASON_CODE=REASON_SAME_ERROR_LIMIT, SEVERITY="WARNING",
ALLOWED_ACTION=ALLOWED_STRATEGY_CHANGE,
ACTION_TEXT=f"same error signature seen {sig_count_total} times "
f"(limit {MAX_SAME_ERROR_SIGNATURE}); strategy change required",
)
return self._persist(out, persist, EVENT_ERROR_SIGNATURE_REPEAT, mission_id, wp_id)
if strategy_label or target_component:
fp_label, fp_hash = strategy_fingerprint(
operation_type="", target_component=target_component or "",
target_files=target_files, actions_class=None,
normalized_strategy=strategy_label or "", intended_change=None,
)
failed_strategy_repeat = any(
a.get("strategy_fingerprint_hash") == fp_hash and a.get("result") == RESULT_FAIL
for a in attempts
)
if failed_strategy_repeat:
out.update(
DECISION="DEBUG", REASON_CODE=REASON_FAILED_STRATEGY_REPEAT,
SEVERITY="WARNING", ALLOWED_ACTION=ALLOWED_STRATEGY_CHANGE,
ACTION_TEXT="same already-failed strategy proposed; blind repeat forbidden",
)
return self._persist(out, persist, EVENT_STRATEGY_REPEAT, mission_id, wp_id)
repair_count = sum(
1 for a in attempts
if (a.get("actor") or "").lower() in ("maker", "repair")
and a.get("result") == RESULT_FAIL
)
if repair_count >= MAX_MAKER_CHECKER_REPAIRS:
out.update(
DECISION="SECOND_OPINION", REASON_CODE=REASON_RETRY_LIMIT, SEVERITY="HIGH",
ALLOWED_ACTION=ALLOWED_SECOND_OPINION,
ACTION_TEXT=f"maker/checker repair limit {MAX_MAKER_CHECKER_REPAIRS} reached; "
f"need fresh second opinion",
)
return self._persist(out, persist, EVENT_RETRY_LIMIT_REACHED, mission_id, wp_id)
# 6) kein messbarer Fortschritt bei wiederholtem FAIL (UNKNOWN != Progress)
no_progress_fails = sum(
1 for a in fail_attempts if a.get("progress") in (PROGRESS_NO, PROGRESS_UNKNOWN)
)
if fail_attempts and no_progress_fails >= NO_PROGRESS_FAILS_LIMIT:
out.update(
DECISION="DEBUG", REASON_CODE=REASON_NO_MEASURABLE_PROGRESS,
SEVERITY="WARNING", ALLOWED_ACTION=ALLOWED_DEBUG,
ACTION_TEXT="repeated failures without measurable progress; "
"UNKNOWN is not progress",
)
return self._persist(out, persist, EVENT_NO_PROGRESS, mission_id, wp_id)
# 7) Retry erlaubt
out.update(
DECISION="RETRY" if is_mutating else "CONTINUE",
REASON_CODE=REASON_RETRY_AVAILABLE, SEVERITY="INFO",
ALLOWED_ACTION=ALLOWED_RETRY if is_mutating else ALLOWED_CONTINUE,
ACTION_TEXT="retry / next action allowed",
)
return self._persist(out, persist, EVENT_RETRY_ALLOWED, mission_id, wp_id)
def _find_open_circuit(
self, mission_id: str, wp_id: Optional[str], target_component: Optional[str]
) -> Optional[Tuple[str, str, str]]:
"""Offene Circuits in Scope-Prioritaet: GLOBAL > MISSION > WP > COMPONENT."""
scopes: List[Tuple[str, str]] = [(SCOPE_GLOBAL, "global")]
if mission_id:
scopes.append((SCOPE_MISSION, mission_id))
if wp_id:
scopes.append((SCOPE_WORK_PACKAGE, wp_id))
if target_component:
scopes.append((SCOPE_COMPONENT, target_component))
conn = _connect(self.db_path)
try:
for st, sid in scopes:
if self._circuit_is_open(conn, st, sid):
return (st, sid, CIRCUIT_OPEN)
finally:
conn.close()
return None
def _persist(
self,
outcome: Dict[str, Any],
persist: bool,
event_type: Optional[str],
mission_id: Optional[str],
wp_id: Optional[str],
) -> Dict[str, Any]:
"""Sichert das Outcome als Safety-Event, falls gewuenscht. Events sind
Beobachtungen; ein Fehler beim Event-Log darf die Entscheidung nicht kippen."""
if persist and event_type:
try:
self.safety_event(
event_type=event_type, severity=outcome["SEVERITY"],
mission_id=mission_id, wp_id=wp_id,
reason=outcome["ACTION_TEXT"], reason_code=outcome["REASON_CODE"],
)
except Exception: # noqa
pass
return outcome