feat(eu-umbau): Bedrock als zweiter KI-Modellweg (Phase 1)
Zweiter Modellweg ueber AWS Bedrock (EU-Inferenzprofile Frankfurt), umschaltbar je Lage (incidents.ai_backend) > Organisation (org_setting ai_backend) > global (ENV AI_BACKEND, Default cli). Neue Datei agents/bedrock_client.py (Converse API, nur eu.-Profile, Token-zu-USD-Umrechnung ueber BEDROCK_PRICING, gleiche Fehlerkategorien wie das CLI). Routing zentral in call_claude; Aufrufe mit WebSearch/WebFetch laufen bis zur staan-Schleife (Phase 2) weiter uebers CLI. Migration: Spalte incidents.ai_backend. requirements: boto3. Hinweis Promote: boto3 muss im Live-venv installiert werden. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Dieser Commit ist enthalten in:
240
src/agents/bedrock_client.py
Normale Datei
240
src/agents/bedrock_client.py
Normale Datei
@@ -0,0 +1,240 @@
|
||||
"""Bedrock-Client. Zweiter Modellweg über AWS Bedrock (EU-Inferenzprofile, Frankfurt).
|
||||
|
||||
Teil des EU/DSGVO-Umbaus (Phase 1). Die Anfragen werden vollständig von AWS in
|
||||
europäischen Regionen verarbeitet, nichts geht an Anthropic. In Phase 1 bedient
|
||||
dieser Weg nur werkzeuglose Aufrufe (tools=None). Aufrufe mit WebSearch/WebFetch
|
||||
laufen weiter über das Claude CLI, bis die staan-Recherche-Schleife (Phase 2)
|
||||
fertig ist. Das Routing sitzt zentral in claude_client.call_claude.
|
||||
|
||||
Kosten. Bedrock meldet nur Token. Die Umrechnung in USD passiert hier über die
|
||||
Preistabelle BEDROCK_PRICING aus config.py, damit ClaudeUsage.cost_usd gefüllt
|
||||
bleibt und Statistik und Credits-Abrechnung unverändert funktionieren.
|
||||
"""
|
||||
import asyncio
|
||||
import logging
|
||||
|
||||
from config import (
|
||||
BEDROCK_CACHE_READ_FACTOR,
|
||||
BEDROCK_CACHE_WRITE_FACTOR,
|
||||
BEDROCK_MAX_TOKENS,
|
||||
BEDROCK_MODEL_MAP,
|
||||
BEDROCK_PRICING,
|
||||
BEDROCK_REGION,
|
||||
CLAUDE_MODEL_STANDARD,
|
||||
CLAUDE_TIMEOUT,
|
||||
)
|
||||
from agents.claude_client import (
|
||||
ClaudeCliError,
|
||||
ClaudeUsage,
|
||||
_cancel_event_var,
|
||||
_sanitize_mdash,
|
||||
)
|
||||
|
||||
logger = logging.getLogger("osint.bedrock_client")
|
||||
|
||||
# Entspricht dem JSON-Zwang des CLI-Wegs bei tools=None (siehe claude_client.call_claude)
|
||||
_JSON_ONLY_SYSTEM = (
|
||||
"CRITICAL: You are a JSON-only output agent. "
|
||||
"Output EXCLUSIVELY a single valid JSON object. "
|
||||
"No explanatory text, no markdown fences, no continuation of previous responses. "
|
||||
"Start your response with { and end with }."
|
||||
)
|
||||
|
||||
# botocore-Fehlercodes, gruppiert nach den error_type-Kategorien des CLI-Wegs,
|
||||
# damit die Retry-Steuerung des Orchestrators unverändert greift.
|
||||
_RATE_LIMIT_CODES = {
|
||||
"ThrottlingException",
|
||||
"TooManyRequestsException",
|
||||
"ServiceUnavailableException",
|
||||
"ServiceQuotaExceededException",
|
||||
"ModelNotReadyException",
|
||||
}
|
||||
_AUTH_ERROR_CODES = {
|
||||
"AccessDeniedException",
|
||||
"UnrecognizedClientException",
|
||||
"ExpiredTokenException",
|
||||
"InvalidSignatureException",
|
||||
}
|
||||
# Transiente botocore-Netzfehler, die als TimeoutError hochgereicht werden
|
||||
# (der Orchestrator behandelt TimeoutError/ConnectionError als retry-fähig).
|
||||
_TRANSIENT_EXC_NAMES = (
|
||||
"ReadTimeoutError",
|
||||
"ConnectTimeoutError",
|
||||
"EndpointConnectionError",
|
||||
"ConnectionClosedError",
|
||||
)
|
||||
|
||||
_client = None
|
||||
|
||||
|
||||
def _get_client():
|
||||
"""Erzeugt den bedrock-runtime-Client einmalig. Lazy Import, damit die App
|
||||
auch ohne installiertes boto3 startet, solange das Backend 'cli' ist."""
|
||||
global _client
|
||||
if _client is None:
|
||||
try:
|
||||
import boto3
|
||||
from botocore.config import Config
|
||||
except ImportError as e:
|
||||
raise ClaudeCliError(
|
||||
"cli_error",
|
||||
f"boto3 ist nicht installiert, Bedrock-Backend nicht nutzbar ({e})",
|
||||
)
|
||||
_client = boto3.client(
|
||||
"bedrock-runtime",
|
||||
region_name=BEDROCK_REGION,
|
||||
config=Config(
|
||||
connect_timeout=10,
|
||||
read_timeout=CLAUDE_TIMEOUT,
|
||||
retries={"max_attempts": 2, "mode": "standard"},
|
||||
),
|
||||
)
|
||||
return _client
|
||||
|
||||
|
||||
def _classify_bedrock_error(exc: Exception) -> str:
|
||||
"""Ordnet einer botocore-Exception eine error_type-Kategorie zu."""
|
||||
code = ""
|
||||
response = getattr(exc, "response", None)
|
||||
if isinstance(response, dict):
|
||||
code = response.get("Error", {}).get("Code", "") or ""
|
||||
if code in _RATE_LIMIT_CODES:
|
||||
return "rate_limit"
|
||||
if code in _AUTH_ERROR_CODES:
|
||||
return "auth_error"
|
||||
return "cli_error"
|
||||
|
||||
|
||||
def _cost_usd(
|
||||
cli_model: str,
|
||||
input_tokens: int,
|
||||
output_tokens: int,
|
||||
cache_write: int,
|
||||
cache_read: int,
|
||||
) -> float:
|
||||
"""Token mal Preis. Preise in BEDROCK_PRICING sind USD je 1 Mio Token."""
|
||||
prices = BEDROCK_PRICING.get(cli_model)
|
||||
if not prices:
|
||||
logger.warning(f"Keine Preise für Modell '{cli_model}' in BEDROCK_PRICING, cost_usd bleibt 0")
|
||||
return 0.0
|
||||
per_in = prices["input"] / 1_000_000
|
||||
per_out = prices["output"] / 1_000_000
|
||||
return (
|
||||
input_tokens * per_in
|
||||
+ output_tokens * per_out
|
||||
+ cache_write * per_in * BEDROCK_CACHE_WRITE_FACTOR
|
||||
+ cache_read * per_in * BEDROCK_CACHE_READ_FACTOR
|
||||
)
|
||||
|
||||
|
||||
async def call_bedrock(
|
||||
prompt: str,
|
||||
model: str | None = None,
|
||||
raw_text: bool = False,
|
||||
timeout: float | None = None,
|
||||
) -> tuple[str, ClaudeUsage]:
|
||||
"""Ruft Claude über AWS Bedrock (Converse API) auf. Gibt (result_text, usage) zurück.
|
||||
|
||||
Gleicher Vertrag wie call_claude mit tools=None. Fehler kommen als
|
||||
ClaudeCliError mit denselben error_type-Kategorien, Timeouts als
|
||||
TimeoutError, Abbrüche als CancelledError.
|
||||
"""
|
||||
cli_model = model or CLAUDE_MODEL_STANDARD
|
||||
bedrock_model = BEDROCK_MODEL_MAP.get(cli_model)
|
||||
if not bedrock_model:
|
||||
raise ClaudeCliError(
|
||||
"cli_error",
|
||||
f"Kein EU-Inferenzprofil für Modell '{cli_model}' in BEDROCK_MODEL_MAP hinterlegt",
|
||||
)
|
||||
if not bedrock_model.startswith("eu."):
|
||||
# Sperre gegen Nicht-EU-Profile. Ein globales Profil würde die
|
||||
# Verarbeitung ausserhalb Europas erlauben, genau das soll dieser Weg
|
||||
# ausschliessen.
|
||||
raise ClaudeCliError(
|
||||
"cli_error",
|
||||
f"Bedrock-Profil '{bedrock_model}' ist kein EU-Profil (eu.-Präfix fehlt)",
|
||||
)
|
||||
|
||||
effective_timeout = timeout if timeout is not None else CLAUDE_TIMEOUT
|
||||
cancel_event = _cancel_event_var.get(None)
|
||||
if cancel_event and cancel_event.is_set():
|
||||
raise asyncio.CancelledError("Cancel angefordert")
|
||||
|
||||
client = _get_client()
|
||||
|
||||
kwargs = {
|
||||
"modelId": bedrock_model,
|
||||
"messages": [{"role": "user", "content": [{"text": prompt}]}],
|
||||
"inferenceConfig": {"maxTokens": BEDROCK_MAX_TOKENS},
|
||||
}
|
||||
if not raw_text:
|
||||
kwargs["system"] = [{"text": _JSON_ONLY_SYSTEM}]
|
||||
|
||||
call_task = asyncio.create_task(asyncio.to_thread(client.converse, **kwargs))
|
||||
wait_tasks = [call_task]
|
||||
cancel_wait_task = None
|
||||
if cancel_event:
|
||||
cancel_wait_task = asyncio.create_task(cancel_event.wait())
|
||||
wait_tasks.append(cancel_wait_task)
|
||||
|
||||
done, pending = await asyncio.wait(
|
||||
wait_tasks, timeout=effective_timeout, return_when=asyncio.FIRST_COMPLETED
|
||||
)
|
||||
for p in pending:
|
||||
p.cancel()
|
||||
|
||||
if call_task not in done:
|
||||
# Der HTTP-Aufruf im Thread lässt sich nicht hart abbrechen, er läuft
|
||||
# aus und sein Ergebnis wird verworfen. Der Callback verhindert
|
||||
# "exception was never retrieved"-Warnungen.
|
||||
call_task.add_done_callback(
|
||||
lambda t: t.exception() if not t.cancelled() else None
|
||||
)
|
||||
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")
|
||||
|
||||
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
|
||||
|
||||
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()
|
||||
|
||||
stop_reason = response.get("stopReason", "")
|
||||
if stop_reason == "max_tokens":
|
||||
logger.warning(
|
||||
f"Bedrock-Antwort bei maxTokens={BEDROCK_MAX_TOKENS} abgeschnitten (Modell {cli_model})"
|
||||
)
|
||||
|
||||
u = response.get("usage", {})
|
||||
input_tokens = int(u.get("inputTokens", 0) or 0)
|
||||
output_tokens = int(u.get("outputTokens", 0) or 0)
|
||||
cache_write = int(u.get("cacheWriteInputTokens", 0) or 0)
|
||||
cache_read = int(u.get("cacheReadInputTokens", 0) or 0)
|
||||
usage = ClaudeUsage(
|
||||
input_tokens=input_tokens,
|
||||
output_tokens=output_tokens,
|
||||
cache_creation_tokens=cache_write,
|
||||
cache_read_tokens=cache_read,
|
||||
cost_usd=_cost_usd(cli_model, input_tokens, output_tokens, cache_write, cache_read),
|
||||
duration_ms=int(response.get("metrics", {}).get("latencyMs", 0) or 0),
|
||||
)
|
||||
logger.info(
|
||||
f"Bedrock [{cli_model} -> {bedrock_model}]: {usage.input_tokens} in / "
|
||||
f"{usage.output_tokens} out / cache {usage.cache_creation_tokens}+{usage.cache_read_tokens} / "
|
||||
f"${usage.cost_usd:.4f} / {usage.duration_ms}ms"
|
||||
)
|
||||
return _sanitize_mdash(result_text), usage
|
||||
@@ -4,12 +4,17 @@ import contextvars
|
||||
import json
|
||||
import logging
|
||||
from dataclasses import dataclass
|
||||
from config import CLAUDE_PATH, CLAUDE_TIMEOUT, CLAUDE_MODEL_FAST, CLAUDE_MODEL_STANDARD
|
||||
from config import CLAUDE_PATH, CLAUDE_TIMEOUT, CLAUDE_MODEL_FAST, CLAUDE_MODEL_STANDARD, AI_BACKEND
|
||||
|
||||
# ContextVar fuer Cancel-Event: Wird vom Orchestrator gesetzt,
|
||||
# call_claude prueft automatisch darauf -- kein Durchreichen noetig.
|
||||
_cancel_event_var: contextvars.ContextVar[asyncio.Event | None] = contextvars.ContextVar("_cancel_event_var", default=None)
|
||||
|
||||
# ContextVar für das KI-Backend des laufenden Refreshs ('cli' oder 'bedrock').
|
||||
# Der Orchestrator setzt es je Lage (Auflösung Lage vor Organisation vor
|
||||
# globalem Default AI_BACKEND). None = globaler Default aus config.
|
||||
_ai_backend_var: contextvars.ContextVar[str | None] = contextvars.ContextVar("_ai_backend_var", default=None)
|
||||
|
||||
logger = logging.getLogger("osint.claude_client")
|
||||
|
||||
|
||||
@@ -88,6 +93,17 @@ async def call_claude(prompt: str, tools: str | None = "WebSearch,WebFetch", mod
|
||||
model: Optionales Modell (z.B. CLAUDE_MODEL_FAST fuer Haiku). None = CLAUDE_MODEL_STANDARD (Opus 4.7).
|
||||
timeout: Override in Sekunden. None = Fallback auf globalen CLAUDE_TIMEOUT (1800s).
|
||||
"""
|
||||
backend = _ai_backend_var.get(None) or AI_BACKEND
|
||||
if backend == "bedrock":
|
||||
if tools:
|
||||
# Phase 1 des EU-Umbaus. Bedrock hat keine eingebaute Websuche,
|
||||
# Aufrufe mit Tools laufen bis zur staan-Schleife (Phase 2)
|
||||
# weiter über das CLI.
|
||||
logger.info(f"Backend 'bedrock' aktiv, Aufruf braucht aber Tools ({tools}) und läuft übers CLI (Phase 1)")
|
||||
else:
|
||||
from agents.bedrock_client import call_bedrock
|
||||
return await call_bedrock(prompt, model=model, raw_text=raw_text, timeout=timeout)
|
||||
|
||||
effective_model = model or CLAUDE_MODEL_STANDARD
|
||||
effective_timeout = timeout if timeout is not None else CLAUDE_TIMEOUT
|
||||
cmd = [CLAUDE_PATH, "-p", "-", "--output-format", "json", "--model", effective_model]
|
||||
|
||||
@@ -11,7 +11,7 @@ from urllib.parse import urlparse, urlunparse, quote_plus
|
||||
|
||||
import httpx
|
||||
|
||||
from agents.claude_client import UsageAccumulator, _cancel_event_var
|
||||
from agents.claude_client import UsageAccumulator, _cancel_event_var, _ai_backend_var
|
||||
from agents.factchecker import find_matching_claim, deduplicate_new_facts, TWOPHASE_MIN_FACTS
|
||||
from source_rules import (
|
||||
_detect_category,
|
||||
@@ -699,6 +699,7 @@ class AgentOrchestrator:
|
||||
self._cancel_events.pop(incident_id, None)
|
||||
self._cancel_requested.discard(incident_id)
|
||||
_cancel_event_var.set(None)
|
||||
_ai_backend_var.set(None)
|
||||
queue.task_done()
|
||||
|
||||
async def _mark_refresh_cancelled(self, incident_id: int):
|
||||
@@ -828,7 +829,21 @@ class AgentOrchestrator:
|
||||
from services.org_settings import (
|
||||
get_org_language, language_display, get_research_language,
|
||||
get_source_language_whitelist, get_translator_enabled,
|
||||
get_org_setting,
|
||||
)
|
||||
# KI-Backend für diesen Refresh auflösen (EU-Umbau Phase 1).
|
||||
# Lage gewinnt vor Organisation, die vor dem globalen Default.
|
||||
# Gilt über das ContextVar für alle call_claude-Aufrufe darunter.
|
||||
from config import AI_BACKEND
|
||||
incident_backend = incident["ai_backend"] if "ai_backend" in incident.keys() else None
|
||||
org_backend = await get_org_setting(db, tenant_id, "ai_backend", default=None) if tenant_id else None
|
||||
ai_backend = (incident_backend or org_backend or AI_BACKEND or "cli").strip().lower()
|
||||
if ai_backend not in ("cli", "bedrock"):
|
||||
logger.warning("Unbekanntes ai_backend '%s' für Lage %s, fallback 'cli'", ai_backend, incident_id)
|
||||
ai_backend = "cli"
|
||||
_ai_backend_var.set(ai_backend)
|
||||
if ai_backend != "cli":
|
||||
logger.info("Lage %s läuft über KI-Backend '%s'", incident_id, ai_backend)
|
||||
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").
|
||||
|
||||
In neuem Issue referenzieren
Einen Benutzer sperren