feat(worker): WU2 §3.7 worker 取 job.userToken + §3.9 额度耗尽/凭据失效分辨

per-user 额度按用户扣的最后一环:三 worker 入口从「读全局 env NEWAPI_KEY」改「优先
job.userToken、缺失回落全局 env key(WARN 脱敏)」;新失败因 quota_exhausted 干净失败。

§3.7 F(取 job token,默认关时字节不变):
- cheap_service_driver._resolve_key(user_token) 优先 per-user token、缺失回落 env+WARN;
  drive_cheap_generation 从 job.userToken 取、透传给凭据体;token 日志脱敏(前后各 4 位)。
- worker_service 整 job(含 userToken)透传给 driver;do_POST 记 userToken 存否(脱敏)。
- wg1 _client.get_api_key/get_client/chat 支持 user_token override、回落 env(WARN 一进程一次)。

§3.9(分辨两类失效,one-api 惯例初值,确切 message 阶段2 待验):
- result_out 新增 quota_exhausted 入枚举 + _map_failure_reason 认 summary.failureReason +
  classify_newapi_failure_reason/newapi_error_status_text(402/403→quota,401→llm_error 可重试)。
- cheap_service_app collector 崩溃路 best-effort 分辨 402/403→写 sidecar failureReason,
  driver→result_out 落 quota_exhausted;凭据失效/其他维持 llm_error,纯旁路不咬生成。
- wg1 _client.chat 额度耗尽抛 NewapiQuotaExhaustedError(不重试);service.handle_job 映射
  quota_exhausted。注:两 worker 主生成经 AgentScope 模型封装调网关,402 埋框架内,主路
  额度耗尽分辨待阶段2 拦模型封装;本次覆盖直连 _client.chat 路 + cheap 崩溃路分辨。

自证:test_result_out 25/25、test_worker_service 17/17(stdlib 直跑);driver/wg1 token 与
分辨逻辑 standalone 断言全过;八文件 py_compile 绿。test_cheap_service_driver 需 pytest
(本机未装),已同步桩签名待阶段2 e2e 真验(per-user token 调网关/used_quota 按用户增/耗尽落 quota_exhausted)。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
lili 2026-07-07 12:44:08 -07:00
parent a069afb649
commit e8f7dc8347
8 changed files with 323 additions and 24 deletions

View File

@ -169,6 +169,9 @@ def _get_collector_cls():
self._tok_in = 0 self._tok_in = 0
self._tok_out = 0 self._tok_out = 0
self._metric_emitted = False # 面二 metric 只发一次(_flush 三条降级路会调多次,防重复计数) self._metric_emitted = False # 面二 metric 只发一次(_flush 三条降级路会调多次,防重复计数)
# WU2 §3.9:模型调用抛错(new-api 非 2xx)时按 one-api 惯例分辨的失败因;仅 quota_exhausted 才写进
# 收口 sidecar 供 driver→result_out 落成 quota_exhausted(凭据失效/其他维持 llm_error,不透传)。默认 None。
self._failure_reason = None
async def on_reply(self, agent, input_kwargs, next_handler): async def on_reply(self, agent, input_kwargs, next_handler):
# C1:框架 _agent.py:615 先 yield ReplyEndEvent 再 yield finish Msg 再收尾;collector 是最外层 on_reply, # C1:框架 _agent.py:615 先 yield ReplyEndEvent 再 yield finish Msg 再收尾;collector 是最外层 on_reply,
@ -194,6 +197,18 @@ def _get_collector_cls():
# 先拿到并 publish 到 bus(异常在下一次推进时才抛),driver 立刻收到回合终结、按盘面组 failed # 先拿到并 publish 到 bus(异常在下一次推进时才抛),driver 立刻收到回合终结、按盘面组 failed
# summary 快速失败;异常原样重抛,不掩盖失败(框架日志照记 exception)。合成失败(极端:事件类 # summary 快速失败;异常原样重抛,不掩盖失败(框架日志照记 exception)。合成失败(极端:事件类
# 构造变化)只告警、退回慢失败路,绝不遮原异常。 # 构造变化)只告警、退回慢失败路,绝不遮原异常。
# WU2 §3.9:run 崩多因 new-api 非 2xx(额度耗尽 402/403 是其一)。best-effort 按 one-api 惯例分辨
# ——仅额度耗尽写 self._failure_reason,随 _flush 落进收口 sidecar,让 driver→result_out 落成
# quota_exhausted(而非含糊 llm_error);凭据失效/其他维持 llm_error。纯旁路,任何异常绝不遮原异常。
try:
import result_out as _ro # noqa: PLC0415 —— 复用单一分辨口径(one-api 惯例,阶段2 精确 message)
_reason = _ro.classify_newapi_failure_reason(*_ro.newapi_error_status_text(e))
if _reason == "quota_exhausted":
self._failure_reason = _reason
print(f"[cheap-service] 模型调用非 2xx 按 §3.9 分辨为额度耗尽 game={self._game_id}"
f"(quota_exhausted):{type(e).__name__}: {str(e)[:160]}", flush=True)
except Exception: # noqa: BLE001 —— 分辨纯旁路,失败不影响崩溃收口与原异常上抛
pass
self._flush() # 崩溃路也先落盘(部分成本/repairs 可回收) self._flush() # 崩溃路也先落盘(部分成本/repairs 可回收)
try: try:
from agentscope.event import ReplyEndEvent # noqa: PLC0415 from agentscope.event import ReplyEndEvent # noqa: PLC0415
@ -230,6 +245,10 @@ def _get_collector_cls():
"budgetSoftTripped": bool(getattr(b, "budget_soft_tripped", False)), "budgetSoftTripped": bool(getattr(b, "budget_soft_tripped", False)),
"traceSummary": self._tracer.summary(), "traceSummary": self._tracer.summary(),
} }
# WU2 §3.9:仅当模型调用非 2xx 被分辨为额度耗尽时才带 failureReason(driver 据它落 quota_exhausted);
# 正常/其他失败不带此键 → driver 不透传 → result_out 走既有 llm_error 兜底,不误标。
if self._failure_reason:
summary["failureReason"] = self._failure_reason
ev = cheap_run.game_dir(self._game_id) / "evidence" ev = cheap_run.game_dir(self._game_id) / "evidence"
ev.mkdir(parents=True, exist_ok=True) ev.mkdir(parents=True, exist_ok=True)
(ev / "service-run-summary.json").write_text( (ev / "service-run-summary.json").write_text(

View File

@ -25,23 +25,46 @@ def _resolve_base() -> str:
return client.resolve_base_url() return client.resolve_base_url()
def _resolve_key() -> str: def _mask_token(tok) -> str:
"""NEWAPI_KEY(_bootstrap 已从凭据档注入 env;client.get_api_key 从 env 读)。""" """token 日志脱敏(§3.8):前后各留 4 位、中间省略;过短/空则整体隐藏。绝不整条打 new-api 凭据。"""
if not tok:
return "<none>"
s = str(tok)
return "****" if len(s) <= 8 else f"{s[:4]}{s[-4:]}"
def _resolve_key(user_token: str | None = None) -> str:
"""解析本次生成用的 new-api 调用凭据(WU2 §3.7 F)。
优先用 job 随派带来的 per-user token(user_token);缺失则回落全局 env NEWAPI_KEY
(_bootstrap 已从凭据档注入 envclient.get_api_key env )回落是系统/编排/bake-off 触发的正常预期
(它们无 game_player 成员身份,后端 dispatch 层已按身份旁路不塞 token);额度门在后端 dispatch 层按身份堵死
(member create 路无 token fail-fast),worker 层只做有则用无则回落不硬拒,以免误伤既有验证路
回落时打一条 WARN:便于发现本该带 token member create job 异常走了共享计量(§3.7)
"""
if user_token:
return user_token
from worker import client # noqa: PLC0415 from worker import client # noqa: PLC0415
return client.get_api_key() key = client.get_api_key()
# 可追溯日志(脱敏):回落全局 key。系统/编排/bake-off 属正常;若为真实 member create job 走到这里则异常。
print(f"[cheap-driver] job 无 userToken → 回落全局 env NEWAPI_KEY({_mask_token(key)})。"
f"系统/编排/bake-off 触发为正常预期;若为真实 member create job 则异常走共享计量(§3.7)。",
flush=True)
return key
def _cheap_credential_payload() -> dict: def _cheap_credential_payload(user_token: str | None = None) -> dict:
"""组注册 cheap M3 凭据(POST /credential/)的请求体 —— OpenAI 兼容路,base 末尾补 /v1(决策②)。 """组注册 cheap M3 凭据(POST /credential/)的请求体 —— OpenAI 兼容路,base 末尾补 /v1(决策②)。
tier2(anthropic_credential,host )不同:cheap OpenAI 兼容 /v1/chat/completions, tier2(anthropic_credential,host )不同:cheap OpenAI 兼容 /v1/chat/completions,
type=openai_credentialbase_url /v1( config.build_model_openai /v1 补全同口径) type=openai_credentialbase_url /v1( config.build_model_openai /v1 补全同口径)
api_key = 本次生成的 new-api 凭据:优先 job.userToken(per-user 额度),缺失回落全局 env(§3.7 F)
""" """
base = _resolve_base().rstrip("/") base = _resolve_base().rstrip("/")
if not base.endswith("/v1"): if not base.endswith("/v1"):
base = base + "/v1" base = base + "/v1"
return {"data": {"type": "openai_credential", "api_key": _resolve_key(), "base_url": base}} return {"data": {"type": "openai_credential", "api_key": _resolve_key(user_token), "base_url": base}}
def _write_session_cfg(session_id: str, *, external_game_id: str, def _write_session_cfg(session_id: str, *, external_game_id: str,
@ -182,6 +205,12 @@ def _build_summary(game_id: str, brief: str, verdict, svc: dict, turn: dict, t0:
} }
# 9d-trace 源维度(verdictFull/driverType/models/stage);对照 None 时各键自然降级(result_out 省略)。 # 9d-trace 源维度(verdictFull/driverType/models/stage);对照 None 时各键自然降级(result_out 省略)。
summary.update(cheap_studio.build_trace_source(verdict, driver_type, attempts, stage, model)) summary.update(cheap_studio.build_trace_source(verdict, driver_type, attempts, stage, model))
# WU2 §3.9:Service 侧在模型调用抛错时(new-api 非 2xx)按 one-api 惯例分辨,把 failureReason 写进收口 sidecar;
# 这里原样透传进 summary,result_out._map_failure_reason 据它把额度耗尽落成 quota_exhausted(而非含糊 llm_error)。
# 仅当 Service 明确判额度耗尽(quota_exhausted)才透传,凭据失效/其他仍走 llm_error(可重试),避免误标。
svc_reason = svc.get("failureReason")
if svc_reason:
summary["failureReason"] = svc_reason
return summary return summary
@ -229,6 +258,12 @@ async def drive_cheap_generation(job: dict, *, base_url: str | None = None, user
brief = job.get("brief") or "" brief = job.get("brief") or ""
t0 = time.time() t0 = time.time()
# WU2 §3.7 F:从 §6.1 job 取 per-user token(后端 dispatchGeneric 对真实 member 装、系统/编排旁路留空);
# 空串归一为 None(缺失 → _resolve_key 回落全局 env key)。整条 job 由 worker_service 透传至此,故直接从 job 取。
user_token = (job.get("userToken") or "").strip() or None
print(f"[cheap-driver] game={game_id} 凭据来源={'per-user token' if user_token else '全局 env key(回落)'} "
f"userToken={_mask_token(user_token)}", flush=True)
# create 路参数:brief→genre 确定性关键词路由(T4 最小品类接线,取代硬编码 scaffold_template=None)。 # create 路参数:brief→genre 确定性关键词路由(T4 最小品类接线,取代硬编码 scaffold_template=None)。
# 命中五类(经营/剧情/TRPG/非遗/解谜)→ 选 per-genre 黄金骨架 + genre 透传丰富度评分(生产 create 路 # 命中五类(经营/剧情/TRPG/非遗/解谜)→ 选 per-genre 黄金骨架 + genre 透传丰富度评分(生产 create 路
# 与 S5 bake_off 复验口径同轨);无命中 → (None, None) 走通用 _template(与旧行为一致,路由只加不减)。 # 与 S5 bake_off 复验口径同轨);无命中 → (None, None) 走通用 _template(与旧行为一致,路由只加不减)。
@ -257,7 +292,7 @@ async def drive_cheap_generation(job: dict, *, base_url: str | None = None, user
headers.update(otel_carrier) # traceparent(/tracestate,若有);出站 :8300 各 POST 都带 headers.update(otel_carrier) # traceparent(/tracestate,若有);出站 :8300 各 POST 都带
async with httpx.AsyncClient() as http: async with httpx.AsyncClient() as http:
# ① 注册 OpenAI 兼容凭据(cheap M3 走 base+/v1)。 # ① 注册 OpenAI 兼容凭据(cheap M3 走 base+/v1)。
cred = (await http.post(f"{base_url}/credential/", json=_cheap_credential_payload(), cred = (await http.post(f"{base_url}/credential/", json=_cheap_credential_payload(user_token),
headers=headers, timeout=30.0)).json() headers=headers, timeout=30.0)).json()
credential_id = cred.get("credential_id") or cred.get("id") credential_id = cred.get("credential_id") or cred.get("id")
# #2 setup fail-fast:Service 未返 id(建凭据/agent/session 任一失败)必须立即诚实失败,不带 id=None 往下走—— # #2 setup fail-fast:Service 未返 id(建凭据/agent/session 任一失败)必须立即诚实失败,不带 id=None 往下走——

View File

@ -17,14 +17,61 @@ import hashlib
import json import json
from pathlib import Path from pathlib import Path
# §6.1 result-out / 契约#6 output.failureReason 七值枚举(对齐后端 FailureReasonEnum,真验改正: # §6.1 result-out / 契约#6 output.failureReason 枚举(对齐后端 FailureReasonEnum,真验改正:
# 原稿用了 content_violation/budget_exceeded/generation_failed 三个不在枚举内的值 → 真 e2e 被后端 # 原稿用了 content_violation/budget_exceeded/generation_failed 三个不在枚举内的值 → 真 e2e 被后端
# DifyCallbackTxService L165 拒「failureReason 不在七值内」、回调事务回滚。真七值见 FailureReasonEnum)。 # DifyCallbackTxService L165 拒「failureReason 不在枚举内」、回调事务回滚。真枚举见 FailureReasonEnum)。
# WU2 §3.9 起新增 quota_exhausted(new-api per-user 额度耗尽/池空,网关 402/403;跨契约共享枚举同批落
# aigc.yaml/dify-workflow-io.json/FailureReasonEnum 三处)。
_FAILURE_REASONS = { _FAILURE_REASONS = {
"unsafe_prompt", "intent_unclear", "no_template_match", "config_invalid", "unsafe_prompt", "intent_unclear", "no_template_match", "config_invalid",
"llm_error", "timeout", "asset_gen_failed", "llm_error", "timeout", "asset_gen_failed", "quota_exhausted",
} }
def newapi_error_status_text(exc) -> tuple:
"""从一个 new-api 调用异常里尽力抽出 (http_status, 错误体文本),供 classify_newapi_failure_reason 分辨(§3.9)。
兼容 openai SDK APIStatusError(.status_code / .response.status_code / .body / .message)与被上层框架
(AgentScope)包装后的普通异常(只剩 str)抽不到 status None,文本恒为 str(exc) 兜底关键词匹配
"""
status = getattr(exc, "status_code", None)
if status is None:
status = getattr(exc, "code", None)
if status is None:
resp = getattr(exc, "response", None)
status = getattr(resp, "status_code", None) if resp is not None else None
parts = []
for attr in ("message", "body"):
v = getattr(exc, attr, None)
if v:
parts.append(str(v))
parts.append(str(exc))
return status, " ".join(parts)
def classify_newapi_failure_reason(status=None, text: str = "") -> str:
"""据 one-api/new-api 惯例把网关非 2xx 分辨成失败因(WU2 §3.9)。
返回 'quota_exhausted'(额度耗尽,干净失败不可重试) 'llm_error'(凭据失效/其他,可重试不塌成额度耗尽)
判据:HTTP 402/403 额度耗尽;401 凭据失效(维持 llm_error);状态码取不到时按错误体关键词兜底
阶段2 待验402/403 与确切错误体 message S0 只坐实了管理面未测 /v1/chat/completions 耗尽错误体,
确切匹配待 S4 真机耗尽/凭据失效测试坐实;此处先给 one-api 惯例初值(状态码优先关键词兜底)
"""
try:
s = int(status) if status is not None else None
except (TypeError, ValueError):
s = None
if s in (402, 403):
return "quota_exhausted"
if s == 401:
return "llm_error" # 凭据失效:维持可重试,不塌成 quota_exhausted(§3.9 ②)
# 状态码取不到时(被框架包装成普通异常)按 one-api 耗尽错误体惯例关键词兜底 —— 确切 message 阶段2 待精确。
t = (text or "").lower()
quota_markers = ("insufficient user quota", "insufficient quota", "余额不足", "额度不足", "quota not enough")
if any(m in t for m in quota_markers):
return "quota_exhausted"
return "llm_error"
# 源工程 2.0 profile 必填三枚举(contracts/agent-loop/source-project.schema.json:39)。 # 源工程 2.0 profile 必填三枚举(contracts/agent-loop/source-project.schema.json:39)。
_REQUIRED_PROFILE_KEYS = ("tickModel", "inputModel", "progressModel") _REQUIRED_PROFILE_KEYS = ("tickModel", "inputModel", "progressModel")
@ -86,14 +133,19 @@ def _read_bundle(game_dir: Path) -> str | None:
def _map_failure_reason(summary: dict, explicit: str | None) -> str: def _map_failure_reason(summary: dict, explicit: str | None) -> str:
"""把 run-summary 的失败信号映射到契约#6 七值 failureReason 枚举。 """把 run-summary 的失败信号映射到契约#6 failureReason 枚举。
显式入参(在七值内)优先;否则统一归 `llm_error`契约#6 catch-all(SaaGraphDispatcher 同口径: 优先级:显式入参(在枚举内)> summary.failureReason(在枚举内)> catch-all `llm_error`
图内部真因[九门不过 / 未收敛 / 预算熔断 / engineBundle ]不在七值内时归生成执行链异常最近桶, summary.failureReason WU2 §3.9 的桥:Service 侧在模型调用抛错时按 one-api 惯例分辨额度耗尽,
真因已落 summary/日志):枚举无 budget/generation 专项值,故不强造统一 llm_error 经收口 sidecar driver._build_summary 透传到 summary,这里据它落成 quota_exhausted(而非含糊 llm_error)
其余图内部真因[九门不过 / 未收敛 / 预算熔断 / engineBundle ]不在枚举内 生成执行链异常最近桶
(SaaGraphDispatcher 同口径,真因已落 summary/日志);枚举无 budget/generation 专项值,不强造统一 llm_error
""" """
if explicit and explicit in _FAILURE_REASONS: if explicit and explicit in _FAILURE_REASONS:
return explicit return explicit
svc_reason = (summary or {}).get("failureReason")
if svc_reason and svc_reason in _FAILURE_REASONS:
return svc_reason
return "llm_error" return "llm_error"

View File

@ -16,13 +16,27 @@ import worker_service as W # noqa: E402
def test_credential_payload_openai_compatible_with_v1(monkeypatch): def test_credential_payload_openai_compatible_with_v1(monkeypatch):
# 便宜档 M3 走 OpenAI 兼容路:type=openai_credential、base_url 带 /v1、api_key 非空。 # 便宜档 M3 走 OpenAI 兼容路:type=openai_credential、base_url 带 /v1、api_key 非空。
monkeypatch.setattr(D, "_resolve_base", lambda: "http://100.64.0.8:3000") monkeypatch.setattr(D, "_resolve_base", lambda: "http://100.64.0.8:3000")
monkeypatch.setattr(D, "_resolve_key", lambda: "sk-test") # WU2 §3.7:_resolve_key 现签名 (user_token=None);缺 token 走 env 回落(此桩恒回同一 key)。
monkeypatch.setattr(D, "_resolve_key", lambda user_token=None: "sk-test")
payload = D._cheap_credential_payload() payload = D._cheap_credential_payload()
assert payload["data"]["type"] == "openai_credential" assert payload["data"]["type"] == "openai_credential"
assert payload["data"]["base_url"].endswith("/v1") assert payload["data"]["base_url"].endswith("/v1")
assert payload["data"]["api_key"] == "sk-test" assert payload["data"]["api_key"] == "sk-test"
def test_resolve_key_prefers_user_token_else_env_fallback(monkeypatch):
# WU2 §3.7 F:job 带 userToken 则用它(不碰 env);缺失回落全局 env key(client.get_api_key)。
import worker.client as _wc
monkeypatch.setattr(_wc, "get_api_key", lambda: "sk-env-fallback")
assert D._resolve_key("sk-user-per") == "sk-user-per" # per-user token 优先
assert D._resolve_key(None) == "sk-env-fallback" # 缺失回落 env
assert D._resolve_key("") == "sk-env-fallback" # 空串视同缺失
# 凭据体据 user_token 选 key(有则 per-user、无则回落 env)。
monkeypatch.setattr(D, "_resolve_base", lambda: "http://100.64.0.8:3000")
assert D._cheap_credential_payload("sk-user-per")["data"]["api_key"] == "sk-user-per"
assert D._cheap_credential_payload()["data"]["api_key"] == "sk-env-fallback"
def test_build_summary_shape_feeds_result_out(tmp_path, monkeypatch): def test_build_summary_shape_feeds_result_out(tmp_path, monkeypatch):
# _build_summary 组的 summary 能被 result_out.build_result_out 吃出 succeeded + trace(costRmb/repairs)。 # _build_summary 组的 summary 能被 result_out.build_result_out 吃出 succeeded + trace(costRmb/repairs)。
monkeypatch.setattr(cheap_run, "game_dir", lambda gid: tmp_path / f"amgen-{gid}") monkeypatch.setattr(cheap_run, "game_dir", lambda gid: tmp_path / f"amgen-{gid}")
@ -136,7 +150,7 @@ def test_drive_cheap_generation_fake_sse(tmp_path, monkeypatch):
import _bootstrap import _bootstrap
monkeypatch.setattr(_bootstrap, "ensure_api_key_env", lambda: None) monkeypatch.setattr(_bootstrap, "ensure_api_key_env", lambda: None)
monkeypatch.setattr(D, "_cheap_credential_payload", monkeypatch.setattr(D, "_cheap_credential_payload",
lambda: {"data": {"type": "openai_credential", "api_key": "sk", "base_url": "http://x/v1"}}) lambda user_token=None: {"data": {"type": "openai_credential", "api_key": "sk", "base_url": "http://x/v1"}})
class _Resp: class _Resp:
def __init__(self, d): self._d = d def __init__(self, d): self._d = d
@ -230,7 +244,7 @@ def _install_common_stubs(tmp_path, monkeypatch):
monkeypatch.setattr("cheap_verify.verify_richness", _fake_richness) monkeypatch.setattr("cheap_verify.verify_richness", _fake_richness)
monkeypatch.setattr(_bootstrap, "ensure_api_key_env", lambda: None) monkeypatch.setattr(_bootstrap, "ensure_api_key_env", lambda: None)
monkeypatch.setattr(D, "_cheap_credential_payload", monkeypatch.setattr(D, "_cheap_credential_payload",
lambda: {"data": {"type": "openai_credential", "api_key": "sk", "base_url": "http://x/v1"}}) lambda user_token=None: {"data": {"type": "openai_credential", "api_key": "sk", "base_url": "http://x/v1"}})
# 有界轮询提速:sidecar 缺失时默认等 5s → 50ms(保留真实 bounded-poll,只缩窗口)。 # 有界轮询提速:sidecar 缺失时默认等 5s → 50ms(保留真实 bounded-poll,只缩窗口)。
_orig_rss = D._read_service_run_summary _orig_rss = D._read_service_run_summary
monkeypatch.setattr(D, "_read_service_run_summary", lambda gid, wait_s=0.05: _orig_rss(gid, wait_s=wait_s)) monkeypatch.setattr(D, "_read_service_run_summary", lambda gid, wait_s=0.05: _orig_rss(gid, wait_s=wait_s))

View File

@ -105,6 +105,77 @@ def test_explicit_failure_reason_wins():
assert out["failureReason"] == "timeout" assert out["failureReason"] == "timeout"
# ── WU2 §3.9:额度耗尽干净失败 + one-api 惯例分辨 ──
def test_quota_exhausted_in_enum():
"""quota_exhausted 已进 failureReason 枚举(跨契约共享枚举同批新增,与 aigc.yaml/dify-workflow-io.json 一致)。"""
assert "quota_exhausted" in R._FAILURE_REASONS
def test_summary_failure_reason_quota_maps_through():
"""WU2 §3.9:summary.failureReason=quota_exhausted(Service 分辨→driver 透传)→ 落 failureReason=quota_exhausted。"""
gd = _make_game_dir(_tmp())
job = {"traceId": "tq", "templateId": "generic", "brief": "x"}
summary = {"finished": False, "verdict": {"pass": False, "failedGates": None},
"failureReason": "quota_exhausted"}
out = R.build_result_out(job, summary, gd)
assert out["status"] == "failed"
assert out["failureReason"] == "quota_exhausted"
def test_summary_failure_reason_ignored_when_not_in_enum():
"""summary.failureReason 非枚举值 → 不透传,回落 catch-all llm_error(防脏值污染契约)。"""
gd = _make_game_dir(_tmp())
job = {"traceId": "tq2", "templateId": "generic", "brief": "x"}
summary = {"finished": False, "verdict": {"pass": False}, "failureReason": "some_garbage"}
out = R.build_result_out(job, summary, gd)
assert out["failureReason"] == "llm_error"
def test_explicit_failure_reason_wins_over_summary():
"""显式入参优先级高于 summary.failureReason。"""
gd = _make_game_dir(_tmp())
job = {"traceId": "tq3", "templateId": "generic", "brief": "x"}
summary = {"verdict": {"pass": False}, "failureReason": "quota_exhausted"}
out = R.build_result_out(job, summary, gd, failure_reason="timeout")
assert out["failureReason"] == "timeout"
def test_classify_newapi_status_codes():
"""one-api 惯例:402/403 → quota_exhausted;401 → llm_error(凭据失效,可重试,不塌成额度耗尽)。"""
assert R.classify_newapi_failure_reason(402) == "quota_exhausted"
assert R.classify_newapi_failure_reason(403) == "quota_exhausted"
assert R.classify_newapi_failure_reason(401) == "llm_error"
assert R.classify_newapi_failure_reason(500) == "llm_error"
assert R.classify_newapi_failure_reason(None) == "llm_error"
def test_classify_newapi_keyword_fallback():
"""状态码取不到(被框架包装)时按 one-api 错误体关键词兜底分辨额度耗尽(阶段2 待精确 message)。"""
assert R.classify_newapi_failure_reason(None, "error: insufficient user quota") == "quota_exhausted"
assert R.classify_newapi_failure_reason(None, "当前分组余额不足,请充值") == "quota_exhausted"
assert R.classify_newapi_failure_reason(None, "connection reset by peer") == "llm_error"
def test_newapi_error_status_text_extraction():
"""从带 status_code/response 的异常对象抽 (status, text);抽不到 status 返 None、文本兜底 str(exc)。"""
class _Err(Exception):
status_code = 402
message = "insufficient user quota"
st, txt = R.newapi_error_status_text(_Err("boom"))
assert st == 402 and "insufficient" in txt
class _Resp:
status_code = 403
class _Wrapped(Exception):
response = _Resp()
st2, _ = R.newapi_error_status_text(_Wrapped("x"))
assert st2 == 403
st3, txt3 = R.newapi_error_status_text(ValueError("plain error"))
assert st3 is None and "plain error" in txt3
def test_traceid_fallback_to_job_id(): def test_traceid_fallback_to_job_id():
"""traceId 缺 → 用 job_id(贯穿键同值)。""" """traceId 缺 → 用 job_id(贯穿键同值)。"""
gd = _make_game_dir(_tmp()) gd = _make_game_dir(_tmp())

View File

@ -46,6 +46,14 @@ def log(msg: str) -> None:
print(f"[cheap-worker-service] {msg}", flush=True) print(f"[cheap-worker-service] {msg}", flush=True)
def _mask_token(tok) -> str:
"""token 日志脱敏(WU2 §3.8):前后各留 4 位、中间省略;过短/空则整体隐藏。绝不整条打 userToken。"""
if not tok:
return "<none>"
s = str(tok)
return "****" if len(s) <= 8 else f"{s[:4]}{s[-4:]}"
# ---------- §6.1 job-in 解析 ---------- # ---------- §6.1 job-in 解析 ----------
def parse_job(raw: bytes) -> dict | None: def parse_job(raw: bytes) -> dict | None:
@ -163,6 +171,9 @@ def _service_run_fn(job: dict):
import cheap_service_driver import cheap_service_driver
base_url = os.environ.get("CHEAP_SERVICE_URL", "http://127.0.0.1:8300") base_url = os.environ.get("CHEAP_SERVICE_URL", "http://127.0.0.1:8300")
# WU2 §3.7 F:整条 job(含 §6.1 追加的 userToken)原样透传给 driver;driver 从 job.get("userToken") 取
# per-user token 装 new-api 凭据、缺失回落全局 env key(见 cheap_service_driver._resolve_key)。
# 此处不拆字段单独传递,避免冗余管道——透传整 job 即满足「worker 取 job token」。
return asyncio.run(cheap_service_driver.drive_cheap_generation(job, base_url=base_url)) return asyncio.run(cheap_service_driver.drive_cheap_generation(job, base_url=base_url))
@ -541,8 +552,10 @@ def make_handler(state: WorkerState):
pass pass
trace_id = job.get("traceId") or job.get("job_id") trace_id = job.get("traceId") or job.get("job_id")
accepted, code = try_enqueue(state, job) accepted, code = try_enqueue(state, job)
# WU2 §3.7/§3.8:记 userToken 是否随 job 带来(脱敏)——便于排「member create 该带 token 却没带」。
log(f"收到 job trace_id={trace_id}, templateId={job.get('templateId')}, " log(f"收到 job trace_id={trace_id}, templateId={job.get('templateId')}, "
f"gameId={job.get('gameId')}{'受理' if accepted else '队满拒'}({code})") f"gameId={job.get('gameId')}, userToken={_mask_token(job.get('userToken'))} "
f"{'受理' if accepted else '队满拒'}({code})")
self._send_json(code, {"accepted": accepted, "job_id": job.get("job_id"), self._send_json(code, {"accepted": accepted, "job_id": job.get("job_id"),
"traceId": trace_id, "queued": state.queue.qsize()}) "traceId": trace_id, "queued": state.queue.qsize()})

View File

@ -43,18 +43,96 @@ for _k in ("NO_PROXY", "no_proxy"):
import openai # noqa: E402 —— 必须在代理旁路之后导入 import openai # noqa: E402 —— 必须在代理旁路之后导入
# WU2 §3.7:回落全局 env key 时打 WARN;但 get_api_key 每次 chat 都会被调,故一进程只 WARN 一次防刷屏。
_ENV_FALLBACK_WARNED = False
def get_api_key():
"""读 NEWAPI_KEY缺失即抛密钥不入库、复现前先设 .env 或 export""" class NewapiQuotaExhaustedError(RuntimeError):
"""new-api 网关按 per-user token 判额度耗尽(WU2 §3.9)。
one-api 惯例:HTTP 402/403 或错误体含额度耗尽关键词**不可重试**干净失败worker 据此回调
failure_reason=quota_exhausted(而非含糊 llm_error);凭据失效(401)不塌成本错误维持可重试 llm_error
阶段2 待验402/403 与确切错误体 message S0 只坐实管理面未测 /v1/chat/completions 耗尽错误体,
确切匹配待 S4 真机耗尽/凭据失效测试;现按 one-api 惯例初值分辨(状态码优先关键词兜底)
"""
def _mask_token(tok) -> str:
"""token 日志脱敏(WU2 §3.8):前后各留 4 位、中间省略;过短/空则整体隐藏。绝不整条打凭据。"""
if not tok:
return "<none>"
s = str(tok)
return "****" if len(s) <= 8 else f"{s[:4]}{s[-4:]}"
def _newapi_error_status_text(exc):
"""从 new-api 调用异常尽力抽 (http_status, 错误体文本);抽不到 status 返 None,文本恒兜底 str(exc)。
兼容 openai SDK APIStatusError(.status_code / .response.status_code / .body / .message)与被上层
框架包装后的普通异常(只剩 str) classify_newapi_failure_reason 分辨
"""
status = getattr(exc, "status_code", None)
if status is None:
status = getattr(exc, "code", None)
if status is None:
resp = getattr(exc, "response", None)
status = getattr(resp, "status_code", None) if resp is not None else None
parts = []
for attr in ("message", "body"):
v = getattr(exc, attr, None)
if v:
parts.append(str(v))
parts.append(str(exc))
return status, " ".join(parts)
def classify_newapi_failure_reason(status=None, text: str = "") -> str:
"""据 one-api 惯例把网关非 2xx 分辨成失败因(WU2 §3.9)。
返回 'quota_exhausted'(额度耗尽,不可重试) 'llm_error'(凭据失效 401/其他,可重试不塌成额度耗尽)
判据:HTTP 402/403 额度耗尽;401 凭据失效;状态码取不到时按错误体关键词兜底
阶段2 待验确切 message/code S4 真机耗尽/凭据失效测试坐实;此处状态码优先关键词兜底为初值
"""
try:
s = int(status) if status is not None else None
except (TypeError, ValueError):
s = None
if s in (402, 403):
return "quota_exhausted"
if s == 401:
return "llm_error" # 凭据失效:维持可重试,不塌成 quota_exhausted(§3.9 ②)
t = (text or "").lower()
quota_markers = ("insufficient user quota", "insufficient quota", "余额不足", "额度不足", "quota not enough")
if any(m in t for m in quota_markers):
return "quota_exhausted"
return "llm_error"
def get_api_key(user_token=None):
"""解析本次调用的 new-api 凭据(WU2 §3.7 F)。
优先用传入的 per-user token(user_token, §6.1 job 派来);缺失回落全局 env NEWAPI_KEY(缺失即抛,
密钥不入库复现前先设 .env export)回落是系统/编排/bake-off/spike 触发的正常预期( game_player
成员身份后端 dispatch 已按身份旁路不塞 token);额度门在后端 dispatch 层按身份堵,worker 层只做有则用
无则回落回落时打一条 WARN(一进程一次,防每次 chat 刷屏)便于发现本该带 token member create job 异常
"""
if user_token:
return user_token
key = os.environ.get("NEWAPI_KEY") key = os.environ.get("NEWAPI_KEY")
if not key: if not key:
raise RuntimeError("NEWAPI_KEY 未设:请在 wg1/gen-worker/.env 写入或 export NEWAPI_KEY=...") raise RuntimeError("NEWAPI_KEY 未设:请在 wg1/gen-worker/.env 写入或 export NEWAPI_KEY=...")
global _ENV_FALLBACK_WARNED
if not _ENV_FALLBACK_WARNED:
_ENV_FALLBACK_WARNED = True
print(f"[wg1-worker] 无 per-user userToken → 回落全局 env NEWAPI_KEY({_mask_token(key)})。"
f"系统/编排/bake-off/spike 为正常预期;若为真实 member create job 则异常走共享计量(§3.7)。",
flush=True)
return key return key
def get_client(): def get_client(user_token=None):
"""构建指向 new-api 的裸 OpenAI 兼容客户端。""" """构建指向 new-api 的裸 OpenAI 兼容客户端(WU2 §3.7:凭据优先 per-user token、缺失回落 env)"""
return openai.OpenAI(api_key=get_api_key(), base_url=BASE_URL, max_retries=2, timeout=180.0) return openai.OpenAI(api_key=get_api_key(user_token), base_url=BASE_URL, max_retries=2, timeout=180.0)
def _extra_body(model, enable_thinking): def _extra_body(model, enable_thinking):
@ -87,14 +165,18 @@ def _cache_hit_from_usage(usage):
return 0 return 0
def chat(model, system, user, max_tokens=16000, temperature=0.0, tries=3, retry_delay=2.0, enable_thinking=False): def chat(model, system, user, max_tokens=16000, temperature=0.0, tries=3, retry_delay=2.0,
enable_thinking=False, user_token=None):
"""单次裸调用 chat.completions。确定性产出 temperature=0瞬时网关错误(502 burst)重试。 """单次裸调用 chat.completions。确定性产出 temperature=0瞬时网关错误(502 burst)重试。
返回 dict{content, prompt_tokens, completion_tokens, prompt_cache_hit_tokens, total_tokens, model, wall_s, id, finish_reason} 返回 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 取证 usage 字段用 openai 2.41.1 口径prompt_tokens / completion_tokens供成本折算与 token 取证
prompt_cache_hit_tokensU4/B9 新增additive = 命中缓存的 prompt token两套字段取 max缺则 0 cost.py 折缓存价 prompt_cache_hit_tokensU4/B9 新增additive = 命中缓存的 prompt token两套字段取 max缺则 0 cost.py 折缓存价
WU2 §3.7user_token 给定则用它调网关(per-user 额度)缺失回落全局 env key
WU2 §3.9网关非 2xx one-api 惯例分辨额度耗尽(402/403)立即抛 NewapiQuotaExhaustedError 不再重试;
凭据失效(401)/瞬时错(502)维持原重试语义(不塌成额度耗尽)
""" """
client = get_client() client = get_client(user_token)
last_err = None last_err = None
for attempt in range(tries): for attempt in range(tries):
try: try:
@ -130,6 +212,11 @@ def chat(model, system, user, max_tokens=16000, temperature=0.0, tries=3, retry_
"id": getattr(resp, "id", None), "id": getattr(resp, "id", None),
} }
except Exception as e: # 瞬时网关错误重试spike 实测 new-api 偶发 502 burst except Exception as e: # 瞬时网关错误重试spike 实测 new-api 偶发 502 burst
# WU2 §3.9:额度耗尽(402/403)重试无意义,立即抛典型错误让 worker 回调 quota_exhausted;
# 凭据失效(401)/瞬时错(502)不塌成额度耗尽,维持下方重试语义(可重试 llm_error)。
if classify_newapi_failure_reason(*_newapi_error_status_text(e)) == "quota_exhausted":
raise NewapiQuotaExhaustedError(
f"new-api 额度耗尽(按 §3.9 分辨,阶段2 待精确 message){type(e).__name__}: {e}") from e
last_err = e last_err = e
if attempt < tries - 1: if attempt < tries - 1:
time.sleep(retry_delay) time.sleep(retry_delay)

View File

@ -54,6 +54,7 @@ if str(WORKER_DIR) not in sys.path:
sys.path.insert(0, str(WORKER_DIR)) sys.path.insert(0, str(WORKER_DIR))
import run # noqa: E402 worker/run.pyGEN_DIR/REPO_ROOT 常量 + scaffold/build/play import run # noqa: E402 worker/run.pyGEN_DIR/REPO_ROOT 常量 + scaffold/build/play
import _client # noqa: E402 裸 new-api 客户端WU2 §3.9NewapiQuotaExhaustedError 分辨额度耗尽)
from agent_loop.studio import run_studio # noqa: E402 L2 agentic 闭环design 展一句话→九门) from agent_loop.studio import run_studio # noqa: E402 L2 agentic 闭环design 展一句话→九门)
from agent_loop import config # noqa: E402 models.yaml 阶段(默认 default_stage from agent_loop import config # noqa: E402 models.yaml 阶段(默认 default_stage
@ -333,6 +334,13 @@ def handle_job(state, job):
log(f"❌ pass 但 bundle 文件缺失:{bundle_path}") log(f"❌ pass 但 bundle 文件缺失:{bundle_path}")
else: else:
failure_reason = "nine_gate_failed" failure_reason = "nine_gate_failed"
except _client.NewapiQuotaExhaustedError as e:
# WU2 §3.9:new-api 按 per-user token 判额度耗尽(402/403)→ 干净失败 quota_exhausted(而非含糊 llm_error),
# 后端据此把任务终态落 quota_exhausted、前端展示「生成额度已用完」。凭据失效(401)不走此路、维持可重试。
# 注:wg1 主生成经 AgentScope 模型封装调网关(config.py credential),402 埋在框架内、此典型错误仅从
# _client.chat 直连路(如 run.py 旧路)propagate;主生成路的额度耗尽分辨待阶段2 拦 AgentScope 模型封装。
failure_reason = "quota_exhausted"
log(f"❌ new-api 额度耗尽 trace_id={trace_id}(§3.9 quota_exhausted): {type(e).__name__}: {e}")
except Exception as e: except Exception as e:
failure_reason = f"{type(e).__name__}: {e}" failure_reason = f"{type(e).__name__}: {e}"
log(f"❌ 真生成异常 trace_id={trace_id}: {failure_reason}") log(f"❌ 真生成异常 trace_id={trace_id}: {failure_reason}")