"""Orchestrierung: fünf Ebenen strikt seriell, je Ebene bauen → QA → Repair-Loop bis 100 %/Stillstand/Limit. Resume ist der Normalfall: der Fortschritt liegt komplett in der DB (topics.status + Item-Status), ein neuer Start überspringt Fertiges. Infra-Erschöpfung pausiert den Lauf (fail-closed), Budget stoppt hart.""" import asyncio import logging import agents import artefakte import auto_loop import db import guide import inventar import korpus import ledger import llm import qa import struktur from config import RUN_BUDGET_TOKENS log = logging.getLogger("creator2.pipeline") EBENEN = [ ("korpus", korpus, "korpus_fertig"), ("inventar", inventar, "inventar_fertig"), ("artefakte", artefakte, "artefakte_fertig"), ("struktur", struktur, "struktur_fertig"), ("guide", guide, "fertig"), ] _ORDNUNG = ["neu", "korpus", "korpus_fertig", "inventar_fertig", "artefakte_fertig", "struktur_fertig", "fertig"] _laeufe: dict[str, asyncio.Task] = {} def _fertig_ab(status: str, marke: str) -> bool: try: return _ORDNUNG.index(status) >= _ORDNUNG.index(marke) except ValueError: return False def lauf_starten(topic: str, budget: int | None = None) -> int: if topic in _laeufe and not _laeufe[topic].done(): raise RuntimeError("Lauf läuft bereits") agents.abbruch_aufheben(f"{topic}-") git_hash, dirty = ledger.git_stand() run_id = db.insert("runs", topic=topic, status="running", git_hash=git_hash, git_dirty=int(dirty), budget_tokens=budget if budget is not None else RUN_BUDGET_TOKENS) _laeufe[topic] = asyncio.ensure_future(_lauf(topic, run_id)) return run_id def _tief_zuruecksetzen(topic: str, ebene: str) -> None: """tief=Re-Verify: Verifiziertes auf 'kandidat' zurückstufen — der BESTEHENDE Verify-Weg beurteilt dann alles neu (verifiziert oder verworfen = Entfernen). Nur wo es einen LLM-Verify gibt; korpus (Soll-Freeze) bleibt bewusst außen vor, inventar/struktur prüfen ohnehin deterministisch alles bei jedem messen.""" if ebene == "artefakte": db.execute("UPDATE artefakte SET status='kandidat' WHERE status='verifiziert'" " AND atom_id IN (SELECT id FROM atome WHERE topic=?)", (topic,)) elif ebene == "guide": db.execute("UPDATE sections SET stage='pruefer', qa_hash='', fix_versuche=0" " WHERE baustein_id IN (SELECT id FROM bausteine WHERE topic=?)", (topic,)) def aktualisierung_starten(topic: str, ebene: str, tief: bool = False) -> int: """Eine Ebene (und bei Auto-Flags die folgenden) erneut fahren — ohne Reset, ohne Status-Rückschritt. bauen füllt Fehlendes, messen+Repair fixt Kaputtes. tief=True: zusätzlich alles Verifizierte erneut durchs Panel (teuer).""" if ebene not in [e[0] for e in EBENEN]: raise ValueError(f"unbekannte Ebene: {ebene}") if topic in _laeufe and not _laeufe[topic].done(): raise RuntimeError("Lauf läuft bereits") agents.abbruch_aufheben(f"{topic}-") if tief: _tief_zuruecksetzen(topic, ebene) git_hash, dirty = ledger.git_stand() run_id = db.insert("runs", topic=topic, status="running", git_hash=git_hash, git_dirty=int(dirty), budget_tokens=RUN_BUDGET_TOKENS, grund=f"aktualisieren:{ebene}" + (":tief" if tief else "")) _laeufe[topic] = asyncio.ensure_future(_lauf(topic, run_id, ab=ebene, erzwingen=True)) return run_id def lauf_pausieren(topic: str) -> None: """Sanfte Pause: Wartende + Neue stoppen, Laufende laufen aus. Der Lauf endet als 'paused', sobald der nächste Call die Pause bemerkt (ManuellePause); Start setzt normal fort (abbruch_aufheben hebt die Pause auf).""" agents.pausieren(f"{topic}-") def lauf_stoppen(topic: str) -> None: agents.abbrechen(f"{topic}-") task = _laeufe.get(topic) if task and not task.done(): task.cancel() run = db.one("SELECT id FROM runs WHERE topic=? AND status='running'" " ORDER BY id DESC", (topic,)) if run: db.update("runs", "id", run["id"], status="stopped", beendet=db.now()) _ROLLBACK = {"guide": "struktur_fertig", "struktur": "artefakte_fertig", "artefakte": "inventar_fertig", "inventar": "korpus_fertig", "korpus": "neu"} def ebenen_entfernen(topic: str, ebene: str) -> None: """Ergebnisse einer Ebene UND aller nachgelagerten verwerfen (rechts hängt an links), Topic-Status entsprechend zurückrollen. Quellen-Snapshots bleiben, außer bei korpus.""" if ebene not in _ROLLBACK: raise ValueError(f"unbekannte Ebene: {ebene}") lauf_stoppen(topic) stufe = ["guide", "struktur", "artefakte", "inventar", "korpus"].index(ebene) db.execute("DELETE FROM sections WHERE baustein_id IN" " (SELECT id FROM bausteine WHERE topic=?)", (topic,)) db.execute("DELETE FROM auftraege WHERE baustein_id IN" " (SELECT id FROM bausteine WHERE topic=?)", (topic,)) db.execute("UPDATE bausteine SET status='neu' WHERE topic=?", (topic,)) db.execute("UPDATE kapitel SET intro='' WHERE topic=?", (topic,)) # Intro ist Guide-Text if stufe >= 1: # struktur db.execute("DELETE FROM bausteine WHERE topic=?", (topic,)) db.execute("DELETE FROM kapitel WHERE topic=?", (topic,)) db.execute("DELETE FROM lernziele WHERE topic=?", (topic,)) # Themen-Schicht fällt mit — das Soll selbst bleibt eingefroren db.execute("DELETE FROM themen WHERE topic=?", (topic,)) db.execute("UPDATE soll SET thema_id=NULL WHERE topic=?", (topic,)) db.execute("UPDATE atome SET ziel_id=NULL, baustein_id=NULL, ord=0 WHERE topic=?", (topic,)) if stufe >= 2: # artefakte db.execute("DELETE FROM leitner WHERE artefakt_id IN (SELECT ar.id FROM artefakte ar" " JOIN atome a ON ar.atom_id=a.id WHERE a.topic=?)", (topic,)) db.execute("DELETE FROM artefakte WHERE atom_id IN" " (SELECT id FROM atome WHERE topic=?)", (topic,)) if stufe >= 3: # inventar db.execute("DELETE FROM anker WHERE atom_id IN (SELECT id FROM atome WHERE topic=?)", (topic,)) db.execute("DELETE FROM kanten WHERE topic=?", (topic,)) db.execute("DELETE FROM atome WHERE topic=?", (topic,)) # atome_stand mit-nullen — sonst überspringt der Resume-Guard die Neu-Extraktion db.execute("UPDATE quellen SET status='extrahiert', atome_stand='' WHERE topic=?", (topic,)) if stufe >= 4: # korpus db.execute("DELETE FROM soll WHERE topic=?", (topic,)) db.execute("DELETE FROM quellen WHERE topic=?", (topic,)) import shutil from config import KORPUS_DIR # Gürtel-und-Hosenträger vor dem destruktiven rmtree: der Zielpfad MUSS # unter KORPUS_DIR liegen (Alt-DBs könnten Namen ohne heutige Validierung # tragen). Sonst nichts löschen. ziel = (KORPUS_DIR / topic).resolve() if ziel != KORPUS_DIR.resolve() and ziel.is_relative_to(KORPUS_DIR.resolve()): shutil.rmtree(ziel, ignore_errors=True) # Token-/Dauer-Anzeige der entfernten Ebenen nullen: Events bleiben (Ledger ist # append-only), aber die Zählung beginnt hinter dem Reset-Marker neu. letzte = db.one("SELECT MAX(id) AS m FROM events")["m"] or 0 resets = db.uj(db.one("SELECT resets FROM topics WHERE name=?", (topic,))["resets"], {}) for name in ["guide", "struktur", "artefakte", "inventar", "korpus"][:stufe + 1]: resets[name] = letzte db.update("topics", "name", topic, status=_ROLLBACK[ebene], resets=db.j(resets)) def topic_loeschen(topic: str) -> None: lauf_stoppen(topic) ebenen_entfernen(topic, "korpus") for run in db.query("SELECT id FROM runs WHERE topic=?", (topic,)): db.execute("DELETE FROM events WHERE run_id=?", (run["id"],)) db.execute("DELETE FROM befunde WHERE run_id=?", (run["id"],)) db.execute("DELETE FROM runs WHERE topic=?", (topic,)) db.execute("DELETE FROM topics WHERE name=?", (topic,)) def voll_reset(topic: str) -> None: """Alles Generierte verwerfen (auch Atome + Artefakte); Quellen-Snapshots und Lauf-Historie bleiben. Für Neuextraktion nach Extraktions-Regeländerungen.""" lauf_stoppen(topic) soll_reset(topic) db.execute("DELETE FROM leitner WHERE artefakt_id IN (SELECT ar.id FROM artefakte ar" " JOIN atome a ON ar.atom_id=a.id WHERE a.topic=?)", (topic,)) db.execute("DELETE FROM artefakte WHERE atom_id IN" " (SELECT id FROM atome WHERE topic=?)", (topic,)) db.execute("DELETE FROM anker WHERE atom_id IN (SELECT id FROM atome WHERE topic=?)", (topic,)) db.execute("DELETE FROM kanten WHERE topic=?", (topic,)) db.execute("DELETE FROM atome WHERE topic=?", (topic,)) # Atome weg → Merker nullen, sonst überspringt der Guard die Neu-Extraktion db.execute("UPDATE quellen SET atome_stand='' WHERE topic=?", (topic,)) db.update("topics", "name", topic, status="neu") def soll_reset(topic: str) -> None: """Soll/Struktur/Guide verwerfen, Atome + Artefakte behalten. Der nächste Lauf baut das Soll neu; der Extraktions-Guard in inventar verhindert Doppel-Atome.""" lauf_stoppen(topic) db.execute("DELETE FROM sections WHERE baustein_id IN" " (SELECT id FROM bausteine WHERE topic=?)", (topic,)) db.execute("DELETE FROM auftraege WHERE baustein_id IN" " (SELECT id FROM bausteine WHERE topic=?)", (topic,)) db.execute("DELETE FROM bausteine WHERE topic=?", (topic,)) db.execute("DELETE FROM kapitel WHERE topic=?", (topic,)) db.execute("DELETE FROM lernziele WHERE topic=?", (topic,)) db.execute("DELETE FROM themen WHERE topic=?", (topic,)) db.execute("DELETE FROM soll WHERE topic=?", (topic,)) db.execute("UPDATE atome SET soll_id=NULL, ziel_id=NULL, baustein_id=NULL, ord=0" " WHERE topic=?", (topic,)) db.execute("UPDATE quellen SET status='neu' WHERE topic=?", (topic,)) db.update("topics", "name", topic, status="korpus") async def _lauf(topic: str, run_id: int, ab: str | None = None, erzwingen: bool = False) -> None: """Normaler Lauf ODER Aktualisierung (ab=Ebene, erzwingen=True): Aktualisieren fährt fertige Ebenen erneut — bauen füllt Fehlendes, messen prüft Form, Repair repariert. Nötig, weil der Resume-Skip fertige Ebenen sonst nie wieder anfasst (Altbestands-Lücke: Katalog-Titel, Aussagen-Nachrüstung, Beleg-Fenster).""" t = db.one("SELECT * FROM topics WHERE name=?", (topic,)) ctx = llm.Kontext(run_id, topic, t["provider"]) start_idx = [e[0] for e in EBENEN].index(ab) if ab else 0 try: for name, modul, marke in EBENEN[start_idx:]: status = db.one("SELECT status FROM topics WHERE name=?", (topic,))["status"] if not erzwingen and _fertig_ab(status, marke): continue db.update("runs", "id", run_id, ebene=name) await modul.bauen(ctx) erst_note, _ = await qa.messen(ctx, name) async def reparieren(n=name, m=modul): offene = db.query("SELECT * FROM befunde WHERE run_id=? AND ebene=?" " AND status='offen'", (run_id, n)) bewegt = await m.reparieren(ctx, offene) neue_note, _ = await qa.messen(ctx, n) return neue_note, bewegt note, grund = await auto_loop.auto_repair_loop(name, erst_note, reparieren) log.info("%s/%s: Note %.1f (%s)", topic, name, note, grund) # Fail-closed-Gate vor dem Ebenen-Marker (nur wo das Modul eins hat): # ein RuntimeError landet im Exception-Handler → Run 'failed', der # Ebenen-Marker wird NICHT gesetzt, die nächste Ebene startet nie. gate = getattr(modul, "gate", None) if gate and (grund_gate := gate(ctx)): raise RuntimeError(f"{name}-Gate: {grund_gate}") # Marker-Schutz: nie zurückdrehen — Aktualisieren einer frühen Ebene # darf z. B. 'fertig' nicht auf 'artefakte_fertig' zurückstufen. if not _fertig_ab(status, marke): db.update("topics", "name", topic, status=marke) flags = db.uj(db.one("SELECT auto FROM topics WHERE name=?", (topic,))["auto"], {}) if not flags.get(name, True) and marke != "fertig": db.update("runs", "id", run_id, status="done", beendet=db.now(), grund=f"Auto-Stopp nach {name}") return db.update("runs", "id", run_id, status="done", beendet=db.now()) except llm.LaufPause as e: agents.abbrechen(f"{topic}-") # Nachzügler-Tasks dürfen keine Tokens mehr ziehen db.update("runs", "id", run_id, status="paused", grund=str(e)[:300], beendet=db.now()) log.warning("%s: Lauf pausiert — %s", topic, e) except ledger.BudgetErschoepft as e: agents.abbrechen(f"{topic}-") db.update("runs", "id", run_id, status="budget", grund=str(e)[:300], beendet=db.now()) log.warning("%s: Budget erreicht", topic) except asyncio.CancelledError: raise except Exception as e: agents.abbrechen(f"{topic}-") db.update("runs", "id", run_id, status="failed", grund=f"{type(e).__name__}: {e}"[:300], beendet=db.now()) log.exception("%s: Lauf fehlgeschlagen", topic)