update
This commit is contained in:
@@ -45,9 +45,11 @@ def _ist_infra(err: str) -> bool:
|
||||
|
||||
|
||||
async def _roher_call(key: str, prompt: str, timeout: int, ctx: Kontext,
|
||||
role: str, caps: str) -> agents.AgentErgebnis:
|
||||
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."""
|
||||
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):
|
||||
@@ -59,18 +61,42 @@ async def _roher_call(key: str, prompt: str, timeout: int, ctx: Kontext,
|
||||
return await haupt
|
||||
fertig, _ = await asyncio.wait({haupt}, timeout=schwelle)
|
||||
if fertig:
|
||||
return haupt.result()
|
||||
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:
|
||||
for aufgabe in asyncio.as_completed([haupt, zwilling]):
|
||||
res = await aufgabe
|
||||
if res.ok:
|
||||
return res
|
||||
return res # beide fertig, keins ok → letztes Ergebnis zur Diagnose
|
||||
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 (haupt, zwilling):
|
||||
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,
|
||||
@@ -83,13 +109,23 @@ async def call(ctx: Kontext, *, stage: str, template: str, werte: dict,
|
||||
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)
|
||||
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
|
||||
@@ -147,13 +183,23 @@ async def alle(coros) -> list:
|
||||
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 asyncio.gather(*(call(ctx, **{**kw, "item": f"{item}-p{i + 1}"})
|
||||
for i in range(groesse)))
|
||||
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"})
|
||||
|
||||
Reference in New Issue
Block a user