Files
creator/backend/bausteine.py
2026-06-22 06:07:22 +02:00

1478 lines
68 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Bausteine-Pipeline: Recherche-Konsens + Klärungs-Loop — reines Inventar, unsortiert.
5x Recherche (min. 3, Grace) → Mapping (Konsens/Rest) → Klärungs-Loop (max.
KONSENS_MAX_RUNDEN Runden): 3 Auswahl-Agenten (min. 2, Grace) entscheiden
über den strittigen Rest, ein Mapping-Agent sortiert in aufnehmen/verwerfen/
weiter strittig. Leerer Rest beendet den Loop; die letzte Runde muss alles
entscheiden. Races nutzen ein Grace-Fenster statt „erste N gewinnen": Nach dem
ersten gültigen Ergebnis dürfen die übrigen Agenten KONSENS_GRACE Sekunden
fertig werden. Der Konsens wird im Code akkumuliert — kein Agent re-emittiert
die Gesamtliste.
"""
import asyncio
import json
import logging
import math
import shutil
import subprocess
from pathlib import Path
from agents import kill_process, cancel_scope, clear_scope
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 lernen import FRAGETYPEN
from paths import arbeit_dir, bausteine_path, frage_muster_path, grafik_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",
"Fragen finden", "Fragen wählen", "Fragen klären", "Fragen prüfen",
"Grafik finden", "Grafik prüfen",
)
def lade_quelle(topic: str) -> dict:
"""Persistierte Quellen-Wahl lesen. Fallback (Alt-Themen ohne quelle.json):
existiert projects/<topic> → 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_frage_muster(topic: str, baustein: str) -> list[dict]:
"""Vordefinierte Frage-Muster eines Bausteins aus dem Sidecar (leer = Fallback auf Live)."""
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())
]
def lade_grafiken(topic: str) -> dict:
"""Alle Baustein-Grafiken eines Themas: {Baustein-Titel: {knoten, kanten, richtung}}.
Leere/ungültige Datei → {} (Frontend zeigt Fallback)."""
g = _json_datei(grafik_path(topic))
return g if isinstance(g, dict) else {}
def lade_grafik(topic: str, baustein: str) -> dict | None:
"""Grafik eines Bausteins (None = keine → Fallback)."""
g = lade_grafiken(topic).get(baustein)
return g if isinstance(g, dict) and g.get("knoten") else None
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",
"Fragen finden", "Fragen wählen", "Fragen klären", "Fragen prüfen",
"Grafik finden", "Grafik 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)
# 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")),
("Grafiken", ("Grafik finden", "Grafik 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),
"grafik": grafik_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-*")) + list(arbeit.glob("grafik-*"))) if arbeit.is_dir() else []
return [
*files["recherche"], files["recherche_mapping"],
*(p for slots in files["auswahl"].values() for p in slots),
*files["mapping"].values(), files["ergaenzung"],
files["sub_roh"], files["sidecar"], files["frage_muster"], files["grafik"], *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 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 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")
# Frage-Muster fertig; nur noch Grafiken offen?
if not _grafik_komplett(topic):
return _step_idx(topic, "Grafik 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
files["grafik"].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 <= 6: # Grafiken
files["grafik"].unlink(missing_ok=True)
glob_del("grafik-*")
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)
for p_alt in _alle_slot_dateien(files):
p_alt.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) -> 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 _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_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 _grafik_schema(data) -> dict | None:
"""{"knoten": [{id, text}], "kanten": [{von, nach, text?}], "richtung"?} → bereinigtes
Grafik-Dict · sonst None. Knoten nicht leer, IDs eindeutig, jede Kante referenziert
existierende Knoten (sonst verworfen). Leere Knotenliste → None."""
if not isinstance(data, dict) or not isinstance(data.get("knoten"), list):
return None
knoten, ids = [], set()
for k in data["knoten"]:
if not isinstance(k, dict):
return None
kid = str(k.get("id", "")).strip()
text = str(k.get("text", "")).strip()
if not kid or not text or kid in ids:
continue
ids.add(kid)
knoten.append({"id": kid, "text": text})
if not knoten:
return None
kanten = []
for e in (data.get("kanten") or []):
if not isinstance(e, dict):
continue
von, nach = str(e.get("von", "")).strip(), str(e.get("nach", "")).strip()
if von not in ids or nach not in ids:
continue # hängende Kante (halluzinierte ID) verwerfen
kante = {"von": von, "nach": nach}
if (t := str(e.get("text", "")).strip()):
kante["text"] = t
kanten.append(kante)
richtung = data.get("richtung") if data.get("richtung") in ("TB", "LR") else "TB"
return {"richtung": richtung, "knoten": knoten, "kanten": kanten}
def _grafik_komplett(topic: str) -> bool:
"""Grafik-Sidecar existiert (Build gelaufen)? Einzelne leere Bausteine fallen im
Frontend auf den Fallback-Hinweis zurück — daher genügt die Datei."""
return isinstance(_json_datei(grafik_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"<!-- baustein: {_titel(entries[num])} -->\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
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: drei Phasen — Finden (1 Generator je Baustein, parallel), Wählen (Code:
relevante Subs + Dedup), Klären (Kritiker je Baustein bereinigt die Tabelle).
Generativ statt Vote: Muster sind Text, kein 3-Rater-Konsens sinnvoll.
{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())
def roh_path(c):
return arbeit / f"frage-muster-c{c}.json"
def final_path(c):
return arbeit / f"frage-muster-final-c{c}.json"
# Phase „Fragen finden": pro Baustein 1 Generator, alle parallel.
async def _finde(c, titel, rel):
fp = roh_path(c)
if _frage_muster_schema(_json_datei(fp)):
return # Resume
sub_block = "\n".join(f"- {s}" for s in rel)
status, _ = await run_single_slot(
ctx, f"Frage-Muster {c}",
key=f"bausteine-{topic}-frage-muster-c{c}",
prompt=_prompt("Frage-Muster-Recherche", topic=topic, baustein=titel,
subbausteine=sub_block, typen=typen_block, out_path=fp, extra=_extra(instructions)),
role="fast", capabilities="files",
payload=lambda result, p=fp: _frage_muster_schema(_json_datei(p)),
timeout=_timeout("frage_muster", len(rel)),
)
if status == FAILED:
_log(topic, f"Frage-Muster Baustein {c} fehlgeschlagen — kein Muster (Fallback Live)")
set_p(f"Fragen finden ({len(bausteine)} Bausteine)…", step=_step_idx(topic, "Fragen finden"))
await asyncio.gather(*[_finde(c, t, r) for c, (t, r) in enumerate(bausteine, 1)], return_exceptions=True)
if is_cancelled():
return None
# Phase „Fragen wählen": Code — Dubletten je Baustein raus, Subbaustein-Titel locker
# auf die relevanten mappen (nichts wegen Titel-Abweichung verwerfen).
def _waehle(c, rel):
eintraege = _frage_muster_schema(_json_datei(roh_path(c))) or []
gesehen, sauber = set(), []
for e in eintraege:
norm = _norm_frage(e["frage"])
if norm in gesehen:
continue
gesehen.add(norm)
sauber.append({**e, "subbaustein": _match_sub(e["subbaustein"], rel)})
return sauber
set_p("Fragen wählen…", step=_step_idx(topic, "Fragen wählen"))
roh_by_c = {c: _waehle(c, rel) for c, (titel, rel) in enumerate(bausteine, 1)}
# Phase „Fragen klären": Kritiker je Baustein bereinigt die Tabelle (eindeutig, distinkt).
async def _klaere(c, titel):
roh = roh_by_c[c]
fp = final_path(c)
if _frage_muster_schema(_json_datei(fp)) or not roh:
return
tabelle = "\n".join(f"{i}. [{e['typ']}] ({e['subbaustein']}) {e['frage']}" for i, e in enumerate(roh, 1))
status, _ = await run_single_slot(
ctx, f"Frage-Muster-Klärung {c}",
key=f"bausteine-{topic}-frage-muster-final-c{c}",
prompt=_prompt("Frage-Muster-Kritik", topic=topic, baustein=titel, tabelle=tabelle, out_path=fp, extra=_extra(instructions)),
role="judge", capabilities="files",
payload=lambda result, p=fp: _frage_muster_schema(_json_datei(p)),
timeout=_timeout("frage_muster_check", len(roh)),
)
if status == FAILED:
_log(topic, f"Frage-Muster-Klärung Baustein {c} fehlgeschlagen — Roh-Muster übernommen")
set_p(f"Fragen klären ({len(bausteine)} Bausteine)…", step=_step_idx(topic, "Fragen klären"))
await asyncio.gather(*[_klaere(c, t) for c, (t, _) in enumerate(bausteine, 1)], return_exceptions=True)
if is_cancelled():
return None
# Geklärte Tabelle je Baustein, Fallback auf Roh-Muster. Subbaustein-Titel locker mappen.
def _finalisiere(c, rel):
final = _frage_muster_schema(_json_datei(final_path(c))) or roh_by_c[c]
return [{**e, "subbaustein": _match_sub(e["subbaustein"], rel)} for e in final]
ergebnis = {titel: _finalisiere(c, rel) for c, (titel, rel) in enumerate(bausteine, 1)}
# Phase „Fragen prüfen": Bausteine mit relevanten Subs aber 0 Mustern eine Runde nachholen.
set_p(f"Fragen prüfen ({len(bausteine)} Bausteine)…", step=_step_idx(topic, "Fragen prüfen"))
leer = [(c, t, r) for c, (t, r) in enumerate(bausteine, 1) if not ergebnis.get(t)]
if leer:
_log(topic, f"Frage-Muster: {len(leer)} Baustein(e) ohne Muster — Nachrunde")
for c, t, r in leer:
roh_path(c).unlink(missing_ok=True)
final_path(c).unlink(missing_ok=True)
await asyncio.gather(*[_finde(c, t, r) for c, t, r in leer], return_exceptions=True)
if is_cancelled():
return None
for c, t, r in leer:
roh_by_c[c] = _waehle(c, r)
await asyncio.gather(*[_klaere(c, t) for c, t, r in leer], return_exceptions=True)
if is_cancelled():
return None
for c, t, r in leer:
ergebnis[t] = _finalisiere(c, r)
rest = [t for c, t, r in leer if not ergebnis.get(t)]
if rest:
_log(topic, f"Frage-Muster: {len(rest)} Baustein(e) bleiben leer (Fallback Live): {rest[:5]}")
return ergebnis
async def _grafik_block(ctx: GenContext, set_p, files: dict, sidecar: dict, instructions: str) -> dict | None:
"""Block F: zwei Phasen — Finden (1 Generator je Baustein, parallel) erzeugt eine
Graph-Grafik (Knoten/Kanten) als JSON, Prüfen (Kritiker je Baustein) fixt Struktur +
Inhalt; Bausteine ohne Grafik eine Runde nachholen.
{Baustein-Titel: {knoten, kanten, richtung}} oder None bei Abbruch."""
topic, is_cancelled = ctx.topic, ctx.is_cancelled
arbeit = files["arbeit"]
beschr = {b["titel"]: b["beschreibung"] for b in lade_uebersicht(topic)}
# Je Baustein die relevanten Subbausteine (relevanz != 'rand') als Graph-Grundlage.
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 {}
def roh_path(c):
return arbeit / f"grafik-c{c}.json"
def final_path(c):
return arbeit / f"grafik-final-c{c}.json"
# Phase „Grafik finden": pro Baustein 1 Generator, alle parallel.
async def _finde(c, titel, rel):
fp = roh_path(c)
if _grafik_schema(_json_datei(fp)):
return # Resume
sub_block = "\n".join(f"- {s}" for s in rel)
status, _ = await run_single_slot(
ctx, f"Grafik {c}",
key=f"bausteine-{topic}-grafik-c{c}",
prompt=_prompt("Baustein-Grafik", topic=topic, baustein=titel,
beschreibung=beschr.get(titel, ""), subbausteine=sub_block,
out_path=fp, extra=_extra(instructions)),
role="guide", capabilities="files",
payload=lambda result, p=fp: _grafik_schema(_json_datei(p)),
timeout=_timeout("grafik", len(rel)),
)
if status == FAILED:
_log(topic, f"Grafik Baustein {c} fehlgeschlagen — keine Grafik (Fallback)")
set_p(f"Grafik finden ({len(bausteine)} Bausteine)…", step=_step_idx(topic, "Grafik finden"))
await asyncio.gather(*[_finde(c, t, r) for c, (t, r) in enumerate(bausteine, 1)], return_exceptions=True)
if is_cancelled():
return None
# Phase „Grafik prüfen": Kritiker je Baustein fixt Struktur/Inhalt; sonst Roh übernehmen.
async def _pruefe(c, titel, rel):
roh = _grafik_schema(_json_datei(roh_path(c)))
fp = final_path(c)
if _grafik_schema(_json_datei(fp)) or not roh:
return
sub_block = "\n".join(f"- {s}" for s in rel)
status, _ = await run_single_slot(
ctx, f"Grafik-Prüfung {c}",
key=f"bausteine-{topic}-grafik-final-c{c}",
prompt=_prompt("Baustein-Grafik-Check", topic=topic, baustein=titel,
subbausteine=sub_block,
grafik=json.dumps(roh, ensure_ascii=False, indent=1),
out_path=fp, extra=_extra(instructions)),
role="judge", capabilities="files",
payload=lambda result, p=fp: _grafik_schema(_json_datei(p)),
timeout=_timeout("grafik_check", len(rel)),
)
if status == FAILED:
_log(topic, f"Grafik-Prüfung Baustein {c} fehlgeschlagen — Roh-Grafik übernommen")
set_p(f"Grafik prüfen ({len(bausteine)} Bausteine)…", step=_step_idx(topic, "Grafik prüfen"))
await asyncio.gather(*[_pruefe(c, t, r) for c, (t, r) in enumerate(bausteine, 1)], return_exceptions=True)
if is_cancelled():
return None
def _finalisiere(c):
return _grafik_schema(_json_datei(final_path(c))) or _grafik_schema(_json_datei(roh_path(c)))
ergebnis = {}
for c, (titel, rel) in enumerate(bausteine, 1):
if (g := _finalisiere(c)):
ergebnis[titel] = g
# Bausteine ohne Grafik eine Runde nachholen.
leer = [(c, t, r) for c, (t, r) in enumerate(bausteine, 1) if t not in ergebnis]
if leer:
_log(topic, f"Grafik: {len(leer)} Baustein(e) ohne Grafik — Nachrunde")
for c, t, r in leer:
roh_path(c).unlink(missing_ok=True)
final_path(c).unlink(missing_ok=True)
await asyncio.gather(*[_finde(c, t, r) for c, t, r in leer], return_exceptions=True)
if is_cancelled():
return None
await asyncio.gather(*[_pruefe(c, t, r) for c, t, r in leer], return_exceptions=True)
if is_cancelled():
return None
for c, t, r in leer:
if (g := _finalisiere(c)):
ergebnis[t] = g
rest = [t for c, t, r in leer if t not in ergebnis]
if rest:
_log(topic, f"Grafik: {len(rest)} Baustein(e) bleiben leer (Fallback): {rest[:5]}")
return ergebnis
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)
# 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)
# 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)
# Block E: Frage-Muster je relevantem Subbaustein × Typ → eigenes Sidecar.
# Zur Prüfungszeit zieht jeder Agent ein Muster ohne Zurücklegen und formuliert
# daraus eine Frage — distinkte Saat verhindert die Doppelfragen der Live-Generierung.
sidecar = _json_datei(files["sidecar"])
if _sidecar_schema(sidecar) is not None and _relevanz_komplett(sidecar) and not _frage_muster_komplett(topic):
muster = await _frage_muster_block(ctx, set_p, files, sidecar, instructions)
if is_cancelled():
abgebrochen()
return
if muster is None:
return # Abbruch
atomic_write_json(files["frage_muster"], muster, indent=1)
# Block F: pro Baustein eine Graph-Grafik (Knoten/Kanten) der relevanten
# Subbausteine → eigenes Sidecar. Von allen Guide-Formaten geteilt; das
# Frontend rendert sie mit Auto-Layout (dagre) + KaTeX.
sidecar = _json_datei(files["sidecar"])
if _sidecar_schema(sidecar) is not None and _relevanz_komplett(sidecar) and _frage_muster_komplett(topic) and not _grafik_komplett(topic):
grafik = await _grafik_block(ctx, set_p, files, sidecar, instructions)
if is_cancelled():
abgebrochen()
return
if grafik is None:
return # Abbruch
atomic_write_json(files["grafik"], grafik, 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)
clear_scope(f"bausteine-{topic}-") # Scope leeren → Neustart blockiert nicht