"""Call-Schicht über agents.py: Template-Rendering, JSON-Parsen, Infra-Behandlung, Hedge, Panels, Ledger, Budget. Fehlerklassen bewusst getrennt (Lektion 1/2/68): - Inhaltsfehler (ungültiges JSON, leere Antwort) → begrenzte Restarts, dann None. - Infra-Fehler (429/Timeout/Netz) → eigener Zähler mit Backoff, erschöpft → LaufPause (fail-closed: der Lauf pausiert fortsetzbar, statt mit Teilergebnis weiterzulaufen).""" import asyncio import logging import time import agents import jsonx import ledger import prompts from config import HEDGE_NACH_S, INFRA_BACKOFF_BASE, INFRA_MAX_RETRIES, timeout_fuer log = logging.getLogger("creator2.llm") MAX_RESTARTS = 2 # inhaltliche Neuversuche pro Call class LaufPause(Exception): """Infrastruktur erschöpft — Lauf pausieren, nicht weiterrechnen.""" class ManuellePause(LaufPause): """Nutzer-Pause: Wartende stoppen, Laufende auslaufen lassen (alle() drainiert statt zu canceln), dann endet der Lauf als 'paused' — Resume via Start.""" class Kontext: """Ein Lauf: wandert durch alle Ebenen, trägt Ledger-Zuordnung.""" def __init__(self, run_id: int, topic: str, provider: str): self.run_id = run_id self.topic = topic self.provider = provider self.ebene = "" _INFRA_MARKER = ("429", "rate limit", "rate_limit", "timeout", "timed out", "connection", "network", "overloaded", "unavailable", "database is locked") def _ist_infra(err: str) -> bool: e = err.lower() return any(m in e for m in _INFRA_MARKER) async def _roher_call(key: str, prompt: str, timeout: int, ctx: Kontext, role: str, caps: str, log_hedge=None) -> agents.AgentErgebnis: """Ein Agent-Call mit Stall-Hedge: läuft der Slot max(HEDGE_NACH_S, timeout/2) ohne Ergebnis, startet genau EIN Zwilling; das erste valide Ergebnis gewinnt. Der zweite (Verlierer-)Call wird über `log_hedge(res|None)` als eigene Ledger-Zeile verbucht — sonst fehlen seine Tokens in Budget und Kennzahlen.""" schwelle = max(HEDGE_NACH_S, timeout / 2) if HEDGE_NACH_S > 0 else 0 async def einer(k: str): return await agents.run_agent(k, prompt, timeout, provider=ctx.provider, role=role, capabilities=caps) haupt = asyncio.ensure_future(einer(key)) try: if not schwelle: return await haupt fertig, _ = await asyncio.wait({haupt}, timeout=schwelle) except asyncio.CancelledError: # Stop cancelt nur den WARTENDEN — der ensure_future-Task liefe sonst # verwaist weiter (aak: 10 Extraktions-Calls überlebten Stop, hielten # API-Slots und zogen Tokens, bis der Timeout sie erlöste) haupt.cancel() raise if fertig: return haupt.result() # kein Zwilling gestartet → kein Hedge-Log zwilling = asyncio.ensure_future(einer(f"{key}-h")) paar = {haupt, zwilling} gewinner = None sonstige: list[agents.AgentErgebnis] = [] # fertig, aber nicht Gewinner letzter_fehler = None try: pending = set(paar) while pending and gewinner is None: done, pending = await asyncio.wait(pending, return_when=asyncio.FIRST_COMPLETED) for t in done: try: res = t.result() except asyncio.TimeoutError as e: letzter_fehler = e # ein Timeout bricht den Hedge NICHT ab — continue # der Zwilling darf noch liefern if res.ok and gewinner is None: gewinner = res else: sonstige.append(res) finally: for t in paar: if not t.done(): t.cancel() if log_hedge is not None: if gewinner is not None: # der Verlierer: fertig (Tokens bekannt) oder gecancelt (0 Tokens) log_hedge(sonstige[0] if sonstige else None) elif len(sonstige) >= 2: # kein Gewinner: call() loggt das erste als Hauptzeile, das zweite hier log_hedge(sonstige[1]) if gewinner is not None: return gewinner if sonstige: return sonstige[0] # keins ok → erstes Ergebnis zur Diagnose raise letzter_fehler if letzter_fehler else asyncio.TimeoutError() async def call(ctx: Kontext, *, stage: str, template: str, werte: dict, role: str = "judge", caps: str = "none", schritt: str | None = None, n: int = 0, item: str = "", erwartet=dict): """Ein LLM-Call → geparstes JSON (erwartet: dict|list) oder roher Text (erwartet=str). None nur nach erschöpften inhaltlichen Restarts. Ledger + Budget immer.""" prompt, thash = prompts.render(template, **werte) timeout = timeout_fuer(schritt or stage, n) infra_rest = INFRA_MAX_RETRIES inhalt_rest = MAX_RESTARTS versuch = 0 def _log_hedge(res) -> None: # Verlierer-/Zweitcall des Hedge als eigene Ledger-Zeile (Budget zählt mit) ledger.log_call(ctx.run_id, ebene=ctx.ebene, stage=stage, item=f"{item}-hedge", template=template, template_hash=thash, role=role, provider=ctx.provider, model=(res.model if res else ""), status="hedge" if res else "hedge_cancel", tokens=(res.tokens if res else None)) while True: versuch += 1 ledger.budget_pruefen(ctx.run_id) # Vorab: nicht erst NACH dem Call stoppen key = f"{ctx.topic}-{ctx.ebene}-{stage}-{item or 'x'}-{versuch}" start = time.monotonic() status, tokens, err, wait_ms, model = "error", None, "", 0, "" try: res = await _roher_call(key, prompt, timeout, ctx, role, caps, _log_hedge) tokens, err = res.tokens, res.err wait_ms = int(res.wait_s * 1000) model = res.model if res.err == "pausiert": # Nutzer-Pause: kein Retry, Lauf sauber beenden status = "pause" raise ManuellePause("manuell pausiert") if res.ok: if erwartet is str: status = "ok" return res.text daten = jsonx.parse(res.text) # Modelle lassen bei EINEM Element die Array-Klammern gern weg — # ein valides Einzel-Dict zählt als 1-elementige Liste (sonst # brannte jeder 1-Item-Chunk 3 Neuversuche durch). if erwartet is list and isinstance(daten, dict) and daten: daten = [daten] # leeres {} bleibt „falscher Typ" (= Ausfall) if daten is not None and isinstance(daten, erwartet): status = "ok" return daten status, err = "parse", f"unparsbar/falscher Typ: {res.text[:200]}" elif "stop=max_tokens" in res.err: status = "cap" # Output-Cap: gleicher Prompt → gleiches Cap, err = res.err # Neuversuch ist deterministisch sinnlos elif _ist_infra(res.err): status = "infra" else: status = "error" err = res.err except asyncio.TimeoutError: status = "infra" err = "timeout" finally: ledger.log_call(ctx.run_id, ebene=ctx.ebene, stage=stage, item=item, template=template, template_hash=thash, role=role, provider=ctx.provider, model=model, status=status, dur_ms=int((time.monotonic() - start) * 1000), wait_ms=wait_ms, tokens=tokens, meta={"err": err[:300]} if err else None) ledger.budget_pruefen(ctx.run_id) if status == "cap": # sofort aufgeben — der Aufrufer reagiert strukturell (Extraktion # halbiert den Chunk); 2 Retries verbrannten sonst ~64k Tokens je Fall log.warning("%s: Output-Cap (stop=max_tokens) — kein Neuversuch", key) return None if status == "infra": infra_rest -= 1 if infra_rest < 0: raise LaufPause(f"{stage}: Infrastruktur erschöpft ({err})") pause = INFRA_BACKOFF_BASE * 2 ** (INFRA_MAX_RETRIES - infra_rest - 1) agents.drossel_melden(pause) # global: auch Timeouts/CLI-429 bremsen alle log.warning("%s: Infra-Fehler (%s), Pause %.0fs", key, err[:80], pause) await asyncio.sleep(pause) continue inhalt_rest -= 1 if inhalt_rest < 0: log.warning("%s: inhaltlich erschöpft (%s)", key, err[:120]) return None log.info("%s: Neuversuch (%s)", key, err[:80]) async def alle(coros) -> list: """gather-Variante, die bei der ersten Exception (LaufPause/Budget) die Geschwister ABBRICHT — sonst brennen Hunderte laufende Repair-Tasks weiter. Ausnahme ManuellePause: sanft — Laufende zu Ende laufen lassen (Ergebnisse landen noch in der DB), Wartende liefern selbst „pausiert".""" tasks = [asyncio.ensure_future(c) for c in coros] try: return await asyncio.gather(*tasks) except ManuellePause: await asyncio.gather(*tasks, return_exceptions=True) raise except BaseException: for t in tasks: if not t.done(): t.cancel() raise def einstimmig(stimmen: list, n: int, urteil) -> bool: """Destruktiver Entscheid (verwerfen/mergen/freisprechen) nur bei VOLLZÄHLIGEM Panel: mindestens n gültige Stimmen, alle True. `urteil(s)` liefert True/False oder None (Ausfall/Parse/fehlendes Feld = keine Stimme → kein Entscheid). Falsch-Verwerfen >> ein überlebendes Item (Lektion 67). Muster von `chunk_urteil` (inventar) verallgemeinert.""" gueltig = [u for u in (urteil(s) for s in stimmen) if u is not None] return len(gueltig) >= n and all(gueltig) async def panel(ctx: Kontext, groesse: int, **kw) -> list: """`groesse` unabhängige Calls, alle Ergebnisse (None-gefiltert). Fällt genau eine Stimme eines 2er-Panels aus, ersetzt EIN Ersatz-Richter (Lektion 20), dann entscheidet der Aufrufer über Konsens.""" item = kw.get("item", "") stimmen = await alle(call(ctx, **{**kw, "item": f"{item}-p{i + 1}"}) for i in range(groesse)) gueltig = [s for s in stimmen if s is not None] if len(gueltig) == groesse - 1 and groesse >= 2: ersatz = await call(ctx, **{**kw, "item": f"{item}-pE"}) if ersatz is not None: gueltig.append(ersatz) return gueltig