"""Bausteine-Pipeline: Recherche-Konsens + Klärungs-Loop — reines Inventar, unsortiert. 5x Recherche (min. 3, Grace) → Mapping (Konsens/Rest) → Klärungs-Loop (max. KONSENS_MAX_RUNDEN Runden): 3 Auswahl-Agenten (min. 2, Grace) entscheiden über den strittigen Rest, ein Mapping-Agent sortiert in aufnehmen/verwerfen/ weiter strittig. Leerer Rest beendet den Loop; die letzte Runde muss alles entscheiden. Races nutzen ein Grace-Fenster statt „erste N gewinnen": Nach dem ersten gültigen Ergebnis dürfen die übrigen Agenten KONSENS_GRACE Sekunden fertig werden. Der Konsens wird im Code akkumuliert — kein Agent re-emittiert die Gesamtliste. """ import asyncio import json import logging import math import re import shutil import subprocess import time from pathlib import Path import database as db import embedding from agents import kill_process, cancel_scope, clear_scope, run_agent from config import KONSENS_GRACE, RECHERCHE_GRACE, KONSENS_MAX_RUNDEN, DEFAULT_PROVIDER, CRAWL_KEEP_PATTERNS, CRAWL_NOISE_PATTERNS, CRAWL_MIN_CHARS, QUELLE_RELEVANZ_CHUNK, QUELLE_RELEVANZ_SNIPPET, EMBEDDING_AKTIV, EMBEDDING_SUB_DUP from fsutil import atomic_write_text, atomic_write_json from jsonio import read_json_file as _json_datei from paths import arbeit_dir, bausteine_path, frage_muster_path, project_dir, subbausteine_path, quelle_path, quelle_crawl_dir, safe_ordner from crawl import crawl from pipeline import ( CANCELLED, FAILED, OK, GenContext, _extra, _gather_fortschritt, _janein_schema, _log, _prompt, _race, _relevanz_schema, _runde_schema, _semaphore, _str_liste, _stufen_schema, _timeout, run_single_slot, ) from textkit import ( _eindeutige_titel, _lade_bausteine, _norm_titel, _parse_auswahl, _parse_subbausteine, _titel, _titel_aufloesen, _titel_index, ) # Subbausteine (Websuche je Baustein) chunken: 1 Agent je ~10 Bausteine, gedeckelt. SUBBAUSTEIN_CHUNK = 10 SUBBAUSTEIN_MAX = 40 # Einstufen ist billig (kurzes Urteil, keine Websuche) → größere Pakete, weniger Dateien/Agenten. STUFE_CHUNK = 100 # Recherche: feste Datei-Batches statt Such-Loop → jede Crawl-Seite genau einmal zugeteilt. RECHERCHE_BATCH = 20 # Crawl-Seiten je Batch RECHERCHE_READERS = 2 # Reader-Agenten je Batch (Konsens ≥2 innerhalb des Batches) RECHERCHE_THEMA_AGENTEN = 5 # Web-Modus (Quelle „thema", kein Crawl-Ordner) RECHERCHE_KAPPE = 1800 # Sicherheits-Deckel je Batch-Agent # uni/projekt: Skript-Text in Abschnitte ~dieser Größe chunken (gegen Lost-in-the-Middle bei # großen Dokumenten). ~12k Zeichen ≈ 3k Token → sicher unter der Recall-Abfall-Schwelle. RECHERCHE_ABSCHNITT_ZEICHEN = 12000 # Sichtung (Content/Noise) ist jetzt ein deterministischer Regel-Filter (config.CRAWL_*). SUBBAUSTEIN_KAPPE = 900 # Subbaustein-Finde-Loop je Chunk (15 min) KONSOLIDIERUNG_CHUNK = 600 # bis hierher EIN globaler Judge (dedupt alles); darüber chunked + Merge-Pass — nur Fallback-Pfad DEDUP_MAX_RUNDEN = 3 # finaler Dedup-Pass: max. Iterationen (die kleinere Liste blockt je Runde neu) DEDUP_MIN_DELTA = 3 # Abbruch, wenn eine Runde weniger als dieses % der Liste entfernt (konvergiert) # Frage-Muster-Chunks per LPT nach Sub-Last balancieren (Makespan), statt nach Baustein-Anzahl. FRAGE_CHUNK_SUBS = 50 # Ziel-Summe relevanter Subs je Chunk FAKTEN_CHUNK_SUBS = 25 # Fakten-Extraktion: kleinere Chunks (Fakten sind umfangreicher als Muster) FAKTEN_CHECK_PANEL = 3 # Judges je Chunk im Fakten-Check (Mehrheit beanstandet) KONSOLIDIERUNG_PANEL = 3 # Mapping-Judges je Chunk (Panel → Reconcile statt Einzel-Judge) SUBBAUSTEIN_PANEL = 3 # Source-Judges in der Subbaustein-Klärung (Mehrheit statt Einzel-Judge) log = logging.getLogger("creator.bausteine") _bausteine_progress: dict[str, str] = {} _bausteine_errors: dict[str, str] = {} _bausteine_cancelled: set[str] = set() _bausteine_step: dict[str, int] = {} BAUSTEINE_STEPS = ( "Recherche", "Konsolidierung", "Klärung", "Subbausteine finden", "Subbausteine wählen", "Subbausteine klären", "Fakten finden", "Fakten prüfen", "Fakten fix", "Stufen finden", "Stufen wählen", "Stufen klären", "Relevanz finden", "Relevanz wählen", "Relevanz klären", "Gliederung", "Fragen finden", "Fragen wählen", "Fragen klären", "Fragen prüfen", "Karteikarten", "Beispiele", ) ARTEFAKT_TYPEN = ("karteikarte", "beispiel") def lade_quelle(topic: str) -> dict: """Persistierte Quellen-Wahl lesen. Fallback (Alt-Themen ohne quelle.json): existiert projects/ → projekt, sonst thema.""" q = _json_datei(quelle_path(topic)) if isinstance(q, dict) and q.get("type") in ("thema", "projekt", "uni", "link"): return q if project_dir(topic).is_dir(): return {"type": "projekt", "ort": f"projects/{topic}", "spec": ""} return {"type": "thema", "ort": "", "spec": ""} def quelle_ordner(topic: str) -> Path | None: """Ordner-Quelle (projekt/uni → Pfad, link → Crawl-Ordner) — sonst None (thema).""" q = lade_quelle(topic) if q["type"] == "link": return quelle_crawl_dir(topic) if q["type"] in ("projekt", "uni"): return safe_ordner(q.get("ort", "")) return None def _crawl_fertig(topic: str) -> bool: return (quelle_crawl_dir(topic) / ".done").exists() # Marker erst bei sauberem Abschluss # Lernpfad-Stufen (anfaenger/fortgeschritten/experte); alte Schwierigkeits-Werte abwärtskompatibel. _STUFEN = ("anfaenger", "fortgeschritten", "experte", "einfach", "mittel", "schwer") async def subbausteine_titel(topic: str, baustein: str) -> list[str]: """Subbaustein-Titel eines Bausteins — DB-first (Konsens), Fallback Sidecar-Datei.""" rows = [s["sub_titel"] for s in await db.list_subbausteine(topic, _norm_titel(baustein)) if s["status"] == "konsens" and s["sub_titel"]] if rows: return rows sc = _json_datei(subbausteine_path(topic)) if not isinstance(sc, dict): return [] return [ t for s in (sc.get(baustein) or []) if isinstance(s, dict) and (t := str(s.get("titel", "")).strip()) ] async def lade_frage_muster(topic: str, baustein: str) -> list[dict]: """Vordefinierte Frage-Muster eines Bausteins — DB-first, Fallback Sidecar (leer = Live).""" rows = await db.list_frage_muster(topic, _norm_titel(baustein)) if rows: return [{"subbaustein": r["sub_titel"], "frage": r["frage"]} for r in rows if r["frage"]] fm = _json_datei(frage_muster_path(topic)) if not isinstance(fm, dict): return [] return [ {"subbaustein": str(e.get("subbaustein", "")).strip(), "frage": frage} for e in (fm.get(baustein) or []) if isinstance(e, dict) and (frage := str(e.get("frage", "")).strip()) ] async def subbausteine_frei(topic: str, baustein: str, max_ebene: int) -> list[str]: """Subbaustein-Titel bis zur freigeschalteten Ebene (≤ max_ebene). Fallback ohne Ebenen-Wissen (Altbestand/Sidecar): alle Subbaustein-Titel.""" rows = await db.subs_mit_ebene(topic, baustein) if not rows: return await subbausteine_titel(topic, baustein) return [s["titel"] for s in rows if s["ebene"] <= max_ebene and s["titel"]] async def lade_frage_muster_frei(topic: str, baustein: str, max_ebene: int) -> list[dict]: """Frage-Muster, gefiltert auf Subbausteine bis zur freigeschalteten Ebene. Ohne Ebenen-Wissen (Altbestand/Sidecar) ungefiltert.""" rows = await db.subs_mit_ebene(topic, baustein) if not rows: return await lade_frage_muster(topic, baustein) frei = {s["norm"] for s in rows if s["ebene"] <= max_ebene} return [m for m in await lade_frage_muster(topic, baustein) if _norm_titel(m["subbaustein"]) in frei] async def lade_uebersicht(topic: str) -> list[dict]: """Strukturierte Baustein-Liste für die Übersicht — DB-first (Konsens + Subs/Stufen/Relevanz), Fallback bausteine.md + Sidecar (Alt-Themen).""" bs = await db.list_bausteine(topic, status="konsens") if bs: out = [] for num, b in enumerate(bs, 1): subs = [s for s in await db.list_subbausteine(topic, b["titel_norm"]) if s["status"] == "konsens"] out.append({ "num": num, "titel": b["titel"], "beschreibung": b["beschreibung"], "subbausteine": [ {"titel": s["sub_titel"], "stufe": s["stufe"] if s["stufe"] in _STUFEN else "fortgeschritten", "relevanz": s["relevanz"] if s["relevanz"] in ("relevant", "rand") else None} for s in subs if s["sub_titel"] ], }) return out entries = _lade_bausteine(_read(bausteine_path(topic))) sidecar = _json_datei(subbausteine_path(topic)) sidecar = sidecar if isinstance(sidecar, dict) else {} out = [] for num, entry in entries.items(): titel = _titel(entry) teile = entry.split(" — ", 1) beschreibung = teile[1].strip() if len(teile) == 2 else "" subbausteine = [ { "titel": t, "stufe": s.get("stufe") if s.get("stufe") in _STUFEN else "fortgeschritten", "relevanz": s.get("relevanz") if s.get("relevanz") in ("relevant", "rand") else None, } for s in (sidecar.get(titel) or []) if isinstance(s, dict) and (t := str(s.get("titel", "")).strip()) ] out.append({"num": num, "titel": titel, "beschreibung": beschreibung, "subbausteine": subbausteine}) return out def _bausteine_steps(topic: str) -> tuple: """Schritte je Quelle: link bekommt vorne „Quelle laden", projekt zusätzlich „Ergänzung". Subbausteine + Stufen sind je drei Phasen (Finden, Wählen, Klären). Pro Phase laufen alle Pakete parallel; der Schritt bleibt, bis das letzte Paket fertig ist. """ q = lade_quelle(topic) base = ("Recherche", "Konsolidierung", "Klärung", "Dedup") rest = ( "Subbausteine finden", "Subbausteine wählen", "Subbausteine klären", "Fakten finden", "Fakten prüfen", "Fakten fix", "Stufen finden", "Stufen wählen", "Stufen klären", "Relevanz finden", "Relevanz wählen", "Relevanz klären", "Gliederung", "Fragen finden", "Fragen wählen", "Fragen klären", "Fragen prüfen", "Karteikarten", "Beispiele", ) mitte = base + (("Ergänzung",) if q["type"] == "projekt" else ()) + rest return (("Quelle aufbereiten",) if q["type"] == "link" else ()) + mitte def _step_idx(topic: str, name: str) -> int: return _bausteine_steps(topic).index(name) def _melde_p(set_p, topic: str, schritt: str): """Async-Melde-Callback für _gather_fortschritt: setzt „ d/t…" + Schritt-Index.""" idx = _step_idx(topic, schritt) async def melde(d, t): set_p(f"{schritt} {d}/{t}…", step=idx) return melde # Grobe Anzeige-Phasen: bündeln die Feinschritte (intern bleibt alles feingranular). # Sonderschritte (Quelle laden, Ergänzung) gehören zur Phase „Inventar". PHASEN = ( ("Quelle", ("Quelle aufbereiten",)), ("Inventar", ("Recherche", "Konsolidierung", "Klärung", "Dedup", "Ergänzung")), ("Subbausteine", ("Subbausteine finden", "Subbausteine wählen", "Subbausteine klären")), ("Fakten", ("Fakten finden", "Fakten prüfen", "Fakten fix")), ("Stufen", ("Stufen finden", "Stufen wählen", "Stufen klären")), ("Relevanz", ("Relevanz finden", "Relevanz wählen", "Relevanz klären")), ("Gliederung", ("Gliederung",)), ("Fragen", ("Fragen finden", "Fragen wählen", "Fragen klären", "Fragen prüfen")), ("Artefakte", ("Karteikarten", "Beispiele")), ) def _phasen(topic: str) -> list[tuple[str, int]]: """[(grob_label, Anzahl vorhandener Feinschritte)] für die aktuelle Quelle.""" feine = _bausteine_steps(topic) return [(label, n) for label, members in PHASEN if (n := sum(f in members for f in feine))] def _phasen_status(topic: str, current: int | None) -> list[dict]: """Grobe Phasen-Zustände aus dem feinen Fortschritt `current` (None = alles pending, len(feine) = alles done). → [{label, state}] mit state done/active/pending.""" out, start = [], 0 for label, n in _phasen(topic): end = start + n if current is None or current < start: state = "pending" elif current >= end: state = "done" else: state = "active" out.append({"label": label, "state": state}) start = end return out def _bausteine_files(topic: str) -> dict: arbeit = arbeit_dir(topic) runden = range(1, KONSENS_MAX_RUNDEN + 1) return { "final": bausteine_path(topic), "arbeit": arbeit, "recherche": [arbeit / f"recherche-{i}.md" for i in (1, 2, 3, 4, 5)], "recherche_mapping": arbeit / "recherche-mapping.json", "auswahl": {n: [arbeit / f"auswahl-r{n}-{i}.json" for i in (1, 2, 3)] for n in runden}, "mapping": {n: arbeit / f"auswahl-mapping-r{n}.json" for n in runden}, "ergaenzung": arbeit / "ergaenzung.json", "sub_roh": arbeit / "subbausteine-roh.json", "fakten": arbeit / "subbausteine-fakten.json", "sidecar": subbausteine_path(topic), "frage_muster": frage_muster_path(topic), "gliederung": arbeit / "gliederung.json", "gliederung_slots": [arbeit / f"gliederung-{i}.json" for i in (1, 2, 3)], "artefakte": arbeit / "artefakte.json", } def _alle_slot_dateien(files: dict) -> list[Path]: arbeit = files["arbeit"] # Subbaustein-/Stufen-Slots sind pro Chunk dynamisch — per Glob einsammeln. dyn = (list(arbeit.glob("subbaustein-*")) + list(arbeit.glob("fakten-*")) + list(arbeit.glob("stufe-*")) + list(arbeit.glob("relevanz-*")) + list(arbeit.glob("frage-muster-*")) + list(arbeit.glob("gliederung-*")) + list(arbeit.glob("artefakt-*")) + list(arbeit.glob("recherche-*")) + list(arbeit.glob("konsolidierung-*")) + list(arbeit.glob("klaerung*")) + list(arbeit.glob("dedup-*"))) if arbeit.is_dir() else [] return [ *files["recherche"], files["recherche_mapping"], *(p for slots in files["auswahl"].values() for p in slots), *files["mapping"].values(), files["ergaenzung"], files["sub_roh"], files["sidecar"], files["frage_muster"], files["fakten"], files["gliederung"], files["artefakte"], *dyn, ] def cancel_bausteine(topic: str) -> bool: if topic not in _bausteine_progress: return False _bausteine_cancelled.add(topic) cancel_scope(f"bausteine-{topic}-") # wartende Agenten bailen vorm Spawn kill_process(f"bausteine-{topic}-") # laufende Subprozesse killen return True def _resume_step(topic: str) -> int: """Erster noch offener Schritt anhand der persistierten Artefakte. Inventar (Recherche→Konsolidierung→Klärung) gilt als fertig, sobald bausteine.md vorliegt.""" files = _bausteine_files(topic) q = lade_quelle(topic) if q["type"] == "link" and not _crawl_fertig(topic): return _step_idx(topic, "Quelle aufbereiten") if not files["final"].exists(): # Inventar (DB-Loop) noch offen return _step_idx(topic, "Recherche") if q["type"] == "projekt" and not files["ergaenzung"].exists(): return _step_idx(topic, "Ergänzung") sidecar = _json_datei(files["sidecar"]) if _sidecar_schema(sidecar) is not None: # Stufen fertig; nur noch Relevanz offen? if not _relevanz_komplett(sidecar): return _step_idx(topic, "Relevanz finden") # Relevanz fertig; Gliederung (Bausteine-Artefakt für den Guide) offen? if not _gliederung_komplett(files): return _step_idx(topic, "Gliederung") # Gliederung fertig; Frage-Muster offen? if not _frage_muster_komplett(topic): return _step_idx(topic, "Fragen finden") # Fragen fertig; Lern-Artefakte (Karteikarten/Beispiele) offen? if not _artefakte_komplett(files): return _step_idx(topic, "Karteikarten") return len(_bausteine_steps(topic)) if _sub_roh_schema(_json_datei(files["sub_roh"])) is None: return _step_idx(topic, "Subbausteine finden") # Subbausteine fertig; Fakten noch offen? (Fakten kommen vor den Stufen.) if not _fakten_komplett(files): return _step_idx(topic, "Fakten finden") return _step_idx(topic, "Stufen finden") def bausteine_status(topic: str) -> dict: # Intern feingranular (Resume/Progress); für die Anzeige zu 5 groben Phasen gebündelt. feine = _bausteine_steps(topic) ready = bausteine_path(topic).exists() generating = topic in _bausteine_progress partial = False if generating: current = _bausteine_step.get(topic) elif ready: current = len(feine) # alles fertig → alle Phasen done else: current = _resume_step(topic) partial = current > 0 return { "ready": ready, "generating": generating, "progress": _bausteine_progress.get(topic), "error": _bausteine_errors.get(topic), "partial": partial, "steps": _phasen_status(topic, current), } def active_bausteine() -> list[dict]: return [{"topic": t, "progress": p} for t, p in _bausteine_progress.items()] def reset_bausteine(topic: str) -> None: """„Entfernen": löscht den GESAMTEN Bausteine-Bereich — Crawl, Sichtung, Inventar … Fragen. BEHÄLT nur die Themen-Config `quelle.json` (Typ/Link/Spec). Re-Generieren crawlt neu. (Crawl/Sichtung gehören zu den Bausteinen; nur die Config ist „Thema".)""" files = _bausteine_files(topic) files["final"].unlink(missing_ok=True) files["sidecar"].unlink(missing_ok=True) files["frage_muster"].unlink(missing_ok=True) shutil.rmtree(quelle_crawl_dir(topic), ignore_errors=True) # Crawl gehört zu Bausteinen shutil.rmtree(files["arbeit"], ignore_errors=True) _bausteine_errors.pop(topic, None) # quelle.json bleibt bewusst stehen — das ist die Themen-Config. def _phase_idx(label: str) -> int: """Index der groben Phase in der kanonischen Reihenfolge (Quelle=0 … Fragen=5).""" order = [l for l, _ in PHASEN] return order.index(label) if label in order else 1 def _reset_ab_phase(topic: str, label: str) -> None: """Datei-Artefakte AB der groben Phase `label` löschen (Quelle/Inventar … Fragen), frühere behalten. Kumulativ. quelle.json + Crawl (.done) bleiben immer (re-crawl nur bei Voll-Reset).""" files = _bausteine_files(topic) arbeit = files["arbeit"] idx = _phase_idx(label) def glob_del(pat: str) -> None: if arbeit.is_dir(): for p in arbeit.glob(pat): p.unlink(missing_ok=True) # Phasen-Index: Quelle=0 · Inventar=1 · Subbausteine=2 · Fakten=3 · Stufen=4 · Relevanz=5 · Gliederung=6 · Fragen=7 · Artefakte=8 if idx <= 8: # Artefakte (Karteikarten/Beispiele) files["artefakte"].unlink(missing_ok=True) glob_del("artefakt-*") if idx <= 7: # Fragen files["frage_muster"].unlink(missing_ok=True) glob_del("frage-muster-*") if idx <= 6: # Gliederung files["gliederung"].unlink(missing_ok=True) glob_del("gliederung-*") if idx <= 5: # Relevanz glob_del("relevanz-*") if idx <= 4: # Stufen + Relevanz teilen die Sidecar → ab Stufen ganz neu files["sidecar"].unlink(missing_ok=True) glob_del("stufe-*") else: # ab Relevanz: Stufen behalten, nur Relevanz-Felder strippen sc = _json_datei(files["sidecar"]) if isinstance(sc, dict): for subs in sc.values(): for s in (subs if isinstance(subs, list) else []): if isinstance(s, dict): s.pop("relevanz", None) atomic_write_json(files["sidecar"], sc, indent=1) if idx <= 3: # Fakten (vor den Stufen) — Fakten-Map + Arbeitsdateien weg files["fakten"].unlink(missing_ok=True) glob_del("fakten-*") if idx <= 2: # Subbausteine files["sub_roh"].unlink(missing_ok=True) glob_del("subbaustein-*") if idx <= 1: # Inventar (und Quelle) = Inventar-Dateien + bausteine.md weg for p_alt in _alle_slot_dateien(files): p_alt.unlink(missing_ok=True) files["final"].unlink(missing_ok=True) def _ergaenzung_schema(data): """{"bausteine": [{"titel", "beschreibung"}]} → Liste (leer erlaubt) · sonst None.""" if not isinstance(data, dict) or not isinstance(data.get("bausteine"), list): return None out = [] for b in data["bausteine"]: if not isinstance(b, dict) or not isinstance(b.get("titel"), str) or not isinstance(b.get("beschreibung"), str): return None titel, beschreibung = b["titel"].strip(), b["beschreibung"].strip() if not titel: return None out.append((titel, beschreibung)) return out def _pdfs_konvertieren(project: Path) -> None: """PDFs im Projekt in .txt wandeln (pdftotext) — Agenten lesen Text statt Seiten-Bildern. Wird vor jeder Projekt-Generierung aufgerufen; konvertiert nur, wenn die .txt fehlt oder älter als das PDF ist. Das Original bleibt unangetastet. Fehlt pdftotext und das Projekt enthält PDFs → harter Fehler statt unzuverlässigem Direkt-Lese-Modus (MiniMax-Bilderlimit, Vision-Kosten). """ pdfs = list(project.rglob("*.pdf")) if not pdfs: return if shutil.which("pdftotext") is None: raise RuntimeError("pdftotext fehlt (poppler-utils installieren) — PDFs im Projekt können nicht gelesen werden") for pdf in pdfs: txt = pdf.with_suffix(".txt") if txt.exists() and txt.stat().st_mtime >= pdf.stat().st_mtime: continue try: subprocess.run(["pdftotext", "-layout", str(pdf), str(txt)], check=True, timeout=120) _log(project.name, f"PDF konvertiert: {pdf.name} → {txt.name}") except Exception as e: raise RuntimeError(f"PDF-Konvertierung fehlgeschlagen ({pdf.name}): {e}") from e _QUELLE_TEMPLATE = {"projekt": "Bausteine-Quelle-Projekt", "uni": "Bausteine-Quelle-Uni", "link": "Bausteine-Quelle-Link"} def _text_abschnitte(text: str, ziel: int = RECHERCHE_ABSCHNITT_ZEICHEN) -> list[str]: """Text an Absatz-/Zeilengrenzen in Abschnitte ~`ziel` Zeichen splitten (gegen Lost-in-the-Middle bei großen Dokumenten). Kleiner Text bleibt EIN Abschnitt. Inhalt bleibt vollständig — nur Trenn-Whitespace fällt weg.""" text = text.strip() if len(text) <= ziel: return [text] if text else [] abschnitte: list[str] = [] buf = "" def flush(): nonlocal buf if buf.strip(): abschnitte.append(buf.strip()) buf = "" for block in re.split(r"\n\s*\n", text): # an Absatz-Grenzen block = block.strip() if not block: continue if len(block) > ziel: # einzelner Riesen-Absatz → hart an Zeilen schneiden flush() for zeile in block.split("\n"): if buf and len(buf) + len(zeile) + 1 > ziel: flush() buf += zeile + "\n" flush() elif buf and len(buf) + len(block) + 2 > ziel: flush() buf = block else: buf = (buf + "\n\n" + block) if buf else block flush() return abschnitte def _build_recherche_prompt(topic: str, out_path: Path, instructions: str, typ: str, ordner: Path | None, fokus: str = "", abschnitt: str = "") -> str: if abschnitt: # Abschnitt-Modus (uni/projekt): Text direkt im Prompt → kleiner Kontext, kein Datei-Lesen. source = abschnitt elif typ in _QUELLE_TEMPLATE: source = _prompt(_QUELLE_TEMPLATE[typ], project=ordner) else: source = _prompt("Bausteine-Quelle-Thema", topic=topic) return _prompt( "Bausteine-Recherche", topic=topic, source=source, bausteine_path=out_path, fokus=fokus, extra=_extra(instructions), ) def _file_payload(path: Path): """Gültig, wenn die Slot-Datei existiert und nummerierte Einträge enthält.""" if not path.exists(): return None text = path.read_text(encoding="utf-8") return text if _parse_auswahl(text) else None def _mapping_schema(data): """{"bausteine": [str, ≥1], "rest": [str]} → (bausteine, rest) · sonst None.""" if not isinstance(data, dict): return None bausteine = _str_liste(data.get("bausteine")) rest = _str_liste(data.get("rest")) if not bausteine or rest is None: return None return bausteine, rest def _sub_roh_schema(data): """{Baustein-Titel: [Subbaustein, …]} → dict · sonst None (Zwischenstand Block B).""" if not isinstance(data, dict) or not data: return None out: dict[str, list[str]] = {} for k, v in data.items(): subs = _str_liste(v) if isinstance(v, list) else None if not isinstance(k, str) or not k.strip() or not subs: return None out[k] = subs return out def _sidecar_schema(data): """{Baustein-Titel: [{titel, stufe}, …]} → dict · sonst None (Sidecar mit Stufen).""" if not isinstance(data, dict) or not data: return None for v in data.values(): if not isinstance(v, list) or not v: return None for s in v: if not isinstance(s, dict) or not str(s.get("titel", "")).strip() or s.get("stufe") not in _STUFEN: return None return data def _relevanz_komplett(data) -> bool: """Jeder Subbaustein der Sidecar trägt eine gültige Relevanz (relevant/rand)?""" if not isinstance(data, dict) or not data: return False return all( isinstance(s, dict) and s.get("relevanz") in ("relevant", "rand") for v in data.values() if isinstance(v, list) for s in v ) def _frage_muster_chunk_schema(data) -> list[dict] | None: """{"muster": [{baustein, subbaustein, frage}, …]} → Liste valider Einträge · sonst None. Ein Muster je Subbaustein (kein Typ-Kreuzprodukt — die Schwierigkeit kommt erst bei der Prüfung aus dem Lerner-Niveau). Ungültige Einzel-Einträge werden übersprungen.""" if not isinstance(data, dict) or not isinstance(data.get("muster"), list): return None out = [] for e in data["muster"]: if not isinstance(e, dict): continue bau = str(e.get("baustein", "")).strip() sub = str(e.get("subbaustein", "")).strip() frage = str(e.get("frage", "")).strip() if not bau or not sub or not frage: continue out.append({"baustein": bau, "subbaustein": sub, "frage": frage}) return out or None def _frage_muster_komplett(topic: str) -> bool: """Frage-Muster-Sidecar existiert (Build gelaufen)? Einzelne leere Bausteine fallen zur Prüfungszeit auf Live-Generierung zurück — daher genügt die Datei.""" return isinstance(_json_datei(frage_muster_path(topic)), dict) def _read(p: Path) -> str: return p.read_text(encoding="utf-8") if p.exists() else "" def _chunk_nums(items: list, n: int) -> list[list]: """Teilt eine flache Liste in n möglichst gleich große Chunks.""" n = max(1, n) size = max(1, math.ceil(len(items) / n)) return [items[i:i + size] for i in range(0, len(items), size)] def _n_chunks(count: int, size: int = SUBBAUSTEIN_CHUNK) -> int: return min(SUBBAUSTEIN_MAX, max(1, math.ceil(count / size))) def _lpt_chunks(gewichte: list[int], target: int) -> list[list[int]]: """Indizes lastbalanciert auf Chunks verteilen (LPT, Makespan-minimal). Gewicht = Kosten je Index. K = ceil(Gesamtgewicht/target); schwerste zuerst in den jeweils leichtesten Bin. → Index-Listen.""" if not gewichte: return [] K = max(1, math.ceil(sum(gewichte) / max(1, target))) bins: list[list[int]] = [[] for _ in range(K)] last = [0] * K for i in sorted(range(len(gewichte)), key=lambda x: gewichte[x], reverse=True): j = min(range(K), key=lambda b: last[b]) bins[j].append(i) last[j] += gewichte[i] return [b for b in bins if b] async def _subbausteine_block(ctx: GenContext, set_p, files: dict, entries: dict, instructions: str) -> dict | None: """Block B (DB + Loop): je Paket Subbausteine in Runden finden (3 Finder, bis 0 neue/Kappe), in der DB sammeln (≥2 Nennungen = Konsens, 1× verworfen), Judge bereinigt je Paket. → {Baustein-Titel: [Subbaustein, …]} (Konsens) oder None. Befüllt DB-Tabelle `subbausteine`.""" topic, provider, is_cancelled = ctx.topic, ctx.provider, ctx.is_cancelled arbeit = files["arbeit"] ordner = quelle_ordner(topic) caps = "files" if ordner else "full" # Quelle für die Beleg-Prüfung im Klär-Schritt (verwirft erfundene/unbelegbare Subs). _typ = lade_quelle(topic).get("type", "thema") source = _prompt(_QUELLE_TEMPLATE[_typ], project=ordner) if _typ in _QUELLE_TEMPLATE else _prompt("Bausteine-Quelle-Thema", topic=topic) nums = list(entries) chunks = _chunk_nums(nums, _n_chunks(len(nums))) n = len(chunks) titel_by_num = {num: _titel(entries[num]) for num in nums} norm_by_num = {num: _norm_titel(titel_by_num[num]) for num in nums} await db.delete_subbausteine(topic) # Frischstart des Blocks (idempotenter Zähler) async def _bekannt_block(chunk): bl = [] for num in chunk: subs = [s["sub_titel"] for s in await db.list_subbausteine(topic, norm_by_num[num])] if subs: bl.append(f"\n" + "\n".join(f"- {s}" for s in subs)) if not bl: return "" # Bekanntes NICHT erneut auflisten lassen (sonst bläht Re-Bestätigung den Mention-Count auf, # Self-Bias/Echo) — nur Fehlendes ergänzen. Der Zähler bleibt so ein ehrliches Konsens-Signal. return ("\n\nBEREITS ERFASST — liste diese NICHT erneut. Finde nur, was FEHLT:\n" + "\n".join(bl)) # Phase „Subbausteine finden": je Paket Loop bis 0 neue Subs / Zeit-Kappe. async def _finde(c, chunk): zuteilung = "\n".join(f"- {entries[num]}" for num in chunk) chunk_idx = _titel_index({num: titel_by_num[num] for num in chunk}) start = time.monotonic() runde = 0 while not is_cancelled(): runde += 1 bekannt = await _bekannt_block(chunk) if runde > 1 else "" paths = [arbeit / f"subbaustein-c{c}-r{runde}-{i}.md" for i in (1, 2, 3)] for p in paths: p.unlink(missing_ok=True) slots = [{ "key": f"bausteine-{topic}-subbaustein-c{c}-r{runde}-{i}", "prompt": _prompt("Subbaustein-Recherche", topic=topic, zuteilung=zuteilung, bekannt=bekannt, out_path=p, extra=_extra(instructions)), "role": "quick", "capabilities": caps, "payload": (lambda result, p=p: _parse_subbausteine(_read(p)) or None), } for i, p in enumerate(paths, 1)] texte = await _race(topic, f"Subbausteine Paket {c} R{runde}", slots, 2, _timeout("subbaustein", len(chunk)), provider, cancelled=is_cancelled, grace=KONSENS_GRACE) if is_cancelled(): return False if not texte: return runde > 1 # Runde 1 ohne Ergebnis = Fehler; spätere = einfach Ende vorhanden = {num: {s["sub_norm"] for s in await db.list_subbausteine(topic, norm_by_num[num])} for num in chunk} neu = 0 for d in texte: for marker, subs in d.items(): num = _titel_aufloesen(chunk_idx, marker) if num is None: continue gesehen = set() for sub in subs: sn = _norm_titel(sub) if not sn or sn in gesehen: continue gesehen.add(sn) if sn not in vorhanden[num]: neu += 1 vorhanden[num].add(sn) await db.upsert_subbaustein(topic, norm_by_num[num], sn, titel_by_num[num], sub) if neu == 0: break if time.monotonic() - start > SUBBAUSTEIN_KAPPE: _log(topic, f"Subbausteine Paket {c}: Zeit-Kappe erreicht (Runde {runde})") break return True oks = await _gather_fortschritt([_finde(c, chunk) for c, chunk in enumerate(chunks, 1)], n, _melde_p(set_p, topic, "Subbausteine finden")) if is_cancelled(): return None if not all(ok is True for ok in oks): _bausteine_errors[topic] = "Subbausteine fehlgeschlagen (Recherche)" return None # Phase „Subbausteine wählen": ≥2 Nennungen = Konsens, 1× verworfen (Code). set_p(f"Subbausteine wählen ({n} Pakete)…", step=_step_idx(topic, "Subbausteine wählen")) for num in nums: for s in await db.list_subbausteine(topic, norm_by_num[num]): await db.set_subbaustein_felder(topic, norm_by_num[num], s["sub_norm"], status=("konsens" if s["nennungen"] >= 2 else "verworfen")) # Phase „Subbausteine klären": Source-Panel (SUBBAUSTEIN_PANEL Judges) prüft Konsens + Unsicher (1×) # gegen die Quelle; Code-Mehrheit je Sub. Externes, mehrstimmiges Gate gegen Einzel-Judge-Bias + Echo. async def _klaere(c, chunk): fp = arbeit / f"subbaustein-final-c{c}.md" if _parse_subbausteine(_read(fp)): return bloecke, hat = [], False konsens_by_num: dict[int, list[str]] = {} for num in chunk: rows = await db.list_subbausteine(topic, norm_by_num[num]) kon = [s["sub_titel"] for s in rows if s["status"] == "konsens"] uns = [s["sub_titel"] for s in rows if s["status"] != "konsens" and s["nennungen"] == 1] konsens_by_num[num] = kon if not kon and not uns: continue hat = True k_zeilen = "\n".join(f"- {s}" for s in kon) if kon else "- (keiner)" u_zeilen = "\n".join(f"- {s}" for s in uns) if uns else "- (keiner)" bloecke.append(f"BAUSTEIN: {titel_by_num[num]}\nKonsens (≥2 Finder):\n{k_zeilen}\nUnsicher (1× — streng gegen Quelle prüfen):\n{u_zeilen}") if not hat: return chunk_idx = _titel_index({num: titel_by_num[num] for num in chunk}) paths = [arbeit / f"subbaustein-final-c{c}-j{j}.md" for j in range(1, SUBBAUSTEIN_PANEL + 1)] offen = [(j, p) for j, p in enumerate(paths, 1) if _parse_subbausteine(_read(p)) is None] for _, p in offen: p.unlink(missing_ok=True) if offen: slots = [{ "key": f"bausteine-{topic}-subbaustein-final-c{c}-j{j}", "prompt": _prompt("Subbaustein-Mapping", topic=topic, source=source, bausteine="\n\n".join(bloecke), out_path=p, extra=_extra(instructions)), "role": "judge", "capabilities": caps, "payload": (lambda result, p=p: _parse_subbausteine(_read(p)) or None), } for j, p in offen] vorhanden = SUBBAUSTEIN_PANEL - len(offen) await _race(topic, f"Subbaustein-Klärung {c}", slots, max(1, 2 - vorhanden), _timeout("subbaustein_check", len(chunk)), provider, cancelled=is_cancelled, grace=KONSENS_GRACE) if is_cancelled(): return outs = [d for p in paths if (d := _parse_subbausteine(_read(p)))] if not outs: # Panel komplett gescheitert → Konsens übernehmen (wie bisher der Fallback) _log(topic, f"Subbaustein-Klärung Paket {c} fehlgeschlagen — Konsens übernommen") text = "\n\n".join(f"\n" + "\n".join(f"- {s}" for s in konsens_by_num[num]) for num in chunk if konsens_by_num[num]) atomic_write_text(fp, text) return # Code-Mehrheit je Baustein/Sub-Norm: behalten wenn Mehrheit der Judges ihn führt (Tie → behalten). bloecke_out = [] for num in chunk: votes: dict[str, int] = {} form: dict[str, str] = {} for d in outs: seen = set() for marker, subs in d.items(): if _titel_aufloesen(chunk_idx, marker) != num: continue for sub in subs: sn = _norm_titel(sub) if not sn or sn in seen: continue seen.add(sn) form.setdefault(sn, sub) votes[sn] = votes.get(sn, 0) + 1 kept = [form[sn] for sn in form if votes[sn] * 2 >= len(outs)] if kept: bloecke_out.append(f"\n" + "\n".join(f"- {s}" for s in kept)) atomic_write_text(fp, "\n\n".join(bloecke_out)) await _gather_fortschritt([_klaere(c, chunk) for c, chunk in enumerate(chunks, 1)], n, _melde_p(set_p, topic, "Subbausteine klären")) if is_cancelled(): return None # Finale Liste je Baustein: Judge-Ausgabe, sonst Konsens-Fallback. DB reconcilen + roh bauen. roh: dict[str, list[str]] = {} for c, chunk in enumerate(chunks, 1): final = _parse_subbausteine(_read(arbeit / f"subbaustein-final-c{c}.md")) or {} chunk_idx = _titel_index({num: titel_by_num[num] for num in chunk}) final_by_num = {_titel_aufloesen(chunk_idx, m): subs for m, subs in final.items() if _titel_aufloesen(chunk_idx, m) is not None} for num in chunk: titel = titel_by_num[num] konsens = [s["sub_titel"] for s in await db.list_subbausteine(topic, norm_by_num[num]) if s["status"] == "konsens"] subs = final_by_num.get(num) or konsens if not subs: continue roh[titel] = subs # DB an die finale Liste angleichen: finale = konsens, Rest verworfen, Neues ergänzen. final_norms = {_norm_titel(s) for s in subs} have = {s["sub_norm"] for s in await db.list_subbausteine(topic, norm_by_num[num])} for s in await db.list_subbausteine(topic, norm_by_num[num]): await db.set_subbaustein_felder(topic, norm_by_num[num], s["sub_norm"], status=("konsens" if s["sub_norm"] in final_norms else "verworfen")) for s in subs: sn = _norm_titel(s) if sn and sn not in have: await db.upsert_subbaustein(topic, norm_by_num[num], sn, titel, s) await db.set_subbaustein_felder(topic, norm_by_num[num], sn, status="konsens") await _dedup_subbausteine(topic, roh) # Near-Dup-Filter je Baustein (deterministisch, kein LLM) if not roh: _bausteine_errors[topic] = "Keine Subbausteine ermittelt" return None return roh async def _dedup_subbausteine(topic: str, roh: dict[str, list[str]]) -> None: """Deterministischer Near-Duplicate-Filter pro Baustein: Subbausteine mit Cosine ≥ EMBEDDING_SUB_DUP sind dieselbe Aussage (im engen Baustein-Kontext zuverlässig — kein LLM nötig). Behält je Dublett-Gruppe den informativsten (längsten); Rest → DB verworfen + aus `roh`. Modell fehlt → still überspringen (wie der übrige Embedding-Fallback).""" if not EMBEDDING_AKTIV or not await asyncio.to_thread(embedding.verfuegbar): return for titel, subs in list(roh.items()): if len(subs) < 2: continue sims = await asyncio.to_thread(embedding.embed_sims, subs) if sims is None: return behalten: list[int] = [] verworfen: list[int] = [] for i in sorted(range(len(subs)), key=lambda x: (-len(subs[x]), x)): # informativster zuerst if any(float(sims[i][j]) >= EMBEDDING_SUB_DUP for j in behalten): verworfen.append(i) else: behalten.append(i) if not verworfen: continue bnorm = _norm_titel(titel) for i in verworfen: await db.set_subbaustein_felder(topic, bnorm, _norm_titel(subs[i]), status="verworfen") roh[titel] = [subs[i] for i in sorted(behalten)] # Original-Reihenfolge der Behaltenen async def _stufen_block(ctx: GenContext, set_p, files: dict, roh: dict, instructions: str) -> dict | None: """Block C: drei Phasen mit Barriere — Finden (einstufen), Wählen (Vote), Klären. Lokale IDs 1..n pro Paket, hinterher auf globale gid gemappt. → {Baustein-Titel: [{titel, stufe}, …]} oder None.""" topic, provider, is_cancelled = ctx.topic, ctx.provider, ctx.is_cancelled arbeit = files["arbeit"] # Kernpunkte je Sub als knapper Kontext (fundiertere Einstufung; Klassifikation braucht wenig). fakten_map = _json_datei(files["fakten"]) fakten_map = fakten_map if isinstance(fakten_map, dict) else {} items = [(titel, sub) for titel, subs in roh.items() for sub in subs] # globale id = index+1 if not items: return {titel: [] for titel in roh} # Chunks aus GANZEN Bausteinen packen (keinen Baustein splitten) → Rater sieht je Baustein # alle Subs und kann relativ einstufen. Item-Indizes je Baustein in roh-Reihenfolge. chunks, cur, i = [], [], 0 for _titel_b, subs in roh.items(): g = list(range(i, i + len(subs))) i += len(subs) if cur and len(cur) + len(g) > STUFE_CHUNK: chunks.append(cur) cur = [] cur.extend(g) if cur: chunks.append(cur) n = len(chunks) def rater_paths(c): return [arbeit / f"stufe-c{c}-{i}.json" for i in (1, 2, 3)] def lset(item_idxs): return set(range(1, len(item_idxs) + 1)) # Phase „Stufen finden": pro Paket 3 Rater (min. 2), lokale IDs. async def _rate(c, item_idxs): local_set = lset(item_idxs) paths = rater_paths(c) vorhanden = sum(1 for p in paths if _stufen_schema(_json_datei(p), local_set)) if vorhanden >= 2: return True enum_zeilen, cur_b = [], None for k, j in enumerate(item_idxs, 1): b, sub = items[j] if b != cur_b: enum_zeilen.append(f"\nBAUSTEIN: {b}") cur_b = b enum_zeilen.append(f"{k}. {sub}") if (kz := _kern_zeile(fakten_map.get(b, {}).get(_norm_titel(sub)))): enum_zeilen.append(f" {kz}") enum = "\n".join(enum_zeilen).strip() offen = [(i, p) for i, p in enumerate(paths, 1) if not _stufen_schema(_json_datei(p), local_set)] slots = [{ "key": f"bausteine-{topic}-stufe-c{c}-{i}", "prompt": _prompt("Stufen-Recherche", topic=topic, subbausteine=enum, out_path=p, extra=_extra(instructions)), "role": "fast", "capabilities": "files", "payload": (lambda result, p=p, ids=local_set: _stufen_schema(_json_datei(p), ids)), } for i, p in offen] neu = await _race(topic, f"Stufen Paket {c}", slots, 2 - vorhanden, _timeout("stufe", len(item_idxs)), provider, cancelled=is_cancelled, grace=KONSENS_GRACE) return not is_cancelled() and neu is not None oks = await _gather_fortschritt([_rate(c, idxs) for c, idxs in enumerate(chunks, 1)], len(chunks), _melde_p(set_p, topic, "Stufen finden")) if is_cancelled(): return None if not all(ok is True for ok in oks): _bausteine_errors[topic] = "Einstufung fehlgeschlagen (Recherche)" return None # Phase „Stufen wählen": Code-Vote je Paket → (ergebnis, strittig). set_p(f"Stufen wählen ({n} Pakete)…", step=_step_idx(topic, "Stufen wählen")) vote_by_c = {} for c, item_idxs in enumerate(chunks, 1): local_set = lset(item_idxs) rater = [d for p in rater_paths(c) if (d := _stufen_schema(_json_datei(p), local_set))] ergebnis: dict[int, str] = {} strittig: dict[int, list[str]] = {} for k in range(1, len(item_idxs) + 1): stimmen = [d[k] for d in rater if k in d] zaehler: dict[str, int] = {} for s in stimmen: zaehler[s] = zaehler.get(s, 0) + 1 best = max(zaehler.values(), default=0) gewinner = [s for s, v in zaehler.items() if v == best] if len(gewinner) == 1 and best >= 2: ergebnis[k] = gewinner[0] else: strittig[k] = stimmen vote_by_c[c] = (ergebnis, strittig) # Phase „Stufen klären": Judge je Paket mit Strittigem, alle parallel. async def _klaere(c, item_idxs): ergebnis, strittig = vote_by_c[c] if strittig: judge_path = arbeit / f"stufe-final-c{c}.json" entsch = _stufen_schema(_json_datei(judge_path), set(strittig)) if entsch is None: strittig_block = "\n".join( f"{k}. [{items[item_idxs[k - 1]][0]}] {items[item_idxs[k - 1]][1]} — Stimmen: {', '.join(stimmen) or 'keine'}" for k, stimmen in strittig.items() ) status, entsch = await run_single_slot( ctx, f"Stufen-Klärung {c}", key=f"bausteine-{topic}-stufe-final-c{c}", prompt=_prompt("Stufen-Mapping", topic=topic, strittig=strittig_block, out_path=judge_path, extra=_extra(instructions)), role="judge", capabilities="files", payload=lambda result, p=judge_path, ids=set(strittig): _stufen_schema(_json_datei(p), ids), timeout=_timeout("stufe_check", len(strittig)), ) if status == FAILED: _log(topic, f"Stufen-Klärung Paket {c} fehlgeschlagen — Default 'fortgeschritten'") entsch = entsch if isinstance(entsch, dict) else {} # Strittige ohne Entscheid → 'fortgeschritten'; Vote-Gewinner bleiben; Judge überschreibt. ergebnis = {**{k: "fortgeschritten" for k in strittig}, **ergebnis, **entsch} return {item_idxs[k - 1] + 1: stufe for k, stufe in ergebnis.items()} parts = await _gather_fortschritt([_klaere(c, idxs) for c, idxs in enumerate(chunks, 1)], len(chunks), _melde_p(set_p, topic, "Stufen klären")) if is_cancelled(): return None stufe_by_id: dict[int, str] = {} for c, part in enumerate(parts, 1): if not isinstance(part, dict): # Klärung ist nicht fatal: Vote-Ergebnis + Default 'fortgeschritten' für Strittige. if isinstance(part, BaseException): _log(topic, f"Stufen-Klärung Paket {c}: {type(part).__name__}: {part}") ergebnis, strittig = vote_by_c[c] item_idxs = chunks[c - 1] merged = {**{k: "fortgeschritten" for k in strittig}, **ergebnis} part = {item_idxs[k - 1] + 1: s for k, s in merged.items()} stufe_by_id.update(part) # Sidecar zusammensetzen — gleiche Reihenfolge wie items → gid stimmt sidecar: dict[str, list[dict]] = {} gid = 0 for titel, subs in roh.items(): lst = [] for sub in subs: gid += 1 lst.append({"titel": sub, "stufe": stufe_by_id.get(gid, "fortgeschritten")}) sidecar[titel] = lst return sidecar _FAKTEN_FELDER = ("kernpunkte", "voraussetzungen", "huerden", "belegte_fakten", "beispiel_idee") def _fakten_schema(data) -> list[dict] | None: """{"fakten": [{baustein, subbaustein, …}]} → valide Liste · sonst None. Trennt belegte_fakten (mit Quelle) hart von beispiel_idee (generativ).""" if not isinstance(data, dict) or not isinstance(data.get("fakten"), list): return None out = [] for e in data["fakten"]: if not isinstance(e, dict): continue bau = str(e.get("baustein", "")).strip() sub = str(e.get("subbaustein", "")).strip() if not bau or not sub: continue bf = [{"text": t, "quelle": str(f.get("quelle", "")).strip()} for f in (e.get("belegte_fakten") or []) if isinstance(f, dict) and (t := str(f.get("text", "")).strip())] out.append({ "baustein": bau, "subbaustein": sub, "kernpunkte": [k for x in (e.get("kernpunkte") or []) if (k := str(x).strip())], "voraussetzungen": str(e.get("voraussetzungen", "")).strip(), "huerden": str(e.get("huerden", "")).strip(), "belegte_fakten": bf, "beispiel_idee": str(e.get("beispiel_idee", "")).strip(), }) return out or None def _fakten_check_schema(data) -> list[tuple[str, bool]] | None: """Fakten-Check → [(sub_norm, verwerfen)] je Beanstandung · {ok:true}→[] · None bei ungültig. verwerfen=True: Sub inhaltlich nicht belegbar (entfernen). verwerfen=False: nur Fakt korrigieren.""" if not isinstance(data, dict): return None if data.get("ok") is True: return [] pr = data.get("probleme") if not isinstance(pr, list): return None return [(sn, bool(p.get("verwerfen"))) for p in pr if isinstance(p, dict) and (sn := _norm_titel(str(p.get("subbaustein", ""))))] def _kern_zeile(fk) -> str: """Knappe Kernpunkt-Zeile für Klassifikation (Stufe/Relevanz) — weniger Kontext genügt dort. Leer, wenn keine Fakten/Kernpunkte (Altbestand).""" if not isinstance(fk, dict) or not fk.get("kernpunkte"): return "" return "Kern: " + " · ".join(str(k) for k in fk["kernpunkte"]) def _fakten_zeilen(fk: dict) -> str: z = [] if fk.get("kernpunkte"): z.append("Kernpunkte: " + " · ".join(str(k) for k in fk["kernpunkte"])) if fk.get("voraussetzungen"): z.append("Voraussetzung: " + fk["voraussetzungen"]) if fk.get("huerden"): z.append("Hürde: " + fk["huerden"]) for bf in fk.get("belegte_fakten", []): z.append(f"FAKT: {bf['text']} (Quelle: {bf.get('quelle', '?')})") if fk.get("beispiel_idee"): z.append("Beispiel: " + fk["beispiel_idee"]) return "\n".join(z) def _fakten_komplett(files: dict) -> bool: """Fakten-Map existiert (Block durch)? {Baustein: {sub_norm: {...}}}.""" d = _json_datei(files["fakten"]) return isinstance(d, dict) and bool(d) async def _fakten_block(ctx, set_p, files: dict, roh: dict, q: dict, ordner, instructions: str) -> tuple | None: """Block: je Sub Quell-Fakten extrahieren (finden) → verifizieren (prüfen) → korrigieren/verwerfen (fix). Extract-once-Grounding: das Ergebnis nährt Stufe/Relevanz/Fragen/Guide. → (fakten_map, verworfen_map) — fakten_map {Baustein: {sub_norm: fakten}}, verworfen_map {Baustein: {sub_norm}} (unbelegbare Subs zum Entfernen) — oder None bei Abbruch/Fehler.""" topic, provider, is_cancelled = ctx.topic, ctx.provider, ctx.is_cancelled arbeit = files["arbeit"] caps = "files" if ordner else "full" typ = q.get("type", "thema") source = _prompt(_QUELLE_TEMPLATE[typ], project=ordner) if typ in _QUELLE_TEMPLATE else _prompt("Bausteine-Quelle-Thema", topic=topic) bausteine = [(titel, [str(s).strip() for s in subs if str(s).strip()]) for titel, subs in roh.items() if subs] if not bausteine: return {}, {} chunks = _lpt_chunks([len(subs) for _, subs in bausteine], FAKTEN_CHUNK_SUBS) def roh_path(ci): return arbeit / f"fakten-c{ci}.json" def erg_path(ci): return arbeit / f"fakten-erg-c{ci}.json" def chk_path(ci, j): return arbeit / f"fakten-check-c{ci}-j{j}.json" def fix_path(ci): return arbeit / f"fakten-fix-c{ci}.json" def ctitel(idxs): return [bausteine[i][0] for i in idxs] def block_text(idxs): return "\n\n".join( f"BAUSTEIN: {bausteine[i][0]}\nSUBBAUSTEINE:\n" + "\n".join(f"- {s}" for s in bausteine[i][1]) for i in idxs) # Roh-Fakten eines Chunks → {Baustein: {sub_norm: {sub, …felder}}}, auf Chunk-Titel gematcht. def roh_map(ci, path): idxs = chunks[ci] rel_by = {bausteine[i][0]: bausteine[i][1] for i in idxs} ct = ctitel(idxs) out: dict[str, dict] = {} for e in _fakten_schema(_json_datei(path)) or []: bt = _match_sub(e["baustein"], ct) if bt not in rel_by: continue sub = _match_sub(e["subbaustein"], rel_by[bt]) out.setdefault(bt, {})[_norm_titel(sub)] = {"sub": sub, **{k: e[k] for k in _FAKTEN_FELDER}} return out # Roh-Fakten + Completeness-Ergänzungen vereinigen (Recall): nur Subs, die in roh existieren. def _chunk_fakten(ci): roh = roh_map(ci, roh_path(ci)) erg = roh_map(ci, erg_path(ci)) if erg_path(ci).exists() else {} if not erg: return roh for bt, fm in roh.items(): ebt = erg.get(bt, {}) for sn, fk in fm.items(): ek = ebt.get(sn) if not ek: continue seen = {str(k).strip().casefold() for k in fk.get("kernpunkte", [])} for k in ek.get("kernpunkte", []): if str(k).strip().casefold() not in seen: seen.add(str(k).strip().casefold()) fk["kernpunkte"].append(k) seent = {bf["text"].strip().casefold() for bf in fk.get("belegte_fakten", [])} for bf in ek.get("belegte_fakten", []): if bf["text"].strip().casefold() not in seent: seent.add(bf["text"].strip().casefold()) fk["belegte_fakten"].append(bf) for f in ("voraussetzungen", "huerden", "beispiel_idee"): if not fk.get(f) and ek.get(f): fk[f] = ek[f] return roh # Phase „Fakten finden": 1 Generator je Chunk. async def _finde(ci, idxs): fp = roh_path(ci) if _fakten_schema(_json_datei(fp)): return True subs_total = sum(len(bausteine[i][1]) for i in idxs) status, _r = await run_single_slot( ctx, f"Fakten {ci}", key=f"bausteine-{topic}-fakten-c{ci}", prompt=_prompt("Fakten-Recherche", topic=topic, source=source, bausteine=block_text(idxs), out_path=fp, extra=_extra(instructions)), role="guide", capabilities=caps, payload=lambda result, p=fp: _fakten_schema(_json_datei(p)), timeout=_timeout("inhalt", subs_total)) return status != FAILED and _fakten_schema(_json_datei(fp)) is not None oks = await _gather_fortschritt([_finde(ci, idxs) for ci, idxs in enumerate(chunks)], len(chunks), _melde_p(set_p, topic, "Fakten finden")) if is_cancelled(): return None if not any(ok is True for ok in oks): _bausteine_errors[topic] = "Fakten-Extraktion fehlgeschlagen" return None # Phase „Fakten ergänzen" (Recall): ein gezielter Gap-Hunt je Chunk sucht source-belegte Fakten, die # der Single-Pass übersah. Best-effort — schlägt nie fehl (keine erg-Datei → Merge nutzt nur roh). async def _ergaenze(ci, idxs): ep = erg_path(ci) if _fakten_schema(_json_datei(ep)): return per = roh_map(ci, roh_path(ci)) if not per: return block = "\n\n".join( f"BAUSTEIN: {bt}\nSUBBAUSTEINE (mit bereits erfassten Fakten):\n" + "\n".join( f"- {fk['sub']}\n Erfasst: " + ("; ".join( list(fk.get("kernpunkte", [])) + [bf["text"] for bf in fk.get("belegte_fakten", [])]) or "(nichts)") for fk in fm.values()) for bt, fm in per.items()) subs_total = sum(len(bausteine[i][1]) for i in idxs) await run_single_slot( ctx, f"Fakten ergänzen {ci}", key=f"bausteine-{topic}-fakten-erg-c{ci}", prompt=_prompt("Fakten-Ergaenzung", topic=topic, source=source, bausteine=block, out_path=ep, extra=_extra(instructions)), role="guide", capabilities=caps, payload=lambda result, p=ep: _fakten_schema(_json_datei(p)), timeout=_timeout("inhalt", subs_total)) set_p("Fakten ergänzen…", step=_step_idx(topic, "Fakten finden")) await _gather_fortschritt([_ergaenze(ci, idxs) for ci, idxs in enumerate(chunks)], len(chunks), _melde_p(set_p, topic, "Fakten finden")) if is_cancelled(): return None # Phase „Fakten prüfen": FAKTEN_CHECK_PANEL Judges je Chunk. Zwei Mehrheits-Mengen: # beanstandet (Fakt ungenau → korrigieren) und verwerfen (Sub nicht belegbar → entfernen). async def _pruefe(ci, idxs): per = _chunk_fakten(ci) # roh + Ergänzungen → Panel verifiziert die Vereinigung if not per: return ci, set(), set() fakten_text = "\n\n".join(f"SUBBAUSTEIN: {fk['sub']}\n{_fakten_zeilen(fk)}" for fm in per.values() for fk in fm.values()) offen = [j for j in (1, 2, 3)[:FAKTEN_CHECK_PANEL] if _fakten_check_schema(_json_datei(chk_path(ci, j))) is None] await asyncio.gather(*[ run_agent(f"bausteine-{topic}-fakten-check-c{ci}-j{j}", _prompt("Fakten-Check", topic=topic, source=source, fakten=fakten_text, out_path=chk_path(ci, j), extra=_extra(instructions)), _timeout("inhalt_check", len(per)), provider=provider, role="judge", capabilities=caps) for j in offen], return_exceptions=True) outs = [s for j in (1, 2, 3)[:FAKTEN_CHECK_PANEL] if (s := _fakten_check_schema(_json_datei(chk_path(ci, j)))) is not None] bvotes: dict[str, int] = {} vvotes: dict[str, int] = {} for s in outs: # s = [(sub_norm, verwerfen)] eines Judges gb, gv = set(), set() for sn, verw in s: if sn not in gb: gb.add(sn); bvotes[sn] = bvotes.get(sn, 0) + 1 if verw and sn not in gv: gv.add(sn); vvotes[sn] = vvotes.get(sn, 0) + 1 schwelle = len(outs) / 2 if outs else 99 beanstandet = {sn for sn, v in bvotes.items() if v > schwelle} # Verwerfen ist irreversibel → strenger als Beanstanden: Mehrheit UND ≥2 zustimmende Judges # (verhindert Löschung durch eine Einzelstimme, wenn das Panel degradiert ist). verwerfen = {sn for sn, v in vvotes.items() if v > schwelle and v >= 2} return ci, beanstandet, verwerfen pruef = await _gather_fortschritt([_pruefe(ci, idxs) for ci, idxs in enumerate(chunks)], len(chunks), _melde_p(set_p, topic, "Fakten prüfen")) if is_cancelled(): return None beanstandet: dict[int, set] = {} verwerfen: dict[int, set] = {} for r in pruef: if isinstance(r, tuple) and len(r) == 3: ci, b, v = r beanstandet[ci] = b verwerfen[ci] = v # Phase „Fakten fix": nur KORRIGIERBARE (beanstandet ohne verwerfen) neu extrahieren. korrigieren = {ci: (beanstandet.get(ci, set()) - verwerfen.get(ci, set())) for ci in beanstandet} n_problem = sum(len(s) for s in korrigieren.values()) if n_problem: set_p(f"Fakten korrigieren ({n_problem})…", step=_step_idx(topic, "Fakten fix")) async def _fix(ci): subs_norm = korrigieren.get(ci, set()) if not subs_norm or _fakten_schema(_json_datei(fix_path(ci))): return idxs = chunks[ci] rel_by = {bausteine[i][0]: bausteine[i][1] for i in idxs} # nur die korrigierbaren Subs als Block ziel = [] for bt, subs in rel_by.items(): betroffen = [s for s in subs if _norm_titel(s) in subs_norm] if betroffen: ziel.append(f"BAUSTEIN: {bt}\nSUBBAUSTEINE:\n" + "\n".join(f"- {s}" for s in betroffen)) if not ziel: return await run_single_slot( ctx, f"Fakten-Fix {ci}", key=f"bausteine-{topic}-fakten-fix-c{ci}", prompt=_prompt("Fakten-Recherche", topic=topic, source=source, bausteine="\n\n".join(ziel), out_path=fix_path(ci), extra=_extra(instructions)), role="guide", capabilities=caps, payload=lambda result, p=fix_path(ci): _fakten_schema(_json_datei(p)), timeout=_timeout("inhalt", len(subs_norm))) await _gather_fortschritt([_fix(ci) for ci in korrigieren], len(korrigieren), _melde_p(set_p, topic, "Fakten fix")) if is_cancelled(): return None # Zusammensetzen: roh + Fix-Overrides für korrigierte. Verworfene Subs raus (+ je Baustein melden). ergebnis: dict[str, dict] = {} verworfen_map: dict[str, set] = {} for ci in range(len(chunks)): per = _chunk_fakten(ci) # roh + Ergänzungen (Recall); Fix überschreibt nur Korrigierte fix = roh_map(ci, fix_path(ci)) if fix_path(ci).exists() else {} verw = verwerfen.get(ci, set()) for bt, fm in per.items(): for sn, fk in fm.items(): if sn in verw: verworfen_map.setdefault(bt, set()).add(sn) continue gewinner = fix.get(bt, {}).get(sn, fk) if sn in korrigieren.get(ci, set()) else fk ergebnis.setdefault(bt, {})[sn] = {k: gewinner[k] for k in _FAKTEN_FELDER} if verworfen_map: _log(topic, f"Fakten-Check verwirft {sum(len(s) for s in verworfen_map.values())} unbelegbare Subbausteine") return ergebnis, verworfen_map async def _relevanz_block(ctx: GenContext, set_p, files: dict, sidecar: dict, instructions: str) -> dict | None: """Block D: drei Phasen mit Barriere — Finden (relevant/rand), Wählen (Vote), Klären. Items aus der Sidecar; lokale IDs 1..n pro Paket → globale gid. → {gid: relevanz} oder None bei Abbruch/Recherche-Fehler. Default bei Lücke/Streit: 'relevant'.""" topic, provider, is_cancelled = ctx.topic, ctx.provider, ctx.is_cancelled arbeit = files["arbeit"] items = [(titel, sub["titel"], sub.get("fakten")) for titel, subs in sidecar.items() for sub in subs] # globale id = index+1 if not items: return {} chunks = _chunk_nums(list(range(len(items))), _n_chunks(len(items), STUFE_CHUNK)) n = len(chunks) def rater_paths(c): return [arbeit / f"relevanz-c{c}-{i}.json" for i in (1, 2, 3)] def lset(item_idxs): return set(range(1, len(item_idxs) + 1)) # Phase „Relevanz finden": pro Paket 3 Rater (min. 2), lokale IDs. async def _rate(c, item_idxs): local_set = lset(item_idxs) paths = rater_paths(c) vorhanden = sum(1 for p in paths if _relevanz_schema(_json_datei(p), local_set)) if vorhanden >= 2: return True enum_zeilen = [] for k, j in enumerate(item_idxs, 1): enum_zeilen.append(f"{k}. [{items[j][0]}] {items[j][1]}") if (kz := _kern_zeile(items[j][2])): enum_zeilen.append(f" {kz}") enum = "\n".join(enum_zeilen) offen = [(i, p) for i, p in enumerate(paths, 1) if not _relevanz_schema(_json_datei(p), local_set)] slots = [{ "key": f"bausteine-{topic}-relevanz-c{c}-{i}", "prompt": _prompt("Relevanz-Recherche", topic=topic, subbausteine=enum, out_path=p, extra=_extra(instructions)), "role": "fast", "capabilities": "files", "payload": (lambda result, p=p, ids=local_set: _relevanz_schema(_json_datei(p), ids)), } for i, p in offen] neu = await _race(topic, f"Relevanz Paket {c}", slots, 2 - vorhanden, _timeout("relevanz", len(item_idxs)), provider, cancelled=is_cancelled, grace=KONSENS_GRACE) return not is_cancelled() and neu is not None oks = await _gather_fortschritt([_rate(c, idxs) for c, idxs in enumerate(chunks, 1)], len(chunks), _melde_p(set_p, topic, "Relevanz finden")) if is_cancelled(): return None if not all(ok is True for ok in oks): _bausteine_errors[topic] = "Relevanz fehlgeschlagen (Recherche)" return None # Phase „Relevanz wählen": Code-Vote je Paket → (ergebnis, strittig). set_p(f"Relevanz wählen ({n} Pakete)…", step=_step_idx(topic, "Relevanz wählen")) vote_by_c = {} for c, item_idxs in enumerate(chunks, 1): local_set = lset(item_idxs) rater = [d for p in rater_paths(c) if (d := _relevanz_schema(_json_datei(p), local_set))] ergebnis: dict[int, str] = {} strittig: dict[int, list[str]] = {} for k in range(1, len(item_idxs) + 1): stimmen = [d[k] for d in rater if k in d] zaehler: dict[str, int] = {} for s in stimmen: zaehler[s] = zaehler.get(s, 0) + 1 best = max(zaehler.values(), default=0) gewinner = [s for s, v in zaehler.items() if v == best] if len(gewinner) == 1 and best >= 2: ergebnis[k] = gewinner[0] else: strittig[k] = stimmen vote_by_c[c] = (ergebnis, strittig) # Phase „Relevanz klären": Judge je Paket mit Strittigem, alle parallel. async def _klaere(c, item_idxs): ergebnis, strittig = vote_by_c[c] if strittig: judge_path = arbeit / f"relevanz-final-c{c}.json" entsch = _relevanz_schema(_json_datei(judge_path), set(strittig)) if entsch is None: strittig_block = "\n".join( f"{k}. [{items[item_idxs[k - 1]][0]}] {items[item_idxs[k - 1]][1]} — Stimmen: {', '.join(stimmen) or 'keine'}" for k, stimmen in strittig.items() ) status, entsch = await run_single_slot( ctx, f"Relevanz-Klärung {c}", key=f"bausteine-{topic}-relevanz-final-c{c}", prompt=_prompt("Relevanz-Mapping", topic=topic, strittig=strittig_block, out_path=judge_path, extra=_extra(instructions)), role="judge", capabilities="files", payload=lambda result, p=judge_path, ids=set(strittig): _relevanz_schema(_json_datei(p), ids), timeout=_timeout("relevanz_check", len(strittig)), ) if status == FAILED: _log(topic, f"Relevanz-Klärung Paket {c} fehlgeschlagen — Default 'relevant'") entsch = entsch if isinstance(entsch, dict) else {} # Strittige ohne Entscheid → 'relevant' (nie versehentlich ausschließen). ergebnis = {**{k: "relevant" for k in strittig}, **ergebnis, **entsch} return {item_idxs[k - 1] + 1: rel for k, rel in ergebnis.items()} parts = await _gather_fortschritt([_klaere(c, idxs) for c, idxs in enumerate(chunks, 1)], len(chunks), _melde_p(set_p, topic, "Relevanz klären")) if is_cancelled(): return None relevanz_by_id: dict[int, str] = {} for c, part in enumerate(parts, 1): if not isinstance(part, dict): # Klärung ist nicht fatal: Vote-Ergebnis + Default 'relevant' für Strittige. if isinstance(part, BaseException): _log(topic, f"Relevanz-Klärung Paket {c}: {type(part).__name__}: {part}") ergebnis, strittig = vote_by_c[c] item_idxs = chunks[c - 1] merged = {**{k: "relevant" for k in strittig}, **ergebnis} part = {item_idxs[k - 1] + 1: s for k, s in merged.items()} relevanz_by_id.update(part) return relevanz_by_id def _norm_frage(t: str) -> str: return " ".join(str(t or "").lower().split()) def _match_sub(agent_sub: str, rel: list[str]) -> str: """Subbaustein-Titel des Agenten auf den passenden relevanten Titel mappen — exakt, dann normalisiert, dann Teilstring (der Agent lässt z. B. das Präfix „Frage: " weg). Kein Treffer → Agent-Titel behalten. So geht KEIN Muster durch Titel-Abweichung verloren.""" if agent_sub in rel: return agent_sub an = _norm_titel(agent_sub) for r in rel: rn = _norm_titel(r) if an and rn and (an == rn or an in rn or rn in an): return r return agent_sub async def _frage_muster_block(ctx: GenContext, set_p, files: dict, sidecar: dict, instructions: str) -> dict | None: """Block E (10er-Chunks): Finden (1 Generator je ~10 Bausteine, parallel), Wählen (Code: je Baustein gruppieren + Dedup), Klären (1 Kritiker je Chunk), Prüfen (Nachrunde). Zuordnung je Eintrag über das `baustein`-Feld (Chunk-Datei trägt mehrere Bausteine). → {Baustein-Titel: [{subbaustein, frage}, …]} oder None bei Abbruch.""" topic, is_cancelled = ctx.topic, ctx.is_cancelled arbeit = files["arbeit"] # ALLE Subbausteine (auch rand) bekommen ein Muster — Rand ist in der FGuide-Ebene prüfbar. # fakten_by: voller Fakten-Kontext je Sub (Generierung profitiert davon — bessere Fragen). bausteine = [] fakten_by: dict[tuple, dict] = {} for titel, subs in sidecar.items(): alle = [] for s in subs: if isinstance(s, dict) and (st := str(s.get("titel", "")).strip()): alle.append(st) if isinstance(s.get("fakten"), dict): fakten_by[(titel, _norm_titel(st))] = s["fakten"] if alle: bausteine.append((titel, alle)) if not bausteine: return {} chunks = _lpt_chunks([len(rel) for _, rel in bausteine], FRAGE_CHUNK_SUBS) # lastbalanciert nach Sub-Zahl def roh_path(ci): return arbeit / f"frage-muster-c{ci}.json" def final_path(ci): return arbeit / f"frage-muster-final-c{ci}.json" def _chunk_titel(idxs): return [bausteine[i][0] for i in idxs] # Phase „Fragen finden": je Chunk 1 Generator, alle parallel. async def _finde(ci, idxs): fp = roh_path(ci) if _frage_muster_chunk_schema(_json_datei(fp)): return # Resume def _sub_zeile(bi, s): zeile = f"- {s}" fk = fakten_by.get((bausteine[bi][0], _norm_titel(s))) if fk and (ft := _fakten_zeilen(fk)): zeile += "\n" + "\n".join(" " + l for l in ft.split("\n")) return zeile block = "\n\n".join( f"BAUSTEIN: {bausteine[i][0]}\nSUBBAUSTEINE:\n" + "\n".join(_sub_zeile(i, s) for s in bausteine[i][1]) for i in idxs ) subs_total = sum(len(bausteine[i][1]) for i in idxs) status, _ = await run_single_slot( ctx, f"Frage-Muster {ci}", key=f"bausteine-{topic}-frage-muster-c{ci}", prompt=_prompt("Frage-Muster-Recherche", topic=topic, bausteine=block, out_path=fp, extra=_extra(instructions)), role="fast", capabilities="files", payload=lambda result, p=fp: _frage_muster_chunk_schema(_json_datei(p)), timeout=_timeout("frage_muster", subs_total), ) if status == FAILED: _log(topic, f"Frage-Muster Chunk {ci} fehlgeschlagen — Bausteine im Fallback (Nachrunde/Live)") async def finde_alle(ci_list): ci_list = list(ci_list) await _gather_fortschritt([_finde(ci, chunks[ci]) for ci in ci_list], len(ci_list), _melde_p(set_p, topic, "Fragen finden")) await finde_alle(range(len(chunks))) if is_cancelled(): return None # Phase „Fragen wählen": Code — Chunk-Dateien je Baustein gruppieren, Dubletten raus, # Baustein-/Subbaustein-Titel locker auf die Vorgaben mappen (nichts wegen Abweichung verwerfen). def _waehle_chunk(ci): idxs = chunks[ci] ctitel = _chunk_titel(idxs) rel_by = {bausteine[i][0]: bausteine[i][1] for i in idxs} out, gesehen = {}, {} for e in _frage_muster_chunk_schema(_json_datei(roh_path(ci))) or []: titel = _match_sub(e["baustein"], ctitel) if titel not in rel_by: continue # nicht zuordenbar → verwerfen sub = _match_sub(e["subbaustein"], rel_by[titel]) seen = gesehen.setdefault(titel, set()) if sub in seen: continue # genau ein Muster je Subbaustein seen.add(sub) out.setdefault(titel, []).append({"subbaustein": sub, "frage": e["frage"]}) return out def _waehle_all(ci_list): roh = {} for ci in ci_list: for titel, eintraege in _waehle_chunk(ci).items(): roh.setdefault(titel, []).extend(eintraege) return roh set_p("Fragen wählen…", step=_step_idx(topic, "Fragen wählen")) roh_by_titel = _waehle_all(range(len(chunks))) # Phase „Fragen klären": je Chunk 1 Kritiker bereinigt die Tabellen (nach Baustein gruppiert). async def _klaere(ci, idxs): fp = final_path(ci) if _frage_muster_chunk_schema(_json_datei(fp)): return # Resume bloecke = [] for i in idxs: t = bausteine[i][0] eintraege = roh_by_titel.get(t) or [] if not eintraege: continue zeilen = "\n".join(f"- ({e['subbaustein']}) {e['frage']}" for e in eintraege) bloecke.append(f"BAUSTEIN: {t}\n{zeilen}") if not bloecke: return # nichts zu klären in diesem Chunk subs_total = sum(len(bausteine[i][1]) for i in idxs) status, _ = await run_single_slot( ctx, f"Frage-Muster-Klärung {ci}", key=f"bausteine-{topic}-frage-muster-final-c{ci}", prompt=_prompt("Frage-Muster-Kritik", topic=topic, tabelle="\n\n".join(bloecke), out_path=fp, extra=_extra(instructions)), role="judge", capabilities="files", payload=lambda result, p=fp: _frage_muster_chunk_schema(_json_datei(p)), timeout=_timeout("frage_muster_check", subs_total), ) if status == FAILED: _log(topic, f"Frage-Muster-Klärung Chunk {ci} fehlgeschlagen — Roh-Muster übernommen") async def klaere_alle(ci_list): ci_list = list(ci_list) await _gather_fortschritt([_klaere(ci, chunks[ci]) for ci in ci_list], len(ci_list), _melde_p(set_p, topic, "Fragen klären")) await klaere_alle(range(len(chunks))) if is_cancelled(): return None # Geklärte Chunk-Tabelle je Baustein, Fallback auf Roh-Muster. Titel locker mappen. def _final_by_titel(ci_list): out = {} for ci in ci_list: idxs = chunks[ci] ctitel = _chunk_titel(idxs) rel_by = {bausteine[i][0]: bausteine[i][1] for i in idxs} for e in _frage_muster_chunk_schema(_json_datei(final_path(ci))) or []: titel = _match_sub(e["baustein"], ctitel) if titel not in rel_by: continue out.setdefault(titel, []).append( {"subbaustein": _match_sub(e["subbaustein"], rel_by[titel]), "frage": e["frage"]}) return out final_by_titel = _final_by_titel(range(len(chunks))) ergebnis = {t: (final_by_titel.get(t) or roh_by_titel.get(t) or []) for t, _ in bausteine} # Phase „Fragen prüfen": Chunks mit ≥1 leeren Baustein eine Runde nachholen. set_p("Fragen prüfen…", step=_step_idx(topic, "Fragen prüfen")) leer_titel = {t for t, _ in bausteine if not ergebnis.get(t)} if leer_titel: nach = [ci for ci, idxs in enumerate(chunks) if any(bausteine[i][0] in leer_titel for i in idxs)] _log(topic, f"Frage-Muster: {len(leer_titel)} Baustein(e) ohne Muster — Nachrunde über {len(nach)} Chunk(s)") for ci in nach: roh_path(ci).unlink(missing_ok=True) final_path(ci).unlink(missing_ok=True) await finde_alle(nach) if is_cancelled(): return None for titel, eintraege in _waehle_all(nach).items(): roh_by_titel[titel] = eintraege await klaere_alle(nach) if is_cancelled(): return None for titel, eintraege in _final_by_titel(nach).items(): final_by_titel[titel] = eintraege for t, _ in bausteine: if not ergebnis.get(t): ergebnis[t] = final_by_titel.get(t) or roh_by_titel.get(t) or [] rest = [t for t, _ in bausteine if not ergebnis.get(t)] if rest: _log(topic, f"Frage-Muster: {len(rest)} Baustein(e) bleiben leer (Fallback Live): {rest[:5]}") return ergebnis # ── Inventar in der DB: Recherche-Loop · Konsolidierung · Klärung ──────────── def _crawl_index(ordner) -> dict[str, str]: """Alias (Dateiname ODER QUELLE:-URL, klein) → kanonischer Seiten-Key (Dateiname).""" idx: dict[str, str] = {} if not ordner or not Path(ordner).is_dir(): return idx for p in sorted(Path(ordner).glob("*.txt")): key = p.name idx[key.lower()] = key try: erste = p.read_text(encoding="utf-8").splitlines()[0] except (OSError, IndexError): erste = "" if erste.startswith("QUELLE:"): url = erste[len("QUELLE:"):].strip() if url: idx[url.lower()] = key idx[url.rstrip("/").lower()] = key return idx async def _set_inventar(topic: str, eintrag: str, status: str) -> None: """Einen Inventar-Eintrag ('Titel — Beschreibung') mit Status in die DB schreiben.""" titel = _titel(eintrag) norm = _norm_titel(titel) if not norm: return teile = [t.strip() for t in eintrag.split(" — ")] besch = teile[1] if len(teile) >= 2 else "" await db.upsert_baustein(topic, norm, titel, besch) await db.set_baustein_status(topic, norm, status) def _sichte_regeln(ordner, pages: list[str]) -> tuple[list[str], list[str]]: """Deterministischer Content/Noise-Filter (config.CRAWL_*). Substring-Match (klein) gegen URL + Dateiname. Reihenfolge: keep > noise > min_chars > behalten. → (content, noise).""" ordner = Path(ordner) content, noise = [], [] for fn in pages: zeilen = _read(ordner / fn).splitlines() url = zeilen[0][len("QUELLE:"):].strip() if zeilen and zeilen[0].startswith("QUELLE:") else "" body = "\n".join(zeilen[1:]).strip() hay = f"{url}\n{fn}".lower() if any(p in hay for p in CRAWL_KEEP_PATTERNS): content.append(fn) elif any(p in hay for p in CRAWL_NOISE_PATTERNS): noise.append(fn) elif len(body) < CRAWL_MIN_CHARS: noise.append(fn) else: content.append(fn) # Default: behalten — alles mit Inhalt bleibt return content, noise def _seite_snippet(ordner, fn: str) -> tuple[str, str]: """(url, snippet) einer Crawl-Seite für das Relevanz-Gate. url aus der QUELLE:-Zeile; snippet = Body-Auszug (Navigations-Boilerplate steht vorn — der Prompt ignoriert es). URL ist das Primärsignal (sprechender Slug), der Snippet stützt nur.""" zeilen = _read(Path(ordner) / fn).splitlines() url = zeilen[0][len("QUELLE:"):].strip() if zeilen and zeilen[0].startswith("QUELLE:") else "" body = "\n".join(zeilen[1:]).strip() snippet = " ".join(body.split())[:QUELLE_RELEVANZ_SNIPPET] return (url or fn), snippet async def _relevanz_sichtung(ctx: GenContext, set_p, files: dict, ordner, content: list[str], spec: str, instructions: str) -> tuple[list[str], list[str]]: """LLM-Themen-Gate nach dem Regel-Filter: jede Content-Seite ja/nein gegen die Spec. Off-topic (anderes Fachgebiet) → raus. Muster wie `_relevanz_block`: kleine Pakete, 3 Rater (`fast`), 2-von-3-Konsens. KONSERVATIV: nur bei klarer „nein"-Mehrheit droppen; Streit/Lücke/ Race-Fehler → behalten. SAFETY: würde das Gate ≥80 % (oder alles) droppen, bleibt alles (Spec-Mismatch/Bug soll die Quelle nicht leeren). → (behalten, raus) als Dateinamen.""" topic, provider, is_cancelled = ctx.topic, ctx.provider, ctx.is_cancelled arbeit = files["arbeit"] pages = sorted(content) if not pages: return content, [] items = [_seite_snippet(ordner, fn) for fn in pages] # Index deckt sich mit `pages` chunks = _chunk_nums(list(range(len(pages))), _n_chunks(len(pages), QUELLE_RELEVANZ_CHUNK)) n = len(chunks) def rater_paths(c): return [arbeit / f"quelle-relevanz-c{c}-{i}.json" for i in (1, 2, 3)] def lset(idxs): return set(range(1, len(idxs) + 1)) async def _rate(c, idxs): local_set = lset(idxs) paths = rater_paths(c) vorhanden = sum(1 for p in paths if _janein_schema(_json_datei(p), local_set)) if vorhanden >= 2: return True enum_zeilen = [] for k, j in enumerate(idxs, 1): url, snip = items[j] enum_zeilen.append(f"{k}. {url}") if snip: enum_zeilen.append(f" {snip}") enum = "\n".join(enum_zeilen) offen = [(i, p) for i, p in enumerate(paths, 1) if not _janein_schema(_json_datei(p), local_set)] slots = [{ "key": f"bausteine-{topic}-quelle-relevanz-c{c}-{i}", "prompt": _prompt("Quelle-Relevanz", topic=topic, spec=spec, seiten=enum, out_path=p, extra=_extra(instructions)), "role": "fast", "capabilities": "files", "payload": (lambda result, p=p, ids=local_set: _janein_schema(_json_datei(p), ids)), } for i, p in offen] neu = await _race(topic, f"Relevanz-Sichtung Paket {c}", slots, 2 - vorhanden, _timeout("relevanz", len(idxs)), provider, cancelled=is_cancelled, grace=KONSENS_GRACE) return not is_cancelled() and neu is not None _qidx = _step_idx(topic, "Quelle aufbereiten") # Gate läuft im Quelle-Schritt (kein eigener Step) set_p(f"Prüfe Relevanz zur Spec ({n} Pakete)…", step=_qidx) async def _melde_sichtung(d, t): set_p(f"Prüfe Relevanz zur Spec {d}/{t}…", step=_qidx) await _gather_fortschritt([_rate(c, idxs) for c, idxs in enumerate(chunks, 1)], n, _melde_sichtung) if is_cancelled(): return content, [] # Abbruch → nichts droppen (Caller bricht ab) # Vote je Seite: nur eine klare „nein"-Mehrheit (≥2 und mehr als „ja") wirft raus. raus: list[str] = [] for c, idxs in enumerate(chunks, 1): local_set = lset(idxs) rater = [d for p in rater_paths(c) if (d := _janein_schema(_json_datei(p), local_set))] for k in range(1, len(idxs) + 1): stimmen = [d[k] for d in rater if k in d] nein, ja = stimmen.count("nein"), stimmen.count("ja") if nein >= 2 and nein > ja: raus.append(pages[idxs[k - 1]]) if raus and len(raus) >= max(1, int(len(pages) * 0.8)): _log(topic, f"Relevanz-Sichtung: würde {len(raus)}/{len(pages)} droppen — verworfen (Spec-Mismatch?), alles behalten") return content, [] raus_set = set(raus) behalten = [fn for fn in pages if fn not in raus_set] return behalten, raus async def _quelle_aufbereiten(ctx: GenContext, set_p, files: dict, q: dict, ordner, instructions: str) -> bool: """Schritt „Quelle aufbereiten": Crawl (link) + PDF-Konvert + Content/Noise-Sichtung. Persistiert die Sichtung in der Coverage-Tabelle (inhalt). → True (ok) / False (Abbruch/Fehler). thema: nichts. projekt/uni: nur PDFs (kuratierter Ordner, keine Sichtung).""" topic, is_cancelled = ctx.topic, ctx.is_cancelled if not ordner: return True # thema → keine Quelle aufzubereiten if q["type"] != "link": await asyncio.to_thread(_pdfs_konvertieren, ordner) # projekt/uni: nur PDFs, keine Sichtung return True if await db.get_step_status(topic, "Quelle aufbereiten") == "fertig": return True if not _crawl_fertig(topic): set_p("Quelle laden (Crawl)…", step=_step_idx(topic, "Quelle aufbereiten")) n = await asyncio.to_thread(crawl, q["ort"], ordner, cancelled=is_cancelled) if is_cancelled(): return False if not n: _bausteine_errors[topic] = "Crawl ergab keine Inhalte — Link/Domain prüfen" return False await asyncio.to_thread(_pdfs_konvertieren, ordner) pages = sorted(set(_crawl_index(ordner).values())) if pages: set_p("Sichte Seiten…", step=_step_idx(topic, "Quelle aufbereiten")) await db.delete_coverage(topic) content, noise = _sichte_regeln(ordner, pages) # deterministischer Regel-Filter if q.get("spec") and content: # Themen-Gate: trennt das Fachgebiet (Regeln können das nicht) content, raus = await _relevanz_sichtung(ctx, set_p, files, ordner, content, q["spec"], instructions) if is_cancelled(): return False if raus: noise = sorted(set(noise) | set(raus)) _log(topic, f"LLM-Relevanz: {len(raus)} Seiten off-topic → Noise") await db.mark_inhalt(topic, sorted(content), sorted(noise)) _log(topic, f"Sichtung: {len(content)} Content / {len(noise)} Noise von {len(pages)} (Regeln + LLM-Gate)") await db.set_step_status(topic, "Quelle aufbereiten", "fertig") return True async def _recherche_batch(ctx: GenContext, set_p, files: dict, q: dict, ordner, instructions: str) -> bool: """Befüllt DB-Tabelle `bausteine` mit Kandidaten (+ Nennungszähler). FESTE Datei-Batches: jede Crawl-Seite wird genau einem Batch zugeteilt und von RECHERCHE_READERS Agenten gelesen (Konsens ≥2 im Batch). Alle zugeteilten Seiten werden als gelesen markiert → 100 % Abdeckung. Ohne Crawl-Ordner (Quelle „thema") → freie Web-Recherche, eine Runde. → True/False.""" topic, provider, is_cancelled = ctx.topic, ctx.provider, ctx.is_cancelled if await db.get_step_status(topic, "Recherche") == "fertig": return True arbeit = files["arbeit"] await db.delete_bausteine(topic) # Coverage/inhalt gehört der Sichtung — NICHT löschen await db.set_step_status(topic, "Recherche", "laufend") async def _ingest(reader_id: str, text: str) -> None: gesehen = set() for eintrag in _parse_auswahl(text).values(): titel = _titel(eintrag) norm = _norm_titel(titel) if not norm or norm in gesehen: continue gesehen.add(norm) # ein Reader = eine Stimme je Konzept teile = [t.strip() for t in eintrag.split(" — ")] besch = teile[1] if len(teile) >= 2 else "" quelle = [teile[2]] if len(teile) >= 3 and teile[2] else [] await db.upsert_baustein(topic, norm, titel, besch, quelle, reader=reader_id) pages = await db.list_content(topic) # von der Sichtung als Content markierte Seiten if not pages and ordner: pages = sorted(set(_crawl_index(ordner).values())) # Fallback (projekt/uni: keine Sichtung) if not pages: # Quelle „thema" (oder kein Crawl): freie Web-Recherche, eine Runde. set_p("Recherche läuft…", step=_step_idx(topic, "Recherche")) caps = "files" if ordner else "full" paths = [arbeit / f"recherche-{i}.md" for i in range(1, RECHERCHE_THEMA_AGENTEN + 1)] for p in paths: p.unlink(missing_ok=True) slots = [{ "key": f"bausteine-{topic}-recherche-{i}", "prompt": _build_recherche_prompt(topic, p, instructions, q["type"], ordner), "role": "quick", "capabilities": caps, "payload": (lambda result, p=p, rid=f"t{i}": ((rid, t) if (t := _file_payload(p)) else None)), } for i, p in enumerate(paths, 1)] texte = await _race(topic, "Recherche", slots, 3, _timeout("recherche"), provider, cancelled=is_cancelled, grace=RECHERCHE_GRACE) if is_cancelled(): return False if not texte: _bausteine_errors[topic] = "Recherche fehlgeschlagen (Minimum nicht erreicht)" return False for rid, text in texte: await _ingest(rid, text) await db.set_step_status(topic, "Recherche", "fertig") return True # uni/projekt: kuratierte, oft GROSSE Dateien (Skript). Statt alle am Stück zu lesen # (Lost-in-the-Middle), in Abschnitte chunken und JEDEN gründlich von 2 Readern lesen — # Text direkt im Prompt (kleiner Kontext), Nennungen akkumulieren zu Konsens. if q["type"] in ("uni", "projekt"): eintraege: list[tuple[str, str]] = [] # (dateiname, abschnitt-text) for fn in sorted(pages): for absch in _text_abschnitte(_read(ordner / fn)): eintraege.append((fn, absch)) if not eintraege: _bausteine_errors[topic] = "Recherche: Quelle leer" return False set_p(f"Recherche ({len(eintraege)} Abschnitte)…", step=_step_idx(topic, "Recherche")) async def _lese_abschnitt(ei: int, fn: str, absch: str) -> None: block = (f"ARBEITE AUSSCHLIESSLICH MIT DIESEM TEXTABSCHNITT (Quelle: {fn}). Lies ihn " f"VOLLSTÄNDIG, überspringe nichts. Notiere `{fn}` als Quelle jedes Bausteins. " f"Suche NICHT im Web — nur dieser Abschnitt zählt.\n\n-----\n{absch}\n-----") paths = [arbeit / f"recherche-a{ei}-{i}.md" for i in range(1, RECHERCHE_READERS + 1)] # Reader-Datei-Wiederverwendung: liegen alle Reader-Outputs valide vor (Resume / # Re-Run ohne Recherche-Änderung), re-ingestieren statt erneut Agenten zu spawnen. vorhanden = [(f"a{ei}-{i}", t) for i, p in enumerate(paths, 1) if (t := _file_payload(p))] if len(vorhanden) == len(paths): for rid, text in vorhanden: await _ingest(rid, text) return for p in paths: p.unlink(missing_ok=True) if is_cancelled(): return slots = [{ "key": f"bausteine-{topic}-recherche-a{ei}-{i}", "prompt": _build_recherche_prompt(topic, p, instructions, q["type"], ordner, abschnitt=block), "role": "quick", "capabilities": "files", "payload": (lambda result, p=p, rid=f"a{ei}-{i}": ((rid, t) if (t := _file_payload(p)) else None)), } for i, p in enumerate(paths, 1)] # Quorum 2: beide Reader pro Abschnitt sollen durch (mehr Augen = mehr Konzepte + # echter Konsens); nach Timeout fällt _race auf das Vorhandene zurück. texte = await _race(topic, f"Recherche Abschnitt {ei}", slots, 2, _timeout("recherche", 1), provider, cancelled=is_cancelled, grace=RECHERCHE_GRACE) for rid, text in (texte or []): await _ingest(rid, text) await _gather_fortschritt([_lese_abschnitt(ei, fn, a) for ei, (fn, a) in enumerate(eintraege, 1)], len(eintraege), _melde_p(set_p, topic, "Recherche")) if is_cancelled(): return False await db.mark_quellen_gelesen(topic, sorted(pages)) gesamt = len(await db.list_bausteine(topic)) _log(topic, f"Recherche (uni/projekt): {gesamt} Kandidaten aus {len(eintraege)} Abschnitten ({len(pages)} Dateien)") if not gesamt: _bausteine_errors[topic] = "Recherche fehlgeschlagen (keine Bausteine)" return False await db.set_step_status(topic, "Recherche", "fertig") return True # Crawl/Link: viele kleine Content-Seiten (Sichtung im Schritt „Quelle aufbereiten"). # Feste Batches, je Batch RECHERCHE_READERS Reader, die GENAU diese Dateien lesen. batches = _chunk_nums(sorted(pages), max(1, math.ceil(len(pages) / RECHERCHE_BATCH))) async def _lese_batch(bi: int, batch: list[str]) -> bool: liste = "\n".join(f"- {p}" for p in batch) fokus = ("WICHTIG — feste Zuteilung: Bearbeite AUSSCHLIESSLICH diese Dateien und lies JEDE " f"vollständig. Ignoriere alle anderen Dateien im Ordner:\n{liste}") paths = [arbeit / f"recherche-b{bi}-{i}.md" for i in range(1, RECHERCHE_READERS + 1)] for p in paths: p.unlink(missing_ok=True) if not is_cancelled(): slots = [{ "key": f"bausteine-{topic}-recherche-b{bi}-{i}", "prompt": _build_recherche_prompt(topic, p, instructions, q["type"], ordner, fokus=fokus), "role": "quick", "capabilities": "files", "payload": (lambda result, p=p, rid=f"b{bi}-{i}": ((rid, t) if (t := _file_payload(p)) else None)), } for i, p in enumerate(paths, 1)] texte = await _race(topic, f"Recherche Batch {bi}", slots, 1, _timeout("recherche", len(batch)), provider, cancelled=is_cancelled, grace=RECHERCHE_GRACE) for rid, text in (texte or []): await _ingest(rid, text) await db.mark_quellen_gelesen(topic, batch) # alle zugeteilten Seiten abhaken (auch ohne Treffer) return not is_cancelled() await _gather_fortschritt([_lese_batch(bi, b) for bi, b in enumerate(batches, 1)], len(batches), _melde_p(set_p, topic, "Recherche")) if is_cancelled(): return False gesamt = len(await db.list_bausteine(topic)) deckung = len(await db.list_coverage(topic)) _log(topic, f"Recherche: {gesamt} Kandidaten, Abdeckung {deckung}/{len(pages)} Seiten ({len(batches)} Batches)") if not gesamt: _bausteine_errors[topic] = "Recherche fehlgeschlagen (keine Bausteine)" return False await db.set_step_status(topic, "Recherche", "fertig") return True def _grp_schema(data, ids: set[int]): """{"gruppen": [[1,3],[2], …]} → Partition von `ids` als Liste von Index-Gruppen. Tolerant: ignoriert Fremd-/Doppel-Nummern; vergessene Kandidaten werden eigenständig (Singleton-Gruppe) ergänzt. None nur bei strukturell kaputtem JSON.""" if not isinstance(data, dict) or not isinstance(data.get("gruppen"), list): return None gruppen, gesehen = [], set() for g in data["gruppen"]: if not isinstance(g, list): return None grp = [] for x in g: try: num = int(x) except (ValueError, TypeError): continue if num in ids and num not in gesehen: gesehen.add(num) grp.append(num) if grp: gruppen.append(grp) gruppen += [[r] for r in sorted(ids - gesehen)] # vergessene Kandidaten bleiben eigenständig return gruppen or None _ASPEKT_MARKER = ("∈ np", "∈np", " in np", "np-schwer", "np-vollständig", "verifizierer", "zertifikat", "ndtm", "nicht-determ", "lower bound", "untere schranke", "bzgl", "als sprache") def _aspekt_marker(titel: str) -> int: """Anzahl Eigenschafts-Marker im Titel (∈NP, NP-schwer, Verifizierer, Lower Bound …). 0 = generisches Hauptkonzept (das Problem selbst); >0 = eine Eigenschaft davon.""" t = titel.casefold() return sum(1 for m in _ASPEKT_MARKER if m in t) def _canonical(kandidaten: list[dict], idxs: list[int], gesehen_norm: set[str]) -> dict: """Repräsentant eines Clusters = das Hauptkonzept (wenigste Eigenschafts-Marker — das Problem selbst, nicht „… ∈ NP"); Tie → häufigster norm-Titel → meiste Reader. Titel global eindeutig (Suffix ' (2)'), damit er als Schlüssel taugt.""" by_norm: dict[str, list[int]] = {} for k in idxs: by_norm.setdefault(_norm_titel(kandidaten[k]["titel"]), []).append(k) def gewicht(nb: str): ms = by_norm[nb] reader = set().union(*[set(kandidaten[m]["reader"]) for m in ms]) if ms else set() return (-_aspekt_marker(nb), len(ms), len(reader)) # aspekt-arm zuerst = Hauptkonzept best = max(by_norm, key=gewicht) k = max(by_norm[best], key=lambda m: len(kandidaten[m]["beschreibung"])) titel = kandidaten[k]["titel"] n = 2 while _norm_titel(titel) in gesehen_norm: titel = f"{kandidaten[k]['titel']} ({n})" n += 1 gesehen_norm.add(_norm_titel(titel)) return {"titel": titel, "beschreibung": kandidaten[k]["beschreibung"]} async def _block_gruppieren(ctx: GenContext, set_p, arbeit: Path, kandidaten: list[dict], blocks: list[list[int]], praefix: str = "konsolidierung", schritt: str = "Konsolidierung") -> list[list[int]]: """Je Ähnlichkeits-Block gruppiert ein Judge die Titel in die echten Bausteine (merge Paraphrasen, split Über-Merges). Singletons direkt. Fehler/Timeout → konservativ jeder Kandidat einzeln (vermeidet fälschliches Über-Mergen). → finale Gruppen (globale Indizes). `praefix`/`schritt` trennen Konsolidierung und Dedup (Artefakte, Race-Key, Fortschritt).""" topic, is_cancelled = ctx.topic, ctx.is_cancelled multi = [(bi, b) for bi, b in enumerate(blocks) if len(b) > 1] ergebnis: list[list[int]] = [list(b) for b in blocks if len(b) == 1] # Singletons direkt def _zeile(k: int, g: int) -> str: b = kandidaten[g] return f"{k}. {b['titel']}" + (f" — {b['beschreibung']}" if b["beschreibung"] else "") async def _grp(bi: int, block: list[int]) -> None: ids = set(range(1, len(block) + 1)) p = arbeit / f"{praefix}-block-c{bi}.json" part = _grp_schema(_json_datei(p), ids) if part is None: # Resume: gültige Datei nicht neu rechnen p.unlink(missing_ok=True) if is_cancelled(): return zeilen = [_zeile(k, block[k - 1]) for k in range(1, len(block) + 1)] status, part = await run_single_slot( ctx, f"Block-Gruppieren {bi}", key=f"bausteine-{topic}-{praefix}-block-c{bi}", prompt=_prompt("Bausteine-Block-Gruppieren", topic=topic, eintraege="\n".join(zeilen), out_path=p), role="judge", capabilities="files", payload=(lambda result, p=p, ids=ids: _grp_schema(_json_datei(p), ids)), timeout=_timeout("recherche_mapping", len(block)), ) part = part if status == OK else None if part is None: # Judge gescheitert → einzeln (kein Über-Merge) ergebnis.extend([idx] for idx in block) else: # lokale Nummern → globale Kandidaten-Indizes ergebnis.extend([block[k - 1] for k in g] for g in part) await _gather_fortschritt([_grp(bi, b) for bi, b in multi], len(multi), _melde_p(set_p, topic, schritt)) return ergebnis async def _konsolidiere_embedding(ctx: GenContext, set_p, files: dict, kandidaten: list[dict]) -> bool: """Zweistufig: Embeddings → grobe Capped-Blocks (High-Recall) → je Multi-Block ein Judge, der die Titel in die echten Bausteine gruppiert → Reader-Union (≥2 = Konsens).""" topic, is_cancelled = ctx.topic, ctx.is_cancelled arbeit = files["arbeit"] texts = [f"{b['titel']} — {b['beschreibung']}" if b["beschreibung"] else b["titel"] for b in kandidaten] sims = await asyncio.to_thread(embedding.embed_sims, texts) if sims is None: # Modell doch nicht verfügbar → Fallback return await _konsolidiere_llm(ctx, set_p, files, kandidaten) # Stufe 1: grobe Ähnlichkeits-Blocks (gedeckelt, kein Giant-Component). blocks = await asyncio.to_thread(embedding.capped_blocks, sims, None, None) # Stufe 2: ein Judge gruppiert JEDEN Multi-Block in die echten Bausteine. gruppen = await _block_gruppieren(ctx, set_p, arbeit, kandidaten, blocks) if is_cancelled(): return False def _min_cos(idxs): # interne Kohärenz zur Kontrolle (Ketten hätten ~0,3) if len(idxs) < 2: return 1.0 return round(min(float(sims[i][j]) for n, i in enumerate(idxs) for j in idxs[n + 1:]), 3) # Konsens = ≥2 distinkte Reader pro Cluster. Legacy-DBs ohne Reader-Tracking (Recherche lief # vor der Migration, kein Re-Ingest) haben leere Reader-Sets → Rückfall auf Titel-Heuristik # (sonst landete ALLES im Rest). hat_reader = any(b["reader"] for b in kandidaten) konsens, rest, debug, gesehen_norm = [], [], [], set() for idxs in gruppen: reader = set().union(*[set(kandidaten[k]["reader"]) for k in idxs]) if idxs else set() if hat_reader: score = len(reader) else: # ohne Reader-Daten: max(Nennungen, Anzahl distinkter Titel-Varianten im Cluster) score = max(max(kandidaten[k]["nennungen"] for k in idxs), len({kandidaten[k]["titel_norm"] for k in idxs})) rep = _canonical(kandidaten, idxs, gesehen_norm) eintrag = f"{rep['titel']} — {rep['beschreibung']}" if rep["beschreibung"] else rep["titel"] (konsens if score >= 2 else rest).append(eintrag) debug.append({"titel": rep["titel"], "reader": sorted(reader), "score": score, "konsens": score >= 2, "min_cos": _min_cos(idxs), "mitglieder": [kandidaten[k]["titel"] for k in idxs]}) atomic_write_json(arbeit / "konsolidierung-cluster.json", debug, indent=1) multi_blocks = sum(1 for b in blocks if len(b) > 1) _log(topic, f"Konsolidierung (Embedding): {len(blocks)} Blocks ({multi_blocks} per LLM gruppiert) " f"→ {len(gruppen)} Cluster aus {len(kandidaten)} Kandidaten " f"→ {len(konsens)} Konsens / {len(rest)} Rest") await db.delete_bausteine(topic) for t in konsens: await _set_inventar(topic, t, "konsens") for t in rest: await _set_inventar(topic, t, "rest") await db.set_step_status(topic, "Konsolidierung", "fertig") return True async def _konsolidiere(ctx: GenContext, set_p, files: dict) -> bool: """Mergt Roh-Kandidaten zu Konsens (≥2 Reader)/Rest. Deterministisch per Embedding-Clustering; fehlt das Modell → Rückfall auf das LLM-Panel (`_konsolidiere_llm`). Status in DB.""" topic = ctx.topic if await db.get_step_status(topic, "Konsolidierung") == "fertig": return True set_p("Konsolidiere Recherche…", step=_step_idx(topic, "Konsolidierung")) kandidaten = await db.list_bausteine(topic) if not kandidaten: _bausteine_errors[topic] = "Konsolidierung: keine Kandidaten" return False if EMBEDDING_AKTIV and await asyncio.to_thread(embedding.verfuegbar): return await _konsolidiere_embedding(ctx, set_p, files, kandidaten) return await _konsolidiere_llm(ctx, set_p, files, kandidaten) async def _konsolidiere_llm(ctx: GenContext, set_p, files: dict, kandidaten: list[dict]) -> bool: """Fallback (nur ohne Embedding-Modell): Panel (KONSOLIDIERUNG_PANEL Judges) mergt Kandidaten semantisch; ein Reconcile-Judge führt die Panel-Ausgaben zur finalen Konsens (≥2)/Rest (1×)-Liste zusammen. Panel statt Einzel-Judge: ein einzelner Judge ist bias-anfällig und instabil.""" topic, provider, is_cancelled = ctx.topic, ctx.provider, ctx.is_cancelled arbeit = files["arbeit"] chunks = _chunk_nums(kandidaten, max(1, math.ceil(len(kandidaten) / KONSOLIDIERUNG_CHUNK))) async def _map_panel(c: int, eintraege: str, anzahl: int): """3 Mapping-Judges über `eintraege` → Reconcile-Judge → (konsens, rest). None bei Abbruch/Fehler.""" paths = [arbeit / f"konsolidierung-c{c}-j{j}.json" for j in range(1, KONSOLIDIERUNG_PANEL + 1)] offen = [(j, p) for j, p in enumerate(paths, 1) if _mapping_schema(_json_datei(p)) is None] for _, p in offen: p.unlink(missing_ok=True) if offen: slots = [{ "key": f"bausteine-{topic}-konsolidierung-c{c}-j{j}", "prompt": _prompt("Bausteine-Recherche-Mapping", topic=topic, n=RECHERCHE_READERS, eintraege=eintraege, out_path=p), "role": "judge", "capabilities": "files", "payload": (lambda result, p=p: _mapping_schema(_json_datei(p))), } for j, p in offen] vorhanden = KONSOLIDIERUNG_PANEL - len(offen) await _race(topic, f"Konsolidierung {c}", slots, max(1, 2 - vorhanden), _timeout("recherche_mapping", anzahl), provider, cancelled=is_cancelled, grace=KONSENS_GRACE) if is_cancelled(): return None outs = [m for p in paths if (m := _mapping_schema(_json_datei(p)))] if not outs: return None # Union der Panel-Titel; je Titel zählen, wie viele Judges ihn als Konsens führen. kvotes: dict[str, int] = {} form: dict[str, str] = {} # norm → Anzeigetitel (erstes Vorkommen) order: list[str] = [] for kk, rr in outs: for t in kk + rr: nt = _norm_titel(_titel(t)) if not nt: continue if nt not in form: form[nt] = t order.append(nt) kvotes.setdefault(nt, 0) for t in kk: nt = _norm_titel(_titel(t)) if nt: kvotes[nt] = kvotes.get(nt, 0) + 1 # Reconcile: ein Merge-Judge über die Union, annotiert mit Judge-Stimmen ("k× genannt"). rp = arbeit / f"konsolidierung-c{c}-reconcile.json" recon = _mapping_schema(_json_datei(rp)) if recon is None: rp.unlink(missing_ok=True) eintraege_r = "\n".join(f"{i}. {form[nt]} ({max(1, kvotes[nt])}× genannt)" for i, nt in enumerate(order, 1)) status, recon = await run_single_slot( ctx, f"Konsolidierung Reconcile {c}", key=f"bausteine-{topic}-konsolidierung-c{c}-reconcile", prompt=_prompt("Bausteine-Recherche-Mapping", topic=topic, n=KONSOLIDIERUNG_PANEL, eintraege=eintraege_r, out_path=rp), role="judge", capabilities="files", payload=lambda result, p=rp: _mapping_schema(_json_datei(p)), timeout=_timeout("recherche_mapping", len(order)), ) if status == CANCELLED: return None recon = recon if status != FAILED else None if recon: return recon # Fallback (Reconcile gescheitert): Code-Mehrheit — Konsens, wenn Mehrheit der Judges Konsens sagt. konsens = [form[nt] for nt in order if kvotes[nt] * 2 >= len(outs) and kvotes[nt] > 0] kset = {_norm_titel(_titel(t)) for t in konsens} return konsens, [form[nt] for nt in order if nt not in kset] konsens, rest = [], [] for c, chunk in enumerate(chunks, 1): eintraege = "\n".join( f"{i}. {b['titel']} — {b['beschreibung']} ({b['nennungen']}× genannt)" for i, b in enumerate(chunk, 1) ) res = await _map_panel(c, eintraege, len(chunk)) if res is None: if is_cancelled(): return False _bausteine_errors[topic] = "Recherche-Mapping fehlgeschlagen" return False k, r = res konsens += k rest += r # Bei mehreren Chunks: ein globaler Merge-Pass über die vereinten Konsens-Einträge, # damit Dubletten über Chunk-Grenzen (DAL×4, PHPUnit×5 …) verschmelzen. if len(chunks) > 1 and konsens: fp = arbeit / "konsolidierung-merge.json" fp.unlink(missing_ok=True) eintraege = "\n".join(f"{i}. {t} (2× genannt)" for i, t in enumerate(konsens, 1)) status, mapping = await run_single_slot( ctx, "Konsolidierung Merge", key=f"bausteine-{topic}-konsolidierung-merge", prompt=_prompt("Bausteine-Recherche-Mapping", topic=topic, n=RECHERCHE_READERS, eintraege=eintraege, out_path=fp), role="judge", capabilities="files", payload=lambda result, p=fp: _mapping_schema(_json_datei(p)), timeout=_timeout("recherche_mapping", len(konsens)), ) if status == CANCELLED: return False if status != FAILED and mapping: konsens, r2 = mapping rest += r2 # vom Merge zurückgestufte Einträge in den Rest # Judge-Ausgabe ist maßgeblich → Inventar in der DB neu setzen. await db.delete_bausteine(topic) for t in konsens: await _set_inventar(topic, t, "konsens") for t in rest: await _set_inventar(topic, t, "rest") await db.set_step_status(topic, "Konsolidierung", "fertig") return True async def _klaere_inventar(ctx: GenContext, set_p, files: dict) -> bool: """Panel (KONSOLIDIERUNG_PANEL Judges) entscheidet über den Rest (1×-Genannte): Mehrheit `aufnehmen` → Konsens, sonst verworfen. Panel statt Einzel-Judge — der Rest-Schnitt ist der schärfste Eingriff; ein einzelner Judge ist hier zu instabil. Konservativer Tie → behalten (nie ein Konzept verlieren).""" topic, provider, is_cancelled = ctx.topic, ctx.provider, ctx.is_cancelled if await db.get_step_status(topic, "Klärung") == "fertig": return True set_p("Klärung läuft…", step=_step_idx(topic, "Klärung")) rest_rows = await db.list_bausteine(topic, status="rest") if rest_rows: arbeit = files["arbeit"] paths = [arbeit / f"klaerung-j{j}.json" for j in range(1, KONSOLIDIERUNG_PANEL + 1)] offen = [(j, p) for j, p in enumerate(paths, 1) if _runde_schema(_json_datei(p), final=True) is None] for _, p in offen: p.unlink(missing_ok=True) if offen: slots = [{ "key": f"bausteine-{topic}-klaerung-j{j}", "prompt": _prompt( "Bausteine-Klaerung", topic=topic, rest="\n".join(f"- {b['titel']}" for b in rest_rows), final="\n- Entscheide JEDEN Eintrag. `rest` MUSS leer sein.", out_path=p, ), "role": "judge", "capabilities": "files", "payload": (lambda result, p=p: _runde_schema(_json_datei(p), final=True)), } for j, p in offen] vorhanden = KONSOLIDIERUNG_PANEL - len(offen) await _race(topic, "Klärung", slots, max(1, 2 - vorhanden), _timeout("auswahl_mapping", len(rest_rows)), provider, cancelled=is_cancelled, grace=KONSENS_GRACE) if is_cancelled(): return False outs = [r for p in paths if (r := _runde_schema(_json_datei(p), final=True))] if not outs: _bausteine_errors[topic] = "Klärung fehlgeschlagen" return False # Mehrheit je Rest-Eintrag (per Norm-Titel). Tie → behalten (votes*2 >= n). votes: dict[str, int] = {} for aufnehmen, _ in outs: for nt in {_norm_titel(_titel(t)) for t in aufnehmen}: votes[nt] = votes.get(nt, 0) + 1 for b in rest_rows: auf = votes.get(b["titel_norm"], 0) * 2 >= len(outs) await db.set_baustein_status(topic, b["titel_norm"], "konsens" if auf else "verworfen") await db.set_step_status(topic, "Klärung", "fertig") return True async def _dedup_inventar(ctx: GenContext, set_p, files: dict) -> bool: """Finaler iterativer Dedup-Pass über die fertige Konsens-Liste: Embedding-Block + LLM- Gruppierung. Fängt Dubletten, die Konsolidierung (Block-Grenzen, Cap) und Klärung (lange Liste, Lost-in-the-Middle) übersehen. Pro Gruppe bleibt EIN Repräsentant, der Rest wird verworfen. Iteriert, weil die kleinere Liste je Runde neu blockt (Cross-Block-Reste).""" topic, is_cancelled = ctx.topic, ctx.is_cancelled if await db.get_step_status(topic, "Dedup") == "fertig": return True if not (EMBEDDING_AKTIV and await asyncio.to_thread(embedding.verfuegbar)): await db.set_step_status(topic, "Dedup", "fertig") # ohne Modell: still überspringen return True set_p("Dedup…", step=_step_idx(topic, "Dedup")) arbeit = files["arbeit"] for runde in range(1, DEDUP_MAX_RUNDEN + 1): konsens = await db.list_bausteine(topic, status="konsens") if len(konsens) < 2: break texts = [f"{b['titel']} — {b['beschreibung']}" if b["beschreibung"] else b["titel"] for b in konsens] sims = await asyncio.to_thread(embedding.embed_sims, texts) if sims is None: break blocks = await asyncio.to_thread(embedding.capped_blocks, sims, None, None) gruppen = await _block_gruppieren(ctx, set_p, arbeit, konsens, blocks, praefix=f"dedup-r{runde}", schritt="Dedup") if is_cancelled(): return False weg = 0 for idxs in gruppen: if len(idxs) < 2: continue # Repräsentant = Hauptkonzept (wenigste Eigenschafts-Marker), dann kürzester Titel; Rest verwerfen. rep = min(idxs, key=lambda k: (_aspekt_marker(konsens[k]["titel"]), len(konsens[k]["titel"]), k)) for k in idxs: if k != rep: await db.set_baustein_status(topic, konsens[k]["titel_norm"], "verworfen") weg += 1 atomic_write_json(arbeit / f"dedup-runde-{runde}.json", {"vorher": len(konsens), "entfernt": weg, "gruppen": [[konsens[k]["titel"] for k in g] for g in gruppen if len(g) > 1]}, indent=1) _log(topic, f"Dedup Runde {runde}: {len(konsens)} → {len(konsens) - weg} (−{weg})") if weg <= max(1, len(konsens) * DEDUP_MIN_DELTA // 100): # konvergiert → Schluss break await db.set_step_status(topic, "Dedup", "fertig") return True # --- Gliederung (Bausteine-Artefakt: Kapitel-Struktur, vom Guide nur gelesen) --- def _gliederung_komplett(files: dict) -> bool: """Gliederung steht (Kapitel-Liste vorhanden)?""" d = _json_datei(files["gliederung"]) return isinstance(d, dict) and isinstance(d.get("kapitel"), list) and bool(d.get("kapitel")) def _gliederung_schema(data, valid: set[int]): """{"kapitel":[{titel,nummern}]} → bereinigt (valide Nummern, je genau einmal) · None bei <80 % Abdeckung (Agent/Judge hat zu viel weggelassen).""" if not isinstance(data, dict) or not isinstance(data.get("kapitel"), list): return None out, seen = [], set() for ch in data["kapitel"]: if not isinstance(ch, dict): continue titel = str(ch.get("titel", "")).strip() or "Kapitel" nums = [] for n in (ch.get("nummern") or []): try: n = int(n) except (ValueError, TypeError): continue if n in valid and n not in seen: seen.add(n) nums.append(n) if nums: out.append({"titel": titel, "nummern": nums}) if not out or len(seen) < 0.8 * len(valid): return None return {"kapitel": out} def _prereq_schema(data, valid: set[int]) -> dict[int, list[int]]: """{"prereqs": {"3": [1, 7]}} → {num: [prereq-nums]} · nur Nummern aus `valid`, ohne Selbstkante. Ungültig/leer → {} (Best-effort: dann Original-Reihenfolge).""" if not isinstance(data, dict) or not isinstance(data.get("prereqs"), dict): return {} out: dict[int, list[int]] = {} for k, v in data["prereqs"].items(): try: num = int(k) except (ValueError, TypeError): continue if num not in valid or not isinstance(v, list): continue pres = [] for p in v: try: p = int(p) except (ValueError, TypeError): continue if p in valid and p != num and p not in pres: pres.append(p) if pres: out[num] = pres return out def _topo_order(nums: list[int], edges: dict[int, list[int]]) -> list[int]: """Kahn-Topo-Sort: Voraussetzungen zuerst. `edges[num]` = Nummern, die VOR num kommen müssen. Stabiler Tie-Break (Original-Reihenfolge von `nums`); Zyklen werden gebrochen (nie Deadlock).""" pos = {n: i for i, n in enumerate(nums)} # Resteingangsgrad nur über gültige Knoten; Selbst-/Fremdkanten ignoriert. pre = {n: [p for p in edges.get(n, []) if p in pos and p != n] for n in nums} fertig: list[int] = [] erledigt: set[int] = set() rest = list(nums) while rest: bereit = [n for n in rest if all(p in erledigt for p in pre[n])] if not bereit: # Zyklus → den in Original-Reihenfolge frühesten Rest-Knoten erzwingen bereit = [min(rest, key=lambda n: pos[n])] nxt = min(bereit, key=lambda n: pos[n]) # stabil: kleinste Original-Position zuerst fertig.append(nxt) erledigt.add(nxt) rest.remove(nxt) return fertig async def _lernreihenfolge(ctx: GenContext, set_p, files: dict, entries: dict, valid: set[int], instructions: str) -> dict: """entries (num→titel) in Lernreihenfolge bringen: LLM extrahiert Prereq-Kanten aus den extrahierten `voraussetzungen`, Code löst per Topo-Sort. Best-effort → sonst entries unverändert.""" if len(entries) < 3: return entries topic = ctx.topic fakten_map = _json_datei(files["fakten"]) fakten_map = fakten_map if isinstance(fakten_map, dict) else {} def _hint(titel): fm = fakten_map.get(titel) or {} vs = [v for fk in fm.values() if isinstance(fk, dict) and (v := str(fk.get("voraussetzungen", "")).strip())] return " · ".join(dict.fromkeys(vs)) pp = files["arbeit"] / "gliederung-prereqs.json" def _payload(result, p=pp): d = _json_datei(p) return d if isinstance(d, dict) and "prereqs" in d else None vorhanden = _json_datei(pp) if not (isinstance(vorhanden, dict) and "prereqs" in vorhanden): zeilen = [f"{n}. {t}" + (f"\n braucht vorher: {h}" if (h := _hint(t)) else "") for n, t in entries.items()] set_p("Gliederung — Lernreihenfolge…", step=_step_idx(topic, "Gliederung")) await run_single_slot( ctx, "Gliederung-Voraussetzungen", key=f"bausteine-{topic}-gliederung-prereqs", prompt=_prompt("Gliederung-Voraussetzungen", topic=topic, bausteine="\n".join(zeilen), out_path=pp, extra=_extra(instructions)), role="guide", capabilities="files", payload=_payload, timeout=_timeout("plan", len(entries))) edges = _prereq_schema(_json_datei(pp), valid) if not edges: return entries # keine/ungültige Kanten → Original-Reihenfolge (kein Regress) ordered = _topo_order(list(entries), edges) return {n: entries[n] for n in ordered} async def _gliederung_block(ctx: GenContext, set_p, files: dict, entries: dict, instructions: str) -> dict: """Format-agnostische Gliederung über ALLE Bausteine — 3 Vorschläge → Judge merged. Bricht nie ab: 0 gültige → ein Kapitel mit allem; fehlende Bausteine landen in „Weitere". → {"kapitel":[{titel,nummern}]} (auch in files["gliederung"]).""" topic, is_cancelled = ctx.topic, ctx.is_cancelled valid = set(entries) step = _step_idx(topic, "Gliederung") # Lernreihenfolge gründen (LLM-Modulo): LLM extrahiert Prereq-Kanten aus den extrahierten # `voraussetzungen`, Code löst per Topo-Sort. Best-effort → sonst Original-Reihenfolge. entries = await _lernreihenfolge(ctx, set_p, files, entries, valid, instructions) liste = "\n".join(f"{n}. {t}" for n, t in entries.items()) set_p("Gliederung — Vorschläge…", step=step) async def _vorschlag(i, path): if _gliederung_schema(_json_datei(path), valid): return True await run_single_slot( ctx, f"Gliederung {i}", key=f"bausteine-{topic}-gliederung-{i}", prompt=_prompt("Guide-Gliederung", topic=topic, bausteine=liste, out_path=path, extra=_extra(instructions)), role="guide", capabilities="files", payload=lambda result, p=path: _gliederung_schema(_json_datei(p), valid), timeout=_timeout("plan", len(entries))) return _gliederung_schema(_json_datei(path), valid) is not None slots = files["gliederung_slots"] await _gather_fortschritt([_vorschlag(i, p) for i, p in enumerate(slots, 1)], len(slots), _melde_p(set_p, topic, "Gliederung")) if is_cancelled(): return {} vorschlaege = [v for p in slots if (v := _gliederung_schema(_json_datei(p), valid))] if not vorschlaege: plan = {"kapitel": [{"titel": "Inhalte", "nummern": list(entries)}]} elif len(vorschlaege) == 1: plan = vorschlaege[0] else: set_p("Gliederung zusammenführen…", step=step) bloecke = "\n\n".join( f"### Vorschlag {i}\n" + "\n".join( f"KAPITEL: {ch['titel']}\n Nummern: {', '.join(str(n) for n in ch['nummern'])}" for ch in v["kapitel"]) for i, v in enumerate(vorschlaege, 1)) await run_single_slot( ctx, "Gliederung-Judge", key=f"bausteine-{topic}-gliederung-judge", prompt=_prompt("Guide-Gliederung-Judge", topic=topic, format_name="den Guide", zweck="alle Bausteine in einem roten Faden", n=len(vorschlaege), bausteine=liste, gliederungen=bloecke, out_path=files["gliederung"], extra=_extra(instructions)), role="judge", capabilities="files", payload=lambda result: _gliederung_schema(_json_datei(files["gliederung"]), valid), timeout=_timeout("plan_judge", len(entries))) plan = _gliederung_schema(_json_datei(files["gliederung"]), valid) or vorschlaege[0] # Vollständigkeit: jeder Baustein kommt vor — fehlende in „Weitere" (gegen weglassende Agenten/Judge). drin = {n for ch in plan["kapitel"] for n in ch["nummern"]} fehlen = [n for n in entries if n not in drin] if fehlen: plan["kapitel"].append({"titel": "Weitere", "nummern": fehlen}) atomic_write_json(files["gliederung"], plan, indent=1) return plan # --- Lern-Artefakte (Karteikarten/Beispiele aus den Fakten) --- def _karten_schema(data): """{"karten":[{baustein,subbaustein,frage,antwort}]} → Liste (auch leer) · None bei kaputt.""" if not isinstance(data, dict) or not isinstance(data.get("karten"), list): return None out = [] for e in data["karten"]: if isinstance(e, dict) and (f := str(e.get("frage", "")).strip()) and (a := str(e.get("antwort", "")).strip()): out.append({"baustein": str(e.get("baustein", "")).strip(), "subbaustein": str(e.get("subbaustein", "")).strip(), "frage": f, "antwort": a}) return out def _beispiel_schema(data): """{"beispiele":[{baustein,subbaustein,problem,schritte,ergebnis}]} → Liste (auch leer) · None bei kaputt.""" if not isinstance(data, dict) or not isinstance(data.get("beispiele"), list): return None out = [] for e in data["beispiele"]: if not isinstance(e, dict): continue problem = str(e.get("problem", "")).strip() schritte = [s for x in (e.get("schritte") or []) if (s := str(x).strip())] if problem and schritte: out.append({"baustein": str(e.get("baustein", "")).strip(), "subbaustein": str(e.get("subbaustein", "")).strip(), "problem": problem, "schritte": schritte, "ergebnis": str(e.get("ergebnis", "")).strip()}) return out def _beispiel_check_schema(data): """Worked-Example-Check → {"ok": true} → set() (alles korrekt); {"probleme":[{"index":N}]} → {N, …} (1-basierte beanstandete Indizes); None bei kaputt.""" if not isinstance(data, dict): return None if data.get("ok") is True: return set() pr = data.get("probleme") if not isinstance(pr, list): return None out: set[int] = set() for p in pr: if isinstance(p, dict): try: out.add(int(p.get("index"))) except (ValueError, TypeError): continue return out _ARTEFAKT_SCHEMA = {"karteikarte": _karten_schema, "beispiel": _beispiel_schema} _ARTEFAKT_PROMPT = {"karteikarte": "Artefakt-Karteikarte", "beispiel": "Artefakt-Beispiel"} _ARTEFAKT_SCHRITT = {"karteikarte": "Karteikarten", "beispiel": "Beispiele"} def _artefakte_komplett(files: dict) -> bool: """Artefakt-Map steht (alle Typen erzeugt)? Werte dürfen leer sein (content-aware).""" d = _json_datei(files["artefakte"]) return isinstance(d, dict) and all(t in d for t in ARTEFAKT_TYPEN) async def _artefakte_block(ctx: GenContext, set_p, files: dict, sidecar: dict, instructions: str) -> dict | None: """Lern-Artefakte je Typ aus den gespeicherten Fakten erzeugen — ein Generierungs-Durchlauf je Typ über Chunks. Worked Examples werden gegen die Fakten verifiziert (falsche verworfen); Karteikarten sind risikoarm und bleiben ungeprüft. → {typ: [eintraege]} (auch in files).""" topic, provider, is_cancelled = ctx.topic, ctx.provider, ctx.is_cancelled arbeit = files["arbeit"] caps = "files" # Bausteine mit Subs + Fakten-Zeilen als Input-Block (extract-once aus den Fakten). bausteine = [] for btitel, subs in sidecar.items(): if not isinstance(subs, list): continue zeilen = [] for s in subs: if not isinstance(s, dict) or not (st := str(s.get("titel", "")).strip()): continue fk = s.get("fakten") if isinstance(s.get("fakten"), dict) else {} zeile = f"- {st}" if fk and (fk_text := _fakten_zeilen(fk)): zeile += "\n" + "\n".join(" " + l for l in fk_text.split("\n")) zeilen.append(zeile) if zeilen: bausteine.append((btitel, zeilen)) if not bausteine: leer = {t: [] for t in ARTEFAKT_TYPEN} atomic_write_json(files["artefakte"], leer, indent=1) return leer chunks = _lpt_chunks([len(z) for _, z in bausteine], FAKTEN_CHUNK_SUBS) def block_text(idxs): return "\n\n".join(f"BAUSTEIN: {bausteine[i][0]}\nSUBBAUSTEINE:\n" + "\n".join(bausteine[i][1]) for i in idxs) # Worked Examples gegen die Fakten prüfen (Panel-Mehrheit) — falsche verwerfen. CoT-Schritte sind # fehleranfällig; ein falsches Beispiel prägt ein fehlerhaftes Schema ein → kein Beispiel > falsches. async def _pruefe_beispiele(ci, idxs, items): if is_cancelled() or not items: return items def cpath(j): return arbeit / f"artefakt-beispiel-check-c{ci}-j{j}.json" beispiele_txt = "\n\n".join( f"{k}. PROBLEM: {e['problem']}\n SCHRITTE: " + " | ".join(e.get("schritte", [])) + (f"\n ERGEBNIS: {e['ergebnis']}" if e.get("ergebnis") else "") for k, e in enumerate(items, 1)) offen = [j for j in (1, 2, 3)[:FAKTEN_CHECK_PANEL] if _beispiel_check_schema(_json_datei(cpath(j))) is None] if offen: await asyncio.gather(*[ run_agent(f"bausteine-{topic}-artefakt-beispiel-check-c{ci}-j{j}", _prompt("Artefakt-Beispiel-Check", topic=topic, fakten=block_text(idxs), beispiele=beispiele_txt, out_path=cpath(j), extra=_extra(instructions)), _timeout("inhalt_check", len(items)), provider=provider, role="judge", capabilities=caps) for j in offen], return_exceptions=True) outs = [s for j in (1, 2, 3)[:FAKTEN_CHECK_PANEL] if (s := _beispiel_check_schema(_json_datei(cpath(j)))) is not None] if not outs: return items # keine Prüfung möglich → behalten (Best-effort) votes: dict[int, int] = {} for s in outs: for idx in s: votes[idx] = votes.get(idx, 0) + 1 schwelle = len(outs) / 2 raus = {idx for idx, v in votes.items() if v > schwelle} # Mehrheit (≥2 von 3) beanstandet → raus if raus: _log(topic, f"Worked-Example-Check Chunk {ci}: {len(raus)}/{len(items)} verworfen") return [e for k, e in enumerate(items, 1) if k not in raus] ergebnis: dict[str, list] = {} for typ in ARTEFAKT_TYPEN: schema = _ARTEFAKT_SCHEMA[typ] def apath(ci, t=typ): return arbeit / f"artefakt-{t}-c{ci}.json" async def _gen(ci, idxs, t=typ, schema=schema): p = apath(ci, t) if schema(_json_datei(p)) is not None: return True await run_single_slot( ctx, f"{_ARTEFAKT_SCHRITT[t]} {ci}", key=f"bausteine-{topic}-artefakt-{t}-c{ci}", prompt=_prompt(_ARTEFAKT_PROMPT[t], topic=topic, bausteine=block_text(idxs), out_path=p, extra=_extra(instructions)), role="guide", capabilities="files", payload=lambda result, p=p, schema=schema: schema(_json_datei(p)), timeout=_timeout("inhalt", sum(len(bausteine[i][1]) for i in idxs))) return schema(_json_datei(p)) is not None await _gather_fortschritt([_gen(ci, idxs) for ci, idxs in enumerate(chunks)], len(chunks), _melde_p(set_p, topic, _ARTEFAKT_SCHRITT[typ])) if is_cancelled(): return None eintraege: list = [] for ci in range(len(chunks)): chunk_items = schema(_json_datei(apath(ci))) or [] if typ == "beispiel" and chunk_items: chunk_items = await _pruefe_beispiele(ci, chunks[ci], chunk_items) eintraege += chunk_items ergebnis[typ] = eintraege atomic_write_json(files["artefakte"], ergebnis, indent=1) return ergebnis async def _mirror_artefakte_db(topic: str, sidecar: dict, artefakte: dict) -> None: """Artefakte in die DB spiegeln. Karteikarte/Beispiel je Sub (sub_norm).""" await db.delete_sub_artefakte(topic) btitel_list = list(sidecar.keys()) for typ in ARTEFAKT_TYPEN: for e in artefakte.get(typ, []): bt = _match_sub(e.get("baustein", ""), btitel_list) bnorm, sn = _norm_titel(bt), _norm_titel(e.get("subbaustein", "")) if not bnorm or not sn: continue daten = json.dumps({k: v for k, v in e.items() if k not in ("baustein", "subbaustein")}, ensure_ascii=False) await db.put_sub_artefakt(topic, bnorm, sn, typ, daten, bt, e.get("subbaustein", "")) async def _mirror_sidecar_db(topic: str, sidecar: dict) -> None: """Sidecar {Baustein-Titel: [{titel, stufe, relevanz}]} in die DB-Tabelle subbausteine spiegeln.""" for btitel, subs in sidecar.items(): bnorm = _norm_titel(btitel) if not bnorm or not isinstance(subs, list): continue for s in subs: if not isinstance(s, dict): continue st = str(s.get("titel", "")).strip() sn = _norm_titel(st) if not sn: continue fakten = json.dumps(s["fakten"], ensure_ascii=False) if isinstance(s.get("fakten"), dict) else None await db.put_subbaustein(topic, bnorm, sn, btitel, st, stufe=s.get("stufe"), relevanz=s.get("relevanz"), fakten=fakten, status="konsens") async def _mirror_frage_muster_db(topic: str, muster: dict) -> None: """Frage-Muster {Baustein-Titel: [{subbaustein, frage}]} in die DB-Tabelle frage_muster spiegeln.""" await db.delete_frage_muster(topic) for btitel, eintraege in muster.items(): bnorm = _norm_titel(btitel) if not bnorm or not isinstance(eintraege, list): continue for e in eintraege: if not isinstance(e, dict): continue sub = str(e.get("subbaustein", "")).strip() sn = _norm_titel(sub) frage = str(e.get("frage", "")).strip() if not (sn and frage): continue await db.upsert_frage_muster(topic, bnorm, sn, btitel, sub, frage) async def _reset_db_ab_phase(topic: str, label: str) -> None: """DB-Inhalt der Phasen ≥ `label` verwerfen (kanonische Reihenfolge Quelle…Artefakte).""" idx = _phase_idx(label) if idx <= 8: # Artefakte (Karteikarten/Beispiele) await db.delete_sub_artefakte(topic) if idx <= 7: # Fragen await db.delete_frage_muster(topic) if idx <= 6: # Gliederung await db.delete_gliederung(topic) if idx <= 2: # Subbausteine (Fakten/Stufen/Relevanz greifen über Sidecar→Mirror) await db.delete_subbausteine(topic) if idx <= 1: # Inventar: Inventar + Recherche-Schritte — Sichtung bleibt await db.delete_bausteine(topic) await db.delete_pipeline_state(topic, ["Recherche", "Konsolidierung", "Klärung", "Dedup"]) if idx <= 0: # Quelle: Sichtung neu (Coverage/inhalt + Schritt) await db.delete_coverage(topic) await db.delete_pipeline_state(topic, ["Quelle aufbereiten"]) async def generate_bausteine(topic: str, instructions: str = "", provider: str = DEFAULT_PROVIDER, ab_phase: int | None = None) -> None: if topic in _bausteine_progress: return _bausteine_progress[topic] = "Wartend…" _bausteine_errors.pop(topic, None) files = _bausteine_files(topic) final_path = files["final"] q = lade_quelle(topic) ordner = quelle_ordner(topic) # projekt/uni/link → Ordner, thema → None instructions = q.get("spec") or instructions # persistierte Spezifikation bevorzugen (auch bei Resume) def set_p(msg: str, step: int | None = None) -> None: _bausteine_progress[topic] = msg if step is not None: _bausteine_step[topic] = step def is_cancelled() -> bool: return topic in _bausteine_cancelled def abgebrochen() -> None: _bausteine_errors[topic] = "Abgebrochen — Fortschritt bleibt erhalten" ctx = GenContext(topic=topic, provider=provider, is_cancelled=is_cancelled) try: async with _semaphore: files["arbeit"].mkdir(parents=True, exist_ok=True) # Re-Run ab gewählter Phase: Artefakte ab dort löschen; der Frischstart-Block # unten wird übersprungen (er würde bei erhaltener Sidecar sonst alles wischen). if ab_phase is not None: phasen = _phasen(topic) label = phasen[ab_phase - 1][0] if 1 <= ab_phase <= len(phasen) else "Inventar" _reset_ab_phase(topic, label) await _reset_db_ab_phase(topic, label) # Schritt „Quelle aufbereiten": Crawl (link) + PDFs + Content/Noise-Sichtung. if not await _quelle_aufbereiten(ctx, set_p, files, q, ordner, instructions): if is_cancelled(): abgebrochen() return # „Neu erstellen": NUR wenn wirklich alles fertig ist (bausteine.md UND # Sidecar) → kompletter Frischstart. Liegt bausteine.md ohne Sidecar vor, # ist das ein Teil-Stand (Block B/C offen) → Resume, nicht wischen. # Bei explizitem Re-Run (ab_phase) hat _reset_ab_phase das schon erledigt. fertig = ab_phase is None and final_path.exists() and _sidecar_schema(_json_datei(files["sidecar"])) is not None if fertig: for p_alt in _alle_slot_dateien(files): p_alt.unlink(missing_ok=True) await db.delete_pipeline_state(topic) await db.delete_bausteine(topic) await db.delete_subbausteine(topic) await db.delete_frage_muster(topic) await db.delete_coverage(topic) await db.delete_gliederung(topic) await db.delete_sub_artefakte(topic) # Inventar (DB): Recherche-Loop → Konsolidierung → Klärung. if not await _recherche_batch(ctx, set_p, files, q, ordner, instructions): if is_cancelled(): abgebrochen() return if not await _konsolidiere(ctx, set_p, files): if is_cancelled(): abgebrochen() return if not await _klaere_inventar(ctx, set_p, files): if is_cancelled(): abgebrochen() return if not await _dedup_inventar(ctx, set_p, files): if is_cancelled(): abgebrochen() return konsens_rows = await db.list_bausteine(topic, status="konsens") entries = { i: (f"{b['titel']} — {b['beschreibung']}" if b["beschreibung"] else b["titel"]) for i, b in enumerate(konsens_rows, 1) } # Nur Projekte: Themenfeld-Ergänzung — Skript/Projekt ist ein Ausschnitt, # ein Web-Agent ergänzt kanonisch fehlende Bausteine, markiert mit [Ergänzung]. if q["type"] == "projekt": set_p("Ergänze Themenfeld…", step=_step_idx(topic, "Ergänzung")) erg_path = files["ergaenzung"] ergaenzungen = _ergaenzung_schema(_json_datei(erg_path)) if ergaenzungen is None: erg_path.unlink(missing_ok=True) status, ergaenzungen = await run_single_slot( ctx, "Ergänzung", key=f"bausteine-{topic}-ergaenzung-1", prompt=_prompt( "Bausteine-Ergaenzung", topic=topic, bausteine="\n".join(f"- {t}" for t in entries.values()), out_path=erg_path, extra=_extra(instructions), ), role="quick", capabilities="full", payload=lambda result: _ergaenzung_schema(_json_datei(erg_path)), timeout=_timeout("ergaenzung"), ) if status == CANCELLED: abgebrochen() return if status == FAILED: _bausteine_errors[topic] = "Ergänzung fehlgeschlagen (kein gültiges Ergebnis)" return idx = _titel_index(entries) neu = [(t, b) for t, b in ergaenzungen if _titel_aufloesen(idx, t) is None] if neu: _log(topic, f"Ergänzung: {len(neu)} Baustein(e) aus dem Themenfeld ergänzt") start = max(entries, default=0) + 1 for off, (t, b) in enumerate(neu): entries[start + off] = f"{t} — {b} [Ergänzung]" # Titel eindeutig machen und unsortiertes Inventar schreiben entries = _eindeutige_titel(entries) atomic_write_text(final_path, "\n".join(f"{i}. {t}" for i, t in entries.items()) + "\n") # Block B + C: Subbausteine je Baustein + Stufen → Sidecar subbausteine.json. # Nicht-destruktiv: bausteine.md steht schon; fehlt die Sidecar, wird beim # nächsten Lauf nur dieser Teil neu versucht. Guide fällt ohne Sidecar zurück. if _sidecar_schema(_json_datei(files["sidecar"])) is None: roh = _sub_roh_schema(_json_datei(files["sub_roh"])) if roh is None: roh = await _subbausteine_block(ctx, set_p, files, entries, instructions) if is_cancelled(): abgebrochen() return if roh is None: return # Fehler ist gesetzt atomic_write_json(files["sub_roh"], roh, indent=1) # Fakten je Sub (VOR der Stufe): Quell-Fakten extrahieren + verifizieren → fakten.json. # Extract-once-Grounding — Stufe/Relevanz/Fragen/Guide nähren sich daraus. if not _fakten_komplett(files): res = await _fakten_block(ctx, set_p, files, roh, q, ordner, instructions) if is_cancelled(): abgebrochen() return if res is None: return # Fehler ist gesetzt fakten_map, verworfen = res # Verworfene (unbelegbare) Subs aus roh streichen — ZUERST (Resume-robust), dann # fakten.json. So sehen Stufen/Relevanz/Gliederung/Fragen/Guide sie nicht mehr. if verworfen: for bt, sns in verworfen.items(): if bt in roh: roh[bt] = [s for s in roh[bt] if _norm_titel(s) not in sns] roh = {bt: subs for bt, subs in roh.items() if subs} # leere Bausteine raus (_sub_roh_schema verlangt ≥1) atomic_write_json(files["sub_roh"], roh, indent=1) atomic_write_json(files["fakten"], fakten_map, indent=1) sidecar = await _stufen_block(ctx, set_p, files, roh, instructions) if is_cancelled(): abgebrochen() return if sidecar is None: return # Fakten in die Sidecar-Subs mergen (DB-Spiegel + Guide-Nutzung). fakten_map = _json_datei(files["fakten"]) if isinstance(fakten_map, dict): for btitel, subs in sidecar.items(): fm = fakten_map.get(btitel, {}) for sub in subs: if (fk := fm.get(_norm_titel(sub["titel"]))): sub["fakten"] = fk atomic_write_json(files["sidecar"], sidecar, indent=1) # Block D: Relevanz je Subbaustein (relevant/rand) → in die Sidecar mergen. # Eigene Phase nach den Stufen; treibt das ProGuide-Format (alle Bausteine # mit ≥1 relevantem Subbaustein) und filtert Rand-Subs aus den Guides. sidecar = _json_datei(files["sidecar"]) if _sidecar_schema(sidecar) is not None and not _relevanz_komplett(sidecar): relevanz_by_id = await _relevanz_block(ctx, set_p, files, sidecar, instructions) if is_cancelled(): abgebrochen() return if relevanz_by_id is None: return # Fehler ist gesetzt gid = 0 for subs in sidecar.values(): for sub in subs: gid += 1 sub["relevanz"] = relevanz_by_id.get(gid, "relevant") atomic_write_json(files["sidecar"], sidecar, indent=1) # Block D.5: Gliederung (Bausteine-Artefakt) — Kapitel-Struktur über ALLE Bausteine, # vom Guide nur noch gelesen. Format-agnostisch; der Guide filtert je Format. if not _gliederung_komplett(files): await _gliederung_block(ctx, set_p, files, entries, instructions) if is_cancelled(): abgebrochen() return # Block E: Frage-Muster je relevantem Subbaustein × Typ → eigenes Sidecar. # Zur Prüfungszeit zieht jeder Agent ein Muster ohne Zurücklegen und formuliert # daraus eine Frage — distinkte Saat verhindert die Doppelfragen der Live-Generierung. sidecar = _json_datei(files["sidecar"]) if _sidecar_schema(sidecar) is not None and _relevanz_komplett(sidecar) and not _frage_muster_komplett(topic): muster = await _frage_muster_block(ctx, set_p, files, sidecar, instructions) if is_cancelled(): abgebrochen() return if muster is None: return # Abbruch atomic_write_json(files["frage_muster"], muster, indent=1) # Block F: Lern-Artefakte (Karteikarten/Beispiele) aus den Fakten — Bonus, # vom Frontend präsentiert. Bricht den Lauf nicht ab (Artefakte sind optional). sidecar = _json_datei(files["sidecar"]) if _sidecar_schema(sidecar) is not None and not _artefakte_komplett(files): artefakte = await _artefakte_block(ctx, set_p, files, sidecar, instructions) if is_cancelled(): abgebrochen() return if artefakte is None: return # Abbruch (Fehler/Cancel) # DB-Spiegel (Brücke): finalen Sidecar- + Frage-Muster-Stand in die DB schreiben. sidecar = _json_datei(files["sidecar"]) if _sidecar_schema(sidecar) is not None: await _mirror_sidecar_db(topic, sidecar) muster = _json_datei(files["frage_muster"]) if isinstance(muster, dict) and muster: await _mirror_frage_muster_db(topic, muster) # Gliederung (Titel-basiert, robust gegen Nummern-Drift) → DB. plan = _json_datei(files["gliederung"]) if isinstance(plan, dict) and plan.get("kapitel"): kapitel = [ {"titel": ch.get("titel", "Kapitel"), "bausteine": [_titel(entries[n]) for n in ch.get("nummern", []) if n in entries]} for ch in plan["kapitel"] ] await db.set_gliederung(topic, json.dumps({"kapitel": kapitel}, ensure_ascii=False)) # Artefakte → DB (Karteikarte/Beispiel je Sub, Diagramm je Baustein). artefakte = _json_datei(files["artefakte"]) if isinstance(artefakte, dict) and _sidecar_schema(sidecar) is not None: await _mirror_artefakte_db(topic, sidecar, artefakte) except Exception as e: log.exception("[%s] Bausteine-Generierung fehlgeschlagen", topic) _bausteine_errors[topic] = str(e)[:2000] finally: # Kein Datei-Cleanup: Zwischendateien bleiben für Resume bzw. Nachvollziehbarkeit. _bausteine_progress.pop(topic, None) _bausteine_step.pop(topic, None) _bausteine_cancelled.discard(topic) clear_scope(f"bausteine-{topic}-") # Scope leeren → Neustart blockiert nicht