"""Task-Engine (R8): DB ist die einzige Wahrheit/Queue. Scheduler blockiert auf einem Wecker-Event + 2-s-Fallback-Poll. Claim per UPDATE…WHERE status='offen'. Resume = derselbe Codepfad (Loop starten, DB lesen); Zombie-Reset beim Start. Ergebnis + status='fertig' + Folgetasks entstehen in EINER Transaktion.""" import asyncio import hashlib import importlib from contextlib import suppress from dataclasses import dataclass, field from . import config, db, graph from .llm import BudgetErschoepft, LaufPause WORKER: dict[str, object] = {} _laeufe: dict[str, asyncio.Task] = {} _wecker: dict[int, asyncio.Event] = {} def worker(name: str, buendel: bool = False): def deco(fn): fn.buendel = buendel # opt-in: Worker verarbeitet t["_buendel"]-Listen WORKER[name] = fn return fn return deco @dataclass class Ergebnis: """Rückgabe jedes Workers. neue_tasks: {knoten, item, payload?, art?, runde?} — Ziel muss über eine Kante vom eigenen Knoten erreichbar sein.""" daten: dict = field(default_factory=dict) neue_tasks: list = field(default_factory=list) befunde: list = field(default_factory=list) # nur Gates: {art, item, detail} gate_status: str = "" # nur Gates: gruen | rot teil_status: dict = field(default_factory=dict) # Bündel: task_id → fertig | neu def module_laden() -> None: """Worker-Module importieren → Registrierung; danach Graph validieren.""" for mod in ("suche", "laden", "korpus", "inventar", "dedup", "struktur", "guide", "redundanz", "montage"): importlib.import_module(f"backend.{mod}") graph.laden(set(WORKER)) def wecken(run_id: int) -> None: ev = _wecker.get(run_id) if ev: ev.set() # ── Lauf-Verwaltung ── def lauf_starten(topic: str, budget: int | None = None) -> int: if topic in _laeufe and not _laeufe[topic].done(): raise RuntimeError("Lauf läuft bereits") run = db.one("SELECT * FROM runs WHERE topic=? AND status IN " "('paused','budget') ORDER BY id DESC LIMIT 1", topic) if run: # Resume: gleicher Run, gleiche Tasks run_id = run["id"] felder = {"status": "running", "grund": ""} if budget is not None: # Budget-Nachschlag beim Fortsetzen felder["budget_tokens"] = budget db.update("runs", "id=?", (run_id,), **felder) else: run_id = db.insert("runs", topic=topic, status="running", gestartet=db.now(), budget_tokens=budget if budget is not None else config.RUN_BUDGET_TOKENS) _laeufe[topic] = asyncio.ensure_future(_lauf(topic, run_id)) return run_id def lauf_pausieren(topic: str) -> None: run = _aktiver_run(topic) if run: db.update("runs", "id=?", (run["id"],), status="paused", grund="manuell") wecken(run["id"]) def lauf_stoppen(topic: str) -> None: run = _aktiver_run(topic) if run: db.update("runs", "id=?", (run["id"],), status="stopped", beendet=db.now()) wecken(run["id"]) task = _laeufe.get(topic) if task and not task.done(): task.cancel() def _aktiver_run(topic: str) -> dict | None: return db.one("SELECT * FROM runs WHERE topic=? AND status='running' " "ORDER BY id DESC LIMIT 1", topic) # ── Kern-Loop ── async def _lauf(topic: str, run_id: int) -> None: g = graph.get() ev = _wecker[run_id] = asyncio.Event() sem_global = asyncio.Semaphore(config.MAX_PARALLEL_GLOBAL) sems = {k.id: asyncio.Semaphore(k.max_parallel) for k in g.knoten.values()} laufende: set[asyncio.Task] = set() # Zombie-Reset: ein Prozess ⇒ beim Start ist garantiert nichts in Arbeit db.execute("UPDATE tasks SET status='offen', fehler='zombie' " "WHERE topic=? AND status='laufend'", topic) _start_tasks(topic, run_id) try: while True: run = db.one("SELECT * FROM runs WHERE id=?", run_id) if run["status"] != "running": break stufe_i = _aktuelle_stufe(g, topic) db.update("runs", "id=? AND stufe!=?", (run_id, g.stufen[stufe_i]), stufe=g.stufen[stufe_i]) _gate_pruefen(g, topic, run_id, stufe_i) gestartet = _dispatch(g, topic, run_id, stufe_i, sems, sem_global, laufende, ev) if not gestartet and not laufende and _leer(topic): _abschliessen(g, topic, run_id) break with suppress(asyncio.TimeoutError): await asyncio.wait_for(ev.wait(), timeout=config.POLL_SEKUNDEN) ev.clear() laufende = {t for t in laufende if not t.done()} # Pause/Stop: Laufende sanft auslaufen lassen if laufende: await asyncio.gather(*laufende, return_exceptions=True) finally: _wecker.pop(run_id, None) def _leer(topic: str) -> bool: return not db.one("SELECT id FROM tasks WHERE topic=? AND " "status IN ('offen','laufend') LIMIT 1", topic) def _aktuelle_stufe(g: graph.Graph, topic: str) -> int: for i, stufe in enumerate(g.stufen): gate = g.gate_der_stufe(stufe) if gate is None: return i gruen = db.one("SELECT id FROM gate_laeufe WHERE topic=? AND knoten=? " "AND status='gruen'", topic, gate.id) if not gruen: return i return len(g.stufen) - 1 def _dispatch(g, topic, run_id, stufe_i, sems, sem_global, laufende, ev) -> int: offen = db.query( "SELECT * FROM tasks WHERE topic=? AND status='offen' ORDER BY runde, id", topic) n = 0 for t in offen: kn = g.knoten.get(t["knoten"]) if kn is None or g.stufen_index(kn.stufe) > stufe_i: continue if sems[kn.id].locked() or sem_global.locked(): continue if kn.barriere and _vorarbeit_offen(g, topic, kn, ausser=t["id"]): continue if kn.typ == "gate" and _vorarbeit_offen(g, topic, kn, ausser=t["id"]): continue if db.execute("UPDATE tasks SET status='laufend', gestartet=? " "WHERE id=? AND status='offen'", db.now(), t["id"]) != 1: continue t = dict(t) t["_buendel"] = _buendel_claimen(g, t, kn) task = asyncio.create_task( _ausfuehren(g, t, kn, sems[kn.id], sem_global, ev)) laufende.add(task) n += 1 return n def _buendel_claimen(g, t: dict, kn: graph.Knoten) -> list[dict]: """Bis zu buendel−1 weitere offene Tasks gleicher (knoten, art) mit-claimen. Nur für Worker mit buendel-Opt-in — sonst blieben Mitglieder 'laufend' hängen.""" fn = WORKER[kn.worker] if kn.buendel <= 1 or not getattr(fn, "buendel", False): return [] mitglieder = [] kandidaten = db.query( "SELECT * FROM tasks WHERE topic=? AND knoten=? AND art=? AND " "status='offen' AND id!=? ORDER BY runde, id LIMIT ?", t["topic"], t["knoten"], t["art"], t["id"], kn.buendel - 1) for k in kandidaten: if db.execute("UPDATE tasks SET status='laufend', gestartet=? " "WHERE id=? AND status='offen'", db.now(), k["id"]) == 1: mitglieder.append(dict(k)) return mitglieder def _vorarbeit_offen(g: graph.Graph, topic: str, kn: graph.Knoten, ausser: int) -> bool: """Barriere/Gate: wartet, bis KEIN offener/laufender Task an Knoten der eigenen oder früherer Stufen existiert (außer sich selbst).""" max_i = g.stufen_index(kn.stufe) ids = [k.id for k in g.knoten.values() if g.stufen_index(k.stufe) <= max_i and k.id != kn.id] qs = ",".join("?" for _ in ids) return bool(db.one( f"SELECT id FROM tasks WHERE topic=? AND status IN ('offen','laufend') " f"AND knoten IN ({qs}) AND id!=? LIMIT 1", topic, *ids, ausser)) def _gate_pruefen(g: graph.Graph, topic: str, run_id: int, stufe_i: int) -> None: """Stufe leer gelaufen → Gate-Task für Runde N anlegen (idempotent).""" gate = g.gate_der_stufe(g.stufen[stufe_i]) if gate is None: return vorhanden = db.one("SELECT id FROM tasks WHERE topic=? AND knoten=? AND " "status IN ('offen','laufend')", topic, gate.id) if vorhanden: return if _vorarbeit_offen(g, topic, gate, ausser=-1): return runde = (db.one("SELECT COUNT(*) c FROM gate_laeufe WHERE topic=? AND knoten=?", topic, gate.id) or {"c": 0})["c"] + 1 if runde > config.GATE_RUNDEN_MAX.get(gate.id, 3): _pausieren(run_id, f"gate_cap:{gate.id}:{runde}") return db.insert("tasks", ignore=True, run_id=run_id, topic=topic, knoten=gate.id, item=f"runde:{runde}", runde=runde, erzeugt_von="engine") async def _ausfuehren(g, t: dict, kn: graph.Knoten, sem, sem_global, ev) -> None: fn = WORKER[kn.worker] alle = [t] + t.get("_buendel", []) try: async with sem_global, sem: erg: Ergebnis = await fn(t) except (LaufPause, BudgetErschoepft) as e: status = "budget" if isinstance(e, BudgetErschoepft) else "paused" for m in alle: db.update("tasks", "id=?", (m["id"],), status="offen", fehler=str(e)[:300]) _pausieren(t["run_id"], f"{status}:{e}", status=status) ev.set() return except asyncio.CancelledError: for m in alle: db.update("tasks", "id=?", (m["id"],), status="offen", fehler="abbruch") raise except Exception as e: # Inhaltsfehler: begrenzte Neuversuche, dann sichtbar for m in alle: _task_neu(m, kn, f"{type(e).__name__}: {e}"[:300]) ev.set() return _abschluss_schreiben(g, t, kn, erg) ev.set() def _task_neu(t: dict, kn: graph.Knoten, fehler: str) -> None: versuch = t["versuch"] + 1 status = "offen" if versuch < kn.max_versuche else "fehler" db.update("tasks", "id=?", (t["id"],), status=status, versuch=versuch, fehler=fehler) def _abschluss_schreiben(g, t: dict, kn: graph.Knoten, erg: Ergebnis) -> None: erlaubt = g.erlaubte_ziele(kn.id) # Bündel-Mitglieder: kaputte Teil-Antworten → neuer Versuch, nie verwerfen for m in t.get("_buendel", []): if erg.teil_status.get(m["id"]) == "neu": _task_neu(m, kn, "buendel_teil_kaputt") else: db.update("tasks", "id=?", (m["id"],), status="fertig", ergebnis=db.j({"buendel_mit": t["id"]}), beendet=db.now()) with db.tx() as c: if erg.teil_status.get(t["id"]) == "neu": c.execute("UPDATE tasks SET status=?, versuch=?, fehler=? WHERE id=?", ("offen" if t["versuch"] + 1 < kn.max_versuche else "fehler", t["versuch"] + 1, "buendel_teil_kaputt", t["id"])) else: c.execute( "UPDATE tasks SET status='fertig', ergebnis=?, beendet=? WHERE id=?", (db.j(erg.daten), db.now(), t["id"])) eingefuegt = 0 for nt in erg.neue_tasks: if nt["knoten"] not in erlaubt: raise graph.GraphFehler( f"{kn.id} darf keine Tasks für {nt['knoten']} erzeugen") cur = c.execute( "INSERT OR IGNORE INTO tasks(run_id, topic, knoten, item, art, " "runde, payload, erzeugt_von) VALUES(?,?,?,?,?,?,?,?)", (t["run_id"], t["topic"], nt["knoten"], nt["item"], nt.get("art", ""), nt.get("runde", t["runde"]), db.j(nt.get("payload", {})), kn.id)) eingefuegt += cur.rowcount if kn.typ == "gate": _gate_abschluss(c, t, kn, erg, eingefuegt) def _gate_abschluss(c, t: dict, kn: graph.Knoten, erg: Ergebnis, neu_eingefuegt: int = 0) -> None: befunde = sorted(f"{b['art']}|{b['item']}|{b.get('detail', '')[:200]}" for b in erg.befunde) fingerprint = hashlib.sha256("\n".join(befunde).encode()).hexdigest()[:16] c.execute("INSERT OR REPLACE INTO gate_laeufe(run_id, topic, knoten, runde, " "status, fingerprint, befunde, ts) VALUES(?,?,?,?,?,?,?,?)", (t["run_id"], t["topic"], kn.id, t["runde"], erg.gate_status, fingerprint, db.j(erg.befunde), db.now())) for b in erg.befunde: # Spiegel-Zeile nur, wenn kein identischer Befund schon offen ist schon = c.execute( "SELECT id FROM befunde WHERE art=? AND item=? AND " "substr(detail,1,200)=? AND status='offen'", (b["art"], b["item"], b.get("detail", "")[:200])).fetchone() if not schon: c.execute("INSERT INTO befunde(run_id, stufe, knoten, art, item, " "detail, runde) VALUES(?,?,?,?,?,?,?)", (t["run_id"], kn.stufe, kn.id, b["art"], b["item"], b.get("detail", "")[:500], t["runde"])) if erg.gate_status == "gruen": # Stufe bestanden → offene Befunde dieser Stufe sind Geschichte c.execute("UPDATE befunde SET status='behoben' WHERE status='offen' AND " "stufe=? AND run_id IN (SELECT id FROM runs WHERE topic=?)", (kn.stufe, t["topic"])) if erg.gate_status == "rot": # Stillstand: identischer Fingerprint n-mal hintereinander UND keine # NEUEN Reparatur-Tasks entstanden (die Eskalations-Kette kann bei # gleicher Befundmenge trotzdem vorankommen: fix → klaerung → neu) alte = [r["fingerprint"] for r in c.execute( "SELECT fingerprint FROM gate_laeufe WHERE topic=? AND knoten=? " "ORDER BY runde DESC LIMIT ?", (t["topic"], kn.id, config.STILLSTAND_N)).fetchall()] if neu_eingefuegt == 0: grund = (f"stillstand:{kn.id}:r{t['runde']}" if len(alte) >= config.STILLSTAND_N and len(set(alte)) == 1 else f"gate_rot_ohne_reparatur:{kn.id}") c.execute("UPDATE runs SET status='paused', grund=? WHERE id=?", (grund, t["run_id"])) def _pausieren(run_id: int, grund: str, status: str = "paused") -> None: db.update("runs", "id=? AND status='running'", (run_id,), status=status, grund=grund[:200]) def _abschliessen(g: graph.Graph, topic: str, run_id: int) -> None: """Kein Task mehr offen: done, wenn alle Gates grün + Montage fertig.""" fertig = db.one("SELECT id FROM tasks WHERE topic=? AND knoten='montage' " "AND status='fertig'", topic) offen_fehler = db.one("SELECT id FROM tasks WHERE topic=? AND status='fehler' " "LIMIT 1", topic) if fertig and not offen_fehler: db.update("runs", "id=?", (run_id,), status="done", beendet=db.now()) else: grund = "tasks_mit_fehler" if offen_fehler else "leerlauf_ohne_montage" _pausieren(run_id, grund) def _start_tasks(topic: str, run_id: int) -> None: """Initiale Tasks der Stufe 1 (idempotent — Resume ist ein No-Op).""" for lens in config.RECHERCHE_LENSES: db.insert("tasks", ignore=True, run_id=run_id, topic=topic, knoten="recherche_plan", item=f"r1:lens:{lens}", runde=1, payload=db.j({"lens": lens}), erzeugt_von="start")