This commit is contained in:
team3
2026-07-13 12:31:25 +02:00
parent 32d7ad9ea1
commit 9cd8e02e22
34 changed files with 747 additions and 606 deletions

View File

@@ -25,6 +25,11 @@ 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."""
@@ -57,9 +62,16 @@ async def _roher_call(key: str, prompt: str, timeout: int, ctx: Kontext,
role=role, capabilities=caps)
haupt = asyncio.ensure_future(einer(key))
if not schwelle:
return await haupt
fertig, _ = await asyncio.wait({haupt}, timeout=schwelle)
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"))
@@ -129,6 +141,9 @@ async def call(ctx: Kontext, *, stage: str, template: str, werte: dict,
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"
@@ -143,6 +158,9 @@ async def call(ctx: Kontext, *, stage: str, template: str, werte: dict,
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:
@@ -159,6 +177,11 @@ async def call(ctx: Kontext, *, stage: str, template: str, werte: dict,
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:
@@ -177,10 +200,15 @@ async def call(ctx: Kontext, *, stage: str, template: str, werte: dict,
async def alle(coros) -> list:
"""gather-Variante, die bei der ersten Exception (LaufPause/Budget) die
Geschwister ABBRICHT — sonst brennen Hunderte laufende Repair-Tasks weiter."""
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():