This commit is contained in:
team3
2026-06-17 22:56:24 +02:00
parent 0ff33271a0
commit dce9156ac8
17 changed files with 743 additions and 42 deletions

View File

@@ -12,20 +12,31 @@ 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
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
from paths import arbeit_dir, bausteine_path, project_dir, subbausteine_path
from pipeline import (
CANCELLED, FAILED, GenContext, _extra, _log, _prompt, _race, _rest_schema,
_runde_schema, _semaphore, _str_liste, _timeout, run_single_slot,
_runde_schema, _semaphore, _str_liste, _stufen_schema, _timeout, run_single_slot,
)
from textkit import _eindeutige_titel, _parse_auswahl, _titel_aufloesen, _titel_index, _vormerge
from textkit import (
_eindeutige_titel, _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")
@@ -34,14 +45,31 @@ _bausteine_errors: dict[str, str] = {}
_bausteine_cancelled: set[str] = set()
_bausteine_step: dict[str, int] = {}
BAUSTEINE_STEPS = ("Recherche", "Konsolidierung", "Klärung")
BAUSTEINE_STEPS = (
"Recherche", "Konsolidierung", "Klärung",
"Subbausteine finden", "Subbausteine wählen", "Subbausteine klären",
"Stufen finden", "Stufen wählen", "Stufen klären",
)
def _bausteine_steps(topic: str) -> tuple:
"""Projekte haben einen zusätzlichen Schritt: Themenfeld-Ergänzung per Web-Recherche."""
"""Projekte haben einen zusätzlichen Schritt (Ergänzung), nach der Klärung eingefügt.
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.
"""
base = ("Recherche", "Konsolidierung", "Klärung")
rest = (
"Subbausteine finden", "Subbausteine wählen", "Subbausteine klären",
"Stufen finden", "Stufen wählen", "Stufen klären",
)
if project_dir(topic).is_dir():
return BAUSTEINE_STEPS + ("Ergänzung",)
return BAUSTEINE_STEPS
return base + ("Ergänzung",) + rest
return base + rest
def _step_idx(topic: str, name: str) -> int:
return _bausteine_steps(topic).index(name)
def _bausteine_files(topic: str) -> dict:
@@ -55,14 +83,20 @@ def _bausteine_files(topic: str) -> dict:
"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-*"))) 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,
]
@@ -89,8 +123,12 @@ def _resume_step(topic: str) -> int:
if not geklaert:
return 2
if project_dir(topic).is_dir() and not files["ergaenzung"].exists():
return 3
return len(_bausteine_steps(topic))
return _step_idx(topic, "Ergänzung")
if _sidecar_schema(_json_datei(files["sidecar"])) is not None:
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:
@@ -127,6 +165,7 @@ def active_bausteine() -> list[dict]:
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/
shutil.rmtree(files["arbeit"], ignore_errors=True)
_bausteine_errors.pop(topic, None)
@@ -200,6 +239,307 @@ def _mapping_schema(data):
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 (finale Sidecar)."""
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 _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 project_dir(topic).is_dir() 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 generate_bausteine(topic: str, instructions: str = "", provider: str = DEFAULT_PROVIDER) -> None:
if topic in _bausteine_progress:
return
@@ -228,9 +568,11 @@ async def generate_bausteine(topic: str, instructions: str = "", provider: str =
files["arbeit"].mkdir(parents=True, exist_ok=True)
if project:
await asyncio.to_thread(_pdfs_konvertieren, project)
# „Neu erstellen": fertige Bausteine → kompletter Frischstart.
# Sonst sind Slot-Dateien Reste eines Abbruchs/Fehlers → Resume.
if final_path.exists():
# „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)
@@ -425,6 +767,27 @@ async def generate_bausteine(topic: str, instructions: str = "", provider: str =
# 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)
except Exception as e:
log.exception("[%s] Bausteine-Generierung fehlgeschlagen", topic)
_bausteine_errors[topic] = str(e)[:2000]