zizi 069ae73e97 feat(aigc/worker): U3 trace 扩展键 + U4 缓存命中落账(B5/B9)
6c6g 后端 Tier0 计划 U3+U4(docs/plans/2026-06-17-001-...)。execution §5.8 trace+cost 落值。
- U3 extractTraceQuietly(Java)+ _extract_trace(worker)additive 抽 modelTier/escalationEvents/giveupDumpPath;
  仅真升档态(stage2/escalationEvents 非空)落 modelTier + 形态校验(非空 str/非空 list)→ 字节兼容(无扩展键键集与扩前一致)、Java/Python 两路同口径
- U4 cost.py 缓存命中折¥:两套字段取 max(DeepSeek prompt_cache_hit_tokens / MiniMax prompt_tokens_details.cached_tokens)
  + cache_ratio 折真实 quota;worker usage 捕获(_safe_int 兜脏不抛);pricing 不可达 tokens-only+costFallback 标记;不阻断仅观测
- 诚实边界:SAA Java 路无 cache 字段源→省略不伪造(follow-up);AgentScope 丢 DeepSeek 顶层字段(待真跑对账);worker 生产路完整落账

验证:mini-desktop Java 12+(SaaGraphDispatcherTraceTest)+ Python 73(test_cost 56/wg1_groupb 17)全绿;
codex 评审 NO-MERGE→2 P0(字节兼容/usage 阻断)+2 P1(pricing 口径/两路一致)全修后绿。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-17 22:26:27 +00:00

