"""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 logging import math import shutil import subprocess import time from pathlib import Path import database as db from agents import kill_process, cancel_scope, clear_scope from config import KONSENS_GRACE, RECHERCHE_GRACE, KONSENS_MAX_RUNDEN, DEFAULT_PROVIDER from fsutil import atomic_write_text, atomic_write_json from jsonio import read_json_file as _json_datei from lernen import FRAGETYPEN 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, GenContext, _extra, _gather_fortschritt, _log, _prompt, _race, _relevanz_schema, _rest_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, _vormerge, ) # 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-Loop: 5 Agenten je Runde, Runden bis volle Crawl-Abdeckung / 0 neue / Zeit-Kappe. RECHERCHE_AGENTEN = 5 RECHERCHE_QUORUM = 3 # je Runde mind. so viele gültige Slot-Ergebnisse RECHERCHE_KAPPE = 1800 # Loop-Gesamt-Sekunden (Notbremse, 30 min) SUBBAUSTEIN_KAPPE = 900 # Subbaustein-Finde-Loop je Chunk (15 min) KONSOLIDIERUNG_CHUNK = 150 # Konsolidierung ab so vielen Kandidaten chunken 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", "Stufen finden", "Stufen wählen", "Stufen klären", "Relevanz finden", "Relevanz wählen", "Relevanz klären", "Fragen finden", "Fragen wählen", "Fragen klären", "Fragen prüfen", ) 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 _STUFEN = ("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"], "typ": r["typ"], "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(), "typ": str(e.get("typ", "")).strip(), "frage": frage} for e in (fm.get(baustein) or []) if isinstance(e, dict) and (frage := str(e.get("frage", "")).strip()) ] 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 "mittel", "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 "mittel", "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") rest = ( "Subbausteine finden", "Subbausteine wählen", "Subbausteine klären", "Stufen finden", "Stufen wählen", "Stufen klären", "Relevanz finden", "Relevanz wählen", "Relevanz klären", "Fragen finden", "Fragen wählen", "Fragen klären", "Fragen prüfen", ) mitte = base + (("Ergänzung",) if q["type"] == "projekt" else ()) + rest return (("Quelle laden",) 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 = ( ("Inventar", ("Quelle laden", "Recherche", "Konsolidierung", "Klärung", "Ergänzung")), ("Subbausteine", ("Subbausteine finden", "Subbausteine wählen", "Subbausteine klären")), ("Stufen", ("Stufen finden", "Stufen wählen", "Stufen klären")), ("Relevanz", ("Relevanz finden", "Relevanz wählen", "Relevanz klären")), ("Fragen", ("Fragen finden", "Fragen wählen", "Fragen klären", "Fragen prüfen")), ) 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", "sidecar": subbausteine_path(topic), "frage_muster": frage_muster_path(topic), } 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("stufe-*")) + list(arbeit.glob("relevanz-*")) + list(arbeit.glob("frage-muster-*"))) 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"], *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 laden") 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; nur noch Frage-Muster offen? if not _frage_muster_komplett(topic): return _step_idx(topic, "Fragen finden") return len(_bausteine_steps(topic)) if _sub_roh_schema(_json_datei(files["sub_roh"])) is None: return _step_idx(topic, "Subbausteine 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: files = _bausteine_files(topic) files["final"].unlink(missing_ok=True) files["sidecar"].unlink(missing_ok=True) # liegt im Themen-Root, nicht in arbeit/ files["frage_muster"].unlink(missing_ok=True) # ebenfalls im Themen-Root quelle_path(topic).unlink(missing_ok=True) shutil.rmtree(quelle_crawl_dir(topic), ignore_errors=True) # gecrawlte Link-Quelle shutil.rmtree(files["arbeit"], ignore_errors=True) _bausteine_errors.pop(topic, None) def _reset_ab_phase(topic: str, phase: int) -> None: """Artefakte AB der groben Phase löschen (1 Inventar … 5 Fragen), frühere behalten. Danach baut der normale Resume-Flow die fehlenden Phasen neu. Kumulativ: Re-Run ab P löscht P…5. quelle.json + Crawl bleiben immer (Quellen-Wahl erhalten).""" files = _bausteine_files(topic) arbeit = files["arbeit"] def glob_del(pat: str) -> None: if arbeit.is_dir(): for p in arbeit.glob(pat): p.unlink(missing_ok=True) if phase <= 5: # Fragen files["frage_muster"].unlink(missing_ok=True) glob_del("frage-muster-*") if phase <= 4: # Relevanz glob_del("relevanz-*") if phase <= 3: # 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 phase <= 2: # Subbausteine files["sub_roh"].unlink(missing_ok=True) glob_del("subbaustein-*") if phase <= 1: # Inventar = kompletter Frischstart (alle Zwischendateien + bausteine.md) 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 _build_recherche_prompt(topic: str, out_path: Path, instructions: str, typ: str, ordner: Path | None, fokus: str = "") -> str: if 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 ("einfach", "mittel", "schwer"): 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_schema(data) -> list[dict] | None: """{"muster": [{subbaustein, typ, frage}, …]} → Liste valider Einträge · sonst None. typ muss ein bekannter Fragetyp sein; subbaustein + frage nicht leer. Leere Liste → None. """ 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): return None sub = str(e.get("subbaustein", "")).strip() typ = str(e.get("typ", "")).strip().casefold() frage = str(e.get("frage", "")).strip() if not sub or typ not in FRAGETYPEN or not frage: return None out.append({"subbaustein": sub, "typ": typ, "frage": frage}) return out or None def _frage_muster_chunk_schema(data) -> list[dict] | None: """{"muster": [{baustein, subbaustein, typ, frage}, …]} → Liste valider Einträge · sonst None. Wie _frage_muster_schema, aber mit `baustein` (Zuordnung im 10er-Chunk). Ungültige Einzel-Einträge werden übersprungen (nicht die ganze Liste verworfen).""" 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() typ = str(e.get("typ", "")).strip().casefold() frage = str(e.get("frage", "")).strip() if not bau or not sub or typ not in FRAGETYPEN or not frage: continue out.append({"baustein": bau, "subbaustein": sub, "typ": typ, "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 _merge_finder(num: int, idx: dict, finder: list[dict]) -> tuple[list[str], list[str]]: """Subbausteine eines Bausteins über die Finder mergen → (konsens ≥2, rest ==1).""" counts: dict[str, int] = {} repr_text: dict[str, str] = {} order: list[str] = [] for d in finder: # die Marker-Liste dieses Finders für genau diesen Baustein subs = next((s for marker, s in d.items() if _titel_aufloesen(idx, marker) == num), []) gesehen: set[str] = set() for sub in subs: key = _norm_titel(sub) if not key or key in gesehen: continue gesehen.add(key) if key not in counts: counts[key], repr_text[key] = 0, sub order.append(key) counts[key] += 1 konsens = [repr_text[k] for k in order if counts[k] >= 2] rest = [repr_text[k] for k in order if counts[k] == 1] return konsens, rest def _final_text(chunk: list[int], entries: dict, daten: dict) -> str: """Marker-Datei aus reinem Konsens (wenn kein Judge nötig).""" teile = [] for num in chunk: konsens = daten[num][0] if konsens: teile.append(f"\n" + "\n".join(f"- {s}" for s in konsens)) return "\n".join(teile) + "\n" def _judge_block(chunk: list[int], entries: dict, daten: dict) -> str: """Eingabe für den Subbaustein-Judge: pro Baustein Konsens + Strittiges.""" lines = [] for num in chunk: konsens, rest = daten[num] lines.append(f"BAUSTEIN: {_titel(entries[num])}") lines.append("Konsens:") lines.extend(f"- {s}" for s in konsens) if not konsens: lines.append("- (noch keiner)") if rest: lines.append("Strittig:") lines.extend(f"- {s}" for s in rest) lines.append("") return "\n".join(lines) 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"] caps = "files" if quelle_ordner(topic) else "full" 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 "" return "\n\nBEREITS GEFUNDEN — bestätige diese Subbausteine erneut UND ergänze fehlende:\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": Judge je Paket bereinigt die Konsens-Liste. async def _klaere(c, chunk): fp = arbeit / f"subbaustein-final-c{c}.md" if _parse_subbausteine(_read(fp)): return bloecke, hat = [], False for num in chunk: subs = [s["sub_titel"] for s in await db.list_subbausteine(topic, norm_by_num[num]) if s["status"] == "konsens"] zeilen = "\n".join(f"- {s}" for s in subs) if subs else "- (keiner)" bloecke.append(f"BAUSTEIN: {titel_by_num[num]}\nKonsens:\n{zeilen}") if subs: hat = True if not hat: return status, _ = await run_single_slot( ctx, f"Subbaustein-Klärung {c}", key=f"bausteine-{topic}-subbaustein-final-c{c}", prompt=_prompt("Subbaustein-Mapping", topic=topic, bausteine="\n\n".join(bloecke), out_path=fp, extra=_extra(instructions)), role="judge", capabilities="files", payload=lambda result, p=fp: _parse_subbausteine(_read(p)) or None, timeout=_timeout("subbaustein_check", len(chunk)), ) if status == FAILED: _log(topic, f"Subbaustein-Klärung Paket {c} fehlgeschlagen — Konsens übernommen") 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") if not roh: _bausteine_errors[topic] = "Keine Subbausteine ermittelt" return None return roh 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"] 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 = _chunk_nums(list(range(len(items))), _n_chunks(len(items), STUFE_CHUNK)) 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 = "\n".join(f"{k}. [{items[j][0]}] {items[j][1]}" for k, j in enumerate(item_idxs, 1)) 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 'mittel'") entsch = entsch if isinstance(entsch, dict) else {} # Strittige ohne Entscheid → 'mittel'; Vote-Gewinner bleiben; Judge überschreibt. ergebnis = {**{k: "mittel" 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 'mittel' 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: "mittel" 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, "mittel")}) sidecar[titel] = lst return sidecar 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"]) 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 = "\n".join(f"{k}. [{items[j][0]}] {items[j][1]}" for k, j in enumerate(item_idxs, 1)) 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, typ, frage}, …]} oder None bei Abbruch.""" topic, is_cancelled = ctx.topic, ctx.is_cancelled arbeit = files["arbeit"] # Je Baustein nur RELEVANTE Subbausteine (relevanz != 'rand'; None zählt als relevant). bausteine = [] for titel, subs in sidecar.items(): rel = [s["titel"] for s in subs if isinstance(s, dict) and s.get("relevanz") != "rand" and str(s.get("titel", "")).strip()] if rel: bausteine.append((titel, rel)) if not bausteine: return {} typen_block = "\n".join(f"- {k}: {v}" for k, v in FRAGETYPEN.items()) chunks = _chunk_nums(list(range(len(bausteine))), _n_chunks(len(bausteine))) # ~10 Bausteine/Chunk 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 block = "\n\n".join( f"BAUSTEIN: {bausteine[i][0]}\nSUBBAUSTEINE:\n" + "\n".join(f"- {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, typen=typen_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 norm = _norm_frage(e["frage"]) seen = gesehen.setdefault(titel, set()) if norm in seen: continue seen.add(norm) out.setdefault(titel, []).append( {"subbaustein": _match_sub(e["subbaustein"], rel_by[titel]), "typ": e["typ"], "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['typ']}] ({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]), "typ": e["typ"], "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 _cited_sources(text: str) -> set[str]: """Zitierte Quellen je Eintrag = 3. ' — '-Segment (URL bzw. Dateiname), kleingeschrieben.""" out = set() for eintrag in _parse_auswahl(text).values(): teile = [t.strip() for t in eintrag.split(" — ")] if len(teile) >= 3 and teile[-1]: out.add(teile[-1].lower()) return out 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 def _abgedeckt(zitiert: set[str], crawl_idx: dict[str, str]) -> set[str]: """Zitierte Quellen → Menge kanonischer Crawl-Seiten-Keys (was nicht passt, fällt weg).""" out = set() for z in zitiert: key = crawl_idx.get(z) or crawl_idx.get(z.rstrip("/")) if key: out.add(key) return out def _fokus_text(bereits: list[str], offen: list[str]) -> str: """Re-Prompt-Block: noch nicht abgedeckte Crawl-Seiten + bereits gefundene Titel.""" teile = [] if offen: liste = "\n".join(f"- {n}" for n in offen[:150]) teile.append( "NOCH NICHT ABGEDECKTE QUELLEN-DATEIEN — lies ZUERST genau diese im Quell-Ordner " f"und ergänze daraus die noch fehlenden Bausteine:\n{liste}" ) if bereits: liste = "\n".join(f"- {t}" for t in bereits) teile.append("BEREITS GEFUNDEN (NICHT wiederholen — liefere nur NEUE Bausteine):\n" + liste) return ("\n\n" + "\n\n".join(teile)) if teile else "" 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) async def _recherche_loop(ctx: GenContext, set_p, files: dict, q: dict, ordner, instructions: str) -> bool: """Füllt DB-Tabelle `bausteine` mit Kandidaten (+ Nennungszähler). Loop: je Runde RECHERCHE_AGENTEN Agenten; re-promptet mit noch offenen Crawl-Seiten, bis alle gelesen / 0 neue / Zeit-Kappe. Resume: fertiger Schritt wird übersprungen. → True (ok) / False (Abbruch).""" 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"] caps = "files" if ordner else "full" crawl_idx = _crawl_index(ordner) alle_seiten = set(crawl_idx.values()) await db.delete_bausteine(topic) await db.delete_coverage(topic) await db.set_step_status(topic, "Recherche", "laufend") set_p("Recherche läuft (Runde 1)…", step=_step_idx(topic, "Recherche")) start = time.monotonic() runde = 0 while not is_cancelled(): runde += 1 vorhandene = {b["titel_norm"] for b in await db.list_bausteine(topic)} gelesen = set(await db.list_coverage(topic)) offen = sorted(alle_seiten - gelesen) if runde > 1 and alle_seiten and not offen: break # volle Abdeckung erreicht fokus = _fokus_text([b["titel"] for b in await db.list_bausteine(topic)], offen) paths = [arbeit / f"recherche-r{runde}-{i}.md" for i in range(1, RECHERCHE_AGENTEN + 1)] for p in paths: p.unlink(missing_ok=True) slots = [{ "key": f"bausteine-{topic}-recherche-r{runde}-{i}", "prompt": _build_recherche_prompt(topic, p, instructions, q["type"], ordner, fokus=fokus), "role": "quick", "capabilities": caps, "payload": (lambda result, p=p: _file_payload(p)), } for i, p in enumerate(paths, 1)] set_p(f"Recherche läuft (Runde {runde})…") texte = await _race(topic, f"Recherche R{runde}", slots, RECHERCHE_QUORUM, _timeout("recherche"), provider, cancelled=is_cancelled, grace=RECHERCHE_GRACE) if is_cancelled(): return False if not texte: if runde == 1: _bausteine_errors[topic] = "Recherche fehlgeschlagen (Minimum nicht erreicht)" return False break # keine neuen Ergebnisse mehr neu, zitiert = 0, set() for text in texte: zitiert |= _cited_sources(text) 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) 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 [] if norm not in vorhandene: neu += 1 vorhandene.add(norm) await db.upsert_baustein(topic, norm, titel, besch, quelle) await db.mark_quellen_gelesen(topic, sorted(_abgedeckt(zitiert, crawl_idx))) gesamt = len(await db.list_bausteine(topic)) deckung = len(await db.list_coverage(topic)) _log(topic, f"Recherche R{runde}: +{neu} neu (gesamt {gesamt}), Abdeckung {deckung}/{len(alle_seiten) or '?'}") if neu == 0: break if time.monotonic() - start > RECHERCHE_KAPPE: _log(topic, f"Recherche: Zeit-Kappe {RECHERCHE_KAPPE}s erreicht (Runde {runde})") break await db.set_step_status(topic, "Recherche", "fertig") return True async def _konsolidiere(ctx: GenContext, set_p, files: dict) -> bool: """Judge mergt Kandidaten semantisch + teilt in Konsens (≥2)/Rest (1×); Status in DB.""" topic, is_cancelled = ctx.topic, ctx.is_cancelled 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 arbeit = files["arbeit"] chunks = _chunk_nums(kandidaten, max(1, math.ceil(len(kandidaten) / KONSOLIDIERUNG_CHUNK))) konsens, rest = [], [] for c, chunk in enumerate(chunks, 1): fp = arbeit / f"konsolidierung-c{c}.json" fp.unlink(missing_ok=True) eintraege = "\n".join( f"{i}. {b['titel']} — {b['beschreibung']} ({b['nennungen']}× genannt)" for i, b in enumerate(chunk, 1) ) status, mapping = await run_single_slot( ctx, f"Konsolidierung {c}", key=f"bausteine-{topic}-konsolidierung-c{c}", prompt=_prompt("Bausteine-Recherche-Mapping", topic=topic, n=RECHERCHE_AGENTEN, eintraege=eintraege, out_path=fp), role="judge", capabilities="files", payload=lambda result, p=fp: _mapping_schema(_json_datei(p)), timeout=_timeout("recherche_mapping", len(chunk)), ) if status == CANCELLED: return False if status == FAILED: _bausteine_errors[topic] = "Recherche-Mapping fehlgeschlagen" return False k, r = mapping konsens += k rest += r # 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: """1 Judge entscheidet über den Rest (1×-Genannte): aufnehmen → Konsens, sonst verworfen.""" topic, is_cancelled = ctx.topic, 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: konsens = [b["titel"] for b in await db.list_bausteine(topic, status="konsens")] fp = files["arbeit"] / "klaerung.json" fp.unlink(missing_ok=True) status, ergebnis = await run_single_slot( ctx, "Klärung", key=f"bausteine-{topic}-klaerung", prompt=_prompt( "Bausteine-Klaerung", topic=topic, konsens="\n".join(f"- {t}" for t in konsens) or "(noch leer)", rest="\n".join(f"- {b['titel']}" for b in rest_rows), final="\n- Entscheide JEDEN Eintrag. `rest` MUSS leer sein.", out_path=fp, ), role="judge", capabilities="files", payload=lambda result, p=fp: _runde_schema(_json_datei(p), final=True), timeout=_timeout("auswahl_mapping", len(rest_rows)), ) if status == CANCELLED: return False if status == FAILED: _bausteine_errors[topic] = "Klärung fehlgeschlagen" return False aufnehmen, _ = ergebnis auf_norm = {_norm_titel(_titel(t)) for t in aufnehmen} for b in rest_rows: await db.set_baustein_status(topic, b["titel_norm"], "konsens" if b["titel_norm"] in auf_norm else "verworfen") await db.set_step_status(topic, "Klärung", "fertig") return True 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 await db.put_subbaustein(topic, bnorm, sn, btitel, st, stufe=s.get("stufe"), relevanz=s.get("relevanz"), status="konsens") async def _mirror_frage_muster_db(topic: str, muster: dict) -> None: """Frage-Muster {Baustein-Titel: [{subbaustein, typ, 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) typ = str(e.get("typ", "")).strip() frage = str(e.get("frage", "")).strip() if not (sn and typ and frage): continue await db.upsert_frage_muster(topic, bnorm, sn, btitel, sub, typ, frage) async def _reset_db_ab_phase(topic: str, ab_phase: int) -> None: """DB-Inhalt der Phasen ≥ ab_phase verwerfen (Inventar=1 … Fragen=5).""" if ab_phase <= 1: await db.delete_bausteine(topic) await db.delete_coverage(topic) await db.delete_pipeline_state(topic, ["Recherche", "Konsolidierung", "Klärung"]) if ab_phase <= 2: await db.delete_subbausteine(topic) if ab_phase <= 5: await db.delete_frage_muster(topic) 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: _reset_ab_phase(topic, ab_phase) await _reset_db_ab_phase(topic, ab_phase) # Link-Quelle: erst crawlen (gleiche Domain, begrenzt) → wird zur Ordner-Quelle. if q["type"] == "link" and not _crawl_fertig(topic): set_p("Quelle laden (Crawl)…", step=_step_idx(topic, "Quelle laden")) n = await asyncio.to_thread(crawl, q["ort"], ordner, cancelled=is_cancelled) if is_cancelled(): abgebrochen() return if not n: _bausteine_errors[topic] = "Crawl ergab keine Inhalte — Link/Domain prüfen" return if ordner: await asyncio.to_thread(_pdfs_konvertieren, ordner) # „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) # Inventar (DB): Recherche-Loop → Konsolidierung → Klärung. if not await _recherche_loop(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 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) sidecar = await _stufen_block(ctx, set_p, files, roh, instructions) if is_cancelled(): abgebrochen() return if sidecar is None: return 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 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) # 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) 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