"""SQLite-Schicht: EINE Writer-Connection (eliminiert 'database is locked' architektonisch, R8), WAL, kurze Transaktionen. Zustand liegt NUR hier — Resume = Prozess neu starten, fertige Tasks überspringen.""" import json import re import sqlite3 import threading from datetime import datetime, timezone from typing import Any, Callable from . import config _conn: sqlite3.Connection | None = None _lock = threading.RLock() on_change: Callable[[str, dict], None] | None = None # ws.Hub.push _LIVE_TABELLEN = {"topics", "runs", "tasks", "gate_laeufe", "befunde"} _RUN_TABELLEN = {"topics", "runs"} # ws-Bereich "run", Rest "graph" # Topic-skopierte Pipeline-Tabellen — EINE Wahrheit für Reset (main.py) PIPELINE_TABELLEN = ("tasks", "gate_laeufe", "quellen", "soll", "atome", "kanten", "lernziele", "bausteine", "kapitel") _TABELLE_RE = re.compile( r"^\s*(?:UPDATE|DELETE\s+FROM|INSERT(?:\s+OR\s+\w+)?\s+INTO)\s+(\w+)", re.IGNORECASE) SCHEMA = """ CREATE TABLE IF NOT EXISTS topics( name TEXT PRIMARY KEY, titel TEXT DEFAULT '', status TEXT DEFAULT 'neu', quellen_max INTEGER, erstellt TEXT DEFAULT (datetime('now'))); CREATE TABLE IF NOT EXISTS runs( id INTEGER PRIMARY KEY, topic TEXT NOT NULL, status TEXT DEFAULT 'running', gestartet TEXT, beendet TEXT, stufe TEXT DEFAULT '', budget_tokens INTEGER DEFAULT 0, grund TEXT DEFAULT '', git_hash TEXT DEFAULT ''); CREATE TABLE IF NOT EXISTS tasks( id INTEGER PRIMARY KEY, run_id INTEGER NOT NULL, topic TEXT NOT NULL, knoten TEXT NOT NULL, item TEXT NOT NULL, status TEXT DEFAULT 'offen', art TEXT DEFAULT '', versuch INTEGER DEFAULT 0, infra_versuch INTEGER DEFAULT 0, runde INTEGER DEFAULT 0, payload TEXT DEFAULT '{}', ergebnis TEXT DEFAULT '{}', fehler TEXT DEFAULT '', erzeugt_von TEXT DEFAULT '', gestartet TEXT, beendet TEXT, UNIQUE(topic, knoten, item)); CREATE INDEX IF NOT EXISTS idx_tasks_disp ON tasks(run_id, status, knoten); CREATE TABLE IF NOT EXISTS gate_laeufe( id INTEGER PRIMARY KEY, run_id INTEGER, topic TEXT, knoten TEXT, runde INTEGER, status TEXT DEFAULT '', fingerprint TEXT DEFAULT '', befunde TEXT DEFAULT '[]', ts TEXT, UNIQUE(topic, knoten, runde)); CREATE TABLE IF NOT EXISTS events( id INTEGER PRIMARY KEY, run_id INTEGER, ts TEXT, stufe TEXT DEFAULT '', knoten TEXT DEFAULT '', item TEXT DEFAULT '', skills TEXT DEFAULT '[]', skill_hash TEXT DEFAULT '', model TEXT DEFAULT '', role TEXT DEFAULT '', status TEXT DEFAULT '', dur_ms INTEGER DEFAULT 0, wait_ms INTEGER DEFAULT 0, tok_in INTEGER DEFAULT 0, tok_out INTEGER DEFAULT 0, tok_cache_read INTEGER DEFAULT 0, tok_cache_write INTEGER DEFAULT 0, meta TEXT DEFAULT '{}'); CREATE INDEX IF NOT EXISTS idx_events_run ON events(run_id); CREATE TABLE IF NOT EXISTS call_texte( event_id INTEGER PRIMARY KEY, prompt TEXT DEFAULT '', antwort TEXT DEFAULT ''); CREATE TABLE IF NOT EXISTS quellen( id INTEGER PRIMARY KEY, topic TEXT, titel TEXT DEFAULT '', url TEXT DEFAULT '', url_norm TEXT DEFAULT '', backend TEXT DEFAULT '', snapshot TEXT DEFAULT '', roh TEXT DEFAULT '', hash TEXT DEFAULT '', runde INTEGER DEFAULT 0, zweck TEXT DEFAULT 'korpus', soll_id INTEGER, status TEXT DEFAULT 'neu', grund TEXT DEFAULT '', atome_stand TEXT DEFAULT '', UNIQUE(topic, url_norm)); CREATE TABLE IF NOT EXISTS soll( id INTEGER PRIMARY KEY, topic TEXT, punkt TEXT, status TEXT DEFAULT 'kandidat', belege TEXT DEFAULT '[]', kapitel_id INTEGER, geprueft TEXT DEFAULT '[]', freispruch TEXT DEFAULT '', verdichtet INTEGER DEFAULT 0); CREATE TABLE IF NOT EXISTS atome( id INTEGER PRIMARY KEY, topic TEXT, titel TEXT, typ TEXT DEFAULT '', definition TEXT DEFAULT '', status TEXT DEFAULT 'neu', soll_id INTEGER, soll_geprueft INTEGER DEFAULT 0, ziel_id INTEGER, baustein_id INTEGER, ord INTEGER DEFAULT 0, merged_into INTEGER); CREATE INDEX IF NOT EXISTS idx_atome_topic ON atome(topic); CREATE TABLE IF NOT EXISTS anker( id INTEGER PRIMARY KEY, atom_id INTEGER, quelle_id INTEGER, start INTEGER DEFAULT -1, ende INTEGER DEFAULT -1, zitat TEXT DEFAULT ''); CREATE INDEX IF NOT EXISTS idx_anker_atom ON anker(atom_id); CREATE TABLE IF NOT EXISTS kanten( id INTEGER PRIMARY KEY, topic TEXT, von_atom INTEGER, zu_atom INTEGER, art TEXT DEFAULT 'verwandt', status TEXT DEFAULT 'aktiv', UNIQUE(von_atom, zu_atom, art)); CREATE TABLE IF NOT EXISTS lernziele( id INTEGER PRIMARY KEY, topic TEXT, titel TEXT DEFAULT '', text TEXT DEFAULT '', soll_id INTEGER, status TEXT DEFAULT 'neu'); CREATE TABLE IF NOT EXISTS bausteine( id INTEGER PRIMARY KEY, topic TEXT, ziel_id INTEGER, titel TEXT DEFAULT '', ord INTEGER DEFAULT 0, kapitel_id INTEGER, status TEXT DEFAULT 'neu'); CREATE TABLE IF NOT EXISTS kapitel( id INTEGER PRIMARY KEY, topic TEXT, titel TEXT DEFAULT '', intro TEXT DEFAULT '', ord INTEGER DEFAULT 0); CREATE TABLE IF NOT EXISTS sections( baustein_id INTEGER PRIMARY KEY, stage TEXT DEFAULT 'writer', text TEXT DEFAULT '', qa_hash TEXT DEFAULT '', fix_versuche INTEGER DEFAULT 0); CREATE TABLE IF NOT EXISTS befunde( id INTEGER PRIMARY KEY, run_id INTEGER, stufe TEXT DEFAULT '', knoten TEXT DEFAULT '', art TEXT DEFAULT '', item TEXT DEFAULT '', detail TEXT DEFAULT '', status TEXT DEFAULT 'offen', quelle TEXT DEFAULT '', runde INTEGER DEFAULT 0, geklaert INTEGER DEFAULT 0); CREATE INDEX IF NOT EXISTS idx_befunde_run ON befunde(run_id, stufe); """ def connect(pfad: str | None = None) -> None: global _conn with _lock: if _conn is not None: _conn.close() ziel = pfad or str(config.DB_PFAD) if ziel != ":memory:": config.DB_PFAD.parent.mkdir(parents=True, exist_ok=True) _conn = sqlite3.connect(ziel, check_same_thread=False) _conn.row_factory = sqlite3.Row _conn.execute("PRAGMA journal_mode=WAL") _conn.execute("PRAGMA synchronous=NORMAL") _conn.execute("PRAGMA busy_timeout=5000") _conn.execute("PRAGMA foreign_keys=ON") _conn.executescript(SCHEMA) for migration in ( # Mini-Migrationen für Bestands-DBs "ALTER TABLE atome ADD COLUMN soll_geprueft INTEGER DEFAULT 0", "ALTER TABLE topics ADD COLUMN quellen_max INTEGER", "ALTER TABLE soll ADD COLUMN verdichtet INTEGER DEFAULT 0", "ALTER TABLE quellen ADD COLUMN soll_id INTEGER"): try: _conn.execute(migration) except sqlite3.OperationalError: pass _conn.commit() def reset_for_tests() -> None: connect(":memory:") def now() -> str: return datetime.now(timezone.utc).strftime("%Y-%m-%d %H:%M:%S") def j(x: Any) -> str: return json.dumps(x, ensure_ascii=False) def uj(s: str, default: Any = None) -> Any: try: return json.loads(s) except (TypeError, ValueError): return default def _notify(tabelle: str, row: dict) -> None: if on_change and tabelle in _LIVE_TABELLEN: on_change(tabelle, row) def query(sql: str, *args) -> list[dict]: with _lock: return [dict(r) for r in _conn.execute(sql, args).fetchall()] def one(sql: str, *args) -> dict | None: rows = query(sql, *args) return rows[0] if rows else None def execute(sql: str, *args) -> int: """Einzelnes Statement + Commit. Rückgabe: rowcount (für Claim-Muster).""" with _lock: cur = _conn.execute(sql, args) _conn.commit() m = _TABELLE_RE.match(sql) # UPDATE/DELETE/INSERT — nicht Positions-Raten _notify(m.group(1).lower() if m else "", {}) return cur.rowcount def insert(tabelle: str, ignore: bool = False, **felder) -> int | None: """INSERT (OR IGNORE); Rückgabe lastrowid oder None bei ignoriertem Duplikat.""" keys = ", ".join(felder) qs = ", ".join("?" for _ in felder) verb = "INSERT OR IGNORE" if ignore else "INSERT" with _lock: cur = _conn.execute(f"{verb} INTO {tabelle}({keys}) VALUES({qs})", tuple(felder.values())) _conn.commit() _notify(tabelle, dict(felder)) return cur.lastrowid if cur.rowcount else None def update(tabelle: str, wo: str, wo_args: tuple, **felder) -> int: setzt = ", ".join(f"{k}=?" for k in felder) with _lock: cur = _conn.execute(f"UPDATE {tabelle} SET {setzt} WHERE {wo}", tuple(felder.values()) + wo_args) _conn.commit() _notify(tabelle, dict(felder)) return cur.rowcount class tx: """Transaktion für atomare Mehrfach-Schreiber (Ergebnis + fertig + Folgetasks). Innerhalb: cursor.execute; niemals LLM-Calls.""" def __enter__(self): _lock.acquire() _conn.execute("BEGIN IMMEDIATE") return _conn def __exit__(self, typ, wert, tb): try: if typ is None: _conn.commit() _notify("tasks", {}) else: _conn.rollback() finally: _lock.release() return False