This commit is contained in:
team3
2026-06-29 00:09:24 +02:00
parent 07b6736ced
commit bf7ebe045b
8 changed files with 483 additions and 40 deletions

View File

@@ -21,14 +21,15 @@ import time
from pathlib import Path
import database as db
import embedding
from agents import kill_process, cancel_scope, clear_scope, run_agent
from config import KONSENS_GRACE, RECHERCHE_GRACE, KONSENS_MAX_RUNDEN, DEFAULT_PROVIDER, CRAWL_KEEP_PATTERNS, CRAWL_NOISE_PATTERNS, CRAWL_MIN_CHARS, QUELLE_RELEVANZ_CHUNK, QUELLE_RELEVANZ_SNIPPET
from config import KONSENS_GRACE, RECHERCHE_GRACE, KONSENS_MAX_RUNDEN, DEFAULT_PROVIDER, CRAWL_KEEP_PATTERNS, CRAWL_NOISE_PATTERNS, CRAWL_MIN_CHARS, QUELLE_RELEVANZ_CHUNK, QUELLE_RELEVANZ_SNIPPET, EMBEDDING_AKTIV
from fsutil import atomic_write_text, atomic_write_json
from jsonio import read_json_file as _json_datei
from paths import arbeit_dir, bausteine_path, frage_muster_path, project_dir, subbausteine_path, quelle_path, quelle_crawl_dir, safe_ordner
from crawl import crawl
from pipeline import (
CANCELLED, FAILED, GenContext, _extra, _gather_fortschritt, _janein_schema, _log, _prompt, _race,
CANCELLED, FAILED, OK, GenContext, _extra, _gather_fortschritt, _janein_schema, _log, _prompt, _race,
_relevanz_schema, _runde_schema, _semaphore, _str_liste, _stufen_schema, _timeout, run_single_slot,
)
from textkit import (
@@ -52,7 +53,9 @@ RECHERCHE_KAPPE = 1800 # Sicherheits-Deckel je Batch-Agent
RECHERCHE_ABSCHNITT_ZEICHEN = 12000
# Sichtung (Content/Noise) ist jetzt ein deterministischer Regel-Filter (config.CRAWL_*).
SUBBAUSTEIN_KAPPE = 900 # Subbaustein-Finde-Loop je Chunk (15 min)
KONSOLIDIERUNG_CHUNK = 600 # bis hierher EIN globaler Judge (dedupt alles); darüber chunked + Merge-Pass
KONSOLIDIERUNG_CHUNK = 600 # bis hierher EIN globaler Judge (dedupt alles); darüber chunked + Merge-Pass — nur Fallback-Pfad
DEDUP_MAX_RUNDEN = 3 # finaler Dedup-Pass: max. Iterationen (die kleinere Liste blockt je Runde neu)
DEDUP_MIN_DELTA = 3 # Abbruch, wenn eine Runde weniger als dieses % der Liste entfernt (konvergiert)
# Frage-Muster-Chunks per LPT nach Sub-Last balancieren (Makespan), statt nach Baustein-Anzahl.
FRAGE_CHUNK_SUBS = 50 # Ziel-Summe relevanter Subs je Chunk
FAKTEN_CHUNK_SUBS = 25 # Fakten-Extraktion: kleinere Chunks (Fakten sind umfangreicher als Muster)
@@ -204,7 +207,7 @@ def _bausteine_steps(topic: str) -> tuple:
laufen alle Pakete parallel; der Schritt bleibt, bis das letzte Paket fertig ist.
"""
q = lade_quelle(topic)
base = ("Recherche", "Konsolidierung", "Klärung")
base = ("Recherche", "Konsolidierung", "Klärung", "Dedup")
rest = (
"Subbausteine finden", "Subbausteine wählen", "Subbausteine klären",
"Fakten finden", "Fakten prüfen", "Fakten fix",
@@ -234,7 +237,7 @@ def _melde_p(set_p, topic: str, schritt: str):
# Sonderschritte (Quelle laden, Ergänzung) gehören zur Phase „Inventar".
PHASEN = (
("Quelle", ("Quelle aufbereiten",)),
("Inventar", ("Recherche", "Konsolidierung", "Klärung", "Ergänzung")),
("Inventar", ("Recherche", "Konsolidierung", "Klärung", "Dedup", "Ergänzung")),
("Subbausteine", ("Subbausteine finden", "Subbausteine wählen", "Subbausteine klären")),
("Fakten", ("Fakten finden", "Fakten prüfen", "Fakten fix")),
("Stufen", ("Stufen finden", "Stufen wählen", "Stufen klären")),
@@ -295,7 +298,7 @@ def _alle_slot_dateien(files: dict) -> list[Path]:
dyn = (list(arbeit.glob("subbaustein-*")) + list(arbeit.glob("fakten-*")) + list(arbeit.glob("stufe-*")) + list(arbeit.glob("relevanz-*"))
+ list(arbeit.glob("frage-muster-*")) + list(arbeit.glob("gliederung-*")) + list(arbeit.glob("artefakt-*"))
+ list(arbeit.glob("recherche-*")) + list(arbeit.glob("konsolidierung-*"))
+ list(arbeit.glob("klaerung*"))) if arbeit.is_dir() else []
+ list(arbeit.glob("klaerung*")) + list(arbeit.glob("dedup-*"))) if arbeit.is_dir() else []
return [
*files["recherche"], files["recherche_mapping"],
*(p for slots in files["auswahl"].values() for p in slots),
@@ -1767,18 +1770,18 @@ async def _recherche_batch(ctx: GenContext, set_p, files: dict, q: dict, ordner,
await db.delete_bausteine(topic) # Coverage/inhalt gehört der Sichtung — NICHT löschen
await db.set_step_status(topic, "Recherche", "laufend")
async def _ingest(text: str) -> None:
async def _ingest(reader_id: str, text: str) -> None:
gesehen = set()
for eintrag in _parse_auswahl(text).values():
titel = _titel(eintrag)
norm = _norm_titel(titel)
if not norm or norm in gesehen:
continue
gesehen.add(norm)
gesehen.add(norm) # ein Reader = eine Stimme je Konzept
teile = [t.strip() for t in eintrag.split("")]
besch = teile[1] if len(teile) >= 2 else ""
quelle = [teile[2]] if len(teile) >= 3 and teile[2] else []
await db.upsert_baustein(topic, norm, titel, besch, quelle)
await db.upsert_baustein(topic, norm, titel, besch, quelle, reader=reader_id)
pages = await db.list_content(topic) # von der Sichtung als Content markierte Seiten
if not pages and ordner:
@@ -1795,7 +1798,7 @@ async def _recherche_batch(ctx: GenContext, set_p, files: dict, q: dict, ordner,
"key": f"bausteine-{topic}-recherche-{i}",
"prompt": _build_recherche_prompt(topic, p, instructions, q["type"], ordner),
"role": "quick", "capabilities": caps,
"payload": (lambda result, p=p: _file_payload(p)),
"payload": (lambda result, p=p, rid=f"t{i}": ((rid, t) if (t := _file_payload(p)) else None)),
} for i, p in enumerate(paths, 1)]
texte = await _race(topic, "Recherche", slots, 3, _timeout("recherche"), provider,
cancelled=is_cancelled, grace=RECHERCHE_GRACE)
@@ -1804,8 +1807,8 @@ async def _recherche_batch(ctx: GenContext, set_p, files: dict, q: dict, ordner,
if not texte:
_bausteine_errors[topic] = "Recherche fehlgeschlagen (Minimum nicht erreicht)"
return False
for text in texte:
await _ingest(text)
for rid, text in texte:
await _ingest(rid, text)
await db.set_step_status(topic, "Recherche", "fertig")
return True
@@ -1827,6 +1830,13 @@ async def _recherche_batch(ctx: GenContext, set_p, files: dict, q: dict, ordner,
f"VOLLSTÄNDIG, überspringe nichts. Notiere `{fn}` als Quelle jedes Bausteins. "
f"Suche NICHT im Web — nur dieser Abschnitt zählt.\n\n-----\n{absch}\n-----")
paths = [arbeit / f"recherche-a{ei}-{i}.md" for i in range(1, RECHERCHE_READERS + 1)]
# Reader-Datei-Wiederverwendung: liegen alle Reader-Outputs valide vor (Resume /
# Re-Run ohne Recherche-Änderung), re-ingestieren statt erneut Agenten zu spawnen.
vorhanden = [(f"a{ei}-{i}", t) for i, p in enumerate(paths, 1) if (t := _file_payload(p))]
if len(vorhanden) == len(paths):
for rid, text in vorhanden:
await _ingest(rid, text)
return
for p in paths:
p.unlink(missing_ok=True)
if is_cancelled():
@@ -1835,14 +1845,14 @@ async def _recherche_batch(ctx: GenContext, set_p, files: dict, q: dict, ordner,
"key": f"bausteine-{topic}-recherche-a{ei}-{i}",
"prompt": _build_recherche_prompt(topic, p, instructions, q["type"], ordner, abschnitt=block),
"role": "quick", "capabilities": "files",
"payload": (lambda result, p=p: _file_payload(p)),
"payload": (lambda result, p=p, rid=f"a{ei}-{i}": ((rid, t) if (t := _file_payload(p)) else None)),
} for i, p in enumerate(paths, 1)]
# Quorum 2: beide Reader pro Abschnitt sollen durch (mehr Augen = mehr Konzepte +
# echter Konsens); nach Timeout fällt _race auf das Vorhandene zurück.
texte = await _race(topic, f"Recherche Abschnitt {ei}", slots, 2, _timeout("recherche", 1),
provider, cancelled=is_cancelled, grace=RECHERCHE_GRACE)
for text in (texte or []):
await _ingest(text)
for rid, text in (texte or []):
await _ingest(rid, text)
await _gather_fortschritt([_lese_abschnitt(ei, fn, a) for ei, (fn, a) in enumerate(eintraege, 1)],
len(eintraege), _melde_p(set_p, topic, "Recherche"))
@@ -1873,12 +1883,12 @@ async def _recherche_batch(ctx: GenContext, set_p, files: dict, q: dict, ordner,
"key": f"bausteine-{topic}-recherche-b{bi}-{i}",
"prompt": _build_recherche_prompt(topic, p, instructions, q["type"], ordner, fokus=fokus),
"role": "quick", "capabilities": "files",
"payload": (lambda result, p=p: _file_payload(p)),
"payload": (lambda result, p=p, rid=f"b{bi}-{i}": ((rid, t) if (t := _file_payload(p)) else None)),
} for i, p in enumerate(paths, 1)]
texte = await _race(topic, f"Recherche Batch {bi}", slots, 1, _timeout("recherche", len(batch)),
provider, cancelled=is_cancelled, grace=RECHERCHE_GRACE)
for text in (texte or []):
await _ingest(text)
for rid, text in (texte or []):
await _ingest(rid, text)
await db.mark_quellen_gelesen(topic, batch) # alle zugeteilten Seiten abhaken (auch ohne Treffer)
return not is_cancelled()
@@ -1896,11 +1906,168 @@ async def _recherche_batch(ctx: GenContext, set_p, files: dict, q: dict, ordner,
return True
def _grp_schema(data, ids: set[int]):
"""{"gruppen": [[1,3],[2], …]} → Partition von `ids` als Liste von Index-Gruppen.
Tolerant: ignoriert Fremd-/Doppel-Nummern; vergessene Kandidaten werden eigenständig
(Singleton-Gruppe) ergänzt. None nur bei strukturell kaputtem JSON."""
if not isinstance(data, dict) or not isinstance(data.get("gruppen"), list):
return None
gruppen, gesehen = [], set()
for g in data["gruppen"]:
if not isinstance(g, list):
return None
grp = []
for x in g:
try:
num = int(x)
except (ValueError, TypeError):
continue
if num in ids and num not in gesehen:
gesehen.add(num)
grp.append(num)
if grp:
gruppen.append(grp)
gruppen += [[r] for r in sorted(ids - gesehen)] # vergessene Kandidaten bleiben eigenständig
return gruppen or None
_ASPEKT_MARKER = ("∈ np", "∈np", " in np", "np-schwer", "np-vollständig", "verifizierer",
"zertifikat", "ndtm", "nicht-determ", "lower bound", "untere schranke",
"bzgl", "als sprache")
def _aspekt_marker(titel: str) -> int:
"""Anzahl Eigenschafts-Marker im Titel (∈NP, NP-schwer, Verifizierer, Lower Bound …).
0 = generisches Hauptkonzept (das Problem selbst); >0 = eine Eigenschaft davon."""
t = titel.casefold()
return sum(1 for m in _ASPEKT_MARKER if m in t)
def _canonical(kandidaten: list[dict], idxs: list[int], gesehen_norm: set[str]) -> dict:
"""Repräsentant eines Clusters = das Hauptkonzept (wenigste Eigenschafts-Marker — das Problem
selbst, nicht „… ∈ NP"); Tie → häufigster norm-Titel → meiste Reader. Titel global eindeutig
(Suffix ' (2)'), damit er als Schlüssel taugt."""
by_norm: dict[str, list[int]] = {}
for k in idxs:
by_norm.setdefault(_norm_titel(kandidaten[k]["titel"]), []).append(k)
def gewicht(nb: str):
ms = by_norm[nb]
reader = set().union(*[set(kandidaten[m]["reader"]) for m in ms]) if ms else set()
return (-_aspekt_marker(nb), len(ms), len(reader)) # aspekt-arm zuerst = Hauptkonzept
best = max(by_norm, key=gewicht)
k = max(by_norm[best], key=lambda m: len(kandidaten[m]["beschreibung"]))
titel = kandidaten[k]["titel"]
n = 2
while _norm_titel(titel) in gesehen_norm:
titel = f"{kandidaten[k]['titel']} ({n})"
n += 1
gesehen_norm.add(_norm_titel(titel))
return {"titel": titel, "beschreibung": kandidaten[k]["beschreibung"]}
async def _block_gruppieren(ctx: GenContext, set_p, arbeit: Path, kandidaten: list[dict],
blocks: list[list[int]], praefix: str = "konsolidierung",
schritt: str = "Konsolidierung") -> list[list[int]]:
"""Je Ähnlichkeits-Block gruppiert ein Judge die Titel in die echten Bausteine (merge
Paraphrasen, split Über-Merges). Singletons direkt. Fehler/Timeout → konservativ jeder
Kandidat einzeln (vermeidet fälschliches Über-Mergen). → finale Gruppen (globale Indizes).
`praefix`/`schritt` trennen Konsolidierung und Dedup (Artefakte, Race-Key, Fortschritt)."""
topic, is_cancelled = ctx.topic, ctx.is_cancelled
multi = [(bi, b) for bi, b in enumerate(blocks) if len(b) > 1]
ergebnis: list[list[int]] = [list(b) for b in blocks if len(b) == 1] # Singletons direkt
def _zeile(k: int, g: int) -> str:
b = kandidaten[g]
return f"{k}. {b['titel']}" + (f"{b['beschreibung']}" if b["beschreibung"] else "")
async def _grp(bi: int, block: list[int]) -> None:
ids = set(range(1, len(block) + 1))
p = arbeit / f"{praefix}-block-c{bi}.json"
part = _grp_schema(_json_datei(p), ids)
if part is None: # Resume: gültige Datei nicht neu rechnen
p.unlink(missing_ok=True)
if is_cancelled():
return
zeilen = [_zeile(k, block[k - 1]) for k in range(1, len(block) + 1)]
status, part = await run_single_slot(
ctx, f"Block-Gruppieren {bi}",
key=f"bausteine-{topic}-{praefix}-block-c{bi}",
prompt=_prompt("Bausteine-Block-Gruppieren", topic=topic, eintraege="\n".join(zeilen), out_path=p),
role="judge", capabilities="files",
payload=(lambda result, p=p, ids=ids: _grp_schema(_json_datei(p), ids)),
timeout=_timeout("recherche_mapping", len(block)),
)
part = part if status == OK else None
if part is None: # Judge gescheitert → einzeln (kein Über-Merge)
ergebnis.extend([idx] for idx in block)
else: # lokale Nummern → globale Kandidaten-Indizes
ergebnis.extend([block[k - 1] for k in g] for g in part)
await _gather_fortschritt([_grp(bi, b) for bi, b in multi],
len(multi), _melde_p(set_p, topic, schritt))
return ergebnis
async def _konsolidiere_embedding(ctx: GenContext, set_p, files: dict, kandidaten: list[dict]) -> bool:
"""Zweistufig: Embeddings → grobe Capped-Blocks (High-Recall) → je Multi-Block ein Judge,
der die Titel in die echten Bausteine gruppiert → Reader-Union (≥2 = Konsens)."""
topic, is_cancelled = ctx.topic, ctx.is_cancelled
arbeit = files["arbeit"]
texts = [f"{b['titel']}{b['beschreibung']}" if b["beschreibung"] else b["titel"] for b in kandidaten]
sims = await asyncio.to_thread(embedding.embed_sims, texts)
if sims is None: # Modell doch nicht verfügbar → Fallback
return await _konsolidiere_llm(ctx, set_p, files, kandidaten)
# Stufe 1: grobe Ähnlichkeits-Blocks (gedeckelt, kein Giant-Component).
blocks = await asyncio.to_thread(embedding.capped_blocks, sims, None, None)
# Stufe 2: ein Judge gruppiert JEDEN Multi-Block in die echten Bausteine.
gruppen = await _block_gruppieren(ctx, set_p, arbeit, kandidaten, blocks)
if is_cancelled():
return False
def _min_cos(idxs): # interne Kohärenz zur Kontrolle (Ketten hätten ~0,3)
if len(idxs) < 2:
return 1.0
return round(min(float(sims[i][j]) for n, i in enumerate(idxs) for j in idxs[n + 1:]), 3)
# Konsens = ≥2 distinkte Reader pro Cluster. Legacy-DBs ohne Reader-Tracking (Recherche lief
# vor der Migration, kein Re-Ingest) haben leere Reader-Sets → Rückfall auf Titel-Heuristik
# (sonst landete ALLES im Rest).
hat_reader = any(b["reader"] for b in kandidaten)
konsens, rest, debug, gesehen_norm = [], [], [], set()
for idxs in gruppen:
reader = set().union(*[set(kandidaten[k]["reader"]) for k in idxs]) if idxs else set()
if hat_reader:
score = len(reader)
else: # ohne Reader-Daten: max(Nennungen, Anzahl distinkter Titel-Varianten im Cluster)
score = max(max(kandidaten[k]["nennungen"] for k in idxs),
len({kandidaten[k]["titel_norm"] for k in idxs}))
rep = _canonical(kandidaten, idxs, gesehen_norm)
eintrag = f"{rep['titel']}{rep['beschreibung']}" if rep["beschreibung"] else rep["titel"]
(konsens if score >= 2 else rest).append(eintrag)
debug.append({"titel": rep["titel"], "reader": sorted(reader), "score": score,
"konsens": score >= 2, "min_cos": _min_cos(idxs),
"mitglieder": [kandidaten[k]["titel"] for k in idxs]})
atomic_write_json(arbeit / "konsolidierung-cluster.json", debug, indent=1)
multi_blocks = sum(1 for b in blocks if len(b) > 1)
_log(topic, f"Konsolidierung (Embedding): {len(blocks)} Blocks ({multi_blocks} per LLM gruppiert) "
f"{len(gruppen)} Cluster aus {len(kandidaten)} Kandidaten "
f"{len(konsens)} Konsens / {len(rest)} Rest")
await db.delete_bausteine(topic)
for t in konsens:
await _set_inventar(topic, t, "konsens")
for t in rest:
await _set_inventar(topic, t, "rest")
await db.set_step_status(topic, "Konsolidierung", "fertig")
return True
async def _konsolidiere(ctx: GenContext, set_p, files: dict) -> bool:
"""Panel (KONSOLIDIERUNG_PANEL Judges) mergt Kandidaten semantisch; ein Reconcile-Judge führt die
Panel-Ausgaben zur finalen Konsens (≥2)/Rest (1×)-Liste zusammen. Status in DB.
Panel statt Einzel-Judge: ein einzelner Judge ist bias-anfällig (Position/Verbosity) und instabil."""
topic, provider, is_cancelled = ctx.topic, ctx.provider, ctx.is_cancelled
"""Mergt Roh-Kandidaten zu Konsens (≥2 Reader)/Rest. Deterministisch per Embedding-Clustering;
fehlt das Modell → Rückfall auf das LLM-Panel (`_konsolidiere_llm`). Status in DB."""
topic = ctx.topic
if await db.get_step_status(topic, "Konsolidierung") == "fertig":
return True
set_p("Konsolidiere Recherche…", step=_step_idx(topic, "Konsolidierung"))
@@ -1908,6 +2075,16 @@ async def _konsolidiere(ctx: GenContext, set_p, files: dict) -> bool:
if not kandidaten:
_bausteine_errors[topic] = "Konsolidierung: keine Kandidaten"
return False
if EMBEDDING_AKTIV and await asyncio.to_thread(embedding.verfuegbar):
return await _konsolidiere_embedding(ctx, set_p, files, kandidaten)
return await _konsolidiere_llm(ctx, set_p, files, kandidaten)
async def _konsolidiere_llm(ctx: GenContext, set_p, files: dict, kandidaten: list[dict]) -> bool:
"""Fallback (nur ohne Embedding-Modell): Panel (KONSOLIDIERUNG_PANEL Judges) mergt Kandidaten
semantisch; ein Reconcile-Judge führt die Panel-Ausgaben zur finalen Konsens (≥2)/Rest (1×)-Liste
zusammen. Panel statt Einzel-Judge: ein einzelner Judge ist bias-anfällig und instabil."""
topic, provider, is_cancelled = ctx.topic, ctx.provider, ctx.is_cancelled
arbeit = files["arbeit"]
chunks = _chunk_nums(kandidaten, max(1, math.ceil(len(kandidaten) / KONSOLIDIERUNG_CHUNK)))
@@ -2026,7 +2203,6 @@ async def _klaere_inventar(ctx: GenContext, set_p, files: dict) -> bool:
set_p("Klärung läuft…", step=_step_idx(topic, "Klärung"))
rest_rows = await db.list_bausteine(topic, status="rest")
if rest_rows:
konsens = [b["titel"] for b in await db.list_bausteine(topic, status="konsens")]
arbeit = files["arbeit"]
paths = [arbeit / f"klaerung-j{j}.json" for j in range(1, KONSOLIDIERUNG_PANEL + 1)]
offen = [(j, p) for j, p in enumerate(paths, 1) if _runde_schema(_json_datei(p), final=True) is None]
@@ -2037,7 +2213,6 @@ async def _klaere_inventar(ctx: GenContext, set_p, files: dict) -> bool:
"key": f"bausteine-{topic}-klaerung-j{j}",
"prompt": _prompt(
"Bausteine-Klaerung", topic=topic,
konsens="\n".join(f"- {t}" for t in konsens) or "(noch leer)",
rest="\n".join(f"- {b['titel']}" for b in rest_rows),
final="\n- Entscheide JEDEN Eintrag. `rest` MUSS leer sein.",
out_path=p,
@@ -2066,6 +2241,52 @@ async def _klaere_inventar(ctx: GenContext, set_p, files: dict) -> bool:
return True
async def _dedup_inventar(ctx: GenContext, set_p, files: dict) -> bool:
"""Finaler iterativer Dedup-Pass über die fertige Konsens-Liste: Embedding-Block + LLM-
Gruppierung. Fängt Dubletten, die Konsolidierung (Block-Grenzen, Cap) und Klärung (lange
Liste, Lost-in-the-Middle) übersehen. Pro Gruppe bleibt EIN Repräsentant, der Rest wird
verworfen. Iteriert, weil die kleinere Liste je Runde neu blockt (Cross-Block-Reste)."""
topic, is_cancelled = ctx.topic, ctx.is_cancelled
if await db.get_step_status(topic, "Dedup") == "fertig":
return True
if not (EMBEDDING_AKTIV and await asyncio.to_thread(embedding.verfuegbar)):
await db.set_step_status(topic, "Dedup", "fertig") # ohne Modell: still überspringen
return True
set_p("Dedup…", step=_step_idx(topic, "Dedup"))
arbeit = files["arbeit"]
for runde in range(1, DEDUP_MAX_RUNDEN + 1):
konsens = await db.list_bausteine(topic, status="konsens")
if len(konsens) < 2:
break
texts = [f"{b['titel']}{b['beschreibung']}" if b["beschreibung"] else b["titel"] for b in konsens]
sims = await asyncio.to_thread(embedding.embed_sims, texts)
if sims is None:
break
blocks = await asyncio.to_thread(embedding.capped_blocks, sims, None, None)
gruppen = await _block_gruppieren(ctx, set_p, arbeit, konsens, blocks,
praefix=f"dedup-r{runde}", schritt="Dedup")
if is_cancelled():
return False
weg = 0
for idxs in gruppen:
if len(idxs) < 2:
continue
# Repräsentant = Hauptkonzept (wenigste Eigenschafts-Marker), dann kürzester Titel; Rest verwerfen.
rep = min(idxs, key=lambda k: (_aspekt_marker(konsens[k]["titel"]), len(konsens[k]["titel"]), k))
for k in idxs:
if k != rep:
await db.set_baustein_status(topic, konsens[k]["titel_norm"], "verworfen")
weg += 1
atomic_write_json(arbeit / f"dedup-runde-{runde}.json",
{"vorher": len(konsens), "entfernt": weg,
"gruppen": [[konsens[k]["titel"] for k in g] for g in gruppen if len(g) > 1]}, indent=1)
_log(topic, f"Dedup Runde {runde}: {len(konsens)}{len(konsens) - weg} ({weg})")
if weg <= max(1, len(konsens) * DEDUP_MIN_DELTA // 100): # konvergiert → Schluss
break
await db.set_step_status(topic, "Dedup", "fertig")
return True
# --- Gliederung (Bausteine-Artefakt: Kapitel-Struktur, vom Guide nur gelesen) ---
def _gliederung_komplett(files: dict) -> bool:
@@ -2459,7 +2680,7 @@ async def _reset_db_ab_phase(topic: str, label: str) -> None:
await db.delete_subbausteine(topic)
if idx <= 1: # Inventar: Inventar + Recherche-Schritte — Sichtung bleibt
await db.delete_bausteine(topic)
await db.delete_pipeline_state(topic, ["Recherche", "Konsolidierung", "Klärung"])
await db.delete_pipeline_state(topic, ["Recherche", "Konsolidierung", "Klärung", "Dedup"])
if idx <= 0: # Quelle: Sichtung neu (Coverage/inhalt + Schritt)
await db.delete_coverage(topic)
await db.delete_pipeline_state(topic, ["Quelle aufbereiten"])
@@ -2534,6 +2755,10 @@ async def generate_bausteine(topic: str, instructions: str = "", provider: str =
if is_cancelled():
abgebrochen()
return
if not await _dedup_inventar(ctx, set_p, files):
if is_cancelled():
abgebrochen()
return
konsens_rows = await db.list_bausteine(topic, status="konsens")
entries = {
i: (f"{b['titel']}{b['beschreibung']}" if b["beschreibung"] else b["titel"])

View File

@@ -21,6 +21,23 @@ LESBARKEIT_MAX = 3.5 # Section zu schwer, wenn der Satz-Schnitt darübe
LESBARKEIT_HART = 4.0 # Einzelsatz ab hier „hart"
LESBARKEIT_HART_ANTEIL = 0.30 # … ODER wenn dieser Anteil der Sätze hart ist
# Bausteine-Konsolidierung: semantisches Embedding-Clustering statt LLM-Listen-Merge.
# Ein kleines mehrsprachiges Satz-Embedding (mean-pool) bildet die Kandidaten-Cluster
# GLOBAL (kein Chunk-Verlust) per Cosine + Union-Find. Titel-Varianten desselben Konzepts
# ("Vertex Cover" / "Vertex Cover Definition") verschmelzen; der Konsens zählt danach die
# echten Reader pro Cluster (≥2 = Konsens). Fehlen transformers/torch oder lädt das Modell
# nicht → Embedding stumm aus, `_konsolidiere` fällt auf den alten Panel-Judge-Pfad zurück.
EMBEDDING_AKTIV = True
EMBEDDING_MODELL = "sentence-transformers/paraphrase-multilingual-MiniLM-L12-v2" # CPU, mehrsprachig, ~470 MB
# Stärkere (größere) CPU-Alternative bei Bedarf: "BAAI/bge-m3".
# Konsolidierung = zweistufig: (1) Embedding bildet GROBE Ähnlichkeits-Blocks (High-Recall),
# (2) ein LLM-Judge gruppiert JEDEN Block in die echten Bausteine (merge Paraphrasen, split
# Über-Merges). Reines Threshold-Blocking erzeugt einen Giant-Component (alles verkettet) →
# darum „Capped-Blocking": greedy nach Cosine mergen, aber Blockgröße deckeln. So bleiben die
# LLM-Listen kurz und stabil (belegt: Embedding-Block + LLM-Judge ≈ 95 % Precision).
EMBEDDING_BLOCK_FLOOR = 0.5 # Mindest-Cosine, damit zwei Kandidaten in EINEN Block dürfen
EMBEDDING_BLOCK_CAP = 25 # max. Titel je Block (LLM-Liste kurz/stabil halten)
# Deckel für gleichzeitige CLI-Agenten-Prozesse (über alle Generierungen hinweg).
# Eigene Spur für interaktive Aufrufe (Chat, Elemente), damit sie nicht hinter
# laufenden Writern in der Warteschlange hängen.

View File

@@ -86,6 +86,7 @@ CREATE TABLE IF NOT EXISTS bausteine (
nennungen INTEGER NOT NULL DEFAULT 1,
status TEXT NOT NULL DEFAULT 'kandidat',
quellen TEXT NOT NULL DEFAULT '[]',
reader TEXT NOT NULL DEFAULT '[]',
updated_at TEXT NOT NULL,
PRIMARY KEY (topic, titel_norm)
)
@@ -266,6 +267,10 @@ async def init_db():
await db.execute("ALTER TABLE subbausteine ADD COLUMN fakten TEXT NOT NULL DEFAULT ''") # DEFAULT '' nötig für NOT NULL beim ADD COLUMN
except aiosqlite.OperationalError:
pass
try: # Migration: bausteine.reader (Reader-Set je Kandidat, JSON) — exakte Konsens-Zählung (≥2 Reader)
await db.execute("ALTER TABLE bausteine ADD COLUMN reader TEXT NOT NULL DEFAULT '[]'")
except aiosqlite.OperationalError:
pass
# Migration: alte vertiefungen-Tabelle → baustein_texte (Bestand = lange Form, art 'deepdive')
cursor = await db.execute("SELECT name FROM sqlite_master WHERE type = 'table' AND name = 'vertiefungen'")
if await cursor.fetchone():
@@ -714,15 +719,31 @@ async def delete_baustein_daten(topic: str) -> None:
# --- Bausteine-Pipeline-Inhalt: Inventar / Subbausteine / Frage-Muster / Coverage / State / Quelle ---
async def upsert_baustein(topic: str, titel_norm: str, titel: str, beschreibung: str = "", quellen: list | None = None) -> None:
"""Kandidat einfügen oder Nennungszähler erhöhen. Erst-Beschreibung bleibt erhalten."""
async def upsert_baustein(topic: str, titel_norm: str, titel: str, beschreibung: str = "",
quellen: list | None = None, reader: str | None = None) -> None:
"""Kandidat einfügen oder Reader-Set vereinigen. Erst-Beschreibung bleibt erhalten.
`reader` = ID des Recherche-Readers (z.B. "a5-1"). Die Vereinigung läuft race-frei in
EINEM Statement (json1), weil mehrere Reader-Coroutinen nebenläufig upserten — ein
read-modify-write über `await` würde Mitglieder verlieren. `nennungen` bleibt synchron
zu `len(reader)`. `reader=None` (z.B. aus `_set_inventar`) lässt das Set unverändert."""
db = await get_db()
rid = reader if isinstance(reader, str) and reader else None
await db.execute(
"""INSERT INTO bausteine (topic, titel_norm, titel, beschreibung, nennungen, status, quellen, updated_at)
VALUES (?, ?, ?, ?, 1, 'kandidat', ?, ?)
"""INSERT INTO bausteine (topic, titel_norm, titel, beschreibung, nennungen, status, quellen, reader, updated_at)
VALUES (?, ?, ?, ?, 1, 'kandidat', ?, ?, ?)
ON CONFLICT(topic, titel_norm) DO UPDATE SET
nennungen = nennungen + 1, quellen = excluded.quellen, updated_at = excluded.updated_at""",
(topic, titel_norm, titel, beschreibung, json.dumps(quellen or [], ensure_ascii=False), _now()),
reader = (SELECT json_group_array(v) FROM (
SELECT value AS v FROM json_each(bausteine.reader)
UNION SELECT ? WHERE ? IS NOT NULL)),
nennungen = (SELECT count(*) FROM (
SELECT value AS v FROM json_each(bausteine.reader)
UNION SELECT ? WHERE ? IS NOT NULL)),
quellen = excluded.quellen, updated_at = excluded.updated_at""",
(topic, titel_norm, titel, beschreibung,
json.dumps(quellen or [], ensure_ascii=False),
json.dumps([rid] if rid else [], ensure_ascii=False), _now(),
rid, rid, rid, rid),
)
await db.commit()
@@ -738,6 +759,7 @@ async def list_bausteine(topic: str, status: str | None = None) -> list[dict]:
for row in rows:
d = _row_to_dict(row, cursor)
d["quellen"] = json.loads(d.get("quellen") or "[]")
d["reader"] = json.loads(d.get("reader") or "[]")
out.append(d)
return out

123
backend/embedding.py Normal file
View File

@@ -0,0 +1,123 @@
"""Semantisches Embedding-Clustering für die Baustein-Konsolidierung.
Mean-Pool-Embeddings eines mehrsprachigen Satz-Modells bilden über Cosine-Blocking +
Union-Find GLOBALE Kandidaten-Cluster (kein Chunk-Verlust). Sichere Paare (Ähnlichkeit
≥ HART) werden ohne LLM gemergt; Grenz-Paare im Band [BAND_LOW, HART) gibt der Aufrufer
einem LLM-Judge zur ja/nein-Entscheidung. Fehlen `transformers`/`torch` oder lädt das
Modell nicht → `embed_sims()` liefert `None`, der Aufrufer fällt auf den alten
Panel-Judge-Pfad zurück (silente Deaktivierung, wie das Lesbarkeits-Gate).
CPU genügt; der Aufrufer wrappt die blockierende Inferenz in `asyncio.to_thread`.
`numpy` ist transitiv über torch vorhanden (bewusst nicht in requirements.txt, analog torch).
"""
import logging
import numpy as np
from config import EMBEDDING_AKTIV, EMBEDDING_MODELL, EMBEDDING_BLOCK_FLOOR, EMBEDDING_BLOCK_CAP
log = logging.getLogger("creator.embedding")
_modell_cache = None # (tokenizer, model, torch) — Singleton
_ladeversuch = False # schon versucht zu laden?
EMBEDDING_BATCH = 32 # Inferenz-Batchgröße (CPU)
EMBEDDING_MAX_LEN = 128 # Titel + Kurzbeschreibung sind kurz → kleiner Truncation-Cap genügt
def _modell():
"""Lädt das Modell einmalig. None = Clustering aus (deaktiviert oder Lade-Fehler)."""
global _modell_cache, _ladeversuch
if _ladeversuch:
return _modell_cache
_ladeversuch = True
if not EMBEDDING_AKTIV:
return None
try:
import torch
from transformers import AutoModel, AutoTokenizer
tok = AutoTokenizer.from_pretrained(EMBEDDING_MODELL)
model = AutoModel.from_pretrained(EMBEDDING_MODELL)
model.eval()
_modell_cache = (tok, model, torch)
log.info("Embedding-Modell geladen: %s", EMBEDDING_MODELL)
except Exception as e:
log.warning("Embedding-Clustering deaktiviert (Modell nicht ladbar): %s", e)
_modell_cache = None
return _modell_cache
def verfuegbar() -> bool:
"""True, wenn das Modell geladen werden konnte. Lädt beim ersten Aufruf (blockierend)."""
return _modell() is not None
def embed(texts: list[str]) -> "np.ndarray | None":
"""Texte → (n, d) L2-normalisierte, mean-gepoolte Embeddings. None = Modell aus."""
if _modell() is None:
return None
tok, model, torch = _modell_cache
out = []
for i in range(0, len(texts), EMBEDDING_BATCH):
batch = texts[i:i + EMBEDDING_BATCH]
enc = tok(batch, return_tensors="pt", truncation=True, max_length=EMBEDDING_MAX_LEN, padding=True)
with torch.no_grad():
hidden = model(**enc).last_hidden_state # (b, t, d)
mask = enc["attention_mask"].unsqueeze(-1).type_as(hidden)
vec = (hidden * mask).sum(1) / mask.sum(1).clamp(min=1e-9) # mean-pool ohne Padding
vec = torch.nn.functional.normalize(vec, p=2, dim=1) # L2 → Cosine = Skalarprodukt
out.append(vec.cpu().numpy())
return np.vstack(out).astype(np.float32)
def _find(parent: list[int], x: int) -> int:
while parent[x] != x:
parent[x] = parent[parent[x]] # Pfad-Kompression
x = parent[x]
return x
def _union(parent: list[int], a: int, b: int) -> None:
ra, rb = _find(parent, a), _find(parent, b)
if ra != rb:
parent[max(ra, rb)] = min(ra, rb) # kleinster Index = Wurzel (deterministisch)
def embed_sims(texts: list[str]):
"""Texte → (n, n) Cosine-Matrix · None = Modell nicht verfügbar (Fallback)."""
embs = embed(texts)
if embs is None:
return None
return embs @ embs.T # (n, n) Cosine, float32 (~2 MB bei n=700)
def capped_blocks(sims, floor: float | None = None, cap: int | None = None) -> list[list[int]]:
"""Grobe Ähnlichkeits-Blocks für den LLM — High-Recall, aber Größe gedeckelt.
Greedy: alle Paare mit Cosine ≥ `floor` nach Cosine absteigend; zwei Blocks werden nur
verschmolzen, wenn der resultierende Block ≤ `cap` bleibt. Verhindert den Giant-Component
(reines Threshold-Blocking verkettet sonst fast alles) und hält die LLM-Listen kurz.
→ Liste von Blocks (Index-Listen), jeder Knoten in genau einem Block.
"""
fl = EMBEDDING_BLOCK_FLOOR if floor is None else floor
cp = EMBEDDING_BLOCK_CAP if cap is None else cap
n = len(sims)
parent = list(range(n))
size = [1] * n
if n >= 2:
iu = np.triu_indices(n, k=1)
s = sims[iu]
kept = np.where(s >= fl)[0]
# höchste Cosine zuerst → engste Paare bilden zuerst Blocks
for k in kept[np.argsort(-s[kept])]:
i, j = int(iu[0][k]), int(iu[1][k])
ri, rj = _find(parent, i), _find(parent, j)
if ri != rj and size[ri] + size[rj] <= cp:
_union(parent, i, j)
r = _find(parent, i)
size[r] = size[ri] + size[rj]
blocks: dict[int, list[int]] = {}
for i in range(n):
blocks.setdefault(_find(parent, i), []).append(i)
return list(blocks.values())

View File

@@ -55,6 +55,14 @@ def _titel_aufloesen(idx: dict[str, int], t: str) -> int | None:
return idx.get(_norm_titel(t)) or idx.get(_norm_titel(_titel(t)))
def _norm_dash(s: str) -> str:
"""Space-umgebene Dash-Varianten (en/em/figure/bar/hyphen) → einheitlicher Trenner ''.
Manche Modelle (v.a. nicht-westliche) setzen statt des Em-Dashs einen En-Dash „–"; ohne
Normalisierung scheitert der ` — `-Split komplett und der ganze Eintrag wird zum Titel.
ASCII-Bindestrich „-" bleibt unangetastet (sonst zerlegt es Formeln wie „n - 1")."""
return re.sub(r"\s+[‒–—―‐]\s+", "", s)
def _parse_auswahl(text: str) -> dict[int, str]:
"""Parst eine Baustein-Liste: `N. Titel — Kurzbeschreibung` pro Zeile."""
entries: dict[int, str] = {}
@@ -63,9 +71,9 @@ def _parse_auswahl(text: str) -> dict[int, str]:
m = re.match(r"\s*(\d+)[.)]\s+(.*\S)", line)
if m:
last = int(m.group(1))
entries[last] = m.group(2)
entries[last] = _norm_dash(m.group(2))
elif last is not None and line.strip():
entries[last] += " " + line.strip()
entries[last] += " " + _norm_dash(line.strip())
return entries