[阶段4.3] 接入自建 AI 网关,修复多模型层的四个真实缺陷

改用甲骨文机上已有的 ai-gateway (129.146.203.203:5100):它本身就
OpenAI 兼容,内部串联 nvidia/gemini/ollama 并轮换 4 个 Gemini key,
比在客户端自己串联更能吸收单厂商的配额和超时。回包里的 provider
字段透传为 meta.upstream,网关侧发生降级时前端也看得见。

fix(ai): 目录里两个 NVIDIA 模型 id 根本不存在
- qwen/qwen2.5-72b-instruct 和 deepseek-ai/deepseek-r1 是我凭印象写的,
  实际 GET /v1/models 里没有,调用一律 404
- 改为该账号清单里确实存在的 nemotron-49b / mistral-large,
  并在注释里写明 id 必须取自实时清单、不能猜

fix(ai): 请求被本机代理劫持导致网关不可达
- requests 默认读 HTTP_PROXY/ALL_PROXY,把发往甲骨文公网 IP 的请求
  也塞进了 127.0.0.1:7897,120s 后超时
- 按 provider 区分:境外厂商(Gemini/NVIDIA)仍走代理,自建网关直连
  (session.trust_env=False)

fix(ai): 承诺的按模型裁剪从未实现
- 模块注释写着 payload 按 (模型窗口, 天数预算) 取小者裁剪,但实际是
  用全局预算构建一次 prompt 发给链上所有模型;365 天数据对 Gemini
  的 1M 窗口无碍,却会撑爆 128k 的模型
- 新增 max_days_for(),在循环内按各模型窗口分别构建 prompt

fix(ai): 推理模型的思考过程吃光输出预算
- 网关首选 nemotron-3-ultra-550b 是推理模型,回答前先输出一段
  chain-of-thought;默认 1024 tokens 全被思考占用,JSON 还没开始
  就被截断
- max_tokens 改为可按 provider 声明,网关条目给 3000

fix(ai): 配置在 import 时被冻结
- DEFAULT_CHAIN/TIMEOUT/DAY_BUDGET 是模块级常量,改环境变量不生效,
  且让开发机 .env 泄漏进测试进程(测试会读到真实 key 和链配置)
- 改为 default_chain()/default_timeout()/default_day_budget() 按调用读取
- conftest 增加 autouse fixture 清空全部 AI_* 变量,测试不再继承 .env

测试 (184 passed, 1 skipped):
- 新增 TestGatewayProvider: 透传 upstream、目标 URL/鉴权头、
  token 失效时继续降级
- 新增 TestProxyPolicy: 境外厂商与自建端点的代理策略相反
- 新增 TestPerModelSizing: 128k 模型收到的 prompt 必须小于 1M 模型
- 新增 TestMaxTokens: 推理端点预算大于默认,且真正写进两种 payload
- 新增 TestLazyConfig: 改环境变量立即生效
- mock 目标从 requests.post 改为 requests.Session.post

实测: 网关链路可返回合法 JSON,但 nemotron-550B 排队较久(约 160s),
故 AI_TIMEOUT_SECONDS 默认调到 180。

Co-Authored-By: Claude Haiku 4.5 <noreply@anthropic.com>
This commit is contained in:
ericwyuan
2026-08-23 17:44:51 +08:00
parent 8616a13525
commit 8882bf44a4
8 changed files with 672 additions and 2019 deletions

View File

