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

368 lines
16 KiB
Python
Raw Permalink 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.
"""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 aus dem Graph ableiten (eine Wahrheit: pipeline.yaml),
importieren → Registrierung; danach voll validieren."""
g = graph.laden(None)
for mod in sorted({k.worker.split(".")[0] for k in g.knoten.values()}):
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 buendel1 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 > gate.runden_max:
_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 die letzte Stufe fertig ist."""
terminal = [k.id for k in g.knoten.values() if k.stufe == g.stufen[-1]]
qs = ",".join("?" * len(terminal))
fertig = db.one(f"SELECT id FROM tasks WHERE topic=? AND knoten IN ({qs}) "
"AND status='fertig'", topic, *terminal)
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_abschluss"
_pausieren(run_id, grund)
def _start_tasks(topic: str, run_id: int) -> None:
"""Initiale Tasks: Knoten der Stufe 1 ohne Normal-Vorgänger (aus dem
Graph abgeleitet, keine Knoten-Namen im Code). Idempotent."""
g = graph.get()
ziele = {k["nach"] for k in g.kanten if k["art"] == "normal"}
for kn in g.knoten.values():
if kn.stufe == g.stufen[0] and kn.id not in ziele:
db.insert("tasks", ignore=True, run_id=run_id, topic=topic,
knoten=kn.id, item="r1:start", runde=1,
payload="{}", erzeugt_von="start")