"""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 from pathlib import Path from agents import kill_process from config import KONSENS_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 paths import arbeit_dir, bausteine_path, project_dir, subbausteine_path, quelle_path, quelle_crawl_dir, safe_ordner from crawl import crawl from pipeline import ( CANCELLED, FAILED, GenContext, _extra, _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 + Stufen entstehen pro Baustein → wie die Writer chunken: # 1 Agent je ~30 Bausteine, gedeckelt. SUBBAUSTEIN_CHUNK = 30 SUBBAUSTEIN_MAX = 20 # Einstufen ist billig (kurzes Urteil, keine Websuche) → größere Pakete, weniger Dateien/Agenten. STUFE_CHUNK = 100 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", ) 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") def subbausteine_titel(topic: str, baustein: str) -> list[str]: """Subbaustein-Titel eines Bausteins aus der Sidecar (leer, wenn keine).""" 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()) ] def lade_uebersicht(topic: str) -> list[dict]: """Strukturierte Baustein-Liste für die Übersicht: Titel + Beschreibung + Subbausteine/Stufen. Verbindet bausteine.md (Nummer/Titel/Beschreibung) mit der Sidecar subbausteine.json (Key = Titel). Fehlt die Sidecar, sind die Subbaustein-Listen leer. """ 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", ) 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 _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), } 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-*"))) 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"], *dyn, ] def cancel_bausteine(topic: str) -> bool: if topic not in _bausteine_progress: return False _bausteine_cancelled.add(topic) kill_process(f"bausteine-{topic}-") return True def _resume_step(topic: str) -> int: """Erster noch offener Schritt anhand der persistierten Zwischendateien.""" files = _bausteine_files(topic) q = lade_quelle(topic) if q["type"] == "link" and not _crawl_fertig(topic): return _step_idx(topic, "Quelle laden") if sum(p.exists() for p in files["recherche"]) < 3: return _step_idx(topic, "Recherche") if not files["recherche_mapping"].exists(): return _step_idx(topic, "Konsolidierung") mapping = _mapping_schema(_json_datei(files["recherche_mapping"])) geklaert = mapping is not None and ( not mapping[1] # kein strittiger Rest or any((r := _runde_schema(_json_datei(p))) is not None and not r[1] for p in files["mapping"].values()) ) if not geklaert: return _step_idx(topic, "Klärung") 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 _relevanz_komplett(sidecar): return len(_bausteine_steps(topic)) return _step_idx(topic, "Relevanz finden") 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: steps = _bausteine_steps(topic) ready = bausteine_path(topic).exists() generating = topic in _bausteine_progress partial = False if generating: current = _bausteine_step.get(topic) states = [ "pending" if current is None else "done" if i < current else "active" if i == current else "pending" for i in range(len(steps)) ] elif ready: states = ["done"] * len(steps) else: nxt = _resume_step(topic) partial = nxt > 0 states = ["done" if i < nxt else "pending" for i in range(len(steps))] return { "ready": ready, "generating": generating, "progress": _bausteine_progress.get(topic), "error": _bausteine_errors.get(topic), "partial": partial, "steps": [{"label": label, "state": s} for label, s in zip(steps, states)], } 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/ 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 _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) -> 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, 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 _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: drei Phasen mit Barriere — Finden, Wählen (Code-Merge), Klären. Pro Phase laufen alle Pakete parallel; der Schritt bleibt, bis das letzte fertig ist. → {Baustein-Titel: [Subbaustein, …]} oder None bei Abbruch/Fehler.""" topic, provider, is_cancelled = ctx.topic, ctx.provider, ctx.is_cancelled arbeit = files["arbeit"] idx = _titel_index(entries) caps = "files" if quelle_ordner(topic) else "full" nums = list(entries) chunks = _chunk_nums(nums, _n_chunks(len(nums))) n = len(chunks) def finder_paths(c): return [arbeit / f"subbaustein-c{c}-{i}.md" for i in (1, 2, 3)] def final_path(c): return arbeit / f"subbaustein-final-c{c}.md" # Phase „Subbausteine finden": pro Paket 3 Finder (min. 2), alle Pakete parallel. async def _finde(c, chunk): paths = finder_paths(c) vorhanden = sum(1 for p in paths if _parse_subbausteine(_read(p))) if vorhanden >= 2: return True zuteilung = "\n".join(f"- {entries[num]}" for num in chunk) offen = [(i, p) for i, p in enumerate(paths, 1) if not _parse_subbausteine(_read(p))] slots = [{ "key": f"bausteine-{topic}-subbaustein-c{c}-{i}", "prompt": _prompt("Subbaustein-Recherche", topic=topic, zuteilung=zuteilung, 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 offen] neu = await _race(topic, f"Subbausteine Paket {c}", slots, 2 - vorhanden, _timeout("subbaustein", len(chunk)), provider, cancelled=is_cancelled, grace=KONSENS_GRACE) return not is_cancelled() and neu is not None set_p(f"Subbausteine finden ({n} Pakete)…", step=_step_idx(topic, "Subbausteine finden")) oks = await asyncio.gather(*[_finde(c, chunk) for c, chunk in enumerate(chunks, 1)], return_exceptions=True) 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": Code-Merge je Paket (instant, kein Agent). set_p(f"Subbausteine wählen ({n} Pakete)…", step=_step_idx(topic, "Subbausteine wählen")) daten_by_c = {} for c, chunk in enumerate(chunks, 1): finder = [d for p in finder_paths(c) if (d := _parse_subbausteine(_read(p)))] daten_by_c[c] = {num: _merge_finder(num, idx, finder) for num in chunk} # Phase „Subbausteine klären": Judge je Paket mit Strittigem, alle parallel. async def _klaere(c, chunk): daten = daten_by_c[c] fp = final_path(c) if _parse_subbausteine(_read(fp)): return if not any(daten[num][1] for num in chunk): atomic_write_text(fp, _final_text(chunk, entries, daten)) 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=_judge_block(chunk, entries, daten), 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 — nur Konsens übernommen") set_p(f"Subbausteine klären ({n} Pakete)…", step=_step_idx(topic, "Subbausteine klären")) await asyncio.gather(*[_klaere(c, chunk) for c, chunk in enumerate(chunks, 1)], return_exceptions=True) if is_cancelled(): return None # Fehlende finale Dateien → Konsens-Fallback; dann alle parsen. roh: dict[str, list[str]] = {} for c, chunk in enumerate(chunks, 1): fp = final_path(c) if not _parse_subbausteine(_read(fp)): atomic_write_text(fp, _final_text(chunk, entries, daten_by_c[c])) for marker, subs in (_parse_subbausteine(_read(fp)) or {}).items(): num = _titel_aufloesen(idx, marker) if num is not None: roh[_titel(entries[num])] = subs 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 set_p(f"Stufen finden ({n} Pakete)…", step=_step_idx(topic, "Stufen finden")) oks = await asyncio.gather(*[_rate(c, idxs) for c, idxs in enumerate(chunks, 1)], return_exceptions=True) 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()} set_p(f"Stufen klären ({n} Pakete)…", step=_step_idx(topic, "Stufen klären")) parts = await asyncio.gather(*[_klaere(c, idxs) for c, idxs in enumerate(chunks, 1)], return_exceptions=True) 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 set_p(f"Relevanz finden ({n} Pakete)…", step=_step_idx(topic, "Relevanz finden")) oks = await asyncio.gather(*[_rate(c, idxs) for c, idxs in enumerate(chunks, 1)], return_exceptions=True) 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()} set_p(f"Relevanz klären ({n} Pakete)…", step=_step_idx(topic, "Relevanz klären")) parts = await asyncio.gather(*[_klaere(c, idxs) for c, idxs in enumerate(chunks, 1)], return_exceptions=True) 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 async def generate_bausteine(topic: str, instructions: str = "", provider: str = DEFAULT_PROVIDER) -> 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) # 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. fertig = 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) # Schritt 1: 5 Recherche-Agenten, min. 3 mit Grace-Fenster — alle gültigen # Slot-Dateien fließen ins Mapping (kein Kappen mehr bei 3) recherchen: list[str] = [] offen = [] for i, path in enumerate(files["recherche"], 1): text = _file_payload(path) if text is not None: recherchen.append(text) else: offen.append((i, path)) vorhanden = len(recherchen) set_p(f"Recherche läuft ({vorhanden} gültig, min. 3)…", step=_step_idx(topic, "Recherche")) if vorhanden < 3: caps = "files" if ordner else "full" slots = [ { "key": f"bausteine-{topic}-recherche-{i}", "prompt": _build_recherche_prompt(topic, path, instructions, q["type"], ordner), "role": "quick", "capabilities": caps, "payload": (lambda result, p=path: _file_payload(p)), } for i, path in offen ] neue = await _race( topic, "Recherche", slots, 3 - vorhanden, _timeout("recherche"), provider, on_update=lambda c: set_p(f"Recherche läuft ({vorhanden + c} gültig, min. 3)…"), cancelled=is_cancelled, grace=KONSENS_GRACE, ) if is_cancelled(): abgebrochen() return if neue is None: _bausteine_errors[topic] = "Recherche fehlgeschlagen (Minimum nicht erreicht)" return recherchen += neue # Schritt 2: Recherche-Mapping — Code-Vormerge (exakte Titel) + 1 Agent # für semantische Dubletten und Konsens/Rest-Teilung (fatal) mapping = _mapping_schema(_json_datei(files["recherche_mapping"])) if mapping is None: set_p("Konsolidiere Recherche…", step=_step_idx(topic, "Konsolidierung")) files["recherche_mapping"].unlink(missing_ok=True) gemergt = _vormerge([_parse_auswahl(t) for t in recherchen]) eintraege = "\n".join(f"{i}. {text} ({n}× genannt)" for i, (text, n) in enumerate(gemergt, 1)) status, mapping = await run_single_slot( ctx, "Recherche-Mapping", key=f"bausteine-{topic}-recherche-mapping", prompt=_prompt( "Bausteine-Recherche-Mapping", topic=topic, n=len(recherchen), eintraege=eintraege, out_path=files["recherche_mapping"], ), role="judge", capabilities="files", payload=lambda result: _mapping_schema(_json_datei(files["recherche_mapping"])), timeout=_timeout("recherche_mapping", len(gemergt)), ) if status == CANCELLED: abgebrochen() return if status == FAILED: _bausteine_errors[topic] = "Recherche-Mapping fehlgeschlagen" return konsens, rest = mapping # Klärungs-Loop: 3 Auswahl-Agenten entscheiden über den Rest, ein # Mapping-Agent sortiert in aufnehmen/verwerfen/weiter strittig. # Leerer Rest beendet den Loop; Runde KONSENS_MAX_RUNDEN muss # alles entscheiden. Der Konsens wächst nur hier im Code. runde = 0 while rest and runde < KONSENS_MAX_RUNDEN: runde += 1 final_runde = runde == KONSENS_MAX_RUNDEN set_p(f"Klärung läuft (Runde {runde}/{KONSENS_MAX_RUNDEN})…", step=_step_idx(topic, "Klärung")) mapping_path = files["mapping"][runde] # Resume: fertiges Runden-Mapping wird direkt übernommen ergebnis = _runde_schema(_json_datei(mapping_path), final=final_runde) if ergebnis is None: mapping_path.unlink(missing_ok=True) konsens_block = "\n".join(f"- {t}" for t in konsens) rest_block = "\n".join(f"- {t}" for t in rest) # 3 Auswahl-Agenten, min. 2 mit Grace-Fenster entscheidungen = [] offen = [] for i, path in enumerate(files["auswahl"][runde], 1): res = _rest_schema(_json_datei(path)) if res is not None: entscheidungen.append(res) else: offen.append((i, path)) if len(entscheidungen) < 2: slots = [ { "key": f"bausteine-{topic}-auswahl-r{runde}-{i}", "prompt": _prompt( "Bausteine-Auswahl", topic=topic, konsens=konsens_block, rest=rest_block, out_path=path, ), "role": "fast", "capabilities": "files", "payload": (lambda result, p=path: _rest_schema(_json_datei(p))), } for i, path in offen ] neue = await _race( topic, f"Auswahl r{runde}", slots, 2 - len(entscheidungen), _timeout("auswahl", len(rest)), provider, cancelled=is_cancelled, grace=KONSENS_GRACE, ) if is_cancelled(): abgebrochen() return if neue is None: _bausteine_errors[topic] = f"Auswahl fehlgeschlagen (Runde {runde}, Minimum nicht erreicht)" return entscheidungen += neue # Votum pro Rest-Eintrag deterministisch zählen indizes = [_titel_index(dict(enumerate(e, 1))) for e in entscheidungen] voten = "\n".join( f"{i}. {text} (von {sum(1 for idx in indizes if _titel_aufloesen(idx, text) is not None)}" f"/{len(entscheidungen)} Agenten übernommen)" for i, text in enumerate(rest, 1) ) final_zusatz = ( "\n- LETZTE RUNDE: Es gibt keine weitere Runde. `rest` MUSS leer sein" " — entscheide JEDEN Eintrag selbst: aufnehmen oder verwerfen." if final_runde else "" ) status, ergebnis = await run_single_slot( ctx, f"Auswahl-Mapping r{runde}", key=f"bausteine-{topic}-auswahl-mapping-r{runde}", prompt=_prompt( "Bausteine-Auswahl-Mapping", topic=topic, n=len(entscheidungen), konsens=konsens_block, rest=voten, final=final_zusatz, out_path=mapping_path, ), role="judge", capabilities="files", payload=lambda result, p=mapping_path, f=final_runde: _runde_schema(_json_datei(p), final=f), timeout=_timeout("auswahl_mapping", len(rest)), ) if status == CANCELLED: abgebrochen() return if status == FAILED: _bausteine_errors[topic] = f"Auswahl-Mapping fehlgeschlagen (Runde {runde})" return aufnehmen, rest = ergebnis _log(topic, f"Klärung Runde {runde}: {len(aufnehmen)} aufgenommen, {len(rest)} weiter strittig") konsens = konsens + aufnehmen entries = {i: t for i, t in enumerate(konsens, 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) 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)