@@ -23,11 +23,30 @@ import re
import requests
DEFAULT_TIMEOUT = float(os.environ.get("AI_TIMEOUT_SECONDS") or 45)
# Tunables are read per call rather than captured at import: module-level
# constants freeze whatever the environment held when the module first loaded,
# which both hides live config changes and leaks a developer's .env into tests.
FALLBACK_TIMEOUT = 60.0
FALLBACK_DAY_BUDGET = 365
# How many days of history to put in the prompt at most. Kept well below the
# model windows so the response always has room.
DEFAULT_DAY_BUDGET = int(os.environ.get("AI_DAY_BUDGET") or 365)
def default_timeout():
return float(os.environ.get("AI_TIMEOUT_SECONDS") or FALLBACK_TIMEOUT)
def default_day_budget():
"""Max days of history to put in a prompt, before per-model trimming."""
return int(os.environ.get("AI_DAY_BUDGET") or FALLBACK_DAY_BUDGET)
# Output cap. Deliberately modest: a long generation is what blows past an
# upstream's own timeout (the self-hosted gateway allows its adapters only
# 30-45s), and the reply here is a short JSON list, not an essay.
FALLBACK_MAX_TOKENS = 1024
def default_max_tokens():
return int(os.environ.get("AI_MAX_TOKENS") or FALLBACK_MAX_TOKENS)
SYSTEM_PROMPT = (
"你是一名严谨的健康数据分析助手负责解读用户的可穿戴设备Garmin数据。\n"
@@ -37,7 +56,8 @@ SYSTEM_PROMPT = (
"3. 给出具体、可执行的建议,而不是泛泛而谈。\n"
"4. 你不是医生,不做诊断;发现明显异常时建议用户咨询专业医师。\n"
"5. 用简体中文回答。\n\n"
"输出严格为 JSON 数组,每个元素形如:\n"
"输出严格为 JSON 数组,最多 5 条,每条 recommendation 不超过 120 字,\n"
"每个元素形如:\n"
'{"category": "睡眠", "recommendation": "……", "priority": "high|medium|low", '
'"basedOn": ["sleep_duration"]}\n'
"不要输出 JSON 以外的任何文字,不要用 markdown 代码块包裹。"
@@ -48,16 +68,54 @@ class AIError(Exception):
"""Raised when a provider cannot produce a completion."""
class Completion:
"""A model reply plus, where the endpoint reports it, the upstream that
actually served the request.
The self-hosted gateway multiplexes over nvidia/gemini/ollama and names
the winner in its response, so `upstream` is what makes a gateway-side
failover visible to the UI instead of silently invisible.
"""
__slots__ = ("text", "upstream")
def __init__(self, text, upstream=None):
self.text = text
self.upstream = upstream
# --- providers --------------------------------------------------------------
class Provider:
"""Base class. Subclasses turn a prompt into text."""
"""Base class. Subclasses turn a prompt into text.
`use_proxy` decides whether HTTP(S)_PROXY / ALL_PROXY from the environment
apply. It matters because the two kinds of endpoint want opposite answers:
overseas vendors (Gemini, NVIDIA) may only be reachable *through* a local
proxy, while a self-hosted box on a public IP is reachable directly and
breaks if forced through one.
"""
name = "base"
def __init__(self, model_id, context_window, api_key_env):
def __init__(
self, model_id, context_window, api_key_env, use_proxy=True, max_tokens=None
):
self.model_id = model_id
self.context_window = context_window
self.api_key_env = api_key_env
self.use_proxy = use_proxy
self._max_tokens = max_tokens
@property
def max_tokens(self):
"""Output cap for this endpoint.
Reasoning models emit a chain-of-thought *before* the answer, so a cap
sized for the answer alone gets spent on the thinking and truncates
before any JSON appears. Those endpoints therefore declare a larger
budget than the default.
"""
return self._max_tokens or default_max_tokens()
@property
def api_key(self):
@@ -66,7 +124,14 @@ class Provider:
def is_configured(self):
return bool(self.api_key)
def generate(self, prompt, timeout=DEFAULT_TIMEOUT):
def _session(self):
session = requests.Session()
# trust_env=False also drops netrc/CA-bundle env lookups, which is the
# intent here: talk to the host directly, exactly as configured.
session.trust_env = self.use_proxy
return session
def generate(self, prompt, timeout=None):
raise NotImplementedError
@@ -76,16 +141,20 @@ class GeminiProvider(Provider):
name = "gemini"
BASE = "https://generativelanguage.googleapis.com/v1beta/models"
def generate(self, prompt, timeout=DEFAULT_TIMEOUT):
def generate(self, prompt, timeout=None):
if not self.is_configured():
raise AIError(f"{self.api_key_env} 未配置")
timeout = timeout or default_timeout()
url = f"{self.BASE}/{self.model_id}:generateContent"
payload = {
"contents": [{"parts": [{"text": prompt}]}],
"generationConfig": {"temperature": 0.4},
"generationConfig": {
"temperature": 0.4,
"maxOutputTokens": self.max_tokens,
},
}
try:
resp = requests.post(
resp = self._session().post(
url,
headers={
"Content-Type": "application/json",
@@ -103,40 +172,70 @@ class GeminiProvider(Provider):
try:
body = resp.json()
parts = body["candidates"][0]["content"]["parts"]
return "".join(p.get("text", "") for p in parts)
return Completion("".join(p.get("text", "") for p in parts))
except (ValueError, KeyError, IndexError) as e:
raise AIError(f"gemini 响应格式异常: {e}") from e
class OpenAICompatProvider(Provider):
"""Any endpoint speaking the OpenAI chat-completions schema (NVIDIA NIM,
Ollama, vLLM, ...)."""
Ollama, vLLM, ...).
`requires_key=False` covers self-hosted runtimes such as Ollama, which
authenticate by network reachability rather than by a token. Those are
opt-in: they count as configured only once their base URL is set, so an
unset OLLAMA_BASE_URL keeps the entry out of the fallback chain.
"""
name = "openai-compat"
def __init__(self, model_id, context_window, api_key_env, base_url_env, default_base_url):
super().__init__(model_id, context_window, api_key_env)
self.base_url = os.environ.get(base_url_env) or default_base_url
def __init__(
self,
model_id,
context_window,
base_url_env,
default_base_url="",
api_key_env=None,
requires_key=True,
use_proxy=True,
max_tokens=None,
):
super().__init__(
model_id, context_window, api_key_env or "", use_proxy, max_tokens
)
self.base_url_env = base_url_env
self.default_base_url = default_base_url
self.requires_key = requires_key
def generate(self, prompt, timeout=DEFAULT_TIMEOUT):
@property
def base_url(self):
return os.environ.get(self.base_url_env) or self.default_base_url
def is_configured(self):
if not self.base_url:
return False
return bool(self.api_key) if self.requires_key else True
def generate(self, prompt, timeout=None):
if not self.is_configured():
raise AIError(f"{self.api_key_env} 未配置")
raise AIError(
f"{self.api_key_env} 未配置" if self.requires_key
else f"{self.base_url_env} 未配置"
)
timeout = timeout or default_timeout()
url = f"{self.base_url.rstrip('/')}/chat/completions"
payload = {
"model": self.model_id,
"messages": [{"role": "user", "content": prompt}],
"temperature": 0.4,
"max_tokens": 2048,
"max_tokens": self.max_tokens,
}
headers = {"Content-Type": "application/json"}
if self.api_key:
headers["Authorization"] = f"Bearer {self.api_key}"
try:
resp = requests.post(
url,
headers={
"Content-Type": "application/json",
"Authorization": f"Bearer {self.api_key}",
},
json=payload,
timeout=timeout,
resp = self._session().post(
url, headers=headers, json=payload, timeout=timeout
)
except requests.RequestException as e:
raise AIError(f"{self.model_id} 请求失败: {e}") from e
@@ -145,56 +244,91 @@ class OpenAICompatProvider(Provider):
raise AIError(f"{self.model_id} HTTP {resp.status_code}: {resp.text[:200]}")
try:
return resp.json()["choices"][0]["message"]["content"]
body = resp.json()
# `provider` is a gateway extension, absent from stock OpenAI
# responses — hence the .get rather than an index.
return Completion(
body["choices"][0]["message"]["content"], body.get("provider")
)
except (ValueError, KeyError, IndexError) as e:
raise AIError(f"{self.model_id} 响应格式异常: {e}") from e
# --- catalog ----------------------------------------------------------------
NVIDIA_BASE = "https://integrate.api.nvidia.com/v1"
def _nvidia(model_id, context_window):
return OpenAICompatProvider(
model_id=model_id,
context_window=context_window,
api_key_env="NVIDIA_API_KEY",
base_url_env="NVIDIA_BASE_URL",
default_base_url=NVIDIA_BASE,
)
def _build_catalog():
"""Model id -> Provider. Text-only models with large context windows."""
"""Model id -> Provider. Text-only models with large context windows.
The NVIDIA model strings below were taken from that account's live
`GET /v1/models` listing. Do not guess them: ids that merely look
plausible (`qwen/qwen2.5-72b-instruct`, `deepseek-ai/deepseek-r1`)
return HTTP 404 from this endpoint.
"""
return {
# Preferred entry: the self-hosted gateway on the Oracle box. It
# multiplexes over nvidia/gemini/ollama behind one OpenAI-compatible
# endpoint and rotates several Gemini keys, so it absorbs the quota
# and timeout failures that a single upstream hits on its own. Its
# reply names the upstream that served the request.
"gateway": OpenAICompatProvider(
model_id=os.environ.get("AI_GATEWAY_MODEL") or "ai-gateway-auto",
context_window=128_000,
api_key_env="AI_GATEWAY_TOKEN",
base_url_env="AI_GATEWAY_BASE_URL",
# Self-hosted and directly reachable: a local proxy would only
# add a hop that times out.
use_proxy=False,
# Its primary upstream is a reasoning model that thinks out loud
# before answering; at the default cap the trace consumed the whole
# budget and the reply was truncated before the JSON began.
max_tokens=3000,
),
# Direct upstreams, for pinning one vendor or for running without the
# gateway. These need their own keys in this app's .env.
"gemini-flash": GeminiProvider(
model_id="gemini-flash-latest",
context_window=1_000_000,
api_key_env="GEMINI_API_KEY",
),
"llama-70b": OpenAICompatProvider(
model_id="meta/llama-3.3-70b-instruct",
context_window=128_000,
api_key_env="NVIDIA_API_KEY",
base_url_env="NVIDIA_BASE_URL",
default_base_url="https://integrate.api.nvidia.com/v1",
),
"qwen-72b": OpenAICompatProvider(
model_id="qwen/qwen2.5-72b-instruct",
context_window=128_000,
api_key_env="NVIDIA_API_KEY",
base_url_env="NVIDIA_BASE_URL",
default_base_url="https://integrate.api.nvidia.com/v1",
),
"deepseek-r1": OpenAICompatProvider(
model_id="deepseek-ai/deepseek-r1",
context_window=128_000,
api_key_env="NVIDIA_API_KEY",
base_url_env="NVIDIA_BASE_URL",
default_base_url="https://integrate.api.nvidia.com/v1",
),
"llama-70b": _nvidia("meta/llama-3.3-70b-instruct", 128_000),
"nemotron-49b": _nvidia("nvidia/llama-3.3-nemotron-super-49b-v1.5", 128_000),
"mistral-large": _nvidia("mistralai/mistral-large-2-instruct", 128_000),
}
CATALOG = _build_catalog()
# Preference order used when no model is requested, and for fallback.
DEFAULT_CHAIN = [
m.strip()
for m in (os.environ.get("AI_MODEL_CHAIN") or "gemini-flash,llama-70b,qwen-72b").split(",")
if m.strip()
]
FALLBACK_CHAIN = "gateway,gemini-flash,llama-70b"
def default_chain():
"""Preference order, read from the environment on every call.
Deliberately not a module-level constant: it is read at request time so a
changed AI_MODEL_CHAIN takes effect without a restart, and so tests can
set it without reaching into module internals.
"""
raw = os.environ.get("AI_MODEL_CHAIN") or FALLBACK_CHAIN
return [m.strip() for m in raw.split(",") if m.strip()]
def list_models():
"""Catalog entries plus whether each one currently has credentials."""
chain = default_chain()
head = chain[0] if chain else None
return [
{
"id": mid,
@@ -202,7 +336,7 @@ def list_models():
"provider": p.name,
"contextWindow": p.context_window,
"configured": p.is_configured(),
"default": mid == DEFAULT_CHAIN[0] if DEFAULT_CHAIN else False,
"default": mid == head,
}
for mid, p in CATALOG.items()
]
@@ -215,7 +349,7 @@ def resolve_chain(preferred=None):
if preferred not in CATALOG:
raise AIError(f"未知模型: {preferred}")
chain.append(preferred)
for mid in DEFAULT_CHAIN:
for mid in default_chain():
if mid in CATALOG and mid not in chain:
chain.append(mid)
configured = [m for m in chain if CATALOG[m].is_configured()]
@@ -237,12 +371,13 @@ _CSV_COLUMNS = [
]
def build_prompt(summary, activities=None, day_budget=DEFAULT_DAY_BUDGET):
def build_prompt(summary, activities=None, day_budget=None):
"""Render health history as a compact CSV prompt.
CSV rather than JSON: roughly 4x fewer tokens for the same numbers, which
is what makes a full year of history practical to send.
"""
day_budget = day_budget if day_budget is not None else default_day_budget()
rows = summary[-day_budget:] if day_budget else summary
header = ",".join(label for _, label in _CSV_COLUMNS) + ",sleep_h,sleep_q"
lines = [header]
@@ -346,26 +481,49 @@ def parse_recommendations(text):
# --- entry point ------------------------------------------------------------
def generate(summary, activities=None, preferred_model=None, day_budget=DEFAULT_DAY_BUDGET):
# One CSV day is ~40 characters ≈ 10 tokens. Half the window is left for the
# system prompt, the activity table and the model's own answer.
_TOKENS_PER_DAY = 10
_WINDOW_UTILISATION = 0.5
def max_days_for(provider, day_budget=None):
"""How many days of history fit in this model's context window.
Models in the chain have windows that differ by more than an order of
magnitude (32k for a local Ollama vs 1M for Gemini), so the payload has to
be sized per model — a prompt that fits Gemini would overflow Ollama.
"""
day_budget = day_budget if day_budget is not None else default_day_budget()
fits = int(provider.context_window * _WINDOW_UTILISATION / _TOKENS_PER_DAY)
return max(1, min(day_budget, fits)) if day_budget else max(1, fits)
def generate(summary, activities=None, preferred_model=None, day_budget=None):
"""Ask the first healthy model in the chain for recommendations.
Returns (recommendations, meta). `meta` records which model answered and
which ones failed, so the UI can show what actually happened.
Returns (recommendations, meta). `meta` records which model answered, how
much history it actually saw, and every model that failed on the way —
the failures are kept even on success so a silent degradation to a weaker
model is still visible.
"""
chain = resolve_chain(preferred_model)
prompt = build_prompt(summary, activities, day_budget)
errors = []
for model_id in chain:
provider = CATALOG[model_id]
days = max_days_for(provider, day_budget)
prompt = build_prompt(summary, activities, days)
try:
raw = provider.generate(prompt)
recs = parse_recommendations(raw)
completion = provider.generate(prompt)
recs = parse_recommendations(completion.text)
return recs, {
"model": model_id,
"provider": provider.name,
"days": min(len(summary), day_budget) if day_budget else len(summary),
"upstream": completion.upstream,
"days": min(len(summary), days),
"fallbackFrom": [e["model"] for e in errors],
"errors": errors,
}
except AIError as e:
errors.append({"model": model_id, "error": str(e)})