Files
creator3/backend/korpus.py
2026-07-24 11:23:18 +02:00

230 lines
11 KiB
Python

"""Stufe 1: Recherche-Planung → Soll-Extraktion (Zitat-Gate) → Konsens
(Judge gruppiert NUR, der Code zählt ≥2 unabhängige Quellen) → hartes Gate
mit Sättigungslogik (Nullrunde nur mit neuen Queries gültig, R11)."""
import hashlib
from . import config, db, engine, graph, laden, llm, textkit
def _thema(topic: str) -> str:
t = db.one("SELECT titel, name FROM topics WHERE name=?", topic)
return (t.get("titel") or t["name"]) if t else topic
@engine.worker("korpus.recherche_plan")
async def recherche_plan(task: dict) -> engine.Ergebnis:
"""EIN Task je Runde plant alle Blickwinkel (das Lens-Auffächern gehört
hierher, nicht in die Engine)."""
lenses = config.RECHERCHE_LENSES
seed = []
if task["runde"] == 1:
# Startquellen-Seed: kuratierter Grundstock ist Pflicht, nicht
# Suchglück — einmalig, idempotent; läuft durchs Inhalts-Urteil
seed = [{"knoten": "suche", "item": "r1:seed",
"payload": {"seed": True}}]
bisherige = [db.uj(t["ergebnis"], {}).get("query", "") for t in db.query(
"SELECT ergebnis FROM tasks WHERE topic=? AND knoten='suche'", task["topic"])]
text = await llm.call(
run_id=task["run_id"], stufe="korpus", knoten="recherche_plan",
item=task["item"], role="judge", n=len(lenses),
skill_namen=graph.skills_von(task["knoten"]),
werte={"thema": _thema(task["topic"]),
"lenses": "\n".join(f"- {l}" for l in lenses),
"runde": task["runde"],
"bisherige_queries": "\n".join(f"- {q}" for q in bisherige if q) or "-"})
if text is None:
raise RuntimeError("recherche_plan ohne Ergebnis")
je_lens: dict[str, list] = {}
for block in textkit.bloecke(text, "QUERY"):
q = block.get("text", "").strip()
if q:
je_lens.setdefault(block.get("lens", "").strip(), []).append(q)
neue = []
for lens in lenses:
queries = je_lens.get(lens, [])[:config.QUERIES_JE_LENS]
if not queries: # Blickwinkel fehlt in der Antwort → Retry, nie stumm
raise RuntimeError(f"recherche_plan: Blickwinkel '{lens}' ohne Query")
for q in queries:
h = hashlib.sha256(q.encode()).hexdigest()[:8]
neue.append({"knoten": "suche", "item": f"r{task['runde']}:q:{h}",
"payload": {"query": q}})
return engine.Ergebnis(daten={"queries": len(neue)}, neue_tasks=seed + neue)
@engine.worker("korpus.soll_extraktion", buendel=True)
async def soll_extraktion(task: dict) -> engine.Ergebnis:
mitglieder = [task] + task.get("_buendel", [])
chunks = {} # idx → (mitglied, quelle, chunk)
for i, m in enumerate(mitglieder, 1):
p = db.uj(m["payload"], {})
quelle = db.one("SELECT * FROM quellen WHERE id=?", p["quelle_id"])
text = laden.snapshot_lesen(quelle)
chunks[i] = (m, quelle, text[p["offset"]:p["offset"] + p["chars"]])
antwort = await llm.call(
run_id=task["run_id"], stufe="korpus", knoten="soll_extraktion",
item=task["item"], role="extraktion", n=len(mitglieder),
skill_namen=graph.skills_von(task["knoten"]),
werte={"thema": _thema(task["topic"]),
"abschnitte": "\n\n".join(f"=== ABSCHNITT {i} ===\n{c}"
for i, (_, _, c) in chunks.items())})
if antwort is None:
raise RuntimeError("soll_extraktion ohne Ergebnis")
je_abschnitt: dict[int, list] = {}
for block in textkit.bloecke(antwort, "PUNKT"):
a = block.get("abschnitt", "").strip()
if a.isdigit():
je_abschnitt.setdefault(int(a), []).append(block)
leer = {int(b["abschnitt"]) for b in textkit.bloecke(antwort, "LEER")
if b.get("abschnitt", "").strip().isdigit()}
n_gesamt, teil = 0, {}
for i, (m, quelle, chunk) in chunks.items():
if i not in je_abschnitt and i not in leer:
teil[m["id"]] = "neu" # Abschnitt fehlt in der Antwort → neuer Versuch
continue
# Resume-Idempotenz: Kandidaten dieses Chunks vorher räumen
db.execute("DELETE FROM soll WHERE topic=? AND status='kandidat' AND "
"belege LIKE ?", task["topic"], f'%"chunk": "{m["item"]}"%')
for block in je_abschnitt.get(i, []):
punkt, beleg = block.get("text", ""), block.get("beleg", "")
# Zitat-Gate: Beleg muss wörtlich in der Quelle stehen
if punkt and beleg and textkit.finde_zitat(chunk, beleg):
db.insert("soll", topic=task["topic"], punkt=punkt[:300],
belege=db.j([{"quelle": quelle["id"],
"zitat": beleg[:500], "chunk": m["item"]}]))
n_gesamt += 1
teil[m["id"]] = "fertig"
return engine.Ergebnis(
daten={"punkte": n_gesamt}, teil_status=teil,
neue_tasks=[{"knoten": "soll_konsens", "item": f"r{task['runde']}:konsens"}])
@engine.worker("korpus.soll_konsens")
async def soll_konsens(task: dict) -> engine.Ergebnis:
"""Judge gruppiert, Code zählt. NEU (Lektion 90): Ziel-Band ≈ √Kandidaten —
ohne Band bestätigte creator3 353 Mikro-Punkte für ein Mini-Thema."""
bestaetigt = 0
for _ in range(2): # 2. Pass faltet zu feine Ergebnisse nochmal
punkte = db.query("SELECT * FROM soll WHERE topic=? AND status IN "
"('kandidat','bestaetigt')", task["topic"])
if not punkte:
return engine.Ergebnis(daten={"bestaetigt": 0})
ziel = max(config.SOLL_PUNKTE_MIN,
round(len(punkte) ** 0.5 * config.SOLL_BAND_FAKTOR))
liste = "\n".join(f"[{k['id']}] {k['punkt']}" for k in punkte)
antwort = await llm.call(
run_id=task["run_id"], stufe="korpus", knoten="soll_konsens",
item=f"{task['item']}:n{len(punkte)}", role="judge", n=len(punkte),
skill_namen=graph.skills_von(task["knoten"]),
werte={"punkte": liste, "ziel": ziel})
if antwort is None:
raise RuntimeError("soll_konsens ohne Ergebnis")
bekannt = {k["id"]: k for k in punkte}
bestaetigt = 0
for gruppe in textkit.bloecke(antwort, "GRUPPE"):
ids = [int(i) for i in gruppe.get("ids", "").replace(" ", "").split(",")
if i.isdigit() and int(i) in bekannt]
if not ids:
continue
belege = []
for i in ids:
belege += db.uj(bekannt[i]["belege"], [])
# Vollständigkeit vor Vorsicht: 1 wörtlicher Beleg genügt (das
# Zitat-Gate schützt vor Halluzination; die ≥2-Quellen-Regel
# verwarf echte Facetten — Lektion 104). Dedup/Verdichtung/
# Prüfungen fangen Überschuss hinten ab.
if belege:
kopf = ids[0]
db.update("soll", "id=?", (kopf,), status="bestaetigt",
punkt=gruppe.get("text", bekannt[kopf]["punkt"])[:300],
belege=db.j(belege))
for i in ids[1:]:
db.update("soll", "id=?", (i,), status="gefaltet")
bestaetigt += 1
if bestaetigt <= ziel * config.SOLL_BAND_TOLERANZ:
break # im Band — fertig
return engine.Ergebnis(daten={"bestaetigt": bestaetigt})
@engine.worker("korpus.gate")
async def gate(task: dict) -> engine.Ergebnis:
topic, runde = task["topic"], task["runde"]
quellen = db.one("SELECT COUNT(*) c FROM quellen WHERE topic=? AND "
"status='geladen'", topic)["c"]
bestaetigt = db.one("SELECT COUNT(*) c FROM soll WHERE topic=? AND "
"status='bestaetigt'", topic)["c"]
vorher = db.uj((db.one(
"SELECT ergebnis FROM tasks WHERE topic=? AND knoten='gate_korpus' AND "
"runde=? AND status='fertig'", topic, runde - 1) or {}).get("ergebnis", ""),
{}).get("bestaetigt", -1)
befunde, neue = [], []
if quellen == 0:
befunde.append({"art": "korpus_leer", "item": topic,
"detail": "keine geladene Quelle"})
if bestaetigt < config.SOLL_PUNKTE_MIN:
befunde.append({"art": "soll_zu_wenig", "item": topic,
"detail": f"{bestaetigt} < {config.SOLL_PUNKTE_MIN}"})
runden_max = graph.get().knoten[task["knoten"]].runden_max
gesaettigt = bestaetigt == vorher # Nullrunde (mit frischen Queries gelaufen)
if not gesaettigt and runde < runden_max:
befunde.append({"art": "unsaettigt", "item": topic,
"detail": f"Runde {runde}: {vorher}{bestaetigt}"})
if befunde:
if runde < runden_max:
neue = [{"knoten": "recherche_plan", "item": f"r{runde + 1}:plan",
"runde": runde + 1, "art": "ruecklauf"}]
return engine.Ergebnis(gate_status="rot", befunde=befunde,
neue_tasks=neue,
daten={"bestaetigt": bestaetigt})
# Harte Vollständigkeits-Stops (Lektion 103): Vollständigkeit ist DAS
# Kernziel — fehlt die Basis, ist die Pipeline kaputt → PAUSE statt
# kleiner Guide. Sichtbar als Run-Grund, Resume prüft erneut.
if config.SEED_PFLICHT and not db.one(
"SELECT id FROM quellen WHERE topic=? AND zweck='seed' AND "
"status='geladen'", topic):
raise llm.LaufPause("vollstaendigkeit:seed_fehlt")
if quellen < config.QUELLEN_MIN:
raise llm.LaufPause(
f"vollstaendigkeit:quellen {quellen}<{config.QUELLEN_MIN}")
# Stufe inventar: Extraktion NUR belegnah (Lektion 92) — 2 Reader je
# Fenster um die Soll-Belege. Quellen ohne Beleg liefern keine Atome.
db.update("topics", "name=?", (topic,), status="korpus_fertig")
belege_je_quelle: dict[int, list] = {}
for s in db.query("SELECT belege FROM soll WHERE topic=? AND "
"status='bestaetigt'", topic):
for b in db.uj(s["belege"], []):
belege_je_quelle.setdefault(b["quelle"], []).append(b["zitat"])
for q in db.query("SELECT * FROM quellen WHERE topic=? AND status='geladen' "
"AND zweck IN ('korpus','seed')", topic):
zitate = belege_je_quelle.get(q["id"])
if not zitate:
continue
text = laden.snapshot_lesen(q)
spans = []
for zitat in zitate:
span = textkit.finde_zitat(text, zitat)
if span:
spans.append([max(0, span[0] - config.EXTRAKT_FENSTER),
min(len(text), span[1] + config.EXTRAKT_FENSTER)])
spans.sort()
gemerged = []
for s in spans:
if gemerged and s[0] <= gemerged[-1][1]:
gemerged[-1][1] = max(gemerged[-1][1], s[1])
else:
gemerged.append(s)
i = 0
for von, bis in gemerged:
for o, c in textkit.abschnitte(text[von:bis]):
for r in range(1, config.READER_JE_ABSCHNITT + 1):
neue.append({"knoten": "atom_extraktion",
"item": f"quelle:{q['id']}:a{i}:r{r}",
"payload": {"quelle_id": q["id"],
"offset": von + o, "chars": len(c)}})
i += 1
return engine.Ergebnis(gate_status="gruen", neue_tasks=neue,
daten={"bestaetigt": bestaetigt})