From a7f1aa29ab366f83617f46c49271773aaf91f829 Mon Sep 17 00:00:00 2001 From: claude-dev Date: Thu, 30 Jul 2026 20:12:31 +0000 Subject: [PATCH] feat(eu-umbau): staan-Recherche-Schleife und Doppelspur (Phase 2) Im EU-Modus laeuft die Recherche jetzt europaeisch: agents/eu_researcher.py steuert eine Schleife aus Bedrock-Planung (Opus, EU-Profile) und staan-Suchen (services/staan_client.py, Suche + Volltexte via full_content=markdown, Pflicht-Domain-Ausschlussliste je Anfrage). Bildet die 4-Phasen-Tiefenrecherche nach und liefert exakt das JSON der bisherigen CLI-WebSearch, Parsing und Filter in researcher.search bleiben unveraendert. Anti-Halluzination: nur URLs aus echten staan-Treffern werden akzeptiert; published_at nur wenn aus Inhalt oder URL ableitbar (staan liefert keine Daten). Doppelspur im Orchestrator: Lane-Schluessel ist jetzt (Organisation, Backend), Aufloesung zentral in _resolve_ai_backend (Lage > Org-Setting > ENV). CLI- und EU-Lauf derselben Organisation laufen gleichzeitig, gleiche Backends seriell. Co-Authored-By: Claude Fable 5 --- CLAUDE.md | 21 +- src/agents/eu_researcher.py | 363 +++++++++++++++++++++++++++++++++++ src/agents/orchestrator.py | 51 +++-- src/agents/researcher.py | 23 ++- src/config.py | 20 ++ src/services/staan_client.py | 131 +++++++++++++ 6 files changed, 589 insertions(+), 20 deletions(-) create mode 100644 src/agents/eu_researcher.py create mode 100644 src/services/staan_client.py diff --git a/CLAUDE.md b/CLAUDE.md index 25c35ed..8763617 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -31,13 +31,18 @@ backend: (EU-Inferenzprofile Frankfurt, agents/bedrock_client.py). Umschaltbar je Lage (incidents.ai_backend) > Organisation (org_setting 'ai_backend') > global (ENV AI_BACKEND, Default 'cli'). Aufloesung je Refresh im - Orchestrator via ContextVar _ai_backend_var. Phase 1 bedient nur - werkzeuglose Aufrufe (tools=None); Aufrufe mit WebSearch/WebFetch laufen - weiter uebers CLI, bis die staan-Recherche-Schleife (Phase 2) fertig ist. - Kosten werden aus Token x BEDROCK_PRICING (config.py) berechnet, damit - das Credits-System unveraendert funktioniert. AWS-Zugang + STAAN_API_KEY - liegen in der Staging-.env. boto3 muss im venv installiert sein (beim - Promote nach Live auch dort nachziehen!). + Orchestrator via ContextVar _ai_backend_var (gemeinsame Aufloesung mit der + Lane-Wahl, siehe _resolve_ai_backend). Seit Phase 2 (2026-07-30) laeuft im + EU-Modus auch die RECHERCHE europaeisch: agents/eu_researcher.py steuert + eine Schleife aus Bedrock-Planung und staan-Suchen (services/staan_client.py) + und liefert dasselbe JSON wie die CLI-WebSearch. Werkzeuglose Aufrufe gehen + direkt ueber bedrock_client; Faktenchecker/Analyzer nutzen im EU-Modus noch + die CLI-WebSearch (kommt in Phase 3). Doppelspur: eine Orchestrator-Lane je + (Organisation, Backend) — CLI- und EU-Lauf derselben Org laufen gleichzeitig, + gleiche Backends seriell. Kosten aus Token x BEDROCK_PRICING plus + STAAN_COST_PER_QUERY_USD, Credits-System unveraendert. AWS-Zugang + + STAAN_API_KEY liegen in der Staging-.env. boto3 muss im venv installiert + sein (beim Promote nach Live auch dort nachziehen!). ki_modelle: schnell: CLAUDE_MODEL_FAST (Haiku) — Feed-Selektion, Topic-Filter, Geoparsing, Uebersetzung, Chat, QC mittel: CLAUDE_MODEL_MEDIUM (Sonnet) — nur Netzwerkanalyse (entity_extractor hat aktuell keinen Aufrufer in der App) @@ -89,6 +94,7 @@ src/: 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" + 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/: rss_parser.py: "RSS-Feed-Parsing (feedparser + httpx), adaptive Keyword-Schwelle, Frische-Bonus, Domain-Cap" @@ -104,6 +110,7 @@ src/: source_suggester.py: "KI-Quellen-Vorschlaege via Haiku + Karteileichen-Heuristik" pdf_ingest.py: "Minutenjob: hochgeladene PDFs einlesen (pdfplumber + OCR-Fallback), uebersetzen, als Pool-Artikel ablegen" org_settings.py: "Key-Value-Einstellungen je Organisation (output_language etc.) mit 60s-Cache" + staan_client.py: "staan.ai Such-Client (EU-Umbau Phase 2): Suche + Volltexte (full_content=markdown) ueber den europaeischen Index, Pflicht-Domain-Ausschlussliste je Anfrage, max 10 Ausschluss-Domains" license_service.py: "Lizenz-Pruefung, Credits-Buchung (charge_usage_to_tenant), Periodenwechsel, Budget-Warnung. expire_licenses() existiert, hat aber KEINEN Scheduler-Job" middleware/: diff --git a/src/agents/eu_researcher.py b/src/agents/eu_researcher.py new file mode 100644 index 0000000..6dc6b9e --- /dev/null +++ b/src/agents/eu_researcher.py @@ -0,0 +1,363 @@ +"""EU-Recherche-Schleife. Gesteuerte Websuche über Bedrock plus staan (Phase 2). + +Ersetzt für den EU-Modellweg die eingebaute WebSearch-Recherche des Claude CLI. +Das Modell (Opus über Bedrock, EU-Profile Frankfurt) schlägt je Runde +Suchanfragen vor, unser Code fragt staan.ai, liefert Treffer und Volltexte +zurück, und das Modell entscheidet, ob es nachhaken will oder genug hat. +Am Ende produziert die Schleife exakt den JSON-Array-Text, den der heutige +WebSearch-Researcher liefert, sodass Parsing, Quellenfilter und die gesamte +nachgelagerte Pipeline unverändert weiterlaufen (Einstieg in +researcher.ResearcherAgent.search). + +Nachgebildet wird auch die vierstufige Tiefenrecherche der Recherche-Lagen +(breite Erfassung, Lückenanalyse, gezielte institutionelle Suche, Vertiefung +per Volltext). Zwei harte Zusagen dieser Schleife gegen Halluzinationen und +für die Aktualität. Erstens akzeptiert die Endauswertung nur URLs, die +wirklich in den staan-Treffern vorkamen (alles andere wird verworfen und +geloggt). Zweitens wird published_at nur gesetzt, wenn das Datum aus Inhalt +oder URL ableitbar ist, staan selbst liefert keine Datumsangaben. +""" +import asyncio +import json +import logging +from datetime import datetime + +from config import ( + CLAUDE_MODEL_STANDARD, + EU_RESEARCH_MAX_ROUNDS_ADHOC, + EU_RESEARCH_MAX_ROUNDS_RESEARCH, + STAAN_COST_PER_QUERY_USD, + TIMEZONE, +) +from agents.claude_client import ClaudeUsage, _cancel_event_var +from agents.bedrock_client import call_bedrock +from services.staan_client import staan_search, market_for_language, StaanError + +logger = logging.getLogger("osint.eu_researcher") + +MAX_QUERIES_PER_ROUND = 4 +MAX_FULLTEXT_QUERIES_PER_ROUND = 2 +MAX_TOTAL_RESULTS = 60 +# Kappung der Textmengen im Abschluss-Prompt (Kosten- und Kontextschutz) +FINAL_FULLTEXT_CHARS = 2500 +FINAL_SNIPPET_CHARS = 400 +ROUND_SNIPPET_CHARS = 250 + +_ROUND_HEADER = """Du steuerst eine OSINT-Recherche über eine europäische Such-API. +Du kannst NICHT selbst suchen. Du schlägst Suchanfragen vor, unser System führt sie aus +und zeigt dir die Treffer. Es ist Runde {round_no} von maximal {max_rounds}. +Heute ist der {today}. Formuliere Suchanfragen so, dass sie AKTUELLE Berichterstattung +finden (konkrete Ereignisse, Akteure, Ortsnamen; keine Jahreszahlen anhängen). + +AUFTRAG: +Titel: {title} +Kontext: {description} +{existing_context}{preferred_sources_block} +SPRACHREGELN: +{lang_instruction} + +{phase_guidance} + +BISHER GESAMMELTE TREFFER ({n_results} Stück): +{results_block} + +VERLAUF DEINER BISHERIGEN SUCHEN: +{rounds_log} + +Antworte NUR mit einem JSON-Objekt in genau diesem Format: +{{"action": "search", "queries": [{{"q": "suchbegriffe", "market": "de-de", "full_content": false}}], "reason": "ein Satz"}} +oder, wenn die Treffer für ein vollständiges Bild reichen: +{{"action": "finish", "reason": "ein Satz"}} + +REGELN FÜR QUERIES: +- Maximal {max_queries} Queries je Runde, jede maximal 400 Zeichen. +- "market" ist "de-de", "en-us" oder "fr-fr" (Standard für diese Lage: "{default_market}"). + Fremdsprachige Suchanfragen (z.B. Russisch, Arabisch, Farsi) funktionieren mit "en-us". +- "full_content": true holt die kompletten Artikeltexte der Treffer dieser Query + (maximal {max_fulltext} Queries je Runde). Nutze das für die wichtigsten Suchen, + deren Artikel du zusammenfassen willst. +- Wiederhole keine Query aus dem Verlauf wortgleich.""" + +_PHASE_GUIDANCE_RESEARCH = { + 1: """PHASE 1, BREITE ERFASSUNG: +Suche nach aktueller Berichterstattung bei Nachrichtenagenturen, Qualitätszeitungen und +öffentlich-rechtlichen Medien. Nutze verschiedene Suchbegriffe und Blickwinkel.""", + 2: """PHASE 2 UND 3, LÜCKENANALYSE UND GEZIELTE TIEFENRECHERCHE: +Prüfe die bisherigen Treffer kritisch. Welche Quellentypen fehlen? Typisch fehlen +Parlamentsdokumente, Behörden-Pressemitteilungen, NGO- und UN-Berichte (ohchr.org, +amnesty.org, hrw.org), Think-Tank-Analysen (IISS, Brookings, SWP, DGAP, Chatham House), +investigative Langform-Berichte und Fachmedien. Suche GEZIELT nach diesen Lücken, +auch mit site:-Operatoren für institutionelle Quellen.""", + 3: """PHASE 4, VERIFIKATION UND VERTIEFUNG: +Fordere für die wichtigsten Suchanfragen jetzt Volltexte an ("full_content": true), +damit die Artikel ausführlich zusammengefasst werden können. Priorisiere Primärquellen +und investigative Berichte. Ergänze nur noch gezielt, was für ein vollständiges Bild fehlt.""", +} + +_PHASE_GUIDANCE_ADHOC = { + 1: """SCHWERPUNKT DIESER RUNDE: +Aktuelle Berichterstattung zur Lage bei seriösen Nachrichtenquellen (Agenturen, +Qualitätszeitungen, öffentlich-rechtliche Medien, Behörden). Verschiedene Blickwinkel.""", + 2: """SCHWERPUNKT DIESER RUNDE: +Lücken schließen und vertiefen. Fordere für die wichtigsten Suchanfragen Volltexte an +("full_content": true), damit die Artikel fundiert zusammengefasst werden können.""", +} + +_FINAL_PROMPT = """Du bist ein OSINT-Recherche-Agent. Die Recherche über die europäische Such-API ist +abgeschlossen. Unten stehen ALLE gesammelten Treffer mit Textauszügen. Heute ist der {today}. + +AUSGABESPRACHE: {output_language} +- KEINE Gedankenstriche verwenden, stattdessen Kommas oder neue Sätze. +- Verwende IMMER echte UTF-8-Umlaute (ä, ö, ü, ß), NIEMALS Umschreibungen. + +AUFTRAG WAR: +Titel: {title} +Kontext: {description} + +WÄHLE aus den Treffern die {target} relevantesten, inhaltlich substanziellen Artikel aus. +Bevorzuge Vielfalt der Quellen und Blickwinkel. Lass Übersichtsseiten, reine Linklisten +und thematisch unpassende Treffer weg. + +GESAMMELTE TREFFER: +{results_block} + +Gib die Ergebnisse AUSSCHLIESSLICH als JSON-Array zurück, ohne Text davor oder danach. +Jedes Element hat diese Felder: +- "headline": Originale Überschrift (aus dem Treffer) +- "headline_de": Übersetzung in die Ausgabesprache (falls Originalsprache abweicht) +- "source": Name der Quelle (z.B. "Reuters", "tagesschau") +- "source_url": Die EXAKTE URL aus dem Treffer. NIEMALS eine URL verändern oder erfinden. +- "content_summary": Zusammenfassung des Inhalts ({summary_len}, in Ausgabesprache, auf Basis der Textauszüge) +- "language": Sprache des Originals (z.B. "de", "en", "fa") +- "published_at": Veröffentlichungsdatum im ISO-Format, NUR wenn es aus Textauszug oder URL + eindeutig hervorgeht (z.B. Datumsangabe im Text oder /2026/07/ im Pfad), sonst null. + +Antworte NUR mit dem JSON-Array.""" + + +def _norm_url(u: str) -> str: + return (u or "").strip().rstrip("/").lower() + + +def _results_block(collected: dict[str, dict], with_texts: bool) -> str: + """Nummerierte Trefferliste für die Prompts. Kompakt in Zwischenrunden, + mit Textauszügen im Abschluss-Prompt.""" + if not collected: + return "(noch keine)" + lines = [] + for i, (url, r) in enumerate(collected.items(), 1): + lines.append(f"{i}. {r['title'] or '(ohne Titel)'} | {r['hostname']}\n URL: {url}") + if with_texts: + if r.get("snippet"): + lines.append(f" Auszug: {r['snippet'][:FINAL_SNIPPET_CHARS]}") + for chunk in (r.get("extra_snippets") or [])[:2]: + lines.append(f" Auszug: {chunk[:FINAL_SNIPPET_CHARS]}") + if r.get("full_text"): + lines.append(f" Volltext (gekürzt): {r['full_text'][:FINAL_FULLTEXT_CHARS]}") + else: + if r.get("snippet"): + lines.append(f" {r['snippet'][:ROUND_SNIPPET_CHARS]}") + return "\n".join(lines) + + +def _parse_action(text: str) -> dict: + """Aktions-JSON des Modells robust parsen. Fallback = finish.""" + from agents.researcher import _extract_json_object + obj = None + try: + parsed = json.loads((text or "").strip()) + if isinstance(parsed, dict): + obj = parsed + except (json.JSONDecodeError, ValueError): + pass + if obj is None: + obj = _extract_json_object(text or "") + if not isinstance(obj, dict) or obj.get("action") not in ("search", "finish"): + logger.warning("EU-Recherche: Aktions-JSON nicht parsebar, beende Suchphase. Sample: %r", (text or "")[:200]) + return {"action": "finish", "queries": []} + queries = [] + for q in obj.get("queries") or []: + if isinstance(q, dict) and str(q.get("q", "")).strip(): + queries.append({ + "q": str(q["q"]).strip()[:400], + "market": str(q.get("market", "")).strip().lower(), + "full_content": bool(q.get("full_content")), + }) + obj["queries"] = queries[:MAX_QUERIES_PER_ROUND] + return obj + + +async def run_eu_research( + *, + title: str, + description: str, + incident_type: str, + lang_instruction: str, + existing_context: str, + preferred_sources_block: str, + output_language: str, + excluded_sources: list[str], + research_language_iso: str, +) -> tuple[str, ClaudeUsage]: + """Führt die komplette EU-Recherche aus. Gibt (json_array_text, usage) zurück. + + Der Rückgabetext wird vom bestehenden ResearcherAgent._parse_response + weiterverarbeitet, der Vertrag entspricht damit exakt dem CLI-Weg. + """ + total = ClaudeUsage() + + def _add(u: ClaudeUsage): + total.input_tokens += u.input_tokens + total.output_tokens += u.output_tokens + total.cache_creation_tokens += u.cache_creation_tokens + total.cache_read_tokens += u.cache_read_tokens + total.cost_usd += u.cost_usd + total.duration_ms += u.duration_ms + + is_research = incident_type == "research" + max_rounds = EU_RESEARCH_MAX_ROUNDS_RESEARCH if is_research else EU_RESEARCH_MAX_ROUNDS_ADHOC + guidance_map = _PHASE_GUIDANCE_RESEARCH if is_research else _PHASE_GUIDANCE_ADHOC + today = datetime.now(TIMEZONE).strftime("%d.%m.%Y") + default_market = market_for_language(research_language_iso) + # Nutzer-Ausschlüsse an staan weiterreichen (die Pflichtliste belegt die + # ersten Plätze, siehe staan_client). Zusätzlich filtert der Aufrufer + # (ResearcherAgent.search) das Endergebnis erneut gegen die volle Liste. + extra_exclude = [d for d in (excluded_sources or []) if d and "." in d][:6] + + collected: dict[str, dict] = {} + rounds_log: list[str] = [] + searches_done = 0 + + for round_no in range(1, max_rounds + 1): + cancel = _cancel_event_var.get(None) + if cancel and cancel.is_set(): + raise asyncio.CancelledError("Cancel angefordert") + + guidance = guidance_map.get(min(round_no, max(guidance_map.keys()))) + prompt = _ROUND_HEADER.format( + round_no=round_no, + max_rounds=max_rounds, + today=today, + title=title, + description=description or "Keine weitere Beschreibung", + existing_context=existing_context or "", + preferred_sources_block=preferred_sources_block or "", + lang_instruction=lang_instruction, + phase_guidance=guidance, + n_results=len(collected), + results_block=_results_block(collected, with_texts=False), + rounds_log="\n".join(rounds_log) or "(noch keine)", + max_queries=MAX_QUERIES_PER_ROUND, + max_fulltext=MAX_FULLTEXT_QUERIES_PER_ROUND, + default_market=default_market, + ) + text, usage = await call_bedrock(prompt, model=CLAUDE_MODEL_STANDARD, timeout=420) + _add(usage) + + action = _parse_action(text) + if action["action"] == "finish": + if round_no == 1 and not collected: + logger.warning("EU-Recherche: Modell wollte ohne eine einzige Suche beenden") + else: + logger.info("EU-Recherche: Suchphase nach Runde %d beendet (%s)", round_no - 1, action.get("reason", "")) + break + + fulltext_used = 0 + for q in action["queries"]: + if len(collected) >= MAX_TOTAL_RESULTS: + logger.info("EU-Recherche: Treffer-Obergrenze %d erreicht", MAX_TOTAL_RESULTS) + break + want_full = q["full_content"] and fulltext_used < MAX_FULLTEXT_QUERIES_PER_ROUND + if want_full: + fulltext_used += 1 + market = q["market"] if q["market"] else default_market + try: + results = await staan_search( + q["q"], market=market, extra_snippets=True, max_snippets=4, + full_content=want_full, extra_exclude=extra_exclude, + ) + searches_done += 1 + except StaanError as e: + logger.warning("EU-Recherche: Suche fehlgeschlagen (%s), überspringe Query", e) + rounds_log.append(f"Runde {round_no}: '{q['q'][:70]}' -> FEHLER") + continue + new_count = 0 + for r in results: + if not r["url"]: + continue + key = r["url"] + if key in collected: + # Volltext/Snippets aus späteren Suchen ergänzen + if r.get("full_text") and not collected[key].get("full_text"): + collected[key]["full_text"] = r["full_text"] + continue + if len(collected) >= MAX_TOTAL_RESULTS: + break + collected[key] = r + new_count += 1 + rounds_log.append( + f"Runde {round_no}: '{q['q'][:70]}' ({market}{', Volltexte' if want_full else ''}) -> {new_count} neue Treffer" + ) + + if not action["queries"]: + break + + total.cost_usd += searches_done * STAAN_COST_PER_QUERY_USD + + if not collected: + logger.warning("EU-Recherche: keine Treffer gesammelt (%d Suchen)", searches_done) + return "[]", total + + target = "15 bis 25" if is_research else "8 bis 15" + summary_len = "5 bis 8 Sätze" if is_research else "3 bis 5 Sätze" + final_prompt = _FINAL_PROMPT.format( + today=today, + output_language=output_language, + title=title, + description=description or "Keine weitere Beschreibung", + target=target, + summary_len=summary_len, + results_block=_results_block(collected, with_texts=True), + ) + text, usage = await call_bedrock(final_prompt, model=CLAUDE_MODEL_STANDARD, timeout=600) + _add(usage) + + from agents.researcher import _extract_json_array + arr = None + try: + parsed = json.loads((text or "").strip()) + if isinstance(parsed, list): + arr = parsed + except (json.JSONDecodeError, ValueError): + pass + if arr is None: + arr = _extract_json_array(text or "") + if not isinstance(arr, list): + # Unparsebarer Abschluss: Rohtext zurückgeben, damit der bestehende + # Recovery-Pfad in _parse_response greift und die Usage verbucht wird. + logger.warning("EU-Recherche: Abschluss-JSON nicht parsebar, reiche Rohtext weiter") + return text or "[]", total + + # Anti-Halluzination. Nur URLs akzeptieren, die wirklich in den staan-Treffern + # vorkamen. published_at, das nicht wie ein Datum aussieht, wird verworfen. + valid = {_norm_url(u) for u in collected} + validated = [] + dropped = 0 + for item in arr: + if not isinstance(item, dict): + continue + if _norm_url(item.get("source_url", "")) not in valid: + dropped += 1 + continue + pub = item.get("published_at") + if pub is not None and not isinstance(pub, str): + item["published_at"] = None + validated.append(item) + if dropped: + logger.warning("EU-Recherche: %d Artikel mit nicht belegten URLs verworfen", dropped) + + logger.info( + "EU-Recherche fertig: %d Artikel aus %d Treffern, %d Suchen, %d Modellrunden, $%.4f", + len(validated), len(collected), searches_done, round_no + 1, total.cost_usd, + ) + return json.dumps(validated, ensure_ascii=False), total diff --git a/src/agents/orchestrator.py b/src/agents/orchestrator.py index 887982c..fa060c2 100644 --- a/src/agents/orchestrator.py +++ b/src/agents/orchestrator.py @@ -470,7 +470,12 @@ class AgentOrchestrator: return False visibility, created_by, tenant_id = await self._get_incident_visibility(incident_id) - lane_key = tenant_id if tenant_id else self.PUBLIC_LANE + # Doppelspur (EU-Umbau Phase 2). Eine Lane je Organisation UND KI-Weg. + # Ein CLI-Lauf und ein EU-Lauf derselben Organisation laufen dadurch + # wirklich gleichzeitig (getrennte Ressourcen, Abo gegen AWS), während + # Läufe über denselben Weg seriell bleiben (Rate-Limit-Schutz). + ai_backend = await self._resolve_ai_backend(incident_id) + lane_key = (tenant_id if tenant_id else self.PUBLIC_LANE, ai_backend) async with self._lanes_lock: queue = self._lanes.get(lane_key) @@ -792,6 +797,35 @@ class AgentOrchestrator: await db.close() return visibility, created_by, tenant_id + async def _resolve_ai_backend(self, incident_id: int) -> str: + """Löst das KI-Backend einer Lage auf (EU-Umbau). + + Lage (incidents.ai_backend) gewinnt vor dem Org-Setting 'ai_backend', + das vor dem globalen Default AI_BACKEND aus config. Wird beim Einreihen + (Lane-Wahl) und beim Refresh-Start (ContextVar) identisch genutzt.""" + from config import AI_BACKEND + from database import get_db + from services.org_settings import get_org_setting + incident_backend = None + tenant_id = None + db = await get_db() + try: + cursor = await db.execute( + "SELECT ai_backend, tenant_id FROM incidents WHERE id = ?", (incident_id,) + ) + row = await cursor.fetchone() + if row: + incident_backend = row["ai_backend"] if "ai_backend" in row.keys() else None + tenant_id = row["tenant_id"] if "tenant_id" in row.keys() else None + org_backend = await get_org_setting(db, tenant_id, "ai_backend", default=None) if tenant_id else None + finally: + await db.close() + 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" + return ai_backend + async def _run_refresh(self, incident_id: int, trigger_type: str = "manual", retry_count: int = 0, user_id: int = None, _suppress_complete: bool = False, _pass_info: dict = None, collect_only: bool = False): """Führt einen kompletten Refresh-Zyklus durch.""" import aiosqlite @@ -829,18 +863,11 @@ 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" + # KI-Backend für diesen Refresh auflösen und per ContextVar setzen + # (EU-Umbau, gemeinsame Auflösung mit der Lane-Wahl beim Einreihen, + # siehe _resolve_ai_backend). + ai_backend = await self._resolve_ai_backend(incident_id) _ai_backend_var.set(ai_backend) if ai_backend != "cli": logger.info("Lage %s läuft über KI-Backend '%s'", incident_id, ai_backend) diff --git a/src/agents/researcher.py b/src/agents/researcher.py index 70e4786..705342d 100644 --- a/src/agents/researcher.py +++ b/src/agents/researcher.py @@ -804,7 +804,28 @@ class ResearcherAgent: ) try: - result, usage = await call_claude(prompt) + # EU-Modellweg (Phase 2). Läuft der Refresh über das Bedrock-Backend, + # übernimmt die gesteuerte staan-Schleife die komplette Recherche und + # liefert denselben JSON-Array-Text wie die CLI-WebSearch. Parsing, + # Quellenfilter und Sprachfilter darunter bleiben unverändert. + from agents.claude_client import _ai_backend_var + from config import AI_BACKEND + _backend = _ai_backend_var.get(None) or AI_BACKEND + if _backend == "bedrock": + from agents.eu_researcher import run_eu_research + result, usage = await run_eu_research( + title=title, + description=description, + incident_type=incident_type, + lang_instruction=lang_instruction, + existing_context=existing_context, + preferred_sources_block=preferred_sources_block, + output_language=output_language, + excluded_sources=await self._get_excluded_sources(user_id=user_id), + research_language_iso=research_language_iso or output_language_iso, + ) + else: + result, usage = await call_claude(prompt) try: articles = self._parse_response(result) except ResearcherParseError as parse_err: diff --git a/src/config.py b/src/config.py index 1fe0c29..cdb628f 100644 --- a/src/config.py +++ b/src/config.py @@ -66,6 +66,26 @@ BEDROCK_PRICING = { BEDROCK_CACHE_WRITE_FACTOR = 1.25 BEDROCK_CACHE_READ_FACTOR = 0.10 +# --- staan.ai Such-API (EU/DSGVO-Umbau, Phase 2) ------------------------------- +# Europäischer Suchindex (European Search Perspective). Ersetzt im EU-Modellweg +# die eingebaute WebSearch/WebFetch des Claude CLI. Client in +# services/staan_client.py, Recherche-Schleife in agents/eu_researcher.py. +STAAN_API_KEY = os.environ.get("STAAN_API_KEY", "") +STAAN_BASE_URL = os.environ.get("STAAN_BASE_URL", "https://api.staan.ai/v2") +# 2 Euro je 1000 Anfragen, als USD-Näherung für die interne Kostenkontrolle. +STAAN_COST_PER_QUERY_USD = float(os.environ.get("STAAN_COST_PER_QUERY_USD", "0.0023")) +# Pflicht-Ausschlussliste je Suchanfrage (Befund aus Phase 0, hebt die +# Trefferqualität deutlich). staan erlaubt maximal 10 Domains je Anfrage, +# Nutzer-Ausschlüsse füllen die restlichen Plätze auf. +STAAN_EXCLUDE_DOMAINS = [ + "youtube.com", "wikipedia.org", "facebook.com", "reddit.com", + "pinterest.com", "tiktok.com", "instagram.com", "x.com", +] +# Rundenbegrenzung der EU-Recherche-Schleife (Planungsaufrufe an das Modell, +# ohne die Abschlussrunde). research bildet die 4-Phasen-Tiefenrecherche nach. +EU_RESEARCH_MAX_ROUNDS_ADHOC = int(os.environ.get("EU_RESEARCH_MAX_ROUNDS_ADHOC", "3")) +EU_RESEARCH_MAX_ROUNDS_RESEARCH = int(os.environ.get("EU_RESEARCH_MAX_ROUNDS_RESEARCH", "5")) + # --- Verbrauchssaetze (Credits je Aktion) ------------------------------------- # Beschlossen 2026-07-23, der Verbrauchsrechner im Verwaltungsportal nutzt # dieselben Werte: Live-Lauf 45, Recherche je Durchlauf 40 (das Anlegen faehrt diff --git a/src/services/staan_client.py b/src/services/staan_client.py new file mode 100644 index 0000000..fd52ed0 --- /dev/null +++ b/src/services/staan_client.py @@ -0,0 +1,131 @@ +"""staan.ai Such-Client. Europäische Websuche für den EU-Modellweg (Phase 2). + +Ein Baustein für Suche UND Volltextabruf. Der Parameter full_content="markdown" +liefert die kompletten Seiteninhalte der Treffer gleich mit, damit braucht die +EU-Recherche-Schleife keinen eigenen Abruf einzelner Seiten. Die Verarbeitung +bleibt vollständig beim europäischen Anbieter (European Search Perspective, +gehostet bei OVHcloud). + +Befund aus Phase 0. Die API liefert praktisch nie ein Veröffentlichungsdatum +und kennt keinen Zeitraum-Filter. Die Aktualität sichern deshalb die +Query-Formulierung (macht die Schleife) und die Datumsextraktion aus den +Inhalten (macht das Modell). Der Domain-Ausschlussfilter ist Pflichtbestandteil +jeder Anfrage, ohne ihn bestehen die Trefferlisten spürbar aus YouTube, +Wikipedia und Social-Media-Seiten. +""" +import asyncio +import logging + +import httpx + +from config import ( + STAAN_API_KEY, + STAAN_BASE_URL, + STAAN_EXCLUDE_DOMAINS, +) + +logger = logging.getLogger("osint.staan_client") + +# Von staan offiziell unterstützte Märkte. Andere Werte fallen auf en-us zurück, +# der Index liefert auch dann fremdsprachige Treffer (in Phase 0 belegt für +# Russisch, Ukrainisch, Arabisch, Farsi, Hebräisch und Japanisch). +ALLOWED_MARKETS = {"de-de", "en-us", "fr-fr"} + +# staan akzeptiert maximal 10 Ausschluss-Domains je Anfrage. Die Pflichtliste +# aus config belegt die ersten Plätze, Nutzer-Ausschlüsse füllen bis 10 auf. +_MAX_EXCLUDE_DOMAINS = 10 + + +class StaanError(RuntimeError): + """Fehler der staan-API (Status, Auth, Netz).""" + + +def market_for_language(lang_iso: str | None) -> str: + """ISO-Sprachcode auf den passendsten staan-Markt abbilden.""" + mapping = {"de": "de-de", "en": "en-us", "fr": "fr-fr"} + return mapping.get((lang_iso or "").lower().strip(), "en-us") + + +async def staan_search( + q: str, + market: str = "en-us", + extra_snippets: bool = True, + max_snippets: int = 5, + full_content: bool = False, + extra_exclude: list[str] | None = None, + timeout: float = 60.0, +) -> list[dict]: + """Eine Suchanfrage gegen staan. Gibt normalisierte Treffer zurück. + + Rückgabe je Treffer: title, url, hostname, snippet, extra_snippets (Liste + von Text-Chunks), full_text (Markdown, gekappt), published_date (fast + immer None, siehe Modul-Docstring). + """ + if not STAAN_API_KEY: + raise StaanError("STAAN_API_KEY ist nicht gesetzt, EU-Suche nicht nutzbar") + + if market not in ALLOWED_MARKETS: + market = "en-us" + + exclude = list(STAAN_EXCLUDE_DOMAINS) + for d in extra_exclude or []: + d = (d or "").strip().lower() + if not d or d in exclude: + continue + if len(exclude) >= _MAX_EXCLUDE_DOMAINS: + break + exclude.append(d) + + body: dict = {"q": (q or "").strip()[:400], "market": market, "exclude_domains": exclude} + if extra_snippets: + body["extra_snippets"] = True + body["max_snippets"] = max(1, min(int(max_snippets), 10)) + if full_content: + body["full_content"] = "markdown" + + headers = {"Authorization": f"Bearer {STAAN_API_KEY}", "Content-Type": "application/json"} + + data = None + async with httpx.AsyncClient(timeout=timeout) as client: + for attempt in (1, 2): + try: + resp = await client.post(f"{STAAN_BASE_URL}/search/web", json=body, headers=headers) + except httpx.TimeoutException as e: + raise StaanError(f"staan Timeout nach {timeout}s ({e})") from e + except httpx.HTTPError as e: + raise StaanError(f"staan Netzwerkfehler ({type(e).__name__}: {e})") from e + if resp.status_code == 429 and attempt == 1: + # Rate-Limit (20 req/s), kurz warten und einmal wiederholen + await asyncio.sleep(1.5) + continue + if resp.status_code not in (200, 201): + # GET antwortet mit 200, POST mit 201 (Created) + raise StaanError(f"staan HTTP {resp.status_code}: {resp.text[:200]}") + try: + data = resp.json() + except ValueError as e: + raise StaanError(f"staan lieferte kein JSON ({e})") from e + break + + results: list[dict] = [] + for r in (data.get("web") or {}).get("results", []): + fc = r.get("full_content") or {} + results.append({ + "title": (r.get("title") or "").strip(), + "url": (r.get("url") or "").strip(), + "hostname": (r.get("hostname") or "").strip(), + "snippet": (r.get("snippet") or "").strip(), + "extra_snippets": [ + (c.get("chunk") or "").strip() + for c in (r.get("extra_snippets") or []) + if (c.get("chunk") or "").strip() + ], + "full_text": ((fc.get("text") or "").strip())[:6000], + "published_date": r.get("published_date"), + }) + + logger.info( + "staan [%s] %r -> %d Treffer%s", + market, body["q"][:60], len(results), " mit Volltexten" if full_content else "", + ) + return results