fix(eu): Kapazitaetsengpaesse bei Bedrock aussitzen statt Schritte verlieren
Am 01.08.2026 lieferte Bedrock waehrend eines Laufs ServiceUnavailableException. Botocore versuchte es dreimal in fuenf Sekunden und gab auf, die Web-Source-Selektion fiel ersatzlos aus. Im Bericht war davon nichts zu sehen, der Lauf galt als vollstaendig. Gemessen sind es zwei betroffene Aufrufe bei rund 340, also 0,6 Prozent, aber jeder Vorfall kostet einen ganzen Pipeline-Schritt. Der Fehler ist auch keiner unserer Quote. ServiceUnavailableException ist ein 503, AWS hat also gerade keine Kapazitaet fuer das Modell. Unsere Meldung sagte trotzdem "Rate-Limit", weil der Code beide Faelle in einen Topf warf. Aendert sich damit: 1. Geduldiges Wiederholen. Kapazitaets- (503), Kontingent- (429) und Verbindungsfehler werden nach 5, 15 und 30 Sekunden erneut versucht. Diese Dellen dauern typischerweise unter einer Minute, unsere bisherigen fuenf Sekunden lagen genau im ungeguenstigsten Fenster. Stufen ueber BEDROCK_RETRY_WAITS einstellbar. 2. Zwei Bremsen gegen lange Laeufe. Gewartet wird nur, wenn die Wartezeit ins Zeitbudget des Aufrufs passt (die Wartezeit zaehlt gegen dasselbe asyncio-Budget, und die Haiku-Planungsaufrufe haben nur 120 Sekunden), und nur solange das Wartebudget des Refreshs reicht (BEDROCK_RETRY_BUDGET_S, Vorgabe 120 Sekunden). Ein Lauf kann sich dadurch um hoechstens zwei Minuten verlaengern. 3. Ursachen werden getrennt benannt. Im Log steht jetzt "hat keine freie Kapazitaet" oder "meldet ausgeschoepftes Kontingent" statt pauschal "Rate-Limit". Nach aussen bleibt beides die Kategorie rate_limit, damit die Retry-Steuerung des Orchestrators unveraendert greift. 4. Stoerungen werden sichtbar. Waehrend einer Wartezeit meldet die Oberflaeche "KI-Dienst hat keine freie Kapazitaet, neuer Versuch in 15 Sekunden", sonst sieht die stehende Anzeige wie ein Haenger aus. Nach dem Lauf steht eine Zusammenfassung im Refresh-Protokoll, auch wenn der Lauf sonst glatt durchlief. 5. botocore laeuft im Modus 'adaptive' statt 'standard'. Es bremst sich bei Drosselung selbst ein und deckt die Sekunden ab, unsere Schleife die Zehnersekunden. Zeitliche Wirkung, gemessen am EU-Lauf der Ceuta-Lage (252 Sekunden, 15 Bedrock-Aufrufe): im stoerungsfreien Fall null, im betroffenen Lauf rund eine Minute mehr. Bei 0,6 Prozent Fehlerrate je Aufruf trifft das etwa jeden elften Lauf, im Mittel ueber alle Laeufe rund zwei Prozent. Gegenprobe mit einem echten Bedrock-Aufruf ueber den geaenderten Pfad: Antwort, Verbrauch und Kosten unveraendert, keine Stoerung protokolliert. Neu sind 23 Pruefungen, insgesamt laufen 179 ohne Netzzugriff und ohne Kosten. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Dieser Commit ist enthalten in:
@@ -98,7 +98,7 @@ src/:
|
||||
entity_extractor.py: "Netzwerkanalyse (Sonnet). ACHTUNG: aktuell ohne Aufrufer in der App, nur regenerate_relations.py im Repo-Root nutzt es"
|
||||
translator.py: "Haiku-Uebersetzung in Batches a 5 (groessere Batches rissen den JSON-Output ab)"
|
||||
claude_client.py: "Shared Claude CLI Client, Usage-Tracking (Token, Kosten), Rate-Limit-Erkennung, Cancel via ContextVar. Routet je nach _ai_backend_var werkzeuglose Aufrufe zu bedrock_client"
|
||||
bedrock_client.py: "EU-Modellweg ueber AWS Bedrock Converse (nur eu.-Inferenzprofile, Token-zu-USD-Umrechnung ueber BEDROCK_PRICING, gleiche Fehlerkategorien wie das CLI). EU-Umbau Phase 1"
|
||||
bedrock_client.py: "EU-Modellweg ueber AWS Bedrock Converse (nur eu.-Inferenzprofile, Token-zu-USD-Umrechnung ueber BEDROCK_PRICING, gleiche Fehlerkategorien wie das CLI). EU-Umbau Phase 1. Wiederholt Kapazitaets- (503) und Kontingentfehler (429) nach BEDROCK_RETRY_WAITS (5/15/30s), begrenzt durch das Zeitbudget des Aufrufs und BEDROCK_RETRY_BUDGET_S je Refresh; Stoerungen landen im Refresh-Protokoll (refresh_log.error_message) und waehrend der Wartezeit als status_update in der Oberflaeche"
|
||||
eu_researcher.py: "EU-Recherche-Schleife (Phase 2): Bedrock-Modell schlaegt Suchanfragen vor, Code fragt staan, Modell entscheidet weiter/fertig. Bildet die 4-Phasen-Tiefenrecherche nach, liefert denselben JSON-Array-Text wie die CLI-WebSearch (Einstieg in researcher.search), akzeptiert NUR URLs aus echten staan-Treffern"
|
||||
|
||||
feeds/:
|
||||
@@ -165,6 +165,7 @@ tests/:
|
||||
test_qc_und_runden.py: "Absicherung der Duplikatpruefung und Abbruch der Suchrunden bei erreichter Treffergrenze"
|
||||
test_mehrfachantwort.py: "Antworten, in denen sich das Modell selbst korrigiert und mehrere Fassungen liefert, die ausfuehrlichste muss gewinnen"
|
||||
test_faktencheck_belegtiefe.py: "Belegzaehlung nach Medienhaus, Nachbedingung fuer bestaetigte Fakten, Sperre gegen das Hochstufen von developing, getrennte Bezeichnungen fuer Widerspruch und Widerlegung"
|
||||
test_bedrock_wiederholung.py: "Wiederholung des EU-Wegs bei Kapazitaets- und Kontingentfehlern, Staffelung, Zeit- und Wartebudget, Sichtbarkeit im Refresh-Protokoll"
|
||||
test_quellenausgabe.py: "Ausgabetreue des Quellenverzeichnisses: Nummern aus dem Lagebild statt Zeilennummern, keine Kappung, ein Verlag ein Name, Statistik deckungsgleich mit dem Verzeichnis"
|
||||
```
|
||||
|
||||
|
||||
@@ -11,7 +11,9 @@ Preistabelle BEDROCK_PRICING aus config.py, damit ClaudeUsage.cost_usd gefüllt
|
||||
bleibt und Statistik und Credits-Abrechnung unverändert funktionieren.
|
||||
"""
|
||||
import asyncio
|
||||
import contextvars
|
||||
import logging
|
||||
import time
|
||||
|
||||
from config import (
|
||||
BEDROCK_CACHE_READ_FACTOR,
|
||||
@@ -20,6 +22,9 @@ from config import (
|
||||
BEDROCK_MODEL_MAP,
|
||||
BEDROCK_PRICING,
|
||||
BEDROCK_REGION,
|
||||
BEDROCK_RETRY_BUDGET_S,
|
||||
BEDROCK_RETRY_MIN_REST_S,
|
||||
BEDROCK_RETRY_WAITS,
|
||||
CLAUDE_MODEL_STANDARD,
|
||||
CLAUDE_TIMEOUT,
|
||||
)
|
||||
@@ -49,6 +54,17 @@ _RATE_LIMIT_CODES = {
|
||||
"ServiceQuotaExceededException",
|
||||
"ModelNotReadyException",
|
||||
}
|
||||
# Innerhalb der Kategorie 'rate_limit' zwei verschiedene Ursachen, die nach
|
||||
# aussen gleich behandelt werden (damit die Retry-Steuerung des Orchestrators
|
||||
# unveraendert greift), im Log aber auseinandergehalten gehoeren:
|
||||
# Kapazitaet heisst, AWS hat gerade keine freie Rechenzeit fuer das Modell.
|
||||
# Kontingent heisst, WIR haben zu viele Anfragen je Minute gestellt.
|
||||
_KAPAZITAETS_CODES = {"ServiceUnavailableException", "ModelNotReadyException"}
|
||||
_KONTINGENT_CODES = {
|
||||
"ThrottlingException",
|
||||
"TooManyRequestsException",
|
||||
"ServiceQuotaExceededException",
|
||||
}
|
||||
_AUTH_ERROR_CODES = {
|
||||
"AccessDeniedException",
|
||||
"UnrecognizedClientException",
|
||||
@@ -66,6 +82,93 @@ _TRANSIENT_EXC_NAMES = (
|
||||
|
||||
_client = None
|
||||
|
||||
# Klartext je Stoerungsart, genutzt in Log, Refresh-Protokoll und Oberflaeche.
|
||||
_ART_TEXT = {
|
||||
"kapazitaet": "hat keine freie Kapazität",
|
||||
"kontingent": "meldet ausgeschöpftes Kontingent",
|
||||
"verbindung": "nicht erreichbar",
|
||||
}
|
||||
|
||||
# --- Wartebudget und Sichtbarkeit je Refresh --------------------------------
|
||||
# Alle drei Variablen gelten fuer den laufenden Refresh (ContextVar, also je
|
||||
# Task getrennt). Ohne gesetzten Kontext verhaelt sich der Client wie bisher,
|
||||
# nur mit Wiederholung: das Budget startet dann beim Standardwert und die
|
||||
# Meldungen gehen ausschliesslich ins Log.
|
||||
_wartebudget_var: contextvars.ContextVar[list] = contextvars.ContextVar("bedrock_wartebudget")
|
||||
_stoerungen_var: contextvars.ContextVar[list] = contextvars.ContextVar("bedrock_stoerungen")
|
||||
# Rueckmeldung an die Oberflaeche waehrend einer Wartezeit. Der Orchestrator
|
||||
# setzt hier eine Funktion (art, sekunden) -> None, damit die Pipeline-Anzeige
|
||||
# nicht wie ein Haenger aussieht.
|
||||
_wartemelder_var: contextvars.ContextVar = contextvars.ContextVar("bedrock_wartemelder")
|
||||
|
||||
|
||||
def refresh_kontext_starten() -> None:
|
||||
"""Setzt Wartebudget und Stoerungsliste fuer einen neuen Refresh zurueck."""
|
||||
_wartebudget_var.set([BEDROCK_RETRY_BUDGET_S])
|
||||
_stoerungen_var.set([])
|
||||
|
||||
|
||||
def stoerungen() -> list[dict]:
|
||||
"""Stoerungen des laufenden Refreshs, aelteste zuerst."""
|
||||
return list(_stoerungen_var.get([]) or [])
|
||||
|
||||
|
||||
def stoerungen_zusammenfassen() -> str:
|
||||
"""Einzeiler ueber die Stoerungen des Laufs, leer wenn es keine gab.
|
||||
|
||||
Landet im Refresh-Protokoll, damit ein Lauf mit Aussetzern nicht wie ein
|
||||
vollstaendiger aussieht.
|
||||
"""
|
||||
eintraege = stoerungen()
|
||||
if not eintraege:
|
||||
return ""
|
||||
je_art: dict[str, int] = {}
|
||||
gewartet = 0.0
|
||||
for e in eintraege:
|
||||
je_art[e["art"]] = je_art.get(e["art"], 0) + 1
|
||||
gewartet += e["wartezeit"]
|
||||
teile = [f"{anzahl}x {_ART_TEXT.get(art, art)}" for art, anzahl in sorted(je_art.items())]
|
||||
return f"KI-Dienst: {', '.join(teile)}, insgesamt {gewartet:.0f}s gewartet"
|
||||
|
||||
|
||||
def wartemelder_setzen(melder) -> None:
|
||||
"""Hinterlegt die Rueckmeldung an die Oberflaeche fuer diesen Refresh."""
|
||||
_wartemelder_var.set(melder)
|
||||
|
||||
|
||||
def _budget_rest() -> float:
|
||||
behaelter = _wartebudget_var.get(None)
|
||||
if behaelter is None:
|
||||
behaelter = [BEDROCK_RETRY_BUDGET_S]
|
||||
_wartebudget_var.set(behaelter)
|
||||
return behaelter[0]
|
||||
|
||||
|
||||
def _budget_verbrauchen(sekunden: float) -> None:
|
||||
behaelter = _wartebudget_var.get(None)
|
||||
if behaelter is not None:
|
||||
behaelter[0] = max(0.0, behaelter[0] - sekunden)
|
||||
|
||||
|
||||
def merke_stoerung(art: str, modell: str, wartezeit: float) -> None:
|
||||
"""Haelt eine wiederholte Stoerung fuer das Refresh-Protokoll fest."""
|
||||
liste = _stoerungen_var.get(None)
|
||||
if liste is None:
|
||||
liste = []
|
||||
_stoerungen_var.set(liste)
|
||||
liste.append({"art": art, "modell": modell, "wartezeit": wartezeit})
|
||||
|
||||
|
||||
def _melde_wartezeit(art: str, sekunden: float) -> None:
|
||||
"""Meldet eine laufende Wartezeit an die Oberflaeche, wenn eingerichtet."""
|
||||
melder = _wartemelder_var.get(None)
|
||||
if melder is None:
|
||||
return
|
||||
try:
|
||||
melder(art, sekunden)
|
||||
except Exception as e: # Anzeige darf den Lauf nie gefaehrden
|
||||
logger.debug("Wartemeldung fehlgeschlagen: %s", e)
|
||||
|
||||
|
||||
def _get_client():
|
||||
"""Erzeugt den bedrock-runtime-Client einmalig. Lazy Import, damit die App
|
||||
@@ -86,18 +189,27 @@ def _get_client():
|
||||
config=Config(
|
||||
connect_timeout=10,
|
||||
read_timeout=CLAUDE_TIMEOUT,
|
||||
retries={"max_attempts": 2, "mode": "standard"},
|
||||
# 'adaptive' bremst sich bei Drosselung selbst ein, statt
|
||||
# stur nachzuschieben. Die Versuchszahl bleibt klein, weil das
|
||||
# geduldige Warten in call_bedrock sitzt: botocore deckt die
|
||||
# Sekunden ab, unsere Schleife die Zehnersekunden.
|
||||
retries={"max_attempts": 3, "mode": "adaptive"},
|
||||
),
|
||||
)
|
||||
return _client
|
||||
|
||||
|
||||
def _classify_bedrock_error(exc: Exception) -> str:
|
||||
"""Ordnet einer botocore-Exception eine error_type-Kategorie zu."""
|
||||
code = ""
|
||||
def _fehlercode(exc: Exception) -> str:
|
||||
"""Botocore-Fehlercode einer Exception, leer wenn keiner vorliegt."""
|
||||
response = getattr(exc, "response", None)
|
||||
if isinstance(response, dict):
|
||||
code = response.get("Error", {}).get("Code", "") or ""
|
||||
return response.get("Error", {}).get("Code", "") or ""
|
||||
return ""
|
||||
|
||||
|
||||
def _classify_bedrock_error(exc: Exception) -> str:
|
||||
"""Ordnet einer botocore-Exception eine error_type-Kategorie zu."""
|
||||
code = _fehlercode(exc)
|
||||
if code in _RATE_LIMIT_CODES:
|
||||
return "rate_limit"
|
||||
if code in _AUTH_ERROR_CODES:
|
||||
@@ -105,6 +217,22 @@ def _classify_bedrock_error(exc: Exception) -> str:
|
||||
return "cli_error"
|
||||
|
||||
|
||||
def stoerungsart(exc: Exception) -> str | None:
|
||||
"""Benennt die Ursache einer wiederholbaren Stoerung, sonst None.
|
||||
|
||||
'kapazitaet' = AWS hat keine freie Rechenzeit fuer das Modell (HTTP 503).
|
||||
'kontingent' = unsere Anfragen je Minute sind ausgeschoepft (HTTP 429).
|
||||
"""
|
||||
code = _fehlercode(exc)
|
||||
if code in _KAPAZITAETS_CODES:
|
||||
return "kapazitaet"
|
||||
if code in _KONTINGENT_CODES:
|
||||
return "kontingent"
|
||||
if type(exc).__name__ in _TRANSIENT_EXC_NAMES:
|
||||
return "verbindung"
|
||||
return None
|
||||
|
||||
|
||||
def _cost_usd(
|
||||
cli_model: str,
|
||||
input_tokens: int,
|
||||
@@ -138,6 +266,10 @@ async def call_bedrock(
|
||||
Gleicher Vertrag wie call_claude mit tools=None. Fehler kommen als
|
||||
ClaudeCliError mit denselben error_type-Kategorien, Timeouts als
|
||||
TimeoutError, Abbrüche als CancelledError.
|
||||
|
||||
Bei Kapazitäts- und Kontingentfehlern wird nach BEDROCK_RETRY_WAITS erneut
|
||||
versucht, begrenzt durch das Zeitbudget dieses Aufrufs und das
|
||||
Wartebudget des laufenden Refreshs.
|
||||
"""
|
||||
cli_model = model or CLAUDE_MODEL_STANDARD
|
||||
bedrock_model = BEDROCK_MODEL_MAP.get(cli_model)
|
||||
@@ -170,6 +302,89 @@ async def call_bedrock(
|
||||
if not raw_text:
|
||||
kwargs["system"] = [{"text": _JSON_ONLY_SYSTEM}]
|
||||
|
||||
beginn = time.monotonic()
|
||||
versuch = 0
|
||||
while True:
|
||||
try:
|
||||
response = await _converse_einmal(
|
||||
client, kwargs, effective_timeout - (time.monotonic() - beginn), cancel_event
|
||||
)
|
||||
break
|
||||
except (asyncio.CancelledError, TimeoutError):
|
||||
raise
|
||||
except Exception as e:
|
||||
art = stoerungsart(e)
|
||||
wartezeit = _wartezeit(
|
||||
art, versuch, rest=effective_timeout - (time.monotonic() - beginn)
|
||||
)
|
||||
if wartezeit is None:
|
||||
_protokolliere_fehler(e, cli_model, versuch)
|
||||
if type(e).__name__ in _TRANSIENT_EXC_NAMES:
|
||||
raise TimeoutError(f"Bedrock Verbindungsfehler ({type(e).__name__}). {e}") from e
|
||||
raise ClaudeCliError(
|
||||
_classify_bedrock_error(e), f"{type(e).__name__}. {e}"
|
||||
) from e
|
||||
versuch += 1
|
||||
_budget_verbrauchen(wartezeit)
|
||||
merke_stoerung(art, cli_model, wartezeit)
|
||||
logger.warning(
|
||||
"Bedrock %s (%s, %s): Versuch %d in %.0fs",
|
||||
_ART_TEXT.get(art, art), type(e).__name__, cli_model, versuch + 1, wartezeit,
|
||||
)
|
||||
_melde_wartezeit(art, wartezeit)
|
||||
await asyncio.sleep(wartezeit)
|
||||
if cancel_event and cancel_event.is_set():
|
||||
raise asyncio.CancelledError("Cancel angefordert")
|
||||
|
||||
return _auswerten(response, cli_model, bedrock_model)
|
||||
|
||||
|
||||
def _wartezeit(art: str | None, versuch: int, rest: float) -> float | None:
|
||||
"""Wartezeit vor dem nächsten Versuch, None wenn nicht wiederholt wird.
|
||||
|
||||
Drei Bedingungen müssen zusammenkommen: die Störung ist wiederholbar, es
|
||||
gibt noch eine Wartestufe, und die Wartezeit passt sowohl in das
|
||||
Zeitbudget dieses Aufrufs als auch in das Wartebudget des Refreshs.
|
||||
"""
|
||||
if not art or versuch >= len(BEDROCK_RETRY_WAITS):
|
||||
return None
|
||||
wartezeit = BEDROCK_RETRY_WAITS[versuch]
|
||||
if rest < wartezeit + BEDROCK_RETRY_MIN_REST_S:
|
||||
logger.warning(
|
||||
"Bedrock: keine Wiederholung, Zeitbudget des Aufrufs reicht nicht "
|
||||
"(%.0fs übrig, %.0fs Wartezeit plus %.0fs Reserve nötig)",
|
||||
max(rest, 0), wartezeit, BEDROCK_RETRY_MIN_REST_S,
|
||||
)
|
||||
return None
|
||||
if _budget_rest() < wartezeit:
|
||||
logger.warning(
|
||||
"Bedrock: keine Wiederholung, Wartebudget des Laufs erschöpft (%.0fs übrig)",
|
||||
_budget_rest(),
|
||||
)
|
||||
return None
|
||||
return wartezeit
|
||||
|
||||
|
||||
def _protokolliere_fehler(e: Exception, cli_model: str, versuch: int) -> None:
|
||||
"""Schreibt den endgültigen Fehler mit seiner Ursache ins Log."""
|
||||
exc_name = type(e).__name__
|
||||
zusatz = f" nach {versuch + 1} Versuchen" if versuch else ""
|
||||
art = stoerungsart(e)
|
||||
if art:
|
||||
logger.warning(
|
||||
"Bedrock %s (%s, %s)%s. %s", _ART_TEXT.get(art, art), exc_name, cli_model, zusatz, e
|
||||
)
|
||||
elif _classify_bedrock_error(e) == "auth_error":
|
||||
logger.error("Bedrock Auth-Fehler (%s). %s", exc_name, e)
|
||||
else:
|
||||
logger.error("Bedrock-Fehler (%s). %s", exc_name, e)
|
||||
|
||||
|
||||
async def _converse_einmal(client, kwargs: dict, rest_timeout: float, cancel_event):
|
||||
"""Führt genau einen Converse-Aufruf aus und gibt die Antwort zurück."""
|
||||
if rest_timeout <= 0:
|
||||
raise TimeoutError("Bedrock Timeout: Zeitbudget vor dem Aufruf aufgebraucht")
|
||||
|
||||
call_task = asyncio.create_task(asyncio.to_thread(client.converse, **kwargs))
|
||||
wait_tasks = [call_task]
|
||||
cancel_wait_task = None
|
||||
@@ -178,7 +393,7 @@ async def call_bedrock(
|
||||
wait_tasks.append(cancel_wait_task)
|
||||
|
||||
done, pending = await asyncio.wait(
|
||||
wait_tasks, timeout=effective_timeout, return_when=asyncio.FIRST_COMPLETED
|
||||
wait_tasks, timeout=rest_timeout, return_when=asyncio.FIRST_COMPLETED
|
||||
)
|
||||
for p in pending:
|
||||
p.cancel()
|
||||
@@ -192,23 +407,15 @@ async def call_bedrock(
|
||||
)
|
||||
if cancel_wait_task is not None and cancel_wait_task in done:
|
||||
raise asyncio.CancelledError("Cancel angefordert")
|
||||
raise TimeoutError(f"Bedrock Timeout nach {effective_timeout}s")
|
||||
raise TimeoutError(f"Bedrock Timeout nach {rest_timeout:.0f}s")
|
||||
|
||||
try:
|
||||
response = call_task.result()
|
||||
except Exception as e:
|
||||
exc_name = type(e).__name__
|
||||
if exc_name in _TRANSIENT_EXC_NAMES:
|
||||
raise TimeoutError(f"Bedrock Verbindungsfehler ({exc_name}). {e}") from e
|
||||
error_type = _classify_bedrock_error(e)
|
||||
if error_type == "rate_limit":
|
||||
logger.warning(f"Bedrock Rate-Limit ({exc_name}). {e}")
|
||||
elif error_type == "auth_error":
|
||||
logger.error(f"Bedrock Auth-Fehler ({exc_name}). {e}")
|
||||
else:
|
||||
logger.error(f"Bedrock-Fehler ({exc_name}). {e}")
|
||||
raise ClaudeCliError(error_type, f"{exc_name}. {e}") from e
|
||||
# Fehler werden hier NICHT behandelt, das entscheidet die Wiederholung in
|
||||
# call_bedrock anhand der Stoerungsart.
|
||||
return call_task.result()
|
||||
|
||||
|
||||
def _auswerten(response: dict, cli_model: str, bedrock_model: str) -> tuple[str, ClaudeUsage]:
|
||||
"""Wandelt die Converse-Antwort in (Text, Verbrauch) um."""
|
||||
out_msg = response.get("output", {}).get("message", {})
|
||||
parts = [c.get("text", "") for c in out_msg.get("content", []) if "text" in c]
|
||||
result_text = "".join(parts).strip()
|
||||
|
||||
@@ -12,6 +12,7 @@ from urllib.parse import urlparse, urlunparse, quote_plus
|
||||
import httpx
|
||||
|
||||
from agents.claude_client import UsageAccumulator, _cancel_event_var, _ai_backend_var
|
||||
from agents import bedrock_client
|
||||
from agents.factchecker import find_matching_claim, deduplicate_new_facts, TWOPHASE_MIN_FACTS
|
||||
from source_rules import (
|
||||
_detect_category,
|
||||
@@ -878,6 +879,24 @@ class AgentOrchestrator:
|
||||
_ai_backend_var.set(ai_backend)
|
||||
if ai_backend != "cli":
|
||||
logger.info("Lage %s läuft über KI-Backend '%s'", incident_id, ai_backend)
|
||||
# Wartebudget des EU-Wegs zuruecksetzen und die Oberflaeche an die
|
||||
# Wartemeldungen anschliessen. Ohne diesen Hinweis sieht eine
|
||||
# Wiederholung nach einem Kapazitaetsengpass wie ein Haenger aus,
|
||||
# weil die Anzeige waehrend der Wartezeit stillsteht.
|
||||
bedrock_client.refresh_kontext_starten()
|
||||
if self._ws_manager:
|
||||
def _warte_melden(art: str, sekunden: float, _iid=incident_id,
|
||||
_vis=visibility, _cb=created_by, _tid=tenant_id) -> None:
|
||||
text = bedrock_client._ART_TEXT.get(art, art)
|
||||
asyncio.create_task(self._ws_manager.broadcast_for_incident({
|
||||
"type": "status_update",
|
||||
"incident_id": _iid,
|
||||
"data": {
|
||||
"status": "running",
|
||||
"detail": f"KI-Dienst {text}, neuer Versuch in {sekunden:.0f} Sekunden...",
|
||||
},
|
||||
}, _vis, _cb, _tid))
|
||||
bedrock_client.wartemelder_setzen(_warte_melden)
|
||||
output_language_iso = await get_org_language(db, tenant_id) if tenant_id else "de"
|
||||
output_language = language_display(output_language_iso)
|
||||
# research_language steuert nur den WebSearch-Prompt ("suche in Sprache X").
|
||||
@@ -2200,18 +2219,24 @@ class AgentOrchestrator:
|
||||
if _missing_key not in _started_keys:
|
||||
await _pipe_skip(_missing_key)
|
||||
|
||||
# Refresh-Log abschließen (mit Token-Statistiken)
|
||||
# Refresh-Log abschließen (mit Token-Statistiken). Störungen des
|
||||
# KI-Dienstes werden auch bei einem erfolgreichen Lauf vermerkt,
|
||||
# sonst sieht ein Durchlauf mit Aussetzern aus wie ein glatter.
|
||||
stoerungshinweis = bedrock_client.stoerungen_zusammenfassen() or None
|
||||
if stoerungshinweis:
|
||||
logger.info("Refresh-Protokoll Lage %s: %s", incident_id, stoerungshinweis)
|
||||
await db.execute(
|
||||
"""UPDATE refresh_log SET
|
||||
completed_at = ?, articles_found = ?, status = 'completed',
|
||||
input_tokens = ?, output_tokens = ?,
|
||||
cache_creation_tokens = ?, cache_read_tokens = ?,
|
||||
total_cost_usd = ?, api_calls = ?
|
||||
total_cost_usd = ?, api_calls = ?, error_message = ?
|
||||
WHERE id = ?""",
|
||||
(datetime.now(TIMEZONE).strftime('%Y-%m-%d %H:%M:%S'), new_count,
|
||||
usage_acc.input_tokens, usage_acc.output_tokens,
|
||||
usage_acc.cache_creation_tokens, usage_acc.cache_read_tokens,
|
||||
round(usage_acc.total_cost_usd, 7), usage_acc.call_count, log_id),
|
||||
round(usage_acc.total_cost_usd, 7), usage_acc.call_count,
|
||||
stoerungshinweis, log_id),
|
||||
)
|
||||
await db.commit()
|
||||
logger.info(
|
||||
|
||||
@@ -70,6 +70,23 @@ BEDROCK_PRICING = {
|
||||
}
|
||||
BEDROCK_CACHE_WRITE_FACTOR = 1.25
|
||||
BEDROCK_CACHE_READ_FACTOR = 0.10
|
||||
# Wartestufen in Sekunden, wenn Bedrock keine Kapazitaet hat
|
||||
# (ServiceUnavailableException, HTTP 503) oder unsere Quote greift
|
||||
# (ThrottlingException, HTTP 429). Botocore gibt nach rund fuenf Sekunden auf,
|
||||
# genau in dem Fenster, in dem sich eine Kapazitaetsdelle meist von selbst
|
||||
# loest. Am 01.08.2026 kostete das einen kompletten Pipeline-Schritt
|
||||
# (Web-Source-Selektion), ohne dass es im Bericht sichtbar wurde.
|
||||
BEDROCK_RETRY_WAITS = [
|
||||
float(x) for x in os.environ.get("BEDROCK_RETRY_WAITS", "5,15,30").split(",") if x.strip()
|
||||
]
|
||||
# Obergrenze fuer die Summe aller Wartezeiten innerhalb eines Refreshs. Ohne
|
||||
# Deckel koennten mehrere Engpaesse einen Lauf spuerbar verlaengern.
|
||||
BEDROCK_RETRY_BUDGET_S = float(os.environ.get("BEDROCK_RETRY_BUDGET_S", "120"))
|
||||
# Sicherheitsabstand: Es wird nur gewartet, wenn danach noch mindestens so
|
||||
# viele Sekunden vom Zeitbudget des Aufrufs uebrig sind. Die Wartezeit zaehlt
|
||||
# gegen dasselbe Budget wie der Aufruf selbst (asyncio.wait), und die
|
||||
# Planungsaufrufe an Haiku haben nur 120 Sekunden.
|
||||
BEDROCK_RETRY_MIN_REST_S = float(os.environ.get("BEDROCK_RETRY_MIN_REST_S", "12"))
|
||||
|
||||
# --- staan.ai Such-API (EU/DSGVO-Umbau, Phase 2) -------------------------------
|
||||
# Europäischer Suchindex (European Search Perspective). Ersetzt im EU-Modellweg
|
||||
|
||||
215
tests/test_bedrock_wiederholung.py
Normale Datei
215
tests/test_bedrock_wiederholung.py
Normale Datei
@@ -0,0 +1,215 @@
|
||||
"""Testet die Wiederholung des EU-Modellwegs bei Kapazitaetsengpaessen.
|
||||
|
||||
Hintergrund: Am 01.08.2026 lieferte Bedrock waehrend eines Laufs
|
||||
ServiceUnavailableException (HTTP 503, keine freie Kapazitaet). Botocore gab
|
||||
nach rund fuenf Sekunden auf, die Web-Source-Selektion fiel ersatzlos aus, und
|
||||
im Bericht war davon nichts zu sehen. Gemessen waren es zwei betroffene
|
||||
Aufrufe bei rund 340, also 0,6 Prozent.
|
||||
|
||||
Geprueft werden die vier Zusagen des Umbaus:
|
||||
1. Kapazitaets- und Kontingentfehler werden mit wachsendem Abstand wiederholt.
|
||||
2. Die Wartezeit passt in das Zeitbudget des Aufrufs und in das des Laufs.
|
||||
3. Beide Ursachen werden im Log auseinandergehalten.
|
||||
4. Ein Lauf mit Aussetzern ist im Refresh-Protokoll erkennbar.
|
||||
|
||||
Laeuft ohne Netzzugriff, ohne Kosten und ohne AWS-Zugang: der Converse-Aufruf
|
||||
ist durch ein Testdouble ersetzt, Wartezeiten werden nicht real abgewartet.
|
||||
|
||||
Aufruf aus dem Projektstamm:
|
||||
venv/bin/python tests/test_bedrock_wiederholung.py
|
||||
"""
|
||||
import asyncio
|
||||
import os
|
||||
import sys
|
||||
|
||||
sys.path.insert(0, os.path.join(os.path.dirname(os.path.abspath(__file__)), "..", "src"))
|
||||
|
||||
import agents.bedrock_client as bc # noqa: E402
|
||||
from agents.claude_client import ClaudeCliError # noqa: E402
|
||||
|
||||
ok = 0
|
||||
fail = 0
|
||||
|
||||
|
||||
def pruefe(name, bedingung, extra=""):
|
||||
global ok, fail
|
||||
if bedingung:
|
||||
ok += 1
|
||||
print(" OK " + name)
|
||||
else:
|
||||
fail += 1
|
||||
print(" FEHL " + name + " " + str(extra))
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Testdoubles
|
||||
# ---------------------------------------------------------------------------
|
||||
class FakeFehler(Exception):
|
||||
"""Botocore-aehnlicher Fehler mit response-Dict."""
|
||||
|
||||
def __init__(self, code):
|
||||
super().__init__(f"{code}: Testfehler")
|
||||
self.response = {"Error": {"Code": code}}
|
||||
|
||||
|
||||
ANTWORT = {
|
||||
"output": {"message": {"content": [{"text": '{"ok": true}'}]}},
|
||||
"stopReason": "end_turn",
|
||||
"usage": {"inputTokens": 10, "outputTokens": 5},
|
||||
"metrics": {"latencyMs": 1200},
|
||||
}
|
||||
|
||||
aufrufe = {"n": 0}
|
||||
gewartet = []
|
||||
antwortfolge = []
|
||||
|
||||
|
||||
class FakeClient:
|
||||
def converse(self, **kwargs):
|
||||
aufrufe["n"] += 1
|
||||
naechste = antwortfolge.pop(0) if antwortfolge else ANTWORT
|
||||
if isinstance(naechste, Exception):
|
||||
raise naechste
|
||||
return naechste
|
||||
|
||||
|
||||
async def fake_sleep(sekunden):
|
||||
gewartet.append(sekunden)
|
||||
|
||||
|
||||
def neu(folge, budget=None, waits=(5, 15, 30)):
|
||||
"""Setzt Testdoubles und Kontext fuer einen Durchgang."""
|
||||
aufrufe["n"] = 0
|
||||
gewartet.clear()
|
||||
antwortfolge[:] = list(folge)
|
||||
bc.BEDROCK_RETRY_WAITS[:] = list(waits)
|
||||
bc.refresh_kontext_starten()
|
||||
if budget is not None:
|
||||
bc._wartebudget_var.set([budget])
|
||||
|
||||
|
||||
bc._get_client = lambda: FakeClient()
|
||||
bc.asyncio.sleep = fake_sleep
|
||||
|
||||
|
||||
def lauf(timeout=420.0):
|
||||
return asyncio.run(bc.call_bedrock("Testauftrag", timeout=timeout))
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# A. Störungsarten auseinanderhalten
|
||||
# ---------------------------------------------------------------------------
|
||||
print("\nA. Ursachen unterscheiden")
|
||||
pruefe("A1 503 ist ein Kapazitaetsengpass",
|
||||
bc.stoerungsart(FakeFehler("ServiceUnavailableException")) == "kapazitaet")
|
||||
pruefe("A2 429 ist unser Kontingent",
|
||||
bc.stoerungsart(FakeFehler("ThrottlingException")) == "kontingent")
|
||||
pruefe("A3 Netzabbruch ist eine Verbindungsstoerung",
|
||||
bc.stoerungsart(type("ReadTimeoutError", (Exception,), {})()) == "verbindung")
|
||||
pruefe("A4 Auth-Fehler wird nicht wiederholt",
|
||||
bc.stoerungsart(FakeFehler("AccessDeniedException")) is None)
|
||||
pruefe("A5 beide Ursachen bleiben nach aussen 'rate_limit'",
|
||||
bc._classify_bedrock_error(FakeFehler("ServiceUnavailableException")) == "rate_limit"
|
||||
and bc._classify_bedrock_error(FakeFehler("ThrottlingException")) == "rate_limit")
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# B. Wiederholung
|
||||
# ---------------------------------------------------------------------------
|
||||
print("\nB. Wiederholung bei Kapazitaetsengpass")
|
||||
neu([FakeFehler("ServiceUnavailableException"), ANTWORT])
|
||||
text, usage = lauf()
|
||||
pruefe("B1 zweiter Versuch liefert das Ergebnis", aufrufe["n"] == 2, aufrufe["n"])
|
||||
pruefe("B2 dazwischen wurde 5 Sekunden gewartet", gewartet == [5], gewartet)
|
||||
pruefe("B3 Verbrauch wird normal zurueckgegeben", usage.input_tokens == 10, usage)
|
||||
|
||||
neu([FakeFehler("ServiceUnavailableException")] * 3 + [ANTWORT])
|
||||
lauf()
|
||||
pruefe("B4 Staffelung 5, 15, 30 Sekunden", gewartet == [5, 15, 30], gewartet)
|
||||
pruefe("B5 vier Versuche insgesamt", aufrufe["n"] == 4, aufrufe["n"])
|
||||
|
||||
neu([FakeFehler("ServiceUnavailableException")] * 5)
|
||||
try:
|
||||
lauf()
|
||||
pruefe("B6 nach der letzten Stufe wird aufgegeben", False, "kein Fehler geworfen")
|
||||
except ClaudeCliError as e:
|
||||
pruefe("B6 nach der letzten Stufe wird aufgegeben", e.error_type == "rate_limit", e.error_type)
|
||||
pruefe("B7 genau vier Versuche, dann Schluss", aufrufe["n"] == 4, aufrufe["n"])
|
||||
|
||||
neu([FakeFehler("AccessDeniedException"), ANTWORT])
|
||||
try:
|
||||
lauf()
|
||||
pruefe("B8 Auth-Fehler wird sofort durchgereicht", False, "kein Fehler geworfen")
|
||||
except ClaudeCliError as e:
|
||||
pruefe("B8 Auth-Fehler wird sofort durchgereicht",
|
||||
e.error_type == "auth_error" and aufrufe["n"] == 1, (e.error_type, aufrufe["n"]))
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# C. Budgets
|
||||
# ---------------------------------------------------------------------------
|
||||
print("\nC. Zeitbudget und Wartebudget")
|
||||
# Planungsaufrufe an Haiku haben nur 120 Sekunden. Nach zwei Wartestufen
|
||||
# (5 + 15) ist zu wenig Rest fuer die dritte (30 + 12 Reserve).
|
||||
neu([FakeFehler("ServiceUnavailableException")] * 5)
|
||||
try:
|
||||
lauf(timeout=40.0)
|
||||
except ClaudeCliError:
|
||||
pass
|
||||
pruefe("C1 Wartezeit passt sich dem Zeitbudget des Aufrufs an",
|
||||
gewartet == [5, 15], gewartet)
|
||||
|
||||
neu([FakeFehler("ServiceUnavailableException")] * 5, budget=10.0)
|
||||
try:
|
||||
lauf()
|
||||
except ClaudeCliError:
|
||||
pass
|
||||
pruefe("C2 Wartebudget des Laufs begrenzt die Wiederholung",
|
||||
gewartet == [5], gewartet)
|
||||
|
||||
neu([FakeFehler("ServiceUnavailableException"), ANTWORT], budget=25.0)
|
||||
lauf()
|
||||
pruefe("C3 verbrauchte Wartezeit wird vom Budget abgezogen",
|
||||
abs(bc._budget_rest() - 20.0) < 0.01, bc._budget_rest())
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# D. Sichtbarkeit
|
||||
# ---------------------------------------------------------------------------
|
||||
print("\nD. Sichtbarkeit im Refresh-Protokoll")
|
||||
neu([FakeFehler("ServiceUnavailableException"),
|
||||
FakeFehler("ThrottlingException"), ANTWORT])
|
||||
meldungen = []
|
||||
bc.wartemelder_setzen(lambda art, sek: meldungen.append((art, sek)))
|
||||
lauf()
|
||||
hinweis = bc.stoerungen_zusammenfassen()
|
||||
pruefe("D1 Hinweis nennt beide Ursachen",
|
||||
"Kapazität" in hinweis and "Kontingent" in hinweis, hinweis)
|
||||
pruefe("D2 Hinweis nennt die gesamte Wartezeit", "20s" in hinweis, hinweis)
|
||||
pruefe("D3 Oberflaeche wird waehrend jeder Wartezeit benachrichtigt",
|
||||
meldungen == [("kapazitaet", 5), ("kontingent", 15)], meldungen)
|
||||
pruefe("D4 Hinweis nutzt echte Umlaute",
|
||||
"Kapazitaet" not in hinweis and "ausgeschoepft" not in hinweis, hinweis)
|
||||
|
||||
neu([ANTWORT])
|
||||
bc.wartemelder_setzen(None)
|
||||
lauf()
|
||||
pruefe("D5 stoerungsfreier Lauf erzeugt keinen Hinweis",
|
||||
bc.stoerungen_zusammenfassen() == "", bc.stoerungen_zusammenfassen())
|
||||
pruefe("D6 stoerungsfreier Lauf wartet nicht", gewartet == [], gewartet)
|
||||
|
||||
# Eine kaputte Anzeige darf den Lauf nicht gefaehrden.
|
||||
neu([FakeFehler("ServiceUnavailableException"), ANTWORT])
|
||||
|
||||
|
||||
def _melder_mit_fehler(art, sek):
|
||||
raise RuntimeError("Anzeige kaputt")
|
||||
|
||||
|
||||
bc.wartemelder_setzen(_melder_mit_fehler)
|
||||
text, _ = lauf()
|
||||
pruefe("D7 Fehler in der Anzeige bricht den Lauf nicht ab", text.startswith("{"), text[:40])
|
||||
bc.wartemelder_setzen(None)
|
||||
|
||||
print("\nErgebnis: " + str(ok) + " bestanden, " + str(fail) + " fehlgeschlagen")
|
||||
sys.exit(1 if fail else 0)
|
||||
In neuem Issue referenzieren
Einen Benutzer sperren