mirror of
https://github.com/ZhuLinsen/daily_stock_analysis
synced 2026-09-20 10:53:33 +08:00
fix: 记录大盘复盘实际生成后端 (#2241)
* fix: report actual market review backend * fix: preserve exhausted fallback diagnostics * fix(review-feedback-2241): update the attempt loop to retain the current model even when the and * fix: preserve resolved market review provider * fix(review-feedback-2241): Resolve aliased LiteLLM providers before recording * fix(review-feedback-2241): preserve the failed LiteLLM route's resolved provider * fix(review-feedback-2241): preserve template fallback for primary LiteLLM exhaustion and Use the * fix(review-feedback-2241): preserve the fallback provider for unqualified response models * fix(review-feedback-2241): Derive exhausted aliases from the router's last deployment and * fix(review-feedback-2241): preserve the provider for gateway-owned slash model IDs * fix(review-feedback-2241): preserve the explicit provider for unqualified failure models * fix(review-feedback-2241): fix several earlier route-alias and router-failure diagnostics gaps, * fix(review-feedback-2241): 评审结论 - 代码检查 :当前整个 PR 仍有 1 个未关闭的高置信度代码 blocker。最新复核摘要:Full * fix(review-feedback-2241): 补一组成功路径回归:同一 alias 下首个 deployment 为 openai/~...、后续 deployment 为直连 * fix(review-feedback-2241): add a regression test that injects an analyzer exposing only is * fix(review-feedback-2241): add a regression that asserts the legacy injected-analyzer path does * fix(review-feedback-2241): 评审结论 - 代码检查 :当前整个 PR 仍有 1 个未关闭的高置信度代码 blocker。最新复核摘要:基于当前 HEAD * fix(review-feedback-2241): 评审结论 - 代码检查 :当前整个 PR 仍有 1 个未关闭的高置信度代码 blocker。最新复核摘要:基于完整 * fix(review-feedback-2241): add a regression test that drives call litellm
This commit is contained in:
323
src/analyzer.py
323
src/analyzer.py
@@ -16,7 +16,7 @@ import math
|
||||
import re
|
||||
import time
|
||||
from dataclasses import dataclass
|
||||
from typing import Optional, Dict, Any, List, Tuple, Callable
|
||||
from typing import Optional, Dict, Any, List, Tuple, Callable, Union
|
||||
|
||||
import litellm
|
||||
from json_repair import repair_json
|
||||
@@ -35,6 +35,7 @@ from src.config import (
|
||||
get_api_keys_for_model,
|
||||
get_config,
|
||||
get_configured_llm_models,
|
||||
get_explicit_llm_channel_model_provider,
|
||||
resolve_news_window_days,
|
||||
)
|
||||
from src.llm.hermes import (
|
||||
@@ -62,6 +63,7 @@ from src.llm.generation_backend import (
|
||||
GenerationBackend,
|
||||
GenerationError,
|
||||
GenerationErrorCode,
|
||||
GenerationResult,
|
||||
)
|
||||
from src.llm.usage import (
|
||||
attach_legacy_message_stability_audit,
|
||||
@@ -275,8 +277,9 @@ class _AllModelsFailedError(Exception):
|
||||
that *did* return a response (but whose JSON could not be validated), so
|
||||
callers can still attempt a best-effort text fallback.
|
||||
|
||||
``last_model`` and ``last_usage`` record the model name and token usage
|
||||
from the last attempt so callers can persist usage even on fallback.
|
||||
``last_model``, ``last_provider`` and ``last_usage`` record the resolved
|
||||
route identity and token usage from the last attempt so callers can persist
|
||||
diagnostics even on fallback.
|
||||
"""
|
||||
|
||||
def __init__(
|
||||
@@ -285,11 +288,13 @@ class _AllModelsFailedError(Exception):
|
||||
*,
|
||||
last_response_text: Optional[str] = None,
|
||||
last_model: Optional[str] = None,
|
||||
last_provider: Optional[str] = None,
|
||||
last_usage: Optional[Dict[str, Any]] = None,
|
||||
):
|
||||
super().__init__(message)
|
||||
self.last_response_text = last_response_text
|
||||
self.last_model = last_model
|
||||
self.last_provider = last_provider
|
||||
self.last_usage = last_usage or {}
|
||||
|
||||
|
||||
@@ -2810,6 +2815,190 @@ class GeminiAnalyzer:
|
||||
return obj.get(key)
|
||||
return getattr(obj, key, None)
|
||||
|
||||
@staticmethod
|
||||
def _resolve_configured_response_provider(
|
||||
configured_model: str,
|
||||
response_model: str,
|
||||
model_list: Optional[List[Dict[str, Any]]] = None,
|
||||
) -> str:
|
||||
"""Match the actual response model against all deployments of one alias."""
|
||||
normalized_configured_model = str(configured_model or "").strip()
|
||||
normalized_response_model = str(response_model or "").strip().lower()
|
||||
if not normalized_configured_model or not normalized_response_model or not model_list:
|
||||
return ""
|
||||
|
||||
for entry in model_list:
|
||||
params = entry.get("litellm_params", {}) or {}
|
||||
model_name = str(entry.get("model_name") or "").strip()
|
||||
if not model_name:
|
||||
model_name = str(params.get("model") or "").strip()
|
||||
if model_name != normalized_configured_model:
|
||||
continue
|
||||
|
||||
deployment_model = str(params.get("model") or "").strip()
|
||||
if deployment_model.lower() != normalized_response_model:
|
||||
continue
|
||||
|
||||
normalized_deployment_model = deployment_model.lower()
|
||||
if normalized_deployment_model.startswith("openai/~") or "openrouter" in normalized_deployment_model:
|
||||
return "openrouter"
|
||||
|
||||
_resolved_model, resolved_provider = resolved_model_provider_identity(
|
||||
deployment_model,
|
||||
)
|
||||
if resolved_provider:
|
||||
return resolved_provider
|
||||
return ""
|
||||
|
||||
def _resolve_response_model_provider(
|
||||
self,
|
||||
response: Any,
|
||||
*,
|
||||
fallback_provider: Optional[str] = None,
|
||||
configured_model: str = "",
|
||||
model_list: Optional[List[Dict[str, Any]]] = None,
|
||||
) -> Tuple[str, str]:
|
||||
"""Return the actual response model/provider when LiteLLM exposes them."""
|
||||
configured_provider = str(fallback_provider or "").strip()
|
||||
normalized_configured_model = str(configured_model or "").strip()
|
||||
if normalized_configured_model:
|
||||
resolved_configured_model, _ = resolved_model_provider_identity(
|
||||
normalized_configured_model,
|
||||
model_list,
|
||||
)
|
||||
configured_route = str(resolved_configured_model or normalized_configured_model).strip().lower()
|
||||
if configured_route.startswith("openai/~") or "openrouter" in configured_route:
|
||||
configured_provider = "openrouter"
|
||||
response_model = str(self._get_response_field(response, "model") or "").strip()
|
||||
if response_model:
|
||||
if "/" not in response_model:
|
||||
return response_model, configured_provider
|
||||
matched_provider = self._resolve_configured_response_provider(
|
||||
normalized_configured_model,
|
||||
response_model,
|
||||
model_list,
|
||||
)
|
||||
if matched_provider:
|
||||
return response_model, matched_provider
|
||||
if configured_provider == "openrouter":
|
||||
return response_model, configured_provider
|
||||
response_provider = get_explicit_llm_channel_model_provider(response_model)
|
||||
if response_provider:
|
||||
return response_model, response_provider
|
||||
return response_model, configured_provider
|
||||
return "", configured_provider
|
||||
|
||||
@staticmethod
|
||||
def _promote_error_identity(details: Any) -> Dict[str, str]:
|
||||
"""Lift route/model diagnostics from nested GenerationError details."""
|
||||
if not isinstance(details, dict):
|
||||
return {}
|
||||
promoted: Dict[str, str] = {}
|
||||
for key in ("last_model", "route_name", "last_provider"):
|
||||
candidate = str(details.get(key) or "").strip()
|
||||
if candidate:
|
||||
promoted[key] = candidate
|
||||
return promoted
|
||||
|
||||
def _resolve_router_failure_identity(
|
||||
self,
|
||||
exc: Any,
|
||||
*,
|
||||
route_name: str,
|
||||
recovery_model_list: List[Dict[str, Any]],
|
||||
) -> Tuple[str, str]:
|
||||
"""Resolve the final Router deployment identity from a transport exception."""
|
||||
normalized_route_name = str(route_name or "").strip()
|
||||
origins = route_deployment_origins(recovery_model_list, normalized_route_name)
|
||||
deployment_count = len(origins.hermes_deployments) + len(origins.non_hermes_deployments)
|
||||
candidate_models: List[str] = []
|
||||
candidate_provider = ""
|
||||
seen_payloads: set[int] = set()
|
||||
|
||||
def _remember_model(value: Any) -> None:
|
||||
normalized = str(value or "").strip()
|
||||
if not normalized:
|
||||
return
|
||||
if deployment_count > 1 and normalized == normalized_route_name:
|
||||
return
|
||||
if normalized not in candidate_models:
|
||||
candidate_models.append(normalized)
|
||||
|
||||
def _remember_provider(value: Any) -> None:
|
||||
nonlocal candidate_provider
|
||||
normalized = str(value or "").strip()
|
||||
if normalized and not candidate_provider:
|
||||
candidate_provider = normalized
|
||||
|
||||
def _walk(payload: Any) -> None:
|
||||
if payload is None:
|
||||
return
|
||||
payload_id = id(payload)
|
||||
if payload_id in seen_payloads:
|
||||
return
|
||||
seen_payloads.add(payload_id)
|
||||
|
||||
if isinstance(payload, dict):
|
||||
params = payload.get("litellm_params")
|
||||
if isinstance(params, dict):
|
||||
_remember_model(params.get("model"))
|
||||
_remember_provider(
|
||||
params.get("custom_llm_provider") or params.get("provider")
|
||||
)
|
||||
for key in (
|
||||
"litellm_model",
|
||||
"response_model",
|
||||
"deployment_model",
|
||||
"model",
|
||||
"model_name",
|
||||
):
|
||||
_remember_model(payload.get(key))
|
||||
_remember_provider(
|
||||
payload.get("llm_provider")
|
||||
or payload.get("litellm_provider")
|
||||
or payload.get("custom_llm_provider")
|
||||
or payload.get("provider")
|
||||
)
|
||||
for key in ("response", "error", "details", "metadata", "body"):
|
||||
_walk(payload.get(key))
|
||||
return
|
||||
|
||||
for key in ("response", "error", "details", "metadata", "body"):
|
||||
nested = getattr(payload, key, None)
|
||||
if nested is not payload:
|
||||
_walk(nested)
|
||||
_remember_provider(
|
||||
getattr(payload, "llm_provider", None)
|
||||
or getattr(payload, "litellm_provider", None)
|
||||
or getattr(payload, "custom_llm_provider", None)
|
||||
or getattr(payload, "provider", None)
|
||||
)
|
||||
for key in (
|
||||
"litellm_model",
|
||||
"response_model",
|
||||
"deployment_model",
|
||||
"model",
|
||||
"model_name",
|
||||
):
|
||||
_remember_model(getattr(payload, key, None))
|
||||
|
||||
_walk(exc)
|
||||
for candidate_model in candidate_models:
|
||||
resolved_model, resolved_provider = resolved_model_provider_identity(
|
||||
candidate_model,
|
||||
recovery_model_list,
|
||||
)
|
||||
normalized_route = str(resolved_model or candidate_model).strip().lower()
|
||||
explicit_provider = get_explicit_llm_channel_model_provider(candidate_model)
|
||||
route_text = f"{candidate_provider} {normalized_route}".strip().lower()
|
||||
if normalized_route.startswith("openai/~") or "openrouter" in route_text:
|
||||
return resolved_model or candidate_model, "openrouter"
|
||||
if candidate_provider and not explicit_provider:
|
||||
return resolved_model or candidate_model, candidate_provider
|
||||
if resolved_provider:
|
||||
return resolved_model or candidate_model, resolved_provider
|
||||
return "", candidate_provider
|
||||
|
||||
def _extract_text_blocks(self, blocks: Any, *, strip: bool = True) -> str:
|
||||
"""Extract final-answer text from OpenAI-compatible content blocks.
|
||||
|
||||
@@ -2980,7 +3169,8 @@ class GeminiAnalyzer:
|
||||
stream_progress_callback: Optional[Callable[[int], None]] = None,
|
||||
response_validator: Optional[Callable[[str], None]] = None,
|
||||
audit_context: Optional[Dict[str, Any]] = None,
|
||||
) -> Tuple[str, str, Dict[str, Any]]:
|
||||
return_generation_result: bool = False,
|
||||
) -> Union[Tuple[str, str, Dict[str, Any]], GenerationResult]:
|
||||
"""Compatibility wrapper around the configured generation backend."""
|
||||
preflight_error = self.get_generation_backend_config_error()
|
||||
if preflight_error is not None and not self._can_use_generation_fallback(preflight_error):
|
||||
@@ -3033,6 +3223,7 @@ class GeminiAnalyzer:
|
||||
except _AllModelsFailedError:
|
||||
raise
|
||||
except GenerationError as fallback_exc:
|
||||
fallback_identity = self._promote_error_identity(fallback_exc.details)
|
||||
raise GenerationError(
|
||||
error_code=fallback_exc.error_code,
|
||||
stage="fallback",
|
||||
@@ -3042,6 +3233,7 @@ class GeminiAnalyzer:
|
||||
provider=fallback_exc.provider,
|
||||
details={
|
||||
"reason": "fallback_backend_failed",
|
||||
**fallback_identity,
|
||||
"primary_error": {
|
||||
"error_code": exc.error_code.value,
|
||||
"backend": exc.backend,
|
||||
@@ -3078,6 +3270,8 @@ class GeminiAnalyzer:
|
||||
"fallback_error": str(fallback_exc),
|
||||
},
|
||||
) from fallback_exc
|
||||
if return_generation_result:
|
||||
return result
|
||||
return result.text, result.model, result.usage
|
||||
|
||||
def _call_litellm_impl(
|
||||
@@ -3128,10 +3322,12 @@ class GeminiAnalyzer:
|
||||
last_error = None
|
||||
last_response_text: Optional[str] = None
|
||||
last_model: Optional[str] = None
|
||||
last_provider: Optional[str] = None
|
||||
last_usage: Dict[str, Any] = {}
|
||||
effective_system_prompt = system_prompt or self.TEXT_SYSTEM_PROMPT
|
||||
router_model_names = set(get_configured_llm_models(config.llm_model_list))
|
||||
for model in models_to_try:
|
||||
last_model = model
|
||||
origins = route_deployment_origins(config.llm_model_list, model)
|
||||
model_stream = bool(stream and not origins.has_hermes)
|
||||
recovery_model_list = config.llm_model_list
|
||||
@@ -3139,6 +3335,8 @@ class GeminiAnalyzer:
|
||||
if legacy_router_model_list and model == config.litellm_model and not use_channel_router:
|
||||
recovery_model_list = legacy_router_model_list
|
||||
usage_model, usage_provider = resolved_model_provider_identity(model, recovery_model_list)
|
||||
if usage_provider:
|
||||
last_provider = usage_provider
|
||||
|
||||
try:
|
||||
def _attach_usage_audit(
|
||||
@@ -3151,7 +3349,9 @@ class GeminiAnalyzer:
|
||||
config,
|
||||
)
|
||||
effective_audit_context = dict(audit_context)
|
||||
effective_audit_context["provider"] = usage_provider
|
||||
effective_audit_context["provider"] = (
|
||||
usage.get("provider") or usage_provider
|
||||
)
|
||||
effective_audit_context["transport"] = (
|
||||
effective_audit_context.get("transport") or "litellm"
|
||||
)
|
||||
@@ -3265,7 +3465,11 @@ class GeminiAnalyzer:
|
||||
if _stream_text is not None:
|
||||
last_response_text = _stream_text
|
||||
last_model = model
|
||||
if usage_provider:
|
||||
_stream_usage["provider"] = usage_provider
|
||||
_stream_usage = _attach_usage_audit(_stream_usage, call_kwargs["messages"])
|
||||
if usage_provider:
|
||||
_stream_usage.setdefault("provider", usage_provider)
|
||||
last_usage = _stream_usage
|
||||
if response_validator is not None:
|
||||
response_validator(_stream_text)
|
||||
@@ -3285,26 +3489,55 @@ class GeminiAnalyzer:
|
||||
logger=logger,
|
||||
)
|
||||
|
||||
response_model, response_provider = self._resolve_response_model_provider(
|
||||
response,
|
||||
fallback_provider=usage_provider,
|
||||
configured_model=model,
|
||||
model_list=recovery_model_list,
|
||||
)
|
||||
actual_model = response_model or model
|
||||
if response_model:
|
||||
last_model = actual_model
|
||||
if response_provider:
|
||||
last_provider = response_provider
|
||||
content = self._extract_completion_text(response)
|
||||
if content:
|
||||
usage_messages = None if audit_context is not None else call_kwargs["messages"]
|
||||
usage = self._normalize_usage(
|
||||
extract_usage_payload(response),
|
||||
model=usage_model or model,
|
||||
provider=usage_provider,
|
||||
model=response_model or usage_model or model,
|
||||
provider=response_provider or usage_provider,
|
||||
messages=usage_messages,
|
||||
)
|
||||
if response_provider or usage_provider:
|
||||
usage["provider"] = response_provider or usage_provider
|
||||
if audit_context is not None:
|
||||
usage = _attach_usage_audit(usage, call_kwargs["messages"])
|
||||
if response_model:
|
||||
usage.setdefault("response_model", response_model)
|
||||
if response_provider or usage_provider:
|
||||
usage.setdefault("provider", response_provider or usage_provider)
|
||||
last_response_text = content
|
||||
last_model = model
|
||||
last_model = actual_model
|
||||
if response_provider:
|
||||
last_provider = response_provider
|
||||
last_usage = usage
|
||||
if response_validator is not None:
|
||||
response_validator(content)
|
||||
return (content, model, usage)
|
||||
return (content, actual_model, usage)
|
||||
raise ValueError("LLM returned empty response")
|
||||
|
||||
except Exception as e:
|
||||
if uses_router:
|
||||
router_model, router_provider = self._resolve_router_failure_identity(
|
||||
e,
|
||||
route_name=model,
|
||||
recovery_model_list=recovery_model_list,
|
||||
)
|
||||
if router_model:
|
||||
last_model = router_model
|
||||
if router_provider:
|
||||
last_provider = router_provider
|
||||
safe_error = self._sanitize_litellm_exception_text(e, config=config, model=model)
|
||||
logger.warning("[LiteLLM] %s failed: %s", model, safe_error)
|
||||
last_error = RuntimeError(f"{type(e).__name__}: {safe_error}")
|
||||
@@ -3314,6 +3547,7 @@ class GeminiAnalyzer:
|
||||
f"All LLM models failed (tried {len(models_to_try)} model(s)). Last error: {last_error}",
|
||||
last_response_text=last_response_text,
|
||||
last_model=last_model,
|
||||
last_provider=last_provider,
|
||||
last_usage=last_usage,
|
||||
)
|
||||
|
||||
@@ -3354,6 +3588,77 @@ class GeminiAnalyzer:
|
||||
logger.error("[generate_text] LLM call failed: %s", exc)
|
||||
return None
|
||||
|
||||
def get_generation_backend_identity(self) -> Tuple[str, str]:
|
||||
"""Return the configured primary backend identity for live diagnostics."""
|
||||
backend_id, _fallback_backend_id = self._resolve_generation_backend_config()
|
||||
if backend_id in LOCAL_CLI_GENERATION_BACKEND_IDS:
|
||||
return backend_id, backend_id
|
||||
config = self._get_runtime_config()
|
||||
return backend_id, str(getattr(config, "litellm_model", "") or "")
|
||||
|
||||
def generate_text_with_metadata(
|
||||
self,
|
||||
prompt: str,
|
||||
max_tokens: int = 2048,
|
||||
temperature: float = 0.7,
|
||||
) -> Optional[GenerationResult]:
|
||||
"""Generate text and return the actual backend/model used for diagnostics."""
|
||||
try:
|
||||
result = self._call_litellm(
|
||||
prompt,
|
||||
generation_config={"max_tokens": max_tokens, "temperature": temperature},
|
||||
return_generation_result=True,
|
||||
)
|
||||
if not isinstance(result, GenerationResult):
|
||||
raise TypeError("generation backend returned an invalid result")
|
||||
if should_persist_usage_telemetry(result.usage):
|
||||
persist_llm_usage(result.usage, result.model, call_type="market_review")
|
||||
return result
|
||||
except GenerationError:
|
||||
raise
|
||||
except _AllModelsFailedError as exc:
|
||||
backend_id, fallback_backend_id = self._resolve_generation_backend_config()
|
||||
if not fallback_backend_id and backend_id == LITELLM_BACKEND_ID:
|
||||
logger.warning(
|
||||
"[generate_text_with_metadata] Primary LiteLLM exhausted all configured models; "
|
||||
"returning empty GenerationResult so caller fallback can continue"
|
||||
)
|
||||
usage = dict(exc.last_usage or {})
|
||||
if exc.last_provider:
|
||||
usage.setdefault("provider", exc.last_provider)
|
||||
return GenerationResult(
|
||||
text="",
|
||||
model=exc.last_model or str(getattr(self._get_runtime_config(), "litellm_model", "") or ""),
|
||||
provider=exc.last_provider or backend_id,
|
||||
backend=backend_id,
|
||||
usage=usage,
|
||||
diagnostics={
|
||||
"reason": "all_models_failed",
|
||||
"configured_primary_backend": backend_id,
|
||||
"configured_fallback_backend": fallback_backend_id,
|
||||
"last_model": exc.last_model,
|
||||
"template_fallback": True,
|
||||
},
|
||||
)
|
||||
failed_backend = fallback_backend_id or backend_id
|
||||
raise GenerationError(
|
||||
error_code=GenerationErrorCode.UNKNOWN_BACKEND_ERROR,
|
||||
stage="fallback" if fallback_backend_id else "generation",
|
||||
retryable=False,
|
||||
fallbackable=False,
|
||||
backend=failed_backend,
|
||||
provider=exc.last_provider or failed_backend,
|
||||
details={
|
||||
"reason": "all_models_failed",
|
||||
"configured_primary_backend": backend_id,
|
||||
"configured_fallback_backend": fallback_backend_id,
|
||||
"last_model": exc.last_model,
|
||||
},
|
||||
) from exc
|
||||
except Exception as exc:
|
||||
logger.error("[generate_text_with_metadata] LLM call failed: %s", exc)
|
||||
raise
|
||||
|
||||
def analyze(
|
||||
self,
|
||||
context: Dict[str, Any],
|
||||
|
||||
@@ -20,16 +20,19 @@ from typing import Optional, Dict, Any, List
|
||||
|
||||
import pandas as pd
|
||||
|
||||
from src.agent.provider_trace import resolved_model_provider_identity
|
||||
from src.config import get_config
|
||||
from src.report_language import normalize_report_language
|
||||
from src.search_service import SearchService
|
||||
from src.core.market_profile import get_profile, MarketProfile
|
||||
from src.core.market_strategy import get_market_strategy_blueprint
|
||||
from src.llm.backend_registry import (
|
||||
LOCAL_CLI_GENERATION_BACKEND_IDS,
|
||||
LITELLM_BACKEND_ID,
|
||||
resolve_generation_backend_id,
|
||||
resolve_generation_fallback_backend_id,
|
||||
)
|
||||
from src.llm.generation_backend import GenerationError
|
||||
from src.llm.generation_backend import GenerationError, GenerationResult
|
||||
from src.schemas.market_light import MARKET_LIGHT_REGIONS, MarketLightSnapshot
|
||||
from src.services.run_diagnostics import record_llm_run, record_llm_run_started
|
||||
from src.services.intelligence_service import IntelligenceService
|
||||
@@ -37,6 +40,7 @@ from data_provider.base import DataFetcherManager
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
_LEGACY_ANALYZER_BACKEND_ID = "legacy_analyzer"
|
||||
|
||||
_ENGLISH_SECTION_PATTERNS = {
|
||||
"market_summary": r"###\s*(?:1\.\s*)?Market Summary",
|
||||
@@ -154,6 +158,130 @@ class MarketAnalyzer:
|
||||
def _log_context(self) -> str:
|
||||
return f"component=market_review region={self.region}"
|
||||
|
||||
@staticmethod
|
||||
def _resolve_configured_response_provider(
|
||||
configured_model: str,
|
||||
response_model: str,
|
||||
model_list: Optional[List[Dict[str, Any]]] = None,
|
||||
) -> str:
|
||||
"""Match the actual response model against all deployments of one alias."""
|
||||
normalized_configured_model = str(configured_model or "").strip()
|
||||
normalized_response_model = str(response_model or "").strip().lower()
|
||||
if not normalized_configured_model or not normalized_response_model or not model_list:
|
||||
return ""
|
||||
|
||||
for entry in model_list:
|
||||
params = entry.get("litellm_params", {}) or {}
|
||||
model_name = str(entry.get("model_name") or "").strip()
|
||||
if not model_name:
|
||||
model_name = str(params.get("model") or "").strip()
|
||||
if model_name != normalized_configured_model:
|
||||
continue
|
||||
|
||||
deployment_model = str(params.get("model") or "").strip()
|
||||
if deployment_model.lower() != normalized_response_model:
|
||||
continue
|
||||
|
||||
normalized_deployment_model = deployment_model.lower()
|
||||
if normalized_deployment_model.startswith("openai/~") or "openrouter" in normalized_deployment_model:
|
||||
return "openrouter"
|
||||
|
||||
_resolved_model, resolved_provider = resolved_model_provider_identity(
|
||||
deployment_model,
|
||||
)
|
||||
if resolved_provider:
|
||||
return resolved_provider
|
||||
return ""
|
||||
|
||||
def _resolve_recorded_provider(
|
||||
self,
|
||||
*,
|
||||
provider: str,
|
||||
model: str,
|
||||
backend: str,
|
||||
usage_provider: str = "",
|
||||
response_model: str = "",
|
||||
) -> str:
|
||||
"""Resolve LiteLLM router aliases before persisting diagnostics."""
|
||||
normalized_backend = str(backend or "").strip().lower()
|
||||
normalized_model = str(model or "").strip()
|
||||
normalized_usage_provider = str(usage_provider or "").strip()
|
||||
normalized_response_model = str(response_model or "").strip()
|
||||
normalized_provider = str(provider or "").strip()
|
||||
resolved_route_provider = ""
|
||||
if normalized_backend == "litellm" and normalized_model:
|
||||
resolved_route_model, resolved_route_provider = resolved_model_provider_identity(
|
||||
normalized_model,
|
||||
getattr(self.config, "llm_model_list", None) or [],
|
||||
)
|
||||
normalized_route = str(resolved_route_model or normalized_model).strip().lower()
|
||||
if normalized_route.startswith("openai/~") or "openrouter" in normalized_route:
|
||||
resolved_route_provider = "openrouter"
|
||||
resolved_response_provider = ""
|
||||
if normalized_response_model:
|
||||
resolved_response_provider = self._resolve_configured_response_provider(
|
||||
normalized_model,
|
||||
normalized_response_model,
|
||||
getattr(self.config, "llm_model_list", None) or [],
|
||||
)
|
||||
if resolved_response_provider == "openrouter":
|
||||
return resolved_response_provider
|
||||
if normalized_usage_provider:
|
||||
return normalized_usage_provider
|
||||
if normalized_response_model:
|
||||
if resolved_response_provider:
|
||||
return resolved_response_provider
|
||||
if resolved_route_provider == "openrouter":
|
||||
return resolved_route_provider
|
||||
_wire_model, resolved_provider = resolved_model_provider_identity(
|
||||
normalized_response_model,
|
||||
)
|
||||
if resolved_provider:
|
||||
return resolved_provider
|
||||
if normalized_backend != "litellm" or not normalized_model:
|
||||
return normalized_provider
|
||||
return resolved_route_provider or normalized_provider or "openai"
|
||||
|
||||
def _resolve_recorded_error_model(
|
||||
self,
|
||||
*,
|
||||
error: Any,
|
||||
fallback_model: str = "",
|
||||
) -> str:
|
||||
"""Preserve route/model diagnostics for LiteLLM configuration failures."""
|
||||
details = getattr(error, "details", None)
|
||||
if isinstance(details, dict):
|
||||
visited: set[int] = set()
|
||||
|
||||
def _find_error_model(payload: Dict[str, Any]) -> str:
|
||||
payload_id = id(payload)
|
||||
if payload_id in visited:
|
||||
return ""
|
||||
visited.add(payload_id)
|
||||
for key in ("last_model", "route_name"):
|
||||
candidate = str(payload.get(key) or "").strip()
|
||||
if candidate:
|
||||
return candidate
|
||||
fallback_error = payload.get("fallback_error")
|
||||
if isinstance(fallback_error, dict):
|
||||
nested_details = fallback_error.get("details")
|
||||
if isinstance(nested_details, dict):
|
||||
candidate = _find_error_model(nested_details)
|
||||
if candidate:
|
||||
return candidate
|
||||
return _find_error_model(fallback_error)
|
||||
return ""
|
||||
|
||||
candidate = _find_error_model(details)
|
||||
if candidate:
|
||||
return candidate
|
||||
backend = str(getattr(error, "backend", "") or "").strip()
|
||||
configured_model = str(getattr(self.config, "litellm_model", "") or "").strip()
|
||||
normalized_fallback_model = str(fallback_model or "").strip()
|
||||
if backend == "litellm":
|
||||
return normalized_fallback_model or configured_model or backend
|
||||
return normalized_fallback_model or backend
|
||||
|
||||
def _get_output_language(self) -> str:
|
||||
"""Return the truthful report language (zh/en/ko) for payload and directives."""
|
||||
return normalize_report_language(
|
||||
@@ -669,8 +797,8 @@ Focus on index trend, liquidity, and sector rotation to shape the next-session t
|
||||
)
|
||||
record_llm_run(
|
||||
success=False,
|
||||
provider="litellm",
|
||||
model=getattr(self.config, "litellm_model", None),
|
||||
provider=backend_error.provider or backend_error.backend,
|
||||
model=self._resolve_recorded_error_model(error=backend_error),
|
||||
call_type="market_review",
|
||||
error_type=type(backend_error).__name__,
|
||||
error_message=backend_error,
|
||||
@@ -688,20 +816,48 @@ Focus on index trend, liquidity, and sector rotation to shape the next-session t
|
||||
prompt = self._build_review_prompt(overview, news)
|
||||
|
||||
logger.info("[大盘] %s action=generate_review status=start", self._log_context())
|
||||
# Use the public generate_text() entry point - never access private analyzer attributes.
|
||||
# Use public analyzer APIs so diagnostics reflect the actual execution backend.
|
||||
llm_started_at = time.perf_counter()
|
||||
provider, model = self._get_analyzer_generation_backend_identity()
|
||||
try:
|
||||
record_llm_run_started(
|
||||
provider="litellm",
|
||||
model=getattr(self.config, "litellm_model", None),
|
||||
provider=provider,
|
||||
model=model,
|
||||
call_type="market_review",
|
||||
)
|
||||
review = self.analyzer.generate_text(prompt, max_tokens=8192, temperature=0.7)
|
||||
generation_result = self._generate_market_review_with_metadata(
|
||||
prompt,
|
||||
provider=provider,
|
||||
model=model,
|
||||
)
|
||||
review = generation_result.text if generation_result else None
|
||||
generation_diagnostics = (
|
||||
getattr(generation_result, "diagnostics", None) if generation_result is not None else None
|
||||
)
|
||||
if generation_result is not None:
|
||||
model = generation_result.model or model
|
||||
usage_payload = getattr(generation_result, "usage", None)
|
||||
provider = self._resolve_recorded_provider(
|
||||
provider=generation_result.provider or generation_result.backend or provider,
|
||||
model=model,
|
||||
backend=generation_result.backend or "",
|
||||
usage_provider=(
|
||||
usage_payload.get("provider") if isinstance(usage_payload, dict) else ""
|
||||
),
|
||||
response_model=(
|
||||
usage_payload.get("response_model") if isinstance(usage_payload, dict) else ""
|
||||
),
|
||||
)
|
||||
except Exception as exc:
|
||||
error_provider = getattr(exc, "provider", None) or getattr(exc, "backend", None) or provider
|
||||
error_model = self._resolve_recorded_error_model(
|
||||
error=exc,
|
||||
fallback_model=model,
|
||||
)
|
||||
record_llm_run(
|
||||
success=False,
|
||||
provider="litellm",
|
||||
model=getattr(self.config, "litellm_model", None),
|
||||
provider=error_provider,
|
||||
model=error_model,
|
||||
call_type="market_review",
|
||||
duration_ms=int((time.perf_counter() - llm_started_at) * 1000),
|
||||
error_type=type(exc).__name__,
|
||||
@@ -709,14 +865,26 @@ Focus on index trend, liquidity, and sector rotation to shape the next-session t
|
||||
)
|
||||
raise
|
||||
|
||||
failed_all_models = (
|
||||
isinstance(generation_diagnostics, dict)
|
||||
and generation_diagnostics.get("reason") == "all_models_failed"
|
||||
)
|
||||
record_llm_run(
|
||||
success=bool(review),
|
||||
provider="litellm",
|
||||
model=getattr(self.config, "litellm_model", None),
|
||||
provider=provider,
|
||||
model=model,
|
||||
call_type="market_review",
|
||||
duration_ms=int((time.perf_counter() - llm_started_at) * 1000),
|
||||
error_type=None if review else "EmptyResponse",
|
||||
error_message=None if review else "empty market review response",
|
||||
error_type=None if review else ("AllModelsFailed" if failed_all_models else "EmptyResponse"),
|
||||
error_message=(
|
||||
None
|
||||
if review
|
||||
else (
|
||||
generation_diagnostics
|
||||
if failed_all_models
|
||||
else "empty market review response"
|
||||
)
|
||||
),
|
||||
)
|
||||
|
||||
if review:
|
||||
@@ -752,6 +920,74 @@ Focus on index trend, liquidity, and sector rotation to shape the next-session t
|
||||
error = method()
|
||||
return error if isinstance(error, GenerationError) else None
|
||||
|
||||
def _get_configured_generation_backend_identity(self) -> tuple[str, str]:
|
||||
"""Best-effort backend identity for legacy analyzers without metadata APIs."""
|
||||
backend_id = str(getattr(self.config, "generation_backend", "") or "").strip().lower()
|
||||
if not backend_id:
|
||||
backend_id = LITELLM_BACKEND_ID
|
||||
if backend_id in LOCAL_CLI_GENERATION_BACKEND_IDS:
|
||||
return backend_id, backend_id
|
||||
return backend_id, str(getattr(self.config, "litellm_model", "") or "")
|
||||
|
||||
def _get_legacy_analyzer_generation_backend_identity(self) -> tuple[str, str]:
|
||||
"""Return a neutral identity for injected analyzers that expose no metadata APIs."""
|
||||
return _LEGACY_ANALYZER_BACKEND_ID, _LEGACY_ANALYZER_BACKEND_ID
|
||||
|
||||
def _get_analyzer_generation_backend_identity(self) -> tuple[str, str]:
|
||||
"""Use analyzer metadata API when available, otherwise fall back to configured identity."""
|
||||
if self.analyzer is None:
|
||||
return self._get_configured_generation_backend_identity()
|
||||
missing = object()
|
||||
if getattr_static(self.analyzer, "get_generation_backend_identity", missing) is not missing:
|
||||
method = getattr(self.analyzer, "get_generation_backend_identity", None)
|
||||
if callable(method):
|
||||
return method()
|
||||
return self._get_legacy_analyzer_generation_backend_identity()
|
||||
|
||||
def _generate_market_review_with_metadata(
|
||||
self,
|
||||
prompt: str,
|
||||
*,
|
||||
provider: str,
|
||||
model: str,
|
||||
) -> Optional[GenerationResult]:
|
||||
"""Support legacy analyzers that only implement the public generate_text() contract."""
|
||||
if self.analyzer is None:
|
||||
return None
|
||||
missing = object()
|
||||
if getattr_static(self.analyzer, "generate_text_with_metadata", missing) is not missing:
|
||||
method = getattr(self.analyzer, "generate_text_with_metadata", None)
|
||||
if callable(method):
|
||||
return method(
|
||||
prompt,
|
||||
max_tokens=8192,
|
||||
temperature=0.7,
|
||||
)
|
||||
|
||||
if getattr_static(self.analyzer, "generate_text", missing) is missing:
|
||||
raise AttributeError(
|
||||
"analyzer must implement generate_text_with_metadata() or generate_text()"
|
||||
)
|
||||
legacy_method = getattr(self.analyzer, "generate_text", None)
|
||||
if not callable(legacy_method):
|
||||
raise AttributeError(
|
||||
"analyzer must implement generate_text_with_metadata() or generate_text()"
|
||||
)
|
||||
review = legacy_method(
|
||||
prompt,
|
||||
max_tokens=8192,
|
||||
temperature=0.7,
|
||||
)
|
||||
if review is None:
|
||||
return None
|
||||
return GenerationResult(
|
||||
text=review,
|
||||
provider=provider,
|
||||
model=model,
|
||||
backend=provider,
|
||||
usage={},
|
||||
)
|
||||
|
||||
def build_market_review_payload(
|
||||
self,
|
||||
overview: MarketOverview,
|
||||
|
||||
Reference in New Issue
Block a user