- JSONFormatter, ContextVars, set_research_run_id() with short logger names - self-built Counter/Histogram/Gauge system (no external deps) - Prometheus text export at /metrics - Request logging middleware with X-Request-ID - Metrics instrumentation: search_queries, sources_fetched, claims, contradictions - Histograms: research_duration, llm_request_duration - Gauge: active_research_runs - 13 + 21 = 34 tests
987 lines
36 KiB
Python
987 lines
36 KiB
Python
"""Research Orchestrator — Hauptklasse für Stage 12.
|
|
|
|
Orchestriert die gesamte Research-Pipeline (PLANNING → SEARCHING → FETCHING →
|
|
EXTRACTING → ANALYZING → COMPARING → SYNTHESIZING → COMPLETED) mit Budget-
|
|
Tracking, State-Machine-Validierung und Fallback-Logik.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import json
|
|
import logging
|
|
import time
|
|
from datetime import datetime, timezone
|
|
from typing import Any
|
|
from uuid import UUID, uuid4
|
|
|
|
from nsct.agents.planner import MockResearchPlanner, ResearchPlanner
|
|
from nsct.config import AppSettings
|
|
from nsct.crawler.manager import CrawlerManager
|
|
from nsct.logging_config import set_research_run_id, get_research_run_id, get_logger
|
|
from nsct.metrics import (
|
|
metrics,
|
|
C_SEARCH_QUERIES_TOTAL,
|
|
C_SOURCES_FETCHED_TOTAL,
|
|
C_CLAIMS_EXTRACTED_TOTAL,
|
|
C_CONTRADICTIONS_DETECTED_TOTAL,
|
|
H_RESEARCH_DURATION,
|
|
H_LLM_REQUEST_DURATION,
|
|
)
|
|
from nsct.models.claim import Claim
|
|
from nsct.orchestration.budget import BudgetExhaustedError, BudgetTracker, HardBudgetConfig
|
|
from nsct.orchestration.models import ResearchRun
|
|
from nsct.orchestration.state import ResearchRunState, StateMachine
|
|
from nsct.providers.abstract import MultiProviderSearch, SearchProvider
|
|
|
|
try:
|
|
from nsct.stages.stage9_synthesis import SynthesisStage
|
|
except ImportError:
|
|
SynthesisStage = None # type: ignore[misc,assignment]
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
class ResearchOrchestrator:
|
|
"""Haupt-Orchestrator für den NSCT Research-Pipeline (Stage 12).
|
|
|
|
Führt die gesamte Pipeline von PLANNING bis COMPLETED sequenziell
|
|
oder schrittweise (run_step) aus. Trackt Budget, validiert States und
|
|
bietet Fallback-Logik für fehlende Dienste.
|
|
"""
|
|
|
|
# ---------------------------------------------------------------
|
|
# Lifecycle
|
|
# ---------------------------------------------------------------
|
|
|
|
def __init__(
|
|
self,
|
|
config: AppSettings,
|
|
research_id: UUID,
|
|
query: str,
|
|
budget_config: HardBudgetConfig | None = None,
|
|
depth: str = "normal",
|
|
) -> None:
|
|
"""Initialisiere den Orchestrator.
|
|
|
|
Parameters
|
|
----------
|
|
config : AppSettings
|
|
Zentrale Anwendungskonfiguration.
|
|
research_id : UUID
|
|
Parent-Research-ID (gruppierung).
|
|
query : str
|
|
Die Forschungsfrage / Query.
|
|
budget_config : HardBudgetConfig | None
|
|
Optionales Budget. Wird aus config abgeleitet, wenn None.
|
|
depth : str
|
|
Suchtiefe ("quick", "normal", "deep").
|
|
"""
|
|
self._config = config
|
|
self._research_id = research_id
|
|
self._query = query
|
|
self._depth = depth
|
|
|
|
# Budget
|
|
if budget_config is not None:
|
|
self._budget_config = budget_config
|
|
else:
|
|
self._budget_config = HardBudgetConfig()
|
|
self._budget_tracker = BudgetTracker(self._budget_config)
|
|
|
|
# State Machine & Run
|
|
self._state_machine = StateMachine(ResearchRunState.CREATED)
|
|
self._run: ResearchRun | None = None
|
|
|
|
# Pipeline-Zwischenspeicher
|
|
self._plan: dict[str, Any] | None = None
|
|
self._search_results: list[dict[str, Any]] = [] # URLs
|
|
self._sources: list[dict[str, Any]] = []
|
|
self._claims: list[Claim] = []
|
|
|
|
# Sub-Components (lazy init)
|
|
self._planner: ResearchPlanner | MockResearchPlanner | None = None
|
|
self._crawler: CrawlerManager | None = None
|
|
self._multi_search: MultiProviderSearch | None = None
|
|
self._llm_provider = None
|
|
|
|
# ---------------------------------------------------------------
|
|
# Properties
|
|
# ---------------------------------------------------------------
|
|
|
|
@property
|
|
def state(self) -> ResearchRunState:
|
|
"""Aktueller Zustand der State Machine."""
|
|
return self._state_machine.current_state
|
|
|
|
@property
|
|
def run_id(self) -> UUID:
|
|
"""UUID des aktuellen Research-Runs."""
|
|
if self._run is None:
|
|
return uuid4()
|
|
return self._run.id
|
|
|
|
@property
|
|
def budget_tracker(self) -> BudgetTracker:
|
|
"""BudgetTracker des aktuellen Runs."""
|
|
return self._budget_tracker
|
|
|
|
@property
|
|
def is_completed(self) -> bool:
|
|
"""True, wenn State == COMPLETED."""
|
|
return self._state_machine.current_state == ResearchRunState.COMPLETED
|
|
|
|
@property
|
|
def is_running(self) -> bool:
|
|
"""True, wenn sich die Maschine in einem nicht-terminalen State befindet."""
|
|
return not self._state_machine.is_terminal()
|
|
|
|
# ---------------------------------------------------------------
|
|
# Private: State Machine
|
|
# ---------------------------------------------------------------
|
|
|
|
def _transition_to(self, state_name: str) -> bool:
|
|
"""Transition via StateMachine + Logging.
|
|
|
|
Parameters
|
|
----------
|
|
state_name : str
|
|
Zielzustand als String (z.B. "planning").
|
|
|
|
Returns
|
|
-------
|
|
bool
|
|
True wenn Übergang erfolgreich, False sonst.
|
|
"""
|
|
result = self._state_machine.transition(state_name)
|
|
if result:
|
|
logger.info("State transition: %s -> %s", self._state_machine.current_state.value, state_name)
|
|
else:
|
|
logger.warning("State transition %s -> %s rejected", self._state_machine.current_state.value, state_name)
|
|
return result
|
|
|
|
def _mark_failed(self, reason: str) -> bool:
|
|
"""Setze State auf FAILED + logge Error.
|
|
|
|
Fällt auf die State Machine zurück, wenn der direkte Übergang
|
|
zum aktuellen State erlaubt ist; sonst erzwingt er FAILED.
|
|
|
|
Parameters
|
|
----------
|
|
reason : str
|
|
Fehlerbeschreibung.
|
|
|
|
Returns
|
|
-------
|
|
bool
|
|
True wenn Zustand erfolgreich auf FAILED gesetzt wurde.
|
|
"""
|
|
logger.error("Research failed: %s (state=%s)", reason, self._state_machine.current_state.value)
|
|
|
|
# Der State Machine _apply erzwingt FAILED direkt,
|
|
# weil _mark_failed als "Sondertransition" gedacht ist.
|
|
if self._state_machine.current_state in (
|
|
ResearchRunState.COMPLETED,
|
|
ResearchRunState.FAILED,
|
|
ResearchRunState.CANCELLED,
|
|
):
|
|
return False
|
|
|
|
self._state_machine.current_state = ResearchRunState.FAILED # type: ignore[assignment]
|
|
return True
|
|
|
|
# ---------------------------------------------------------------
|
|
# Private: Budget
|
|
# ---------------------------------------------------------------
|
|
|
|
def _check_budget(self) -> None:
|
|
"""Budget-Check vor jedem LLM-/Download-Aufruf.
|
|
|
|
Raises
|
|
------
|
|
BudgetExhaustedError
|
|
Wenn ein Limit erreicht ist.
|
|
"""
|
|
self._budget_tracker.check_budget()
|
|
|
|
# ---------------------------------------------------------------
|
|
# Private: Sub-Component Init
|
|
# ---------------------------------------------------------------
|
|
|
|
def _get_planner(self) -> ResearchPlanner | MockResearchPlanner:
|
|
"""Planner instantiieren (LLM) oder Fallback (Mock)."""
|
|
if self._planner is not None:
|
|
return self._planner
|
|
|
|
try:
|
|
from nsct.providers.llm import get_provider
|
|
from nsct.providers.metrics import ProviderMetrics
|
|
|
|
metrics = ProviderMetrics()
|
|
llm_prov = get_provider(self._config, metrics)
|
|
self._planner = ResearchPlanner(config=self._config, llm_provider=llm_prov, metrics=metrics)
|
|
logger.info("Using LLM-based ResearchPlanner")
|
|
except Exception as exc:
|
|
logger.warning("LLM Planner unavailable (%s) — falling back to MockResearchPlanner", exc)
|
|
self._planner = MockResearchPlanner(config=self._config)
|
|
return self._planner
|
|
|
|
def _get_crawler(self) -> CrawlerManager:
|
|
"""CrawlerManager instantiieren oder Fallback."""
|
|
if self._crawler is not None:
|
|
return self._crawler
|
|
|
|
try:
|
|
self._crawler = CrawlerManager()
|
|
except Exception as exc:
|
|
logger.warning("Crawler init failed (%s) — crawler will return empty docs", exc)
|
|
self._crawler = CrawlerManager()
|
|
return self._crawler
|
|
|
|
def _get_multi_search(self) -> MultiProviderSearch:
|
|
"""MultiProviderSearch mit den konfigurierten Providern."""
|
|
if self._multi_search is not None:
|
|
return self._multi_search
|
|
|
|
providers: list[SearchProvider] = []
|
|
try:
|
|
from nsct.providers.searxng import SearXNGProvider
|
|
if self._config.searxng_base_url:
|
|
provider = SearXNGProvider(base_url=self._config.searxng_base_url)
|
|
providers.append(provider)
|
|
except ImportError:
|
|
pass
|
|
except Exception as exc:
|
|
logger.warning("SearXNG provider init failed (%s) — search will return empty", exc)
|
|
|
|
self._multi_search = MultiProviderSearch(providers=providers if providers else None)
|
|
if not providers:
|
|
logger.warning("No search providers configured — search will return empty results")
|
|
return self._multi_search
|
|
|
|
def _get_llm_provider(self):
|
|
"""LLM-Provider für Claim Extraction."""
|
|
if self._llm_provider is not None:
|
|
return self._llm_provider
|
|
|
|
try:
|
|
from nsct.providers.llm import get_provider
|
|
from nsct.providers.metrics import ProviderMetrics
|
|
|
|
metrics = ProviderMetrics()
|
|
self._llm_provider = get_provider(self._config, metrics)
|
|
except Exception as exc:
|
|
logger.warning("LLM provider unavailable (%s) — extraction will produce no claims", exc)
|
|
self._llm_provider = None
|
|
return self._llm_provider
|
|
|
|
# ---------------------------------------------------------------
|
|
# Public: Pipeline
|
|
# ---------------------------------------------------------------
|
|
|
|
async def start(self) -> ResearchRun:
|
|
"""Erstelle den ResearchRun und initialisiere die Pipeline.
|
|
|
|
Setzt State auf CREATED, erstellt ResearchRun-Instanz,
|
|
initialisiert Budget-Tracker und State Machine.
|
|
|
|
Returns
|
|
-------
|
|
ResearchRun
|
|
Das erstellte Run-Objekt.
|
|
"""
|
|
# Budget init
|
|
self._budget_tracker.increment_llm_requests(1) # planner call
|
|
self._budget_tracker.increment_search(1) # initial search plan
|
|
|
|
# ResearchRun erstellen
|
|
self._run = self._create_research_run()
|
|
self._transition_to("created")
|
|
|
|
logger.info(
|
|
"Research run started: id=%s, query=%s, depth=%s",
|
|
self._run.id,
|
|
self._query[:60],
|
|
self._depth,
|
|
)
|
|
return self._run
|
|
|
|
async def run(self) -> dict[str, Any]:
|
|
"""Führe die gesamte Pipeline sequenziell aus.
|
|
|
|
Pipeline:
|
|
CREATED → PLANNING → SEARCHING → FETCHING → EXTRACTING
|
|
→ ANALYZING → COMPARING → SYNTHESIZING → COMPLETED / FAILED
|
|
|
|
Returns
|
|
-------
|
|
dict
|
|
Das Ergebnis der Pipeline (Report oder Error-Info).
|
|
"""
|
|
# Start falls noch nicht passiert
|
|
if self._run is None:
|
|
await self.start()
|
|
|
|
# Record research start timestamp for duration histogram
|
|
_start = time.monotonic()
|
|
|
|
# Pipeline-Schritte nacheinander
|
|
steps = [
|
|
"planning",
|
|
"searching",
|
|
"fetching",
|
|
"analyzing",
|
|
"comparing",
|
|
"synthesizing",
|
|
]
|
|
|
|
# Gap-Analysis & iterative Suche (Stage 13)
|
|
gap_results = await self._run_gap_analysis_loop(
|
|
claims=self._claims,
|
|
sources=self._sources,
|
|
max_iterations=2,
|
|
)
|
|
if gap_results.get("gap_queries"):
|
|
self._search_results.extend(gap_results.get("gap_search_results", []))
|
|
logger.info("Gap iteration complete: %d gap queries, %d additional results",
|
|
len(gap_results.get("gap_queries", [])),
|
|
len(gap_results.get("gap_search_results", [])))
|
|
if gap_results.get("gap_claims"):
|
|
self._claims.extend(gap_results["gap_claims"])
|
|
|
|
for step_name in steps:
|
|
try:
|
|
# Budget prüfen vor jedem Schritt
|
|
self._check_budget()
|
|
|
|
result = await self.run_step(step_name)
|
|
|
|
# Wenn ein Schritt FAILED meldet, Pipeline brechen
|
|
if result.get("success") is False:
|
|
error_msg = result.get("error", f"Step {step_name} failed")
|
|
self._mark_failed(error_msg)
|
|
return {
|
|
"success": False,
|
|
"error": error_msg,
|
|
"failed_step": step_name,
|
|
"report": None,
|
|
}
|
|
|
|
# Zeit-Tracking
|
|
self._budget_tracker.record_time_elapsed()
|
|
|
|
except BudgetExhaustedError:
|
|
logger.error("Budget exhausted during step %s", step_name)
|
|
self._mark_failed(f"Budget exhausted at step: {step_name}")
|
|
raise
|
|
|
|
except Exception as exc:
|
|
logger.error("Unhandled error in step %s: %s", step_name, exc)
|
|
self._mark_failed(str(exc))
|
|
return {
|
|
"success": False,
|
|
"error": str(exc),
|
|
"failed_step": step_name,
|
|
"report": None,
|
|
}
|
|
|
|
# Alle Schritte erfolgreich → COMPLETED
|
|
self._transition_to("completed")
|
|
self._budget_tracker.record_time_elapsed()
|
|
|
|
report = {
|
|
"success": True,
|
|
"report": self._get_report(),
|
|
"state": self._state_machine.current_state.value,
|
|
"budget_usage": self._budget_tracker.get_usage(),
|
|
}
|
|
|
|
elapsed = time.monotonic() - _start
|
|
metrics.observe(H_RESEARCH_DURATION, elapsed)
|
|
|
|
logger.info("Research pipeline completed: %d claims, state=%s, %.1fs",
|
|
len(self._claims), report["state"], elapsed)
|
|
return report
|
|
|
|
async def run_step(self, step_name: str) -> dict[str, Any]:
|
|
"""Führe EINEN einzelnen Pipeline-Schritt aus.
|
|
|
|
Parameters
|
|
----------
|
|
step_name : str
|
|
Name des Schritts (z.B. "planning", "searching").
|
|
|
|
Returns
|
|
-------
|
|
dict
|
|
Ergebnis-Dict mit success, data, error.
|
|
"""
|
|
step_name = step_name.lower().strip()
|
|
|
|
try:
|
|
handler = getattr(self, f"_step_{step_name}", None)
|
|
if handler is None:
|
|
logger.warning("Unknown step: %s", step_name)
|
|
return {"success": False, "error": f"Unknown step: {step_name}", "data": {}}
|
|
|
|
self._check_budget()
|
|
result = await handler()
|
|
return result
|
|
|
|
except BudgetExhaustedError:
|
|
self._mark_failed(f"Budget exhausted during {step_name}")
|
|
raise
|
|
|
|
except Exception as exc:
|
|
logger.error("Error in step %s: %s", step_name, exc)
|
|
self._mark_failed(str(exc))
|
|
return {"success": False, "error": str(exc), "data": {}}
|
|
|
|
# ---------------------------------------------------------------
|
|
# Private: Pipeline Steps
|
|
# ---------------------------------------------------------------
|
|
|
|
async def _step_planning(self) -> dict[str, Any]:
|
|
"""PLANNING: Erstelle Recherchestrategie mit dem Planner."""
|
|
self._transition_to("planning")
|
|
|
|
llm_duration_start = time.monotonic()
|
|
try:
|
|
planner = self._get_planner()
|
|
if isinstance(planner, MockResearchPlanner):
|
|
logger.info("Using MockResearchPlanner (LLM not available)")
|
|
plan = planner.plan(self._query, language="de")
|
|
else:
|
|
self._budget_tracker.increment_llm_requests(1)
|
|
plan = await planner.plan(self._query, language="de")
|
|
# Track LLM request duration
|
|
elapsed = time.monotonic() - llm_duration_start
|
|
metrics.observe(H_LLM_REQUEST_DURATION, elapsed)
|
|
|
|
self._plan = plan
|
|
|
|
logger.info("Planning complete: topic=%s, queries=%d",
|
|
plan.get("topic", ""), len(plan.get("queries", [])))
|
|
|
|
# Update run metadata (create a new frozen instance)
|
|
if self._run is not None:
|
|
self._run = self._run.model_copy(
|
|
update={
|
|
"research_plan": plan,
|
|
"updated_at": datetime.utcnow(),
|
|
}
|
|
)
|
|
|
|
return {"success": True, "data": {"plan": plan}, "plan": plan}
|
|
|
|
except Exception as exc:
|
|
logger.error("Planning failed: %s", exc)
|
|
# Fallback: Create a minimal plan
|
|
fallback = self._get_fallback_plan()
|
|
self._plan = fallback
|
|
# Store error in _run metadata via new instance
|
|
if self._run is not None:
|
|
self._run = self._run.model_copy(
|
|
update={
|
|
"plan_error": str(exc),
|
|
"updated_at": datetime.utcnow(),
|
|
}
|
|
)
|
|
return {"success": True, "data": {"plan": fallback}, "plan": fallback, "warning": "Used fallback plan"}
|
|
|
|
async def _step_searching(self) -> dict[str, Any]:
|
|
"""SEARCHING: Führe Suchanfragen durch und sammle URLs."""
|
|
self._transition_to("searching")
|
|
|
|
try:
|
|
multi_search = self._get_multi_search()
|
|
plan_queries = self._plan.get("queries", []) if self._plan else []
|
|
|
|
if not plan_queries:
|
|
# Fallback: verwende die ursprüngliche Query
|
|
plan_queries = [{"query": self._query, "purpose": "general", "category": "general", "language": "de"}]
|
|
|
|
all_results = []
|
|
self._budget_tracker.increment_search(len(plan_queries))
|
|
|
|
# Führe jede Query parallel aus
|
|
search_tasks = []
|
|
for q in plan_queries:
|
|
query_text = q.get("query", self._query)
|
|
search_tasks.append(multi_search.search(query_text, language="de", max_results=5))
|
|
|
|
raw_results = await asyncio.gather(*search_tasks, return_exceptions=True)
|
|
|
|
for raw in raw_results:
|
|
if isinstance(raw, Exception):
|
|
logger.warning("Search query failed: %s", raw)
|
|
continue
|
|
if isinstance(raw, list):
|
|
all_results.extend(raw)
|
|
|
|
# Dedupliziere und extrahiere URLs
|
|
seen_urls = set()
|
|
self._search_results = []
|
|
for r in all_results:
|
|
url = r.url if hasattr(r, "url") else r.get("url", "")
|
|
if url and url not in seen_urls:
|
|
seen_urls.add(url)
|
|
result_dict = {
|
|
"url": url,
|
|
"title": getattr(r, "title", "") or r.get("title", ""),
|
|
"snippet": getattr(r, "snippet", "") or r.get("snippet", ""),
|
|
"provider": getattr(r, "provider", "") or r.get("provider", ""),
|
|
}
|
|
self._search_results.append(result_dict)
|
|
|
|
logger.info("Search complete: %d unique URLs collected", len(self._search_results))
|
|
metrics.increment(C_SEARCH_QUERIES_TOTAL, len(plan_queries))
|
|
return {"success": True, "data": {"urls": self._search_results}, "url_count": len(self._search_results)}
|
|
|
|
except Exception as exc:
|
|
logger.warning("Searching failed (graceful fallback): %s — returning empty results", exc)
|
|
self._search_results = []
|
|
return {"success": True, "data": {"urls": []}, "url_count": 0, "warning": str(exc)}
|
|
|
|
async def _step_fetching(self) -> dict[str, Any]:
|
|
"""FETCHING: Crawle gesammelte URLs und extrahiere Inhalt."""
|
|
self._transition_to("fetching")
|
|
|
|
try:
|
|
if not self._search_results:
|
|
logger.warning("No URLs to fetch — skipping fetching step")
|
|
self._sources = []
|
|
return {"success": True, "data": {"sources": []}, "source_count": 0}
|
|
|
|
crawler = self._get_crawler()
|
|
urls = [r["url"] for r in self._search_results]
|
|
|
|
# Budget: Abschätzung der Download-Größe
|
|
self._budget_tracker.increment_download_bytes(len(urls) * 500_000) # ~500KB pro Seite
|
|
self._budget_tracker.increment_sources(len(urls))
|
|
|
|
docs = await crawler.fetch_and_extract_many(urls)
|
|
self._budget_tracker.record_time_elapsed()
|
|
|
|
# In sources-Dicts umwandeln
|
|
self._sources = []
|
|
for i, doc in enumerate(docs):
|
|
src = {
|
|
"id": str(uuid4()),
|
|
"url": doc.url or "",
|
|
"title": doc.title or "",
|
|
"domain": doc.url.split("//")[-1].split("/")[0] if doc.url else "",
|
|
"content": doc.text or "",
|
|
"links": doc.links or [],
|
|
"metadata": doc.metadata or {},
|
|
"error": doc.metadata.get("error", "") if doc.metadata else "",
|
|
}
|
|
self._sources.append(src)
|
|
|
|
# Track bytes
|
|
content_len = len(doc.text or "")
|
|
self._budget_tracker.increment_download_bytes(content_len)
|
|
|
|
successful = [s for s in self._sources if not s.get("error")]
|
|
metrics.increment(C_SOURCES_FETCHED_TOTAL, len(self._sources))
|
|
logger.info("Fetching complete: %d/%d sources successfully fetched", len(successful), len(self._sources))
|
|
return {
|
|
"success": True,
|
|
"data": {"sources": self._sources},
|
|
"source_count": len(self._sources),
|
|
"successful_count": len(successful),
|
|
}
|
|
|
|
except Exception as exc:
|
|
logger.warning("Fetching failed (graceful fallback): %s — returning empty sources", exc)
|
|
self._sources = []
|
|
return {"success": True, "data": {"sources": []}, "source_count": 0, "warning": str(exc)}
|
|
|
|
async def _step_extracting(self) -> dict[str, Any]:
|
|
"""EXTRACTING: Extrahiere Claims aus den extrahierten Quellen."""
|
|
self._transition_to("extracting")
|
|
|
|
try:
|
|
llm_provider = self._get_llm_provider()
|
|
|
|
if llm_provider is None or not self._sources:
|
|
logger.warning("No LLM provider or no sources — extracting produces no claims")
|
|
self._claims = []
|
|
return {"success": True, "data": {"claims": []}, "claim_count": 0}
|
|
|
|
self._budget_tracker.increment_llm_requests(len(self._sources))
|
|
|
|
# Stage 5: Claim Extraction
|
|
try:
|
|
from nsct.stages.stage5_extract_claims import Stage5Extractor
|
|
|
|
extractor = Stage5Extractor(
|
|
llm_provider=llm_provider,
|
|
config=self._config,
|
|
research_run_id=self._run.id if self._run else uuid4(),
|
|
sources=self._sources,
|
|
)
|
|
self._claims = await extractor.extract()
|
|
except NameError:
|
|
# stage5_extract_claims nicht importierbar
|
|
logger.warning("Stage5Extractor not available — skipping claim extraction")
|
|
self._claims = []
|
|
|
|
# Update Run
|
|
if self._run is not None:
|
|
self._run = self._run.model_copy(
|
|
update={
|
|
"claim_count": len(self._claims),
|
|
"updated_at": datetime.utcnow(),
|
|
}
|
|
)
|
|
|
|
logger.info("Extraction complete: %d claims", len(self._claims))
|
|
metrics.increment(C_CLAIMS_EXTRACTED_TOTAL, len(self._claims))
|
|
return {"success": True, "data": {"claims": [c.model_dump() for c in self._claims]}, "claim_count": len(self._claims)}
|
|
|
|
except Exception as exc:
|
|
logger.warning("Claim extraction failed (graceful fallback): %s", exc)
|
|
self._claims = []
|
|
return {"success": True, "data": {"claims": []}, "claim_count": 0, "warning": str(exc)}
|
|
|
|
async def _step_analyzing(self) -> dict[str, Any]:
|
|
"""ANALYZING: Platzhalter — Stage 6/7/8 werden später eingebunden.
|
|
|
|
Derzeit: Validiere Claims und sammle Metadaten.
|
|
"""
|
|
self._transition_to("analyzing")
|
|
|
|
logger.info("ANALYZING step: placeholder — Stage 6/7/8 pending integration")
|
|
|
|
# Grundlegende Claim-Validierung
|
|
claim_stats = {
|
|
"total": len(self._claims),
|
|
"by_type": {},
|
|
"avg_confidence": 0.0,
|
|
}
|
|
|
|
if self._claims:
|
|
type_counts: dict[str, int] = {}
|
|
total_conf = 0.0
|
|
for c in self._claims:
|
|
t = c.claim_type.value if hasattr(c.claim_type, "value") else str(c.claim_type)
|
|
type_counts[t] = type_counts.get(t, 0) + 1
|
|
total_conf += c.confidence
|
|
claim_stats["avg_confidence"] = round(total_conf / len(self._claims), 3)
|
|
claim_stats["by_type"] = type_counts
|
|
else:
|
|
claim_stats["by_type"] = {}
|
|
|
|
self._run_metadata = claim_stats
|
|
if self._run is not None:
|
|
self._run = self._run.model_copy(
|
|
update={
|
|
"metadata": {**self._run.metadata, "analysis": claim_stats},
|
|
"updated_at": datetime.utcnow(),
|
|
}
|
|
)
|
|
|
|
return {
|
|
"success": True,
|
|
"data": {"claim_stats": claim_stats},
|
|
"claim_count": len(self._claims),
|
|
}
|
|
|
|
async def _step_comparing(self) -> dict[str, Any]:
|
|
"""COMPARING: Platzhalter — Stage 8 Evidence Scoring.
|
|
|
|
Derzeit: Leere Comparison, nur Logging.
|
|
"""
|
|
self._transition_to("comparing")
|
|
|
|
logger.info("COMPARING step: placeholder — Stage 8 Evidence Scoring pending integration")
|
|
|
|
# Placeholder: keine Comparison-Daten
|
|
self._comparison_data: dict[str, Any] = {
|
|
"total_claims": len(self._claims),
|
|
"comparisons": [],
|
|
"evidence_scores": {},
|
|
}
|
|
|
|
if self._run is not None:
|
|
self._run = self._run.model_copy(
|
|
update={
|
|
"metadata": {**self._run.metadata, "comparison": self._comparison_data},
|
|
"updated_at": datetime.utcnow(),
|
|
}
|
|
)
|
|
|
|
return {
|
|
"success": True,
|
|
"data": self._comparison_data,
|
|
}
|
|
|
|
async def _step_synthesizing(self) -> dict[str, Any]:
|
|
"""SYNTHESIZING: Erzeuge Synthese-Bericht mit SynthesisStage."""
|
|
self._transition_to("synthesizing")
|
|
|
|
try:
|
|
# LLM-Anfrage zählen
|
|
self._budget_tracker.increment_llm_requests(1)
|
|
|
|
if SynthesisStage is None:
|
|
logger.warning("SynthesisStage not available — generating fallback report")
|
|
return self._fallback_synthesis()
|
|
|
|
llm_provider = self._get_llm_provider()
|
|
if llm_provider is None:
|
|
logger.warning("No LLM provider for synthesis — generating fallback report")
|
|
return self._fallback_synthesis()
|
|
|
|
# Claims als Dicts für SynthesisStage
|
|
claims_dicts = [c.model_dump() for c in self._claims]
|
|
|
|
stage = SynthesisStage(
|
|
research_run_id=self._run.id if self._run else uuid4(),
|
|
llm_provider=llm_provider,
|
|
config=self._config,
|
|
claims=claims_dicts,
|
|
)
|
|
|
|
result = await stage.execute()
|
|
|
|
if result.success:
|
|
logger.info("Synthesis complete: report generated successfully")
|
|
report = result.data if result.data else {}
|
|
self._budget_tracker.record_time_elapsed()
|
|
return {
|
|
"success": True,
|
|
"data": report,
|
|
"report": report,
|
|
}
|
|
else:
|
|
logger.warning("Synthesis returned success=False — generating fallback")
|
|
return self._fallback_synthesis_data(result.errors)
|
|
|
|
except Exception as exc:
|
|
logger.error("Synthesis failed: %s — generating fallback", exc)
|
|
return self._fallback_synthesis()
|
|
|
|
def _fallback_synthesis(self) -> dict[str, Any]:
|
|
"""Fallback Synthese wenn SynthesisStage nicht verfügbar ist."""
|
|
logger.info("Generating fallback synthesis report")
|
|
|
|
fallback_report = {
|
|
"success": True,
|
|
"data": self._get_report(),
|
|
"report": self._get_report(),
|
|
"warning": "Fallback synthesis (LLM unavailable)",
|
|
}
|
|
self._budget_tracker.record_time_elapsed()
|
|
return fallback_report
|
|
|
|
def _fallback_synthesis_data(self, errors: list[str]) -> dict[str, Any]:
|
|
"""Fallback report aus gescheiterter Synthese."""
|
|
report = self._get_report()
|
|
fallback = {
|
|
"success": True,
|
|
"data": report,
|
|
"report": report,
|
|
"synthesis_errors": errors,
|
|
"warning": "Synthesis degraded to fallback",
|
|
}
|
|
self._budget_tracker.record_time_elapsed()
|
|
return fallback
|
|
|
|
# ---------------------------------------------------------------
|
|
# Private: Helpers
|
|
# ---------------------------------------------------------------
|
|
|
|
def _create_research_run(self) -> ResearchRun:
|
|
"""Erstelle eine ResearchRun-Instanz mit allen relevanten Feldern.
|
|
|
|
Returns
|
|
-------
|
|
ResearchRun
|
|
Das neue Run-Objekt.
|
|
"""
|
|
budget_json = None
|
|
if self._budget_config is not None:
|
|
budget_json = self._budget_config.model_dump_json()
|
|
|
|
run = ResearchRun(
|
|
id=uuid4(),
|
|
research_id=self._research_id,
|
|
query=self._query,
|
|
state="created",
|
|
budget_config_json=budget_json,
|
|
metadata={
|
|
"depth": self._depth,
|
|
"created_by": "orchestrator",
|
|
},
|
|
)
|
|
self._run = run
|
|
return run
|
|
|
|
def _get_report(self) -> dict[str, Any]:
|
|
"""Erzeuge den finalen Report aus allen Zwischenspeichern."""
|
|
report: dict[str, Any] = {
|
|
"run_id": self._run.id if self._run else None,
|
|
"query": self._query,
|
|
"state": self._state_machine.current_state.value,
|
|
"plan": self._plan,
|
|
"search_results": self._search_results,
|
|
"sources": self._sources,
|
|
"claims": [c.model_dump() for c in self._claims],
|
|
"claim_count": len(self._claims),
|
|
"source_count": len(self._sources),
|
|
"budget_usage": self._budget_tracker.get_usage(),
|
|
"duration_seconds": (
|
|
time.monotonic() - self._budget_tracker._start_time
|
|
if hasattr(self._budget_tracker, "_start_time")
|
|
else 0
|
|
),
|
|
}
|
|
|
|
if self._comparison_data:
|
|
report["comparison"] = self._comparison_data
|
|
if self._run is not None:
|
|
report["metadata"] = self._run.metadata
|
|
|
|
return report
|
|
|
|
def _get_fallback_plan(self) -> dict[str, Any]:
|
|
"""Minimaler Fallback-Plan falls der Planner komplett versagt."""
|
|
return {
|
|
"topic": self._query[:120],
|
|
"queries": [
|
|
{
|
|
"query": self._query,
|
|
"purpose": "Falls Planner ausgefallen: Basis-Query",
|
|
"category": "general",
|
|
"language": "de",
|
|
}
|
|
],
|
|
"fallback": True,
|
|
"error": "Planner completely failed — using minimal fallback",
|
|
}
|
|
|
|
# ---------------------------------------------------------------
|
|
# Lifecycle
|
|
# ---------------------------------------------------------------
|
|
|
|
async def reset(self) -> None:
|
|
"""Setze den Orchestrator zurück (zustand → CREATED)."""
|
|
self._state_machine.reset(ResearchRunState.CREATED)
|
|
self._plan = None
|
|
self._search_results = []
|
|
self._sources = []
|
|
self._claims = []
|
|
self._planner = None
|
|
self._crawler = None
|
|
self._multi_search = None
|
|
self._llm_provider = None
|
|
self._budget_tracker = BudgetTracker(self._budget_config)
|
|
logger.info("Orchestrator reset to CREATED state")
|
|
|
|
# ---------------------------------------------------------------
|
|
# Private: Gap Analysis Loop (Stage 13)
|
|
# ---------------------------------------------------------------
|
|
|
|
async def _run_gap_analysis_loop(
|
|
self,
|
|
claims: list[Claim],
|
|
sources: list[dict[str, Any]],
|
|
max_iterations: int,
|
|
) -> dict[str, Any]:
|
|
"""Gap-Analyse durchführen und bei Bedarf iterative Suchanfragen generieren.
|
|
|
|
Parameters
|
|
----------
|
|
claims : list[Claim]
|
|
Extrahierte Claims.
|
|
sources : list[dict]
|
|
Gesammelte Quellen.
|
|
max_iterations : int
|
|
Maximale Anzahl Iterationen.
|
|
|
|
Returns
|
|
-------
|
|
dict
|
|
Gap-Ergebnisse mit gap_queries, gap_search_results, gap_claims.
|
|
"""
|
|
from nsct.models.gap_analysis import IterationReport
|
|
from nsct.stages.stage13_gap_analysis import GapAnalysisEngine
|
|
|
|
result = {
|
|
"gap_queries": [],
|
|
"gap_search_results": [],
|
|
"gap_claims": [],
|
|
"report": None,
|
|
}
|
|
|
|
try:
|
|
engine = GapAnalysisEngine(config=self._config, max_iterations=max_iterations)
|
|
|
|
for iteration in range(1, max_iterations + 1):
|
|
report = engine.analyze(
|
|
claims=claims,
|
|
sources=sources,
|
|
iteration_number=iteration,
|
|
research_run_id=str(self._run.id) if self._run else "",
|
|
)
|
|
|
|
result["report"] = report
|
|
|
|
if not report.has_gaps:
|
|
logger.info("Gap analysis: no more gaps at iteration %d", iteration)
|
|
break
|
|
|
|
logger.info("Gap analysis iteration %d: %d findings, %d queries",
|
|
iteration, len(report.findings), len(report.gap_search_queries))
|
|
|
|
result["gap_queries"].extend(report.gap_search_queries)
|
|
|
|
# Führe Gap-Suchen aus
|
|
if report.gap_search_queries:
|
|
gap_results = await self._execute_gap_searches(
|
|
report.gap_search_queries,
|
|
max_iterations,
|
|
)
|
|
result["gap_search_results"].extend(gap_results.get("search_results", []))
|
|
if gap_results.get("claims"):
|
|
result["gap_claims"].extend(gap_results["claims"])
|
|
|
|
return result
|
|
|
|
except Exception as exc:
|
|
logger.warning("Gap analysis failed: %s", exc)
|
|
return result
|
|
|
|
async def _execute_gap_searches(
|
|
self,
|
|
gap_queries: list[Any],
|
|
_max_iterations: int,
|
|
) -> dict[str, Any]:
|
|
"""Gap-Suchanfragen ausführen und Ergebnisse sammeln."""
|
|
search_results = []
|
|
claims = []
|
|
|
|
try:
|
|
multi_search = self._get_multi_search()
|
|
search_tasks = []
|
|
|
|
for q in gap_queries:
|
|
query_text = q.query if hasattr(q, "query") else q.get("query", "")
|
|
language = q.language if hasattr(q, "language") else "de"
|
|
search_tasks.append(multi_search.search(query_text, language=language, max_results=3))
|
|
|
|
raw_results = await asyncio.gather(*search_tasks, return_exceptions=True)
|
|
|
|
for raw in raw_results:
|
|
if isinstance(raw, Exception):
|
|
logger.warning("Gap search failed: %s", raw)
|
|
continue
|
|
if isinstance(raw, list):
|
|
search_results.extend(raw)
|
|
|
|
except Exception as exc:
|
|
logger.warning("Gap search execution failed: %s", exc)
|
|
|
|
return {"search_results": search_results, "claims": claims} |