135 lines
6.0 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""_client.py —— W-G1 L1 裸路模型客户端(复用 spike 代理旁路 + 裸 openai 客户端 + usage 读取)。
L1 纪律(spike VERDICT §3/§6):裸 llm_client,不引 AgentScope(框架 token 膨胀+本地开销吃便宜档单价)。
唯一模型出口 = new-api OpenAI 兼容网关(BASE_URL/v1),绝不直连厂商 SDK;密钥只走环境变量/本地 .env,绝不入库。
"""
import os
import time
from pathlib import Path
from urllib.parse import urlparse
# ── 本地 .env 加载(wg1/gen-worker/.env,gitignored):密钥外置,不落代码 ──
_ENV_PATH = Path(__file__).resolve().parent.parent / ".env"
def _load_env():
"""把 wg1/gen-worker/.env 里的 KEY=VAL 注入环境(已存在的不覆盖)。"""
if not _ENV_PATH.exists():
return
for line in _ENV_PATH.read_text(encoding="utf-8").splitlines():
line = line.strip()
if not line or line.startswith("#") or "=" not in line:
continue
k, v = line.split("=", 1)
os.environ.setdefault(k.strip(), v.strip())
_load_env()
# 网关地址(可被环境变量覆盖;默认 new-api Tailscale 网关 OpenAI 兼容路径)。
BASE_URL = os.environ.get("NEWAPI_BASE_URL", "http://100.64.0.8:3000/v1")
# L1 目标模型(可被 WG1_MODEL 覆盖;冒烟默认最便宜档 deepseek-v4-flash)。
DEFAULT_MODEL = os.environ.get("WG1_MODEL", "deepseek-v4-flash")
# ── 代理旁路(spike VERDICT §4-坑1,必须在 import openai 前)──
# 本机 clash 代理(HTTP(S)_PROXY=127.0.0.1:7897) 会把发往 100.64.0.8 的请求经代理转发→502;
# httpx/openai 默认 trust_env=True,须把网关 host 并入大小写两个 NO_PROXY 才直连。
_GW_HOST = urlparse(BASE_URL).hostname or "100.64.0.8"
for _k in ("NO_PROXY", "no_proxy"):
_cur = os.environ.get(_k, "")
if _GW_HOST not in _cur:
os.environ[_k] = (_cur + "," + _GW_HOST) if _cur else _GW_HOST
import openai # noqa: E402 —— 必须在代理旁路之后导入
def get_api_key():
"""读 NEWAPI_KEY(缺失即抛,密钥不入库、复现前先设 .env 或 export)。"""
key = os.environ.get("NEWAPI_KEY")
if not key:
raise RuntimeError("NEWAPI_KEY 未设:请在 wg1/gen-worker/.env 写入或 export NEWAPI_KEY=...")
return key
def get_client():
"""构建指向 new-api 的裸 OpenAI 兼容客户端。"""
return openai.OpenAI(api_key=get_api_key(), base_url=BASE_URL, max_retries=2, timeout=180.0)
def _cache_hit_from_usage(usage):
"""从 openai 原始 usage 抽命中缓存的 prompt token(兼容两套字段取 max,U4/B9)。
DeepSeek `usage.prompt_cache_hit_tokens` / MiniMax `usage.prompt_tokens_details.cached_tokens`,取 max;
字段全缺/异常 → 0(按 miss 计,不报错)。逻辑与 cost.cache_hit_tokens 同口径,但本处只依赖 openai 原始对象、
不反向 import cost(避免 _client→cost 循环依赖;cost 已 import _client)。
"""
if usage is None:
return None
try:
cands = []
ds = getattr(usage, "prompt_cache_hit_tokens", None)
if ds is not None:
cands.append(int(ds))
details = getattr(usage, "prompt_tokens_details", None)
cached = getattr(details, "cached_tokens", None) if details is not None else None
if cached is not None:
cands.append(int(cached))
return max(cands) if cands else 0
except Exception:
return 0
def chat(model, system, user, max_tokens=16000, temperature=0.0, tries=3, retry_delay=2.0):
"""单次裸调用 chat.completions。确定性产出 temperature=0;瞬时网关错误(502 burst)重试。
返回 dict:{content, prompt_tokens, completion_tokens, prompt_cache_hit_tokens, total_tokens, model, wall_s, id, finish_reason}。
usage 字段用 openai 2.41.1 口径:prompt_tokens / completion_tokens(供成本折算与 token 取证);
prompt_cache_hit_tokens(U4/B9 新增,additive 键)= 命中缓存的 prompt token(两套字段取 max;缺则 0),供 cost.py 折缓存价。
"""
client = get_client()
last_err = None
for attempt in range(tries):
try:
t0 = time.perf_counter()
resp = client.chat.completions.create(
model=model,
messages=[
{"role": "system", "content": system},
{"role": "user", "content": user},
],
temperature=temperature,
max_tokens=max_tokens,
stream=False,
)
wall = time.perf_counter() - t0
choice = resp.choices[0]
usage = resp.usage
return {
"content": choice.message.content or "",
"finish_reason": getattr(choice, "finish_reason", None),
"prompt_tokens": getattr(usage, "prompt_tokens", None),
"completion_tokens": getattr(usage, "completion_tokens", None),
# U4/B9:命中缓存的 prompt token(additive 键;旧 dict 消费方不受影响,新口径供 cost 折缓存价)。
"prompt_cache_hit_tokens": _cache_hit_from_usage(usage),
"total_tokens": getattr(usage, "total_tokens", None),
"model": model,
"wall_s": round(wall, 3),
"id": getattr(resp, "id", None),
}
except Exception as e: # 瞬时网关错误重试(spike 实测 new-api 偶发 502 burst)
last_err = e
if attempt < tries - 1:
time.sleep(retry_delay)
raise RuntimeError(f"模型调用失败({tries} 次重试后):{last_err}")
if __name__ == "__main__":
# 自测:列模型(零 token 成本)确认 客户端+代理旁路+密钥 三件齐活。
c = get_client()
ids = sorted(m.id for m in c.models.list().data)
print(f"[_client] BASE_URL={BASE_URL} NO_PROXY 含网关={_GW_HOST in os.environ.get('NO_PROXY','')}")
print(f"[_client] 可用模型 {len(ids)} 个:")
for i in ids:
print(" ", i)