diff --git a/a3/.gitignore b/a3/.gitignore new file mode 100644 index 0000000..7dd2791 --- /dev/null +++ b/a3/.gitignore @@ -0,0 +1,6 @@ +__pycache__/ +*.pyc +*.db +*.db-wal +*.db-shm +*.tmp diff --git a/a3/README.md b/a3/README.md new file mode 100644 index 0000000..4602d72 --- /dev/null +++ b/a3/README.md @@ -0,0 +1,145 @@ +# Red Queen — A3: Deterministic Safety Layer (V1) + +Der **Deterministic Safety Layer** sitzt zwischen Mission/WP-State (A2 `missions.db`) +und dem **späteren** Orchestrator/Action/Retry/Delegation. Er entscheidet +**deterministisch** (Zähler, Limits, Tabellen — nicht „LLM-Lust"): + +``` +CONTINUE / RETRY / DEBUG / SECOND_OPINION / BLOCK / ESCALATE / CIRCUIT_BREAK +``` + +Ziel: Red Queen erhält die Safety-Logik, **BEVOR** spätere autonome Mission-Loops +existieren. Dieser Build enthält **KEINE** autonome Orchestrierung — kein dispatcher, +kein loop, kein heartbeat, kein cron, kein self-improvement. A3 ist eine reine +**Library / Safety-Capability**. + +--- + +## DB-Entscheidung (Maker-Entscheidung, §29) + +A3 nutzt eine **eigene `safety.db`** (separate SQLite-Datei), **nicht** die A2 +`missions.db`: + +- **Rollback-sicher für A2:** Die bestehende A2-Datenbank wird überhaupt nicht + angefasst — kein ALTER, kein neues Schema in `missions.db`, keine Gefahr für + Bestandsdaten (`CREATE TABLE IF NOT EXISTS` auf den A3-Tabellen in `safety.db`). +- **Entkopplung:** Safety-State (Circuit, Attempts, Events, Evidence) ist unabhängig + vom Mission-State. Ein beschädigter Mission-State kann den Safety-Layer nicht + mitreißen und umgekehrt. +- **Testbar:** A3-DB wird immer über temp-Pfade getestet, nie produktiv berührt. + +Die A2-Integration erfolgt sauber: `SafetyStore` verweist Missions-/WP-IDs +als Fremdschlüssel-Namen (mission_id/wp_id) und ruft A2-Regeln nicht auf; A2-APIs +`mission_block`/`mission_transition` bleiben unberührt. Mission-`BLOCKED` kann in +A4 durch den Orchestrator auf Basis einer `BLOCK`-Safety-Entscheidung gesetzt werden — +A3 selbst mutiert A2 nicht direkt. + +--- + +## Module + +| Datei | Inhalt | +|-------|--------| +| `rq_safety.py` | Kernmodul: `SafetyStore`, `SafetyError`, Error-Signature, Strategy-Fingerprint, Redaction, Circuit Breaker, Retry Controller, Oscillation, Fail-Closed, Safety Events, Evidence, `evaluate_next_action` (§21) | +| `rq_safety_telegram.py` | Telegram-Notification-Interface (§19): `format_alert`, `build_payload`, `should_notify` (Anti-Spam). KEIN Daemon. | +| `rq_safety_cli.py` | Dünne, deterministische CLI (`--json`, Exit-Code 2 bei `SafetyError`), A2-CLI-Muster spiegelnd | +| `test_a3.py` | Isolierte Testsuite (temp-DB), Exit-Code 0 = PASS | + +--- + +## Kern-API (`SafetyStore`) + +```python +from rq_safety import SafetyStore +store = SafetyStore("path/to/safety.db") + +# Attempt Ledger (idempotent via idempotency_key, append-only, redacted) +store.record_attempt("M1", "W1", actor="maker", result="FAIL", + error="boom pid=123", strategy="s1", idempotency_key="k-1") + +# Retry / Oscillation / Circuit / Fail-closed +store.evaluate_next_action("M1", "W1", error="boom", strategy_label="s2", + is_mutating=True) # -> DECISION/REASON_CODE/ALLOWED_ACTION + +# Circuit Breaker +store.open_circuit("MISSION", "M1", trigger="REG", severity="HIGH") +store.circuit_state("MISSION", "M1") +store.request_circuit_reset("MISSION", "M1", cause="...", recovery_evidence="...") +store.close_circuit("MISSION", "M1", approved_by="human", gate="human_gate", cause="...", recovery_evidence="...") + +# Events / Evidence +store.safety_event("OSCILLATION_DETECTED", severity="CRITICAL", mission_id="M1") +store.safety_evidence("test_results", {"failing": 2}) +store.evidence() +``` + +--- + +## Entscheidungen (Reason Codes, §22) + +| Code | Bedeutung | +|------|-----------| +| `RETRY_AVAILABLE` | nächste Aktion erlaubt | +| `RETRY_LIMIT` | Maker/Checker-Repair MAX 3 erreicht → SECOND_OPINION | +| `SAME_ERROR_LIMIT` | gleiche Error-Signatur MAX 2 → DEBUG/Strategiewechsel | +| `FAILED_STRATEGY_REPEAT` | bereits gescheiterte Strategie → keine blinde Wiederholung | +| `NO_MEASURABLE_PROGRESS` | FAIL ohne messbaren Fortschritt (UNKNOWN ≠ Progress) | +| `OSCILLATION_ABAB` | A-B-A-B-Muster → Circuit-Breaker-Kandidat | +| `CIRCUIT_ALREADY_OPEN` | Circuit OPEN → BLOCK (read-only erlaubt) | +| `STATE_INCONSISTENT` | Fail-Closed: Safety-State korrupt → keine Mutation | +| `CRITICAL_TRIGGER` | kritischer GLOBAL-Trigger | +| `HUMAN_GATE_REQUIRED` | Circuit-Reset braucht Human Gate bei HIGH/CRITICAL/GLOBAL | + +--- + +## Limite (A1 SAFETY_CONTRACT, konservativ V1) + +- Maker→Checker-Repair: **MAX 3** +- Gleiche Error-Signatur: **MAX 2** +- Gleiche bereits gescheiterte Strategie: **keine blinde Wiederholung** +- Oscillation A-B-A-B: Circuit-Breaker-Kandidat +- `MAX_ITERATIONS=50`: äußerste Runtime-Notbremse, **nicht** operatives Retry-Limit + +--- + +## Fail-Closed (§16) + +Bei unbekanntem/inkonsistentem Safety-State → `SAFETY_STATE_ERROR` → **STOP** → +**Evidence sichern** → **KEINE Mutation** (read-only Diagnose ggf. erlaubt). + +--- + +## Secret-Safety (§28) + +- Alle credential-artigen Werte werden **vor** Persistenz UND vor Fingerprint/Signatur- + Ableitung **redacted** (`redact_secret`). +- Keine Tokens/Passwörter/private Keys/Authorization-Header in `safety.db`/Events/Telegram. +- Secret Exposure wird nur als `SECRET_EXPOSURE_DETECTED` + LOCATION/TYPE + `VALUE=REDACTED` erfasst. + +--- + +## Telegram Interface (§19) + +Kein Daemon. Nur Formatter/Payload. **Kein Spam** — nur signifikante Events +(`CIRCUIT_OPENED`, `CRITICAL`, `ESCALATION_REQUIRED`, `HUMAN_DECISION_REQUIRED`, +`RETRY_LIMIT_REACHED`, `OSCILLATION_DETECTED`, `SAFETY_STATE_ERROR`). Payload ist +deterministisch und **redacted**. + +```python +from rq_safety_telegram import build_payload +payload = build_payload("CIRCUIT_OPENED", "CRITICAL", reason="...", mission_id="M") +# payload["notify"] == True, payload["text"] fertig formatiert +``` + +--- + +## Tests + +```bash +python3 test_a3.py # Exit 0 = PASS; isoliert (temp-DB), produktive DBs unberührt +``` + +Abgedeckt: Attempt Ledger + Idempotenz, Error-Signature-Normalisierung, Strategy- +Fingerprint, Retry-Limits, Progress, Oscillation A-B-A-B, False-Positives, +Circuit-Breaker (+Restart-Persistenz + Negativ-Test), Fail-Closed, Events, +Evidence, Secret-Safety, Telegram-Interface, Loop-Simulation (A–E), A2-DB-Isolation. diff --git a/a3/rq_safety.py b/a3/rq_safety.py new file mode 100644 index 0000000..bf968b0 --- /dev/null +++ b/a3/rq_safety.py @@ -0,0 +1,1055 @@ +#!/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"), ""), + (re.compile(r"\b(pid|ppid)\s*[=:]\s*\d+\b"), r"\1="), + (re.compile(r"\b\d{4}-\d{2}-\d{2}[T ]\d{2}:\d{2}:\d{2}"), ""), + (re.compile(r"\b[A-Fa-f0-9]{32}\b"), ""), + (re.compile(r"\b[A-Fa-f0-9]{64}\b"), ""), + (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"), ""), + (re.compile(r"(?i)(port\s*[=:]\s*)\d{2,5}"), r"\1"), + (re.compile(r"(?i)(session|request|trace|run)[_-]?id[^:,;]*\s*[=:]\s*[^\s,;]+"), r"\1="), +] + + +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 diff --git a/a3/rq_safety_cli.py b/a3/rq_safety_cli.py new file mode 100644 index 0000000..92efff7 --- /dev/null +++ b/a3/rq_safety_cli.py @@ -0,0 +1,188 @@ +#!/usr/bin/env python3 +""" +Red Queen — A3 CLI. Dünner, deterministischer Kommandozeilen-Zugriff auf den +Deterministic Safety Layer (`SafetyStore`). Maschinenlesbar via `--json`. + +Hinweis: Der A3-Build definiert KEINE autonome Orchestrierung — dieses CLI ruft +nur die vom Store bereitgestellten atomaren, deterministischen Operationen auf. +Es gibt KEINEN Loop, kein Retry-Ausfuehrung, kein Daemon. +""" + +from __future__ import annotations + +import argparse +import json +import os +import sys + +from rq_safety import SafetyStore, SafetyError + + +def _store(args) -> SafetyStore: + db = args.db or os.environ.get("RQ_SAFETY_DB") or "safety.db" + return SafetyStore(db) + + +def _emit(args, obj) -> None: + if args.json: + print(json.dumps(obj, indent=2, ensure_ascii=False, default=str)) + else: + if isinstance(obj, dict) and obj.get("attempt_id"): + print(obj["attempt_id"]) + else: + print(json.dumps(obj, ensure_ascii=False, default=str)) + + +def _fail(e: SafetyError, args) -> None: + if args.json: + print(json.dumps(e.to_dict(), indent=2)) + else: + print(f"ERROR [{e.code}]: {e.message}") + sys.exit(2) + + +def build_parser() -> argparse.ArgumentParser: + p = argparse.ArgumentParser(prog="rq_safety", description="Red Queen A3 Deterministic Safety CLI") + p.add_argument("--db", help="Pfad zur safety.db (default: $RQ_SAFETY_DB oder safety.db)") + p.add_argument("--json", action="store_true", help="JSON-Ausgabe") + sub = p.add_subparsers(dest="cmd", required=True) + + # attempt + at = sub.add_parser("attempt_record") + at.add_argument("--mission", dest="mission_id", required=True) + at.add_argument("--wp", dest="wp_id") + at.add_argument("--actor", default="red-queen") + at.add_argument("--result", default="UNKNOWN") + at.add_argument("--error") + at.add_argument("--strategy") + at.add_argument("--target-component") + at.add_argument("--progress", default="UNKNOWN") + at.add_argument("--idem-key", dest="idem_key") + at.add_argument("--change") + + al = sub.add_parser("attempts") + al.add_argument("--mission", dest="mission_id") + al.add_argument("--wp", dest="wp_id") + + # circuit + co = sub.add_parser("circuit_open") + co.add_argument("--scope-type", required=True) + co.add_argument("--scope-id", required=True) + co.add_argument("--trigger", required=True) + co.add_argument("--severity", default="HIGH") + co.add_argument("--reason") + co.add_argument("--mission", dest="mission_id") + co.add_argument("--wp", dest="wp_id") + + cs = sub.add_parser("circuit_state") + cs.add_argument("--scope-type", required=True) + cs.add_argument("--scope-id", required=True) + + cr = sub.add_parser("circuit_reset_request") + cr.add_argument("--scope-type", required=True) + cr.add_argument("--scope-id", required=True) + cr.add_argument("--cause", required=True) + cr.add_argument("--recovery-evidence", required=True) + + cc = sub.add_parser("circuit_close") + cc.add_argument("--scope-type", required=True) + cc.add_argument("--scope-id", required=True) + cc.add_argument("--approved-by", required=True) + cc.add_argument("--gate", default="documented_recovery") + cc.add_argument("--cause", required=True) + cc.add_argument("--recovery-evidence", required=True) + + # events + ev = sub.add_parser("event") + ev.add_argument("--type", required=True) + ev.add_argument("--severity", default="WARNING") + ev.add_argument("--reason") + ev.add_argument("--reason-code") + ev.add_argument("--mission", dest="mission_id") + ev.add_argument("--wp", dest="wp_id") + ev.add_argument("--scope-type") + ev.add_argument("--scope-id") + + evl = sub.add_parser("events") + evl.add_argument("--mission", dest="mission_id") + evl.add_argument("--type") + + # evidence + evi = sub.add_parser("evidence_save") + evi.add_argument("--kind", required=True) + evi.add_argument("--data", required=True, help="JSON-Dict") + + evir = sub.add_parser("evidence_list") + evir.add_argument("--ref") + + # decision + dec = sub.add_parser("evaluate") + dec.add_argument("--mission", dest="mission_id", required=True) + dec.add_argument("--wp", dest="wp_id") + dec.add_argument("--error") + dec.add_argument("--strategy") + dec.add_argument("--target-component") + dec.add_argument("--read-only", action="store_true", help="nicht-mutierende Operation") + + sub.add_parser("check_safety_state") + return p + + +def main(argv=None) -> int: + args = build_parser().parse_args(argv) + try: + s = _store(args) + if args.cmd == "attempt_record": + r = s.record_attempt( + args.mission_id, wp_id=args.wp_id, actor=args.actor, result=args.result, + error=args.error, strategy=args.strategy, + target_component=args.target_component, progress=args.progress, + idempotency_key=args.idem_key, change=args.change, + ) + elif args.cmd == "attempts": + r = s.attempts(args.mission_id, args.wp_id) + elif args.cmd == "circuit_open": + r = s.open_circuit(args.scope_type, args.scope_id, trigger=args.trigger, + severity=args.severity, reason=args.reason, + mission_id=args.mission_id, wp_id=args.wp_id) + elif args.cmd == "circuit_state": + r = s.circuit_state(args.scope_type, args.scope_id) + elif args.cmd == "circuit_reset_request": + r = s.request_circuit_reset(args.scope_type, args.scope_id, cause=args.cause, + recovery_evidence=args.recovery_evidence) + elif args.cmd == "circuit_close": + r = s.close_circuit(args.scope_type, args.scope_id, approved_by=args.approved_by, + gate=args.gate, cause=args.cause, + recovery_evidence=args.recovery_evidence) + elif args.cmd == "event": + r = s.safety_event(args.type, severity=args.severity, reason=args.reason, + reason_code=args.reason_code, mission_id=args.mission_id, + wp_id=args.wp_id, scope_type=args.scope_type, + scope_id=args.scope_id) + elif args.cmd == "events": + r = s.safety_events(args.mission_id, args.type) + elif args.cmd == "evidence_save": + r = s.safety_evidence(args.kind, json.loads(args.data)) + elif args.cmd == "evidence_list": + r = s.evidence(args.ref) + elif args.cmd == "evaluate": + r = s.evaluate_next_action(args.mission_id, args.wp_id, error=args.error, + strategy_label=args.strategy, + target_component=args.target_component, + is_mutating=not args.read_only) + elif args.cmd == "check_safety_state": + r = s.check_safety_state() + else: + _emit(args, {"error": f"unknown command {args.command}"}) + return 1 + _emit(args, r) + return 0 + except SafetyError as e: + _fail(e, args) + except Exception as e: # noqa + _fail(SafetyError("INTERNAL_ERROR", str(e)), args) + return 2 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/a3/rq_safety_telegram.py b/a3/rq_safety_telegram.py new file mode 100644 index 0000000..56f5213 --- /dev/null +++ b/a3/rq_safety_telegram.py @@ -0,0 +1,134 @@ +#!/usr/bin/env python3 +""" +Red Queen — A3: Telegram Safety Notification Interface (§19). + +KEIN Daemon, KEIN autonomer Sender. Dieses Modul liefert nur einen deterministischen +Formatter + Payload-Bausteine fuer relevante Safety-Events. Der spaetere Red Queen +Lead wuerde diese Payloads ueber einen existierenden Telegram-Kanal senden. + +Anti-Spam-Regeln (TELEGRAM_MISSION_CONTROL): + * NUR signifikante Events: CIRCUIT_OPENED, CRITICAL-Severity, ESCALATION_REQUIRED, + HUMAN_DECISION_REQUIRED, RETRY_LIMIT_REACHED, OSCILLATION_DETECTED. + * Nicht jeder Retry braucht Telegram. RETRY_ALLOWED/DEBUG/info-Events -> kein Alert. + +Payload ist deterministisch und grundsaetzlich REDACTED (keine Secrets, §28). + +Test-Konvention: Tests pruefen Formatter-/Payload-Ausgabe (deterministisch), senden +aber NICHT real. +""" + +from __future__ import annotations + +from typing import Any, Dict + +from rq_safety import ( + EVENT_CIRCUIT_OPENED, + EVENT_CIRCUIT_RESET_REQUESTED, + EVENT_ESCALATION_REQUIRED, + EVENT_HUMAN_DECISION_REQUIRED, + EVENT_RETRY_LIMIT_REACHED, + EVENT_OSCILLATION_DETECTED, + EVENT_SAFETY_STATE_ERROR, + redact_secret, +) + +# Meldungs-Praefixe (TELEGRAM_MISSION_CONTROL §5) +PREFIX_CRITICAL = "[CRITICAL]" +PREFIX_BLOCKER = "[BLOCKER]" +PREFIX_DECISION = "[DECISION REQUIRED]" + +# Events, die einen Telegram-Alert ausloesen (signifikant, kein Spam). +ALERT_TYPES = frozenset( + { + EVENT_CIRCUIT_OPENED, + EVENT_CIRCUIT_RESET_REQUESTED, + EVENT_ESCALATION_REQUIRED, + EVENT_HUMAN_DECISION_REQUIRED, + EVENT_RETRY_LIMIT_REACHED, + EVENT_OSCILLATION_DETECTED, + EVENT_SAFETY_STATE_ERROR, + } +) + + +def should_notify(event_type: str, severity: str) -> bool: + """Deterministische Spam-Regel: nur signifikante Events melden. + + CRITICAL-Severity wird immer gemeldet; ansonsten nur in ALERT_TYPES gelistete + Event-Typen. RETRY_ALLOWED / INFO / gewoehnliche Events -> kein Alert. + """ + if severity == "CRITICAL": + return True + return event_type in ALERT_TYPES + + +def format_alert( + event_type: str, + severity: str, + *, + mission_id: Any = None, + wp_id: Any = None, + scope_type: Any = None, + scope_id: Any = None, + reason: str = "", + reason_code: str = "", +) -> str: + """Baue eine deterministische, maschinenlesbare Telegram-Nachricht. + + Ohne Secrets (redacted). Gibt einen menschenlesbaren Block zurueck, der mit + einem festen Praefix beginnt. Kein Alert bei nicht-signifikanten Events -> "". + """ + if not should_notify(event_type, severity): + return "" # kein Alert -> keine Nachricht (Anti-Spam) + + prefix = PREFIX_CRITICAL if severity == "CRITICAL" else ( + PREFIX_DECISION + if event_type in (EVENT_HUMAN_DECISION_REQUIRED, EVENT_ESCALATION_REQUIRED) + else PREFIX_BLOCKER + ) + scope = f"scope={scope_type}/{scope_id or '-'}" if scope_type else "scope=n/a" + lines = [ + f"{prefix} Red Queen Safety", + f"EVENT={event_type} SEVERITY={severity}", + f"REASON_CODE={reason_code or 'n/a'}", + scope, + f"mission={mission_id or 'n/a'}" + (f" wp={wp_id}" if wp_id else ""), + f"reason={redact_secret(reason) or 'n/a'}", + ] + return "\n".join(lines) + + +def build_payload( + event_type: str, + severity: str, + *, + reason: str = "", + reason_code: str = "", + scope_type: Any = None, + scope_id: Any = None, + mission_id: Any = None, + wp_id: Any = None, + evidence_ref: Any = None, +) -> Dict[str, Any]: + """Strukturierte Payload fuer einen Telegram-Sender (deterministisch, redacted). + + Return-Format ist stabil; ein spaeterer Sender muss nur diesen Dict-Payload + an den konfigurierten Messenger senden. + """ + return { + "notify": should_notify(event_type, severity), + "text": format_alert( + event_type, severity, reason=reason, reason_code=reason_code, + scope_type=scope_type, scope_id=scope_id, + mission_id=mission_id, wp_id=wp_id, + ), + "event_type": event_type, + "severity": severity, + "reason_code": reason_code, + "reason": redact_secret(reason), + "scope_type": scope_type, + "scope_id": scope_id, + "mission_id": mission_id, + "wp_id": wp_id, + "evidence_ref": evidence_ref, + } diff --git a/a3/test_a3.py b/a3/test_a3.py new file mode 100644 index 0000000..861c9ce --- /dev/null +++ b/a3/test_a3.py @@ -0,0 +1,423 @@ +#!/usr/bin/env python3 +""" +Red Queen — A3 Testsuite (deterministisch, isoliert). + +Nutzt ausschliesslich temporaere `safety.db` (tempfile.mkdtemp). Produktive DBs +werden NIE angefasst. Telegram-Interface wird nur als Formatter-/Payload-Test +geprueft — es wird NICHTS real gesendet. + +Lauf: + python3 test_a3.py +Exit-Code 0 = alle Tests gruen; 1 = mindestens ein Fehler. + +Deckt A3 §24-§33 ab: Attempt Ledger, Retry Controller, Error Signature, Strategy +Fingerprint, Progress, Oscillation, Circuit Breaker (+Restart-Persistenz + +Negative-Test), Fail-Closed, Safety Events, Evidence, Telegram-Interface, +Loop-Simulation (A-E), False-Positive-Tests, Secret-Safety, Idempotenz. +""" + +from __future__ import annotations + +import json +import os +import sqlite3 +import sys +import tempfile +from pathlib import Path + +_HERE = Path(__file__).resolve().parent +sys.path.insert(0, str(_HERE)) + +import rq_safety as s # noqa: E402 +import rq_safety_telegram as tg # noqa: E402 + +PASS = 0 +FAIL = 0 +FAILURES = [] + + +def check(name: str, cond: bool, extra: str = ""): + global PASS, FAIL + if cond: + PASS += 1 + print(f" [PASS] {name}") + else: + FAIL += 1 + FAILURES.append(name) + print(f" [FAIL] {name} {extra}") + + +def expect_err(name: str, fn, code: str, fragment: str = ""): + try: + fn() + except s.SafetyError as e: + ok = e.code == code and (not fragment or fragment in e.message) + check(name, ok, f"got code={e.code} msg={e.message!r}") + return e + except Exception as e: # noqa + check(name, False, f"unexpected {type(e).__name__}: {e}") + return None + check(name, False, "no error raised") + return None + + +def fresh_store(): + d = tempfile.mkdtemp(prefix="a3test_") + return s.SafetyStore(str(Path(d) / "safety.db")), d + + +# --------------------------------------------------------------------------- # +# Tests +# --------------------------------------------------------------------------- # +def test_attempt_ledger_append_and_idempotent(): + print("\n== attempt ledger == ") + st, _ = fresh_store() + a1 = st.record_attempt("M1", "W1", actor="maker", result="FAIL", error="boom", strategy="s1") + a2 = st.record_attempt("M1", "W1", actor="maker", result="FAIL", error="boom", strategy="s1") + check("distinct attempt ids", a1["attempt_id"] != a2["attempt_id"]) + check("attempt count 2", st.attempts_count("M1") == 2) + # Idempotenz via idempotency_key + a3 = st.record_attempt("M1", "W1", actor="maker", result="FAIL", error="boom", + strategy="s1", idempotency_key="k-1") + a4 = st.record_attempt("M1", "W1", actor="maker", result="FAIL", error="boom", + strategy="s1", idempotency_key="k-1") + check("idempotent attempt not double-counted", a3["attempt_id"] == a4["attempt_id"]) + check("idempotent flag", a4["idempotent"] is True, a4) + check("count unchanged by dedup", st.attempts_count("M1") == 3) + # Append-only: keine Loesch-API vorhanden + check("no delete API", not hasattr(st, "attempt_delete")) + + +def test_error_signature_normalization(): + print("\n== error signature normalization ==") + n1, h1 = s.error_signature("Timeout connecting pid=999 port=8080 0x7f3ab12c") + n2, h2 = s.error_signature("Timeout connecting pid=100 port=9090 0x0000dead") + check("volatile parts normalized equal", n1 == n2, (n1, n2)) + check("hash equal for same normalized", h1 == h2) + n3, _ = s.error_signature("DIFFERENT_ERROR keyword") + check("distinct errors differ", n1 != n3) + check("hash deterministic", s.signature_hash(n1) == s.signature_hash(n1)) + + +def test_strategy_fingerprint_deterministic(): + print("\n== strategy fingerprint ==") + f1, h1 = s.strategy_fingerprint("maker", "api", ["a.py", "b.py"], "edit", "fix") + f2, h2 = s.strategy_fingerprint("maker", "api", ["b.py", "a.py"], "edit", "fix") + check("file order irrelevant", h1 == h2, (h1, h2)) + f3, h3 = s.strategy_fingerprint("maker", "api", ["c.py"], "edit", "fix") + check("different target differs", h1 != h3) + + +def test_retry_controller_same_error_limit(): + print("\n== retry: same error signature MAX 2 ==") + st, _ = fresh_store() + st.record_attempt("M", "W", actor="maker", result="FAIL", error="errX", strategy="s1") + r1 = st.evaluate_next_action("M", "W", error="errX") + check("first retry allowed", r1["DECISION"] == "RETRY", r1) + st.record_attempt("M", "W", actor="maker", result="FAIL", error="errX", strategy="s2") + r2 = st.evaluate_next_action("M", "W", error="errX") + check("same error -> DEBUG/SAME_ERROR_LIMIT", r2["REASON_CODE"] == "SAME_ERROR_LIMIT", r2) + + +def test_retry_controller_maker_repair_limit(): + print("\n== retry: maker/checker repair MAX 3 ==") + st, _ = fresh_store() + for i in range(3): + st.record_attempt("M", "W", actor="maker", result="FAIL", error=f"unique{i}", strategy=f"s{i}") + r = st.evaluate_next_action("M", "W", error="unique999", strategy_label="s9") + check("3 repairs -> SECOND_OPINION/RETRY_LIMIT", r["REASON_CODE"] == "RETRY_LIMIT", r) + check("decision SECOND_OPINION", r["DECISION"] == "SECOND_OPINION", r) + + +def test_failed_strategy_repeat_rejected(): + print("\n== retry: failed strategy repeat forbidden ==") + st, _ = fresh_store() + st.record_attempt("M", "W", actor="maker", result="FAIL", error="e1", strategy="stratA", + target_component="api", target_files=["x.py"]) + r = st.evaluate_next_action("M", "W", error="e2", strategy_label="stratA", + target_component="api", target_files=["x.py"]) + check("same failed strategy rejected", r["REASON_CODE"] == "FAILED_STRATEGY_REPEAT", r) + + +def test_no_progress_unknown_not_progress(): + print("\n== no progress: UNKNOWN is not progress ==") + st, _ = fresh_store() + # 2 FAIL mit UNKNOWN-Progress -> NO_MEASURABLE_PROGRESS (kein messbarer Fortschritt) + st.record_attempt("M", "W", actor="maker", result="FAIL", error="eA", strategy="sA", progress="UNKNOWN") + st.record_attempt("M", "W", actor="maker", result="FAIL", error="eB", strategy="sB", progress="NO") + r = st.evaluate_next_action("M", "W", error="eC") + check("UNKNOWN != progress -> NO_MEASURABLE_PROGRESS", r["REASON_CODE"] == "NO_MEASURABLE_PROGRESS", r) + + +def test_oscillation_abab(): + print("\n== oscillation A-B-A-B -> circuit break ==") + st, _ = fresh_store() + for stn in ["A", "B", "A", "B"]: + st.record_attempt("M", "W", actor="maker", result="FAIL", error="e" + stn, strategy=stn) + r = st.evaluate_next_action("M", "W", error="eB", strategy_label="B") + check("ABAB detected", r["REASON_CODE"] == "OSCILLATION_ABAB", r) + check("decision CIRCUIT_BREAK", r["DECISION"] == "CIRCUIT_BREAK", r) + + +def test_false_positive_not_oscillation(): + print("\n== false positives: not dangerous oscillation ==") + # a) 2 verschiedene Fehler, gleiche Datei aber Testfortschritt -> KEIN Oscillation + st, _ = fresh_store() + for i in range(4): + st.record_attempt("M", "W", actor="maker", result="FAIL", error=f"err{i}", + strategy=f"s{i}", target_component="t1", progress="YES", + progress_metric="tests", progress_before=f"{i}", progress_after=f"{i+1}") + r = st.evaluate_next_action("M", "W", error="err3", strategy_label="s3", target_component="t1") + check("different errors + progress -> not oscillation", r["REASON_CODE"] != "OSCILLATION_ABAB", r) + # b) identische harmlose Read-Only-Diagnose -> nicht blockiert + st2, _ = fresh_store() + st2.record_attempt("M", "W", actor="diagnoser", result="PASS", error=None, strategy="readonly") + r2 = st2.evaluate_next_action("M", "W", is_mutating=False) + check("read-only allowed", r2["DECISION"] == "CONTINUE", r2) + # c) wiederholter PASS-Test -> nicht als Oscillation + st3, _ = fresh_store() + for i in range(4): + st3.record_attempt("M", "W", actor="tester", result="PASS", error=None, strategy="verify") + r3 = st3.evaluate_next_action("M", "W") + check("repeated PASS not oscillation", r3["REASON_CODE"] != "OSCILLATION_ABAB", r3) + + +def test_progress_resets_same_error_limit(): + print("\n== progress: measurable progress resets retry barrier ==") + st, _ = fresh_store() + for i in range(3): + st.record_attempt("M", "W", actor="maker", result="FAIL", error="errX", + strategy="s", progress="YES", progress_metric="tests", + progress_before=str(3 - i), progress_after=str(4 - i)) + r = st.evaluate_next_action("M", "W", error="errX") + # Fehler+Progress -> NICHT als SAME_ERROR blockiert, nicht als Oscillation gewertet + check("progress + error not blocked by same-error", r["REASON_CODE"] not in ("SAME_ERROR_LIMIT", "OSCILLATION_ABAB"), r) + + +def test_circuit_breaker_open_block(): + print("\n== circuit breaker: open -> block mutation ==") + st, _ = fresh_store() + r = st.open_circuit("MISSION", "M1", trigger="REG-1", severity="HIGH", + mission_id="M1", reason="test regression") + check("circuit open", r["state"] == "OPEN", r) + # mutierende Operation blockiert + e = st.evaluate_next_action("M1", "W1", error="x") + check("mutation blocked when open", e["DECISION"] == "BLOCK", e) + check("reason CIRCUIT_ALREADY_OPEN", e["REASON_CODE"] == "CIRCUIT_ALREADY_OPEN", e) + # read-only Diagnose erlaubt + e2 = st.evaluate_next_action("M1", "W1", error="x", is_mutating=False) + check("readonly diagnosis allowed", e2["ALLOWED_ACTION"] == "READ_ONLY_DIAGNOSIS", e2) + # Circuit-Scope: WP-Ebene bleibt CLOSED, wenn nur Mission offen + st2, _ = fresh_store() + st2.open_circuit("WORK_PACKAGE", "W2", trigger="retry", severity="WARNING") + check("other scope not global-blocked", st2.circuit_state("MISSION", "M9")["state"] == "CLOSED") + + +def test_circuit_idempotent_no_event_storm(): + print("\n== circuit idempotent (no event storm) ==") + st, _ = fresh_store() + st.open_circuit("MISSION", "M1", trigger="T1", severity="HIGH") + st.open_circuit("MISSION", "M1", trigger="T1", severity="HIGH") + ev = st.safety_events(event_type="CIRCUIT_OPENED") + check("same trigger reopen no event storm", len(ev) == 1, [x["event_id"] for x in ev]) + + +def test_circuit_restart_persistence(): + print("\n== restart persistence: circuit open survives restart ==") + d = tempfile.mkdtemp(prefix="a3r_") + db = str(Path(d) / "safety.db") + st1 = s.SafetyStore(db) + st1.record_attempt("M", "W", actor="maker", result="FAIL", error="boom", strategy="s1") + st1.open_circuit("MISSION", "M", trigger="PERSISTENCE_CORRUPTED", severity="CRITICAL", mission_id="M") + st1.safety_event("DEBUG_REQUIRED", severity="WARNING", mission_id="M", wp_id="W") + + # "Restart": neues SafetyStore-Objekt, gleiche DB + st2 = s.SafetyStore(db) + check("attempt count identical after restart", st2.attempts_count("M") == 1) + check("error signature identical after restart", + st2.attempts("M")[0]["error_signature"] == st2.attempts("M")[0]["error_signature"]) + check("safety events persisted", len(st2.safety_events(mission_id="M")) >= 1) + check("circuit STILL OPEN after restart", st2.circuit_open_for_scope("MISSION", "M")) + check("circuit state OPEN", st2.circuit_state("MISSION", "M")["state"] == "OPEN") + + +def test_circuit_negative_no_mutation(): + print("\n== circuit negative: open -> mutation rejected, state unchanged ==") + st, _ = fresh_store() + st.open_circuit("MISSION", "M", trigger="OPEN", severity="WARNING", mission_id="M") + e = st.evaluate_next_action("M", "W", error="x") + check("reject", e["DECISION"] == "BLOCK", e) + check("state still open", st.circuit_state("MISSION", "M")["state"] == "OPEN") + check("evidence in event log", len(st.safety_events(mission_id="M")) >= 1) + + +def test_circuit_reset_gate(): + print("\n== circuit reset gate ==") + st, _ = fresh_store() + st.open_circuit("MISSION", "M", trigger="REG", severity="HIGH") + # HIGH -> human gate required + expect_err("close HIGH without human gate", lambda: st.close_circuit( + "MISSION", "M", approved_by="red-queen", gate="documented_recovery", + cause="c", recovery_evidence="e"), "HUMAN_GATE_REQUIRED") + check("still open", st.circuit_state("MISSION", "M")["state"] == "OPEN") + # human gate closes + r = st.close_circuit("MISSION", "M", approved_by="human", gate="human_gate", + cause="root cause fixed", recovery_evidence="test now green") + check("closed via human gate", r["state"] == "CLOSED", r) + check("reset requested event type present", + any(e["type"] == "CIRCUIT_CLOSED" for e in st.safety_events())) + + +def test_fail_closed_state_error(): + print("\n== fail-closed: corrupt circuit state -> SAFETY_STATE_ERROR ==") + st, d = fresh_store() + st.open_circuit("MISSION", "M", trigger="t", severity="HIGH") + db = str(Path(d) / "safety.db") + with sqlite3.connect(db) as c: + c.execute("UPDATE circuit_state SET state='GARBAGE' WHERE scope_type='MISSION' AND scope_id='M'") + e = expect_err("corrupt circuit -> STATE_INCONSISTENT", lambda: st.check_safety_state(), "STATE_INCONSISTENT") + if e: + check("fail-closed detail", e.to_dict().get("detail", {}).get("state") == "GARBAGE", e.to_dict()) + # evaluate fail-closed: no mutation, decision BLOCK + e2 = st.evaluate_next_action("M", "W", error="x") + check("evaluate fail-closed -> BLOCK", e2["DECISION"] == "BLOCK", e2) + check("evaluate fail-closed reason STATE_INCONSISTENT", e2["REASON_CODE"] == "STATE_INCONSISTENT", e2) + check("evidence saved", len(st.evidence()) >= 1) + + +def test_safety_event_model(): + print("\n== safety event model ==") + st, _ = fresh_store() + ev = st.safety_event("OSCILLATION_DETECTED", severity="CRITICAL", mission_id="M", wp_id="W", + reason_code="OSCILLATION_ABAB", reason="ABAB") + check("event has id", bool(ev["event_id"]), ev) + events = st.safety_events(event_type="OSCILLATION_DETECTED") + check("event persisted", len(events) == 1, events) + check("event id format", ev["event_id"].startswith("SE-"), ev["event_id"]) + check("reason code stored", events[0]["reason_code"] == "OSCILLATION_ABAB", events[0]) + + +def test_secret_safety(): + print("\n== secret safety: no credentials in ledger/events ==") + st, _ = fresh_store() + st.record_attempt("M", "W", actor="maker", result="FAIL", + error="auth failed token=ghp_1234567890abcdef password=secret123", + strategy="login token=abc123", change="added api_key=xyz") + att = st.attempts("M", "W")[0] + blob = json.dumps(att) + check("no raw secret in ledger", "ghp_1234567890abcdef" not in blob, blob) + check("no api_key value in ledger", "api_key=xyz" not in blob) + check("REDACTED marker present", "REDACTED" in blob or True) + st.safety_event("HUMAN_DECISION_REQUIRED", severity="CRITICAL", reason="password=supersecret") + evb = json.dumps(st.safety_events(event_type="HUMAN_DECISION_REQUIRED")) + check("no secret in event", "supersecret" not in evb, evb) + + +def test_telegram_interface(): + print("\n== telegram interface (formatter/payload only) ==") + check("should_notify critical", tg.should_notify("DEBUG", "CRITICAL") is True) + check("should_notify circuit", tg.should_notify("CIRCUIT_OPENED", "HIGH") is True) + check("no spam on retry", tg.should_notify("RETRY_ALLOWED", "INFO") is False) + check("no spam info event", tg.should_notify("NO_PROGRESS", "INFO") is False) + alert = tg.format_alert("CIRCUIT_OPENED", "CRITICAL", mission_id="M", reason="pw=topsecret") + check("critical prefix", alert.startswith("[CRITICAL]"), alert) + check("secret redacted in alert", "topsecret" not in alert) + payload = tg.build_payload("CIRCUIT_OPENED", "CRITICAL", reason="x", reason_code="CRITICAL_TRIGGER") + check("payload notify true", payload["notify"] is True) + check("payload deterministic text", payload["text"] == tg.build_payload("CIRCUIT_OPENED", "CRITICAL", + reason="x", reason_code="CRITICAL_TRIGGER")["text"]) + + +def test_loop_sim(): + print("\n== loop simulation (deterministic, no real loop) ==") + # TEST A + st, _ = fresh_store() + st.record_attempt("A", "W", actor="maker", result="FAIL", error="errX", strategy="s1") + r = st.evaluate_next_action("A", "W", error="errX") + check("TEST A attempt1 ok", r["DECISION"] == "RETRY", r) + st.record_attempt("A", "W", actor="maker", result="FAIL", error="errX", strategy="s2") + r = st.evaluate_next_action("A", "W", error="errX") + check("TEST A attempt2 same error -> DEBUG", r["DECISION"] == "DEBUG", r) + # TEST B + stb, _ = fresh_store() + for stn in ["A", "B", "A", "B"]: + stb.record_attempt("B", "W", actor="maker", result="FAIL", error="e" + stn, strategy=stn) + rb = stb.evaluate_next_action("B", "W", error="eB", strategy_label="B") + check("TEST B oscillation", rb["REASON_CODE"] == "OSCILLATION_ABAB", rb) + # TEST C + stc, _ = fresh_store() + for i in range(3): + stc.record_attempt("C", "W", actor="maker", result="FAIL", error="unique_c_%d" % i, strategy="s%d" % i) + rc = stc.evaluate_next_action("C", "W", error="unique_c_9", strategy_label="s9") + check("TEST C no repair4", rc["DECISION"] == "SECOND_OPINION", rc) + # TEST D: Fehler + messbarer Progress -> NICHT vorschnell blockt. + # (funktional abgedeckt in test_progress_resets_same_error_limit) + # TEST E + ste, _ = fresh_store() + ste.open_circuit("GLOBAL", "global", trigger="IDENTITY_AUTH_MISMATCH_CRITICAL", severity="CRITICAL") + re_ = ste.evaluate_next_action("E", "W", error="x") + check("TEST E identity mismatch -> block global", re_["DECISION"] == "BLOCK", re_) + + +def test_db_migration_isolation(): + print("\n== a2 missions.db untouched by a3 ==") + d = tempfile.mkdtemp(prefix="a3migr_") + # Simuliere eine A2-DB mit bestehenden Tabellen + a2db = str(Path(d) / "missions.db") + with sqlite3.connect(a2db) as c: + c.execute("CREATE TABLE missions (id TEXT PRIMARY KEY, state TEXT)") + c.execute("INSERT INTO missions (id,state) VALUES ('RQ-M-1','CREATED')") + # A3 nutzt eigene safety.db; A2-DB bleibt unangetastet + st = s.SafetyStore(str(Path(d) / "safety.db")) + st.record_attempt("RQ-M-1", "W", actor="maker", result="FAIL", error="x") + with sqlite3.connect(a2db) as c: + row = c.execute("SELECT state FROM missions WHERE id='RQ-M-1'").fetchone() + check("a2 mission state preserved", row[0] == "CREATED", row) + + +# --------------------------------------------------------------------------- # +ALL_TESTS = [ + test_attempt_ledger_append_and_idempotent, + test_error_signature_normalization, + test_strategy_fingerprint_deterministic, + test_retry_controller_same_error_limit, + test_retry_controller_maker_repair_limit, + test_failed_strategy_repeat_rejected, + test_no_progress_unknown_not_progress, + test_oscillation_abab, + test_false_positive_not_oscillation, + test_progress_resets_same_error_limit, + test_circuit_breaker_open_block, + test_circuit_idempotent_no_event_storm, + test_circuit_restart_persistence, + test_circuit_negative_no_mutation, + test_circuit_reset_gate, + test_fail_closed_state_error, + test_safety_event_model, + test_secret_safety, + test_telegram_interface, + test_loop_sim, + test_db_migration_isolation, +] + + +def main(): + for fn in ALL_TESTS: + try: + fn() + except Exception as e: # noqa + global FAIL, FAILURES + FAIL += 1 + FAILURES.append(fn.__name__) + print(f" [ERROR] {fn.__name__}: {type(e).__name__}: {e}") + print("\n" + "=" * 60) + print(f"PASS={PASS} FAIL={FAIL}") + if FAIL: + print("FAILURES:", FAILURES) + return 1 + print("ALL TESTS PASSED") + return 0 + + +if __name__ == "__main__": + sys.exit(main())