feat(ai): add guarded Claude fallback route
This commit is contained in:
@@ -8,8 +8,9 @@ AI Router - Phase 13.3 #87
|
||||
延遲目標: < 50ms (規則引擎優先)
|
||||
|
||||
唯一 production 路由(所有風險、複雜度與工具路徑共用):
|
||||
GCP-A Ollama -> GCP-B Ollama -> host111 Ollama -> Gemini API。
|
||||
Claude/NVIDIA/Nemo/OpenClaw 僅保留 non-executable shadow metadata。
|
||||
GCP-A Ollama -> GCP-B Ollama -> host111 Ollama -> Anthropic Claude API
|
||||
-> Gemini API。Claude/Gemini 都必須通過 persistent enablement、原子成本閘與
|
||||
per-provider circuit breaker;NVIDIA/Nemo/OpenClaw 僅保留 shadow metadata。
|
||||
|
||||
版本: v5.0
|
||||
建立: 2026-03-26 (台北時區)
|
||||
@@ -28,7 +29,7 @@ Claude/NVIDIA/Nemo/OpenClaw 僅保留 non-executable shadow metadata。
|
||||
| v4.2 | 2026-04-04 | Claude Code | Phase 25 P0 實測修正: _local_fallback_chain 移除 Nemotron(雲端),僅留 Ollama(本地); timeout 依實測調整(NIM 60s/Ollama 200s) |
|
||||
| v4.3 | 2026-04-05 | Claude Code | Phase 25 P0 架構修正: 實測 Ollama CPU ~238s(不可用); NIM 實測 2-27s avg 10.6s; DIAGNOSE 改走 _full_fallback_chain(NIM 主力); _local_fallback_chain 廢棄 |
|
||||
| v4.4 | 2026-04-27 | Claude Sonnet 4.6 | A2 INC-20260425: DIAGNOSE fallback chain 移除 Ollama (CPU 238s 二次 timeout); 新增 _diagnose_fallback_chain (NEMO→GEMINI→CLAUDE); 新增 aiops_diagnose_fallback_total metric |
|
||||
| v5.0 | 2026-07-14 | Codex | 舊路由全部 superseded;所有路徑固定 GCP-A→GCP-B→host111→Gemini,legacy providers 僅 shadow metadata |
|
||||
| v5.1 | 2026-07-15 | Codex | Claude 以 paid/cost/canary gate 插入 host111 與 Gemini 之間;legacy providers 維持 shadow metadata |
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -36,8 +37,10 @@ from __future__ import annotations
|
||||
import asyncio
|
||||
import hashlib
|
||||
import json as _json
|
||||
import re
|
||||
import time
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import UTC, datetime
|
||||
from enum import Enum
|
||||
from typing import TYPE_CHECKING, Protocol
|
||||
|
||||
@@ -50,6 +53,8 @@ if TYPE_CHECKING:
|
||||
from src.core.config import get_settings
|
||||
from src.services.ai_provider_policy import (
|
||||
NON_EXECUTABLE_SHADOW_PROVIDERS,
|
||||
PAID_PROVIDER_ORDER,
|
||||
PRODUCTION_OLLAMA_ORDER,
|
||||
PRODUCTION_PROVIDER_ORDER,
|
||||
normalize_production_execution_order,
|
||||
)
|
||||
@@ -70,6 +75,24 @@ from src.services.model_registry import get_model_registry
|
||||
|
||||
logger = structlog.get_logger(__name__)
|
||||
|
||||
_DURABLE_OLLAMA_UNAVAILABLE_RECEIPT_SCHEMA = "ollama_unavailable_receipt_v1"
|
||||
_DURABLE_OLLAMA_UNAVAILABLE_RECEIPT_PREFIX = "ai:ollama:unavailable_receipt:"
|
||||
_DURABLE_OLLAMA_UNAVAILABLE_RECEIPT_ID = re.compile(r"^[A-Za-z0-9:_-]{1,128}$")
|
||||
_DURABLE_OLLAMA_UNAVAILABLE_MAX_SECONDS = 300.0
|
||||
|
||||
|
||||
def _is_paid_provider_identity(provider: object) -> bool:
|
||||
"""Treat decorated paid identities as paid and never cache-executable."""
|
||||
|
||||
normalized = str(provider or "").strip().lower()
|
||||
return any(
|
||||
normalized == paid
|
||||
or normalized.startswith(f"{paid}_")
|
||||
or normalized.startswith(f"{paid}:")
|
||||
or normalized.startswith(f"{paid}-")
|
||||
for paid in PAID_PROVIDER_ORDER
|
||||
)
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# Provider 定義
|
||||
@@ -214,7 +237,7 @@ class AIRouter:
|
||||
動態選擇最適合的 AI Provider 和模型。
|
||||
|
||||
Intent/risk/complexity only choose an Ollama model. They never choose a
|
||||
provider. Provider placement is the global four-hop production contract.
|
||||
provider. Provider placement is the global five-hop production contract.
|
||||
"""
|
||||
|
||||
def __init__(self):
|
||||
@@ -252,6 +275,7 @@ class AIRouter:
|
||||
(AIProviderEnum.OLLAMA_GCP_A, self._ollama_default),
|
||||
(AIProviderEnum.OLLAMA_GCP_B, self._ollama_default),
|
||||
(AIProviderEnum.OLLAMA_LOCAL, self._ollama_default),
|
||||
(AIProviderEnum.CLAUDE, self._claude_default),
|
||||
(AIProviderEnum.GEMINI, self._gemini_default),
|
||||
]
|
||||
|
||||
@@ -487,7 +511,7 @@ class AIRouter:
|
||||
選擇 Provider 和模型 (Phase 13.3 #87 核心邏輯)
|
||||
|
||||
所有 intent / risk / complexity 只選擇 Ollama model;provider
|
||||
execution order 永遠是 GCP-A、GCP-B、host111,最後才是 Gemini。
|
||||
execution order 永遠是 GCP-A、GCP-B、host111、Claude,最後才是 Gemini。
|
||||
|
||||
Args:
|
||||
intent: 正規化後的意圖
|
||||
@@ -521,7 +545,7 @@ class AIRouter:
|
||||
model = self._ollama_summary if score <= 1 else self._ollama_default
|
||||
reason = (
|
||||
f"複雜度={score}/5, 風險={risk.value};全域 production chain "
|
||||
"GCP-A -> GCP-B -> host111 -> Gemini"
|
||||
"GCP-A -> GCP-B -> host111 -> Claude -> Gemini"
|
||||
)
|
||||
return provider, model, reason
|
||||
|
||||
@@ -556,7 +580,7 @@ class AIRouter:
|
||||
# DEPRECATED 2026-04-28 — 已由 _build_fallback_chain_for_intent 取代,無呼叫方
|
||||
建立 Fallback 鏈 (排除已選 Provider)
|
||||
|
||||
Fallback 順序: GCP-A → GCP-B → host111 → Gemini
|
||||
Fallback 順序: GCP-A → GCP-B → host111 → Claude → Gemini
|
||||
|
||||
Args:
|
||||
selected_provider: 已選擇的 Provider
|
||||
@@ -719,7 +743,7 @@ class AIRouter:
|
||||
"condition": "all intents, risks, complexities, tools and background jobs",
|
||||
"provider": "ollama_gcp_a",
|
||||
"provider_order": list(PRODUCTION_PROVIDER_ORDER),
|
||||
"reason": "single global production route; legacy providers are shadow metadata only",
|
||||
"reason": "single global production route; paid providers require cost gates and legacy providers remain shadow metadata only",
|
||||
}
|
||||
]
|
||||
|
||||
@@ -917,6 +941,10 @@ class AIRouterExecutor:
|
||||
self._circuit_breakers[name] = _SimpleCircuitBreaker(
|
||||
name, failure_threshold=10, recovery_timeout=30.0
|
||||
)
|
||||
elif name == "claude":
|
||||
self._circuit_breakers[name] = _SimpleCircuitBreaker(
|
||||
name, failure_threshold=3, recovery_timeout=120.0
|
||||
)
|
||||
else:
|
||||
self._circuit_breakers[name] = _SimpleCircuitBreaker(name)
|
||||
return self._circuit_breakers[name]
|
||||
@@ -926,10 +954,203 @@ class AIRouterExecutor:
|
||||
"""生成 Cache Key (與 openclaw.py 相容)"""
|
||||
ctx_hash = ""
|
||||
if context:
|
||||
ctx_hash = f":{context.get('alert_type', '')}:{context.get('target_resource', '')}"
|
||||
cloud_context = context.get("paid_cloud_context") or {}
|
||||
receipt = (
|
||||
cloud_context.get("cloud_sanitization_receipt") or {}
|
||||
if isinstance(cloud_context, dict)
|
||||
else {}
|
||||
)
|
||||
ctx_hash = (
|
||||
f":{context.get('alert_type', '')}:{context.get('target_resource', '')}:"
|
||||
f"{(cloud_context or {}).get('data_classification', context.get('data_classification', ''))}:"
|
||||
f"{receipt.get('receipt_id', '') if isinstance(receipt, dict) else ''}"
|
||||
)
|
||||
content = f"{prompt}{ctx_hash}"
|
||||
return f"llm_cache:{hashlib.sha256(content.encode()).hexdigest()[:16]}"
|
||||
|
||||
@staticmethod
|
||||
def _paid_alert_execution_inputs(
|
||||
prompt: str,
|
||||
context: dict | None,
|
||||
) -> tuple[str | None, dict | None, str | None]:
|
||||
"""Validate an allowlist sanitizer receipt and isolate cloud inputs."""
|
||||
|
||||
is_alert = bool(context) and any(
|
||||
key in context
|
||||
for key in ("alert_type", "alertname", "alert_name", "incident_id", "signals")
|
||||
)
|
||||
if not is_alert:
|
||||
return prompt, context, None
|
||||
cloud_context = (context or {}).get("paid_cloud_context")
|
||||
paid_prompt = (context or {}).get("paid_cloud_prompt")
|
||||
if not isinstance(cloud_context, dict) or not isinstance(paid_prompt, str):
|
||||
return None, None, "alert_cloud_sanitization_receipt_missing"
|
||||
receipt = cloud_context.get("cloud_sanitization_receipt")
|
||||
payload = cloud_context.get("sanitized_payload")
|
||||
if not isinstance(receipt, dict) or not isinstance(payload, dict):
|
||||
return None, None, "alert_cloud_sanitization_receipt_missing"
|
||||
if (
|
||||
receipt.get("schema_version") != "openclaw_paid_cloud_sanitization_v1"
|
||||
or receipt.get("status") != "verified"
|
||||
or receipt.get("sanitizer") != "src.services.sanitization_service.sanitize"
|
||||
or receipt.get("classifier") != "allowlisted_structured_alert_fields_v1"
|
||||
or receipt.get("dlp_policy") != "paid_cloud_allowlist_dlp_v1"
|
||||
or receipt.get("raw_payload_forwarded") is not False
|
||||
or receipt.get("raw_log_payload_forwarded") is not False
|
||||
or receipt.get("secret_value_exposed") is not False
|
||||
or receipt.get("private_network_value_exposed") is not False
|
||||
):
|
||||
return None, None, "alert_cloud_sanitization_receipt_invalid"
|
||||
canonical = _json.dumps(
|
||||
payload,
|
||||
ensure_ascii=False,
|
||||
sort_keys=True,
|
||||
separators=(",", ":"),
|
||||
)
|
||||
payload_sha256 = hashlib.sha256(canonical.encode("utf-8")).hexdigest()
|
||||
prompt_sha256 = hashlib.sha256(paid_prompt.encode("utf-8")).hexdigest()
|
||||
if receipt.get("payload_sha256") != payload_sha256:
|
||||
return None, None, "alert_cloud_sanitization_payload_mismatch"
|
||||
if cloud_context.get("paid_prompt_sha256") != prompt_sha256:
|
||||
return None, None, "alert_cloud_sanitized_prompt_mismatch"
|
||||
# Receipt fields are caller-controlled until they are rebound to the
|
||||
# original alert through the canonical allowlist sanitizer. A matching
|
||||
# self-signed hash alone is not an authority boundary.
|
||||
from src.services.openclaw import build_paid_cloud_alert_context
|
||||
|
||||
rebuilt_context, rebuilt_reason = build_paid_cloud_alert_context(context)
|
||||
if rebuilt_context is None:
|
||||
return None, None, rebuilt_reason
|
||||
rebuilt_receipt = rebuilt_context.get("cloud_sanitization_receipt") or {}
|
||||
if (
|
||||
rebuilt_context.get("sanitized_payload") != payload
|
||||
or rebuilt_receipt.get("receipt_id") != receipt.get("receipt_id")
|
||||
):
|
||||
return None, None, "alert_cloud_sanitization_receipt_untrusted"
|
||||
restricted_context = dict(cloud_context)
|
||||
for correlation_field in ("trace_id", "run_id", "work_item_id"):
|
||||
if (context or {}).get(correlation_field):
|
||||
restricted_context[correlation_field] = (context or {})[
|
||||
correlation_field
|
||||
]
|
||||
return paid_prompt, restricted_context, None
|
||||
|
||||
@staticmethod
|
||||
async def _bounded_ollama_unavailable_receipts(
|
||||
context: dict | None,
|
||||
) -> set[str]:
|
||||
"""Verify current-run free-lane unavailability against durable truth.
|
||||
|
||||
Caller-provided JSON is only a lookup hint. The full receipt must exist
|
||||
independently in Redis, carry a short expiry, and acknowledge a
|
||||
completed failover-manager verifier before it can replace an attempt.
|
||||
"""
|
||||
|
||||
run_id = str((context or {}).get("run_id") or "")
|
||||
if not run_id:
|
||||
return set()
|
||||
raw_receipts = (context or {}).get("ollama_unavailable_receipts") or []
|
||||
if isinstance(raw_receipts, dict):
|
||||
raw_receipts = list(raw_receipts.values())
|
||||
valid: set[str] = set()
|
||||
allowed_reasons = {
|
||||
"endpoint_unreachable",
|
||||
"health_probe_failed",
|
||||
"provider_not_registered",
|
||||
"provider_disabled",
|
||||
"bounded_deadline_exhausted",
|
||||
}
|
||||
try:
|
||||
from src.core.redis_client import get_redis
|
||||
|
||||
redis = get_redis()
|
||||
except Exception as redis_error:
|
||||
logger.warning(
|
||||
"ai_router_ollama_unavailable_receipt_store_unavailable",
|
||||
error_type=type(redis_error).__name__,
|
||||
)
|
||||
return set()
|
||||
now = datetime.now(UTC)
|
||||
bounded_receipts = (
|
||||
raw_receipts[: len(PRODUCTION_OLLAMA_ORDER)]
|
||||
if isinstance(raw_receipts, list)
|
||||
else []
|
||||
)
|
||||
for receipt in bounded_receipts:
|
||||
if not isinstance(receipt, dict):
|
||||
continue
|
||||
provider = str(receipt.get("provider") or "")
|
||||
receipt_id = str(receipt.get("receipt_id") or "")
|
||||
if not (
|
||||
provider in PRODUCTION_OLLAMA_ORDER
|
||||
and receipt.get("schema_version")
|
||||
== _DURABLE_OLLAMA_UNAVAILABLE_RECEIPT_SCHEMA
|
||||
and receipt.get("source") == "ollama_failover_manager"
|
||||
and receipt.get("status") == "bounded_unavailable"
|
||||
and str(receipt.get("run_id") or "") == run_id
|
||||
and receipt.get("check_completed") is True
|
||||
and str(receipt.get("reason") or "") in allowed_reasons
|
||||
and _DURABLE_OLLAMA_UNAVAILABLE_RECEIPT_ID.fullmatch(receipt_id)
|
||||
):
|
||||
continue
|
||||
try:
|
||||
durable_raw = await redis.get(
|
||||
f"{_DURABLE_OLLAMA_UNAVAILABLE_RECEIPT_PREFIX}{receipt_id}"
|
||||
)
|
||||
if isinstance(durable_raw, bytes):
|
||||
durable_raw = durable_raw.decode("utf-8")
|
||||
durable = _json.loads(durable_raw) if durable_raw else None
|
||||
except Exception as receipt_error:
|
||||
logger.warning(
|
||||
"ai_router_ollama_unavailable_receipt_read_failed",
|
||||
provider=provider,
|
||||
error_type=type(receipt_error).__name__,
|
||||
)
|
||||
continue
|
||||
if not isinstance(durable, dict):
|
||||
continue
|
||||
contract_fields = (
|
||||
"schema_version",
|
||||
"receipt_id",
|
||||
"source",
|
||||
"provider",
|
||||
"run_id",
|
||||
"status",
|
||||
"check_completed",
|
||||
"reason",
|
||||
"observed_at",
|
||||
"expires_at",
|
||||
)
|
||||
if any(durable.get(field) != receipt.get(field) for field in contract_fields):
|
||||
continue
|
||||
if (
|
||||
durable.get("durable_write_ack") is not True
|
||||
or durable.get("verifier_status") != "verified"
|
||||
or durable.get("verified_by")
|
||||
!= "ollama_failover_manager.health_probe"
|
||||
):
|
||||
continue
|
||||
try:
|
||||
observed_at = datetime.fromisoformat(
|
||||
str(durable.get("observed_at") or "").replace("Z", "+00:00")
|
||||
)
|
||||
expires_at = datetime.fromisoformat(
|
||||
str(durable.get("expires_at") or "").replace("Z", "+00:00")
|
||||
)
|
||||
except ValueError:
|
||||
continue
|
||||
if observed_at.tzinfo is None or expires_at.tzinfo is None:
|
||||
continue
|
||||
lifetime = (expires_at - observed_at).total_seconds()
|
||||
if not (
|
||||
0 < lifetime <= _DURABLE_OLLAMA_UNAVAILABLE_MAX_SECONDS
|
||||
and observed_at <= now
|
||||
and now < expires_at
|
||||
):
|
||||
continue
|
||||
valid.add(provider)
|
||||
return valid
|
||||
|
||||
async def execute(
|
||||
self,
|
||||
prompt: str,
|
||||
@@ -990,39 +1211,89 @@ class AIRouterExecutor:
|
||||
error="production_route_missing_ollama_lane",
|
||||
)
|
||||
|
||||
# Paid-provider runtime control is evaluated before cache lookup. This
|
||||
# prevents a previously cached Gemini response from bypassing a newly
|
||||
# asserted durable disable and prevents Redis read failures from
|
||||
# silently re-adding the paid lane.
|
||||
# Paid-provider runtime control is evaluated before cache lookup. This
|
||||
# prevents a cached cloud response from bypassing a newly asserted
|
||||
# durable disable and prevents Redis read failures from silently
|
||||
# re-adding either paid lane.
|
||||
preflight_errors: list[str] = []
|
||||
if "gemini" in provider_order:
|
||||
from src.services.ai_providers.interfaces import (
|
||||
is_provider_enabled_by_env,
|
||||
from src.services.ai_control import is_provider_disabled
|
||||
from src.services.ai_providers.interfaces import (
|
||||
cloud_context_block_reason,
|
||||
is_provider_enabled_by_env,
|
||||
)
|
||||
|
||||
paid_prompt, paid_context, paid_input_block = self._paid_alert_execution_inputs(
|
||||
prompt,
|
||||
context,
|
||||
)
|
||||
cloud_privacy_block = paid_input_block or cloud_context_block_reason(
|
||||
paid_context,
|
||||
require_local=require_local,
|
||||
)
|
||||
if cloud_privacy_block:
|
||||
blocked_paid_providers = [
|
||||
provider
|
||||
for provider in PAID_PROVIDER_ORDER
|
||||
if provider in provider_order
|
||||
]
|
||||
provider_order = [
|
||||
provider
|
||||
for provider in provider_order
|
||||
if provider not in PAID_PROVIDER_ORDER
|
||||
]
|
||||
preflight_errors.extend(
|
||||
f"{provider}: {cloud_privacy_block}"
|
||||
for provider in blocked_paid_providers
|
||||
)
|
||||
|
||||
if not is_provider_enabled_by_env("gemini"):
|
||||
provider_order = [
|
||||
provider for provider in provider_order if provider != "gemini"
|
||||
]
|
||||
preflight_errors.append("gemini: disabled_by_environment")
|
||||
logger.info("ai_router_provider_environment_disabled_pre_cache", provider="gemini")
|
||||
else:
|
||||
logger.warning(
|
||||
"ai_router_cloud_context_blocked_pre_cache",
|
||||
reason=cloud_privacy_block,
|
||||
blocked_providers=blocked_paid_providers,
|
||||
)
|
||||
if any(provider in provider_order for provider in PAID_PROVIDER_ORDER):
|
||||
for paid_provider in PAID_PROVIDER_ORDER:
|
||||
if paid_provider not in provider_order:
|
||||
continue
|
||||
if not is_provider_enabled_by_env(paid_provider):
|
||||
provider_order = [
|
||||
provider
|
||||
for provider in provider_order
|
||||
if provider != paid_provider
|
||||
]
|
||||
preflight_errors.append(
|
||||
f"{paid_provider}: disabled_by_environment"
|
||||
)
|
||||
logger.info(
|
||||
"ai_router_provider_environment_disabled_pre_cache",
|
||||
provider=paid_provider,
|
||||
)
|
||||
continue
|
||||
try:
|
||||
from src.services.ai_control import is_provider_disabled
|
||||
|
||||
if await is_provider_disabled("gemini"):
|
||||
if await is_provider_disabled(paid_provider):
|
||||
provider_order = [
|
||||
provider for provider in provider_order if provider != "gemini"
|
||||
provider
|
||||
for provider in provider_order
|
||||
if provider != paid_provider
|
||||
]
|
||||
preflight_errors.append("gemini: disabled_by_runtime_control")
|
||||
logger.info("ai_router_provider_disabled_pre_cache", provider="gemini")
|
||||
preflight_errors.append(
|
||||
f"{paid_provider}: disabled_by_runtime_control"
|
||||
)
|
||||
logger.info(
|
||||
"ai_router_provider_disabled_pre_cache",
|
||||
provider=paid_provider,
|
||||
)
|
||||
except Exception as disable_check_error:
|
||||
provider_order = [
|
||||
provider for provider in provider_order if provider != "gemini"
|
||||
provider
|
||||
for provider in provider_order
|
||||
if provider != paid_provider
|
||||
]
|
||||
preflight_errors.append("gemini: disabled_state_unavailable")
|
||||
preflight_errors.append(
|
||||
f"{paid_provider}: disabled_state_unavailable"
|
||||
)
|
||||
logger.warning(
|
||||
"ai_router_gemini_disable_state_unavailable_pre_cache_blocked",
|
||||
"ai_router_paid_disable_state_unavailable_pre_cache_blocked",
|
||||
provider=paid_provider,
|
||||
error_type=type(disable_check_error).__name__,
|
||||
)
|
||||
|
||||
@@ -1035,6 +1306,14 @@ class AIRouterExecutor:
|
||||
if cached:
|
||||
data = _json.loads(cached)
|
||||
cached_provider = data.get("provider", "cache")
|
||||
if _is_paid_provider_identity(cached_provider):
|
||||
logger.info(
|
||||
"ai_router_paid_cache_bypassed",
|
||||
cache_key=cache_key[:30],
|
||||
cached_provider=cached_provider,
|
||||
reason="current_policy_and_ordered_attempt_recheck_required",
|
||||
)
|
||||
raise ValueError("paid provider cache is non-executable")
|
||||
provider_allowed = cached_provider in provider_order
|
||||
ollama_first_required = (
|
||||
bool(context)
|
||||
@@ -1082,20 +1361,27 @@ class AIRouterExecutor:
|
||||
from_cache=True,
|
||||
)
|
||||
except Exception as e:
|
||||
logger.debug("ai_router_cache_read_failed", error=str(e))
|
||||
logger.debug(
|
||||
"ai_router_cache_read_failed",
|
||||
error_type=type(e).__name__,
|
||||
)
|
||||
|
||||
# ③ 遍歷 Provider + 閘門 (D3)
|
||||
# 2026-04-02 ogt: C1 修復 — 建立 Langfuse Trace (D5)
|
||||
# 包住整個執行鏈,記錄每個 Provider 的 generation
|
||||
try:
|
||||
from src.services.langfuse_client import langfuse_trace
|
||||
alert_type_value = str((context or {}).get("alert_type", ""))
|
||||
_lf_trace_ctx = langfuse_trace(
|
||||
"ai_router_execute",
|
||||
metadata={
|
||||
"provider_order": provider_order,
|
||||
"prompt_length": len(prompt),
|
||||
"require_local": require_local,
|
||||
"alert_type": (context or {}).get("alert_type", ""),
|
||||
"alert_type_length": len(alert_type_value),
|
||||
"alert_type_sha256": hashlib.sha256(
|
||||
alert_type_value.encode("utf-8")
|
||||
).hexdigest(),
|
||||
},
|
||||
)
|
||||
_lf_trace_ctx.__enter__()
|
||||
@@ -1104,6 +1390,10 @@ class AIRouterExecutor:
|
||||
|
||||
errors: list[str] = list(preflight_errors)
|
||||
attempted_providers: set[str] = set()
|
||||
bounded_unavailable_providers = (
|
||||
await self._bounded_ollama_unavailable_receipts(context)
|
||||
)
|
||||
circuit_open_ollama_providers: set[str] = set()
|
||||
provider_timeout_seconds: float | None = None
|
||||
try:
|
||||
configured_provider_timeout = float(
|
||||
@@ -1170,6 +1460,40 @@ class AIRouterExecutor:
|
||||
errors.append(f"{provider_name}: not_registered")
|
||||
continue
|
||||
|
||||
if provider_name in PAID_PROVIDER_ORDER:
|
||||
unresolved_ollama_hops = [
|
||||
hop
|
||||
for hop in PRODUCTION_OLLAMA_ORDER
|
||||
if hop not in attempted_providers
|
||||
and hop not in bounded_unavailable_providers
|
||||
]
|
||||
if circuit_open_ollama_providers:
|
||||
errors.append(
|
||||
f"{provider_name}: cloud_blocked_ollama_circuit_open("
|
||||
f"{','.join(sorted(circuit_open_ollama_providers))})"
|
||||
)
|
||||
logger.warning(
|
||||
"ai_router_cloud_blocked_by_ollama_circuit_open",
|
||||
provider=provider_name,
|
||||
circuit_open_providers=sorted(circuit_open_ollama_providers),
|
||||
)
|
||||
continue
|
||||
if unresolved_ollama_hops:
|
||||
errors.append(
|
||||
f"{provider_name}: cloud_blocked_ollama_attempt_receipts_missing("
|
||||
f"{','.join(unresolved_ollama_hops)})"
|
||||
)
|
||||
logger.warning(
|
||||
"ai_router_cloud_blocked_missing_ollama_attempt_receipts",
|
||||
provider=provider_name,
|
||||
unresolved_ollama_hops=unresolved_ollama_hops,
|
||||
attempted_providers=sorted(attempted_providers),
|
||||
bounded_unavailable_providers=sorted(
|
||||
bounded_unavailable_providers
|
||||
),
|
||||
)
|
||||
continue
|
||||
|
||||
# 隱私過濾 (D7)
|
||||
# 2026-04-27 Claude Sonnet 4.6: F6 — privacy_skip 不設 _last_attempted_provider(未嘗試)
|
||||
if require_local and provider.privacy_level != "local":
|
||||
@@ -1190,6 +1514,8 @@ class AIRouterExecutor:
|
||||
# 閘門 1: Circuit Breaker (per-provider, C2 修復)
|
||||
cb = self._get_circuit_breaker(provider_name)
|
||||
if cb.is_open():
|
||||
if provider_name in PRODUCTION_OLLAMA_ORDER:
|
||||
circuit_open_ollama_providers.add(provider_name)
|
||||
if alert_requires_ollama_before_cloud and provider_name.startswith("ollama"):
|
||||
logger.warning(
|
||||
"ai_router_alert_ollama_circuit_bypassed",
|
||||
@@ -1204,9 +1530,10 @@ class AIRouterExecutor:
|
||||
|
||||
# 閘門 2: Rate Limiter
|
||||
# 2026-04-02 Claude Code: Phase 24 B3 + C1 修復 — Rate Limiter (含 openclaw_nemo)
|
||||
# Gemini owns an atomic requests/tokens/cost reservation inside the
|
||||
# provider so every generation path is guarded exactly once.
|
||||
if provider_name in ("openclaw_nemo", "nemotron", "claude"):
|
||||
# Paid providers own atomic requests/tokens/cost reservations inside
|
||||
# their provider implementations; legacy free/provider lanes retain
|
||||
# this compatibility rate gate.
|
||||
if provider_name in ("openclaw_nemo", "nemotron"):
|
||||
try:
|
||||
from src.services.ai_rate_limiter import get_ai_rate_limiter
|
||||
rate_limiter = get_ai_rate_limiter()
|
||||
@@ -1223,7 +1550,27 @@ class AIRouterExecutor:
|
||||
async with sem:
|
||||
try:
|
||||
attempted_providers.add(provider_name)
|
||||
provider_call = provider.analyze(prompt, context)
|
||||
provider_prompt = (
|
||||
paid_prompt
|
||||
if provider_name in PAID_PROVIDER_ORDER
|
||||
else prompt
|
||||
)
|
||||
provider_context = (
|
||||
paid_context
|
||||
if provider_name in PAID_PROVIDER_ORDER
|
||||
else context
|
||||
)
|
||||
if provider_name in PAID_PROVIDER_ORDER and (
|
||||
provider_prompt is None or provider_context is None
|
||||
):
|
||||
errors.append(
|
||||
f"{provider_name}: paid_cloud_execution_inputs_missing"
|
||||
)
|
||||
continue
|
||||
provider_call = provider.analyze(
|
||||
provider_prompt,
|
||||
provider_context,
|
||||
)
|
||||
result = (
|
||||
await asyncio.wait_for(
|
||||
provider_call,
|
||||
@@ -1238,24 +1585,30 @@ class AIRouterExecutor:
|
||||
cb.record_success()
|
||||
|
||||
# 記錄費用
|
||||
if result.cost_usd > 0 and provider_name != "gemini":
|
||||
if (
|
||||
result.cost_usd > 0
|
||||
and provider_name not in PAID_PROVIDER_ORDER
|
||||
):
|
||||
try:
|
||||
rate_limiter = get_ai_rate_limiter()
|
||||
await rate_limiter.record_cost(provider_name, result.cost_usd)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
# 寫入 Cache (D4)
|
||||
try:
|
||||
redis = get_redis()
|
||||
cache_data = _json.dumps({
|
||||
"response": result.raw_response,
|
||||
"provider": result.provider,
|
||||
"cached_at": time.strftime("%Y-%m-%dT%H:%M:%S+08:00"),
|
||||
})
|
||||
await redis.set(cache_key, cache_data, ex=cache_ttl)
|
||||
except Exception:
|
||||
pass
|
||||
# Paid responses are deliberately non-cacheable: every
|
||||
# paid candidate must re-run current durable policy,
|
||||
# privacy, canary, cost, and ordered-attempt gates.
|
||||
if provider_name not in PAID_PROVIDER_ORDER:
|
||||
try:
|
||||
redis = get_redis()
|
||||
cache_data = _json.dumps({
|
||||
"response": result.raw_response,
|
||||
"provider": result.provider,
|
||||
"cached_at": time.strftime("%Y-%m-%dT%H:%M:%S+08:00"),
|
||||
})
|
||||
await redis.set(cache_key, cache_data, ex=cache_ttl)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
logger.info(
|
||||
"ai_router_execute_success",
|
||||
@@ -1267,13 +1620,46 @@ class AIRouterExecutor:
|
||||
# D5: 記錄 Langfuse generation
|
||||
if _lf_trace_ctx:
|
||||
try:
|
||||
paid_provider = provider_name in PAID_PROVIDER_ORDER
|
||||
prompt_hash = hashlib.sha256(
|
||||
provider_prompt.encode("utf-8")
|
||||
).hexdigest()
|
||||
response_hash = hashlib.sha256(
|
||||
result.raw_response.encode("utf-8")
|
||||
).hexdigest()
|
||||
receipt = (
|
||||
(provider_context or {}).get(
|
||||
"cloud_sanitization_receipt"
|
||||
)
|
||||
or {}
|
||||
)
|
||||
_lf_trace_ctx.generation(
|
||||
name=f"{provider_name}_call",
|
||||
model=provider_name,
|
||||
input=prompt[:500],
|
||||
output=result.raw_response[:500],
|
||||
input=None if paid_provider else provider_prompt[:500],
|
||||
output=(
|
||||
None
|
||||
if paid_provider
|
||||
else result.raw_response[:500]
|
||||
),
|
||||
usage={"total": result.tokens} if result.tokens else None,
|
||||
metadata={"cost_usd": result.cost_usd, "latency_ms": round(result.latency_ms, 1)},
|
||||
metadata={
|
||||
"cost_usd": result.cost_usd,
|
||||
"latency_ms": round(result.latency_ms, 1),
|
||||
"content_redacted": paid_provider,
|
||||
"prompt_length": len(provider_prompt),
|
||||
"prompt_sha256": prompt_hash,
|
||||
"response_length": len(result.raw_response),
|
||||
"response_sha256": response_hash,
|
||||
"sanitization_receipt_id": (
|
||||
receipt.get("receipt_id")
|
||||
if isinstance(receipt, dict)
|
||||
else None
|
||||
),
|
||||
"generation_receipt_id": result.audit_metadata.get(
|
||||
"generation_receipt_id"
|
||||
),
|
||||
},
|
||||
)
|
||||
_lf_trace_ctx.__exit__(None, None, None)
|
||||
except Exception:
|
||||
@@ -1282,14 +1668,53 @@ class AIRouterExecutor:
|
||||
|
||||
# Provider 回傳 success=False
|
||||
cb.record_failure()
|
||||
errors.append(f"{provider_name}: {result.error}")
|
||||
logger.warning("ai_router_provider_failed", provider=provider_name, error=result.error)
|
||||
if provider_name in PAID_PROVIDER_ORDER:
|
||||
error_text = str(result.error or "")
|
||||
error_sha256 = hashlib.sha256(
|
||||
error_text.encode("utf-8")
|
||||
).hexdigest()
|
||||
errors.append(
|
||||
f"{provider_name}: provider_failed(error_sha256={error_sha256})"
|
||||
)
|
||||
logger.warning(
|
||||
"ai_router_paid_provider_failed",
|
||||
provider=provider_name,
|
||||
error_length=len(error_text),
|
||||
error_sha256=error_sha256,
|
||||
)
|
||||
else:
|
||||
errors.append(f"{provider_name}: {result.error}")
|
||||
logger.warning(
|
||||
"ai_router_provider_failed",
|
||||
provider=provider_name,
|
||||
error=result.error,
|
||||
)
|
||||
# 2026-04-27 A2: 記錄失敗的 provider,供下輪迭代計算 fallback metric
|
||||
_last_attempted_provider = provider_name
|
||||
|
||||
except Exception as e:
|
||||
errors.append(f"{provider_name}: {e}")
|
||||
logger.warning("ai_router_provider_exception", provider=provider_name, error=str(e))
|
||||
if provider_name in PAID_PROVIDER_ORDER:
|
||||
error_text = str(e)
|
||||
error_sha256 = hashlib.sha256(
|
||||
error_text.encode("utf-8")
|
||||
).hexdigest()
|
||||
errors.append(
|
||||
f"{provider_name}: provider_exception(error_sha256={error_sha256})"
|
||||
)
|
||||
logger.warning(
|
||||
"ai_router_paid_provider_exception",
|
||||
provider=provider_name,
|
||||
error_type=type(e).__name__,
|
||||
error_length=len(error_text),
|
||||
error_sha256=error_sha256,
|
||||
)
|
||||
else:
|
||||
errors.append(f"{provider_name}: {e}")
|
||||
logger.warning(
|
||||
"ai_router_provider_exception",
|
||||
provider=provider_name,
|
||||
error=str(e),
|
||||
)
|
||||
# 2026-04-05 Claude Code: v4.3 — Timeout 不計 CB 失敗
|
||||
# NIM 偶爾 GPU 忙碌導致 27s,timeout 不代表 NIM 故障
|
||||
# 只有明確連線錯誤(非 timeout)才累積 CB 失敗次數
|
||||
|
||||
Reference in New Issue
Block a user