From e8f7dc83477fcec5282282987505908f4f4f335e Mon Sep 17 00:00:00 2001 From: lili Date: Tue, 7 Jul 2026 12:44:08 -0700 Subject: [PATCH] =?UTF-8?q?feat(worker):=20WU2=20=C2=A73.7=20worker=20?= =?UTF-8?q?=E5=8F=96=20job.userToken=20+=20=C2=A73.9=20=E9=A2=9D=E5=BA=A6?= =?UTF-8?q?=E8=80=97=E5=B0=BD/=E5=87=AD=E6=8D=AE=E5=A4=B1=E6=95=88?= =?UTF-8?q?=E5=88=86=E8=BE=A8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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) --- cheap-worker/cheap_service_app.py | 19 ++++ cheap-worker/cheap_service_driver.py | 47 ++++++-- cheap-worker/result_out.py | 66 ++++++++++-- .../tests/test_cheap_service_driver.py | 20 +++- cheap-worker/tests/test_result_out.py | 71 ++++++++++++ cheap-worker/worker_service.py | 15 ++- wg1/gen-worker/worker/_client.py | 101 ++++++++++++++++-- wg1/gen-worker/worker/service.py | 8 ++ 8 files changed, 323 insertions(+), 24 deletions(-) diff --git a/cheap-worker/cheap_service_app.py b/cheap-worker/cheap_service_app.py index 888fb578..ddea391c 100644 --- a/cheap-worker/cheap_service_app.py +++ b/cheap-worker/cheap_service_app.py @@ -169,6 +169,9 @@ def _get_collector_cls(): self._tok_in = 0 self._tok_out = 0 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): # C1:框架 _agent.py:615 先 yield ReplyEndEvent 再 yield finish Msg 再收尾;collector 是最外层 on_reply, @@ -194,6 +197,18 @@ def _get_collector_cls(): # 先拿到并 publish 到 bus(异常在下一次推进时才抛),driver 立刻收到回合终结、按盘面组 failed # 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 可回收) try: from agentscope.event import ReplyEndEvent # noqa: PLC0415 @@ -230,6 +245,10 @@ def _get_collector_cls(): "budgetSoftTripped": bool(getattr(b, "budget_soft_tripped", False)), "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.mkdir(parents=True, exist_ok=True) (ev / "service-run-summary.json").write_text( diff --git a/cheap-worker/cheap_service_driver.py b/cheap-worker/cheap_service_driver.py index 85486553..7315d83b 100644 --- a/cheap-worker/cheap_service_driver.py +++ b/cheap-worker/cheap_service_driver.py @@ -25,23 +25,46 @@ def _resolve_base() -> str: return client.resolve_base_url() -def _resolve_key() -> str: - """NEWAPI_KEY(_bootstrap 已从凭据档注入 env;client.get_api_key 从 env 读)。""" +def _mask_token(tok) -> str: + """token 日志脱敏(§3.8):前后各留 4 位、中间省略;过短/空则整体隐藏。绝不整条打 new-api 凭据。""" + if not tok: + return "" + 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 已从凭据档注入 env、client.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 - 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(决策②)。 与 tier2(anthropic_credential,host 根)不同:cheap 走 OpenAI 兼容 /v1/chat/completions,故 type=openai_credential、base_url 带 /v1(与 config.build_model_openai 的 /v1 补全同口径)。 + api_key = 本次生成的 new-api 凭据:优先 job.userToken(per-user 额度),缺失回落全局 env(§3.7 F)。 """ base = _resolve_base().rstrip("/") if not base.endswith("/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, @@ -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 省略)。 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 @@ -229,6 +258,12 @@ async def drive_cheap_generation(job: dict, *, base_url: str | None = None, user brief = job.get("brief") or "" 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)。 # 命中五类(经营/剧情/TRPG/非遗/解谜)→ 选 per-genre 黄金骨架 + genre 透传丰富度评分(生产 create 路 # 与 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 都带 async with httpx.AsyncClient() as http: # ① 注册 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() credential_id = cred.get("credential_id") or cred.get("id") # #2 setup fail-fast:Service 未返 id(建凭据/agent/session 任一失败)必须立即诚实失败,不带 id=None 往下走—— diff --git a/cheap-worker/result_out.py b/cheap-worker/result_out.py index 5afa56eb..90c1cd56 100644 --- a/cheap-worker/result_out.py +++ b/cheap-worker/result_out.py @@ -17,14 +17,61 @@ import hashlib import json 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 被后端 -# 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 = { "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)。 _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: - """把 run-summary 的失败信号映射到契约#6 七值 failureReason 枚举。 + """把 run-summary 的失败信号映射到契约#6 failureReason 枚举。 - 显式入参(在七值内)优先;否则统一归 `llm_error`——契约#6 catch-all(SaaGraphDispatcher 同口径: - 图内部真因[九门不过 / 未收敛 / 预算熔断 / engineBundle 缺]不在七值内时归「生成执行链异常」最近桶, - 真因已落 summary/日志)。注:枚举无 budget/generation 专项值,故不强造、统一 llm_error。 + 优先级:显式入参(在枚举内)> summary.failureReason(在枚举内)> catch-all `llm_error`。 + summary.failureReason 是 WU2 §3.9 的桥:Service 侧在模型调用抛错时按 one-api 惯例分辨额度耗尽, + 经收口 sidecar → driver._build_summary 透传到 summary,这里据它落成 quota_exhausted(而非含糊 llm_error)。 + 其余图内部真因[九门不过 / 未收敛 / 预算熔断 / engineBundle 缺]不在枚举内 → 归「生成执行链异常」最近桶 + (SaaGraphDispatcher 同口径,真因已落 summary/日志);枚举无 budget/generation 专项值,不强造、统一 llm_error。 """ if explicit and explicit in _FAILURE_REASONS: return explicit + svc_reason = (summary or {}).get("failureReason") + if svc_reason and svc_reason in _FAILURE_REASONS: + return svc_reason return "llm_error" diff --git a/cheap-worker/tests/test_cheap_service_driver.py b/cheap-worker/tests/test_cheap_service_driver.py index aa6a0876..82271fbd 100644 --- a/cheap-worker/tests/test_cheap_service_driver.py +++ b/cheap-worker/tests/test_cheap_service_driver.py @@ -16,13 +16,27 @@ import worker_service as W # noqa: E402 def test_credential_payload_openai_compatible_with_v1(monkeypatch): # 便宜档 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_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() assert payload["data"]["type"] == "openai_credential" assert payload["data"]["base_url"].endswith("/v1") 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): # _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}") @@ -136,7 +150,7 @@ def test_drive_cheap_generation_fake_sse(tmp_path, monkeypatch): import _bootstrap monkeypatch.setattr(_bootstrap, "ensure_api_key_env", lambda: None) 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: 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(_bootstrap, "ensure_api_key_env", lambda: None) 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,只缩窗口)。 _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)) diff --git a/cheap-worker/tests/test_result_out.py b/cheap-worker/tests/test_result_out.py index 4510db36..cef6ae7c 100644 --- a/cheap-worker/tests/test_result_out.py +++ b/cheap-worker/tests/test_result_out.py @@ -105,6 +105,77 @@ def test_explicit_failure_reason_wins(): 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(): """traceId 缺 → 用 job_id(贯穿键同值)。""" gd = _make_game_dir(_tmp()) diff --git a/cheap-worker/worker_service.py b/cheap-worker/worker_service.py index 9a31d36e..f0ac2882 100644 --- a/cheap-worker/worker_service.py +++ b/cheap-worker/worker_service.py @@ -46,6 +46,14 @@ def log(msg: str) -> None: print(f"[cheap-worker-service] {msg}", flush=True) +def _mask_token(tok) -> str: + """token 日志脱敏(WU2 §3.8):前后各留 4 位、中间省略;过短/空则整体隐藏。绝不整条打 userToken。""" + if not tok: + return "" + s = str(tok) + return "****" if len(s) <= 8 else f"{s[:4]}…{s[-4:]}" + + # ---------- §6.1 job-in 解析 ---------- def parse_job(raw: bytes) -> dict | None: @@ -163,6 +171,9 @@ def _service_run_fn(job: dict): import cheap_service_driver 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)) @@ -541,8 +552,10 @@ def make_handler(state: WorkerState): pass trace_id = job.get("traceId") or job.get("job_id") 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')}, " - 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"), "traceId": trace_id, "queued": state.queue.qsize()}) diff --git a/wg1/gen-worker/worker/_client.py b/wg1/gen-worker/worker/_client.py index eb7fa54d..67ca585a 100644 --- a/wg1/gen-worker/worker/_client.py +++ b/wg1/gen-worker/worker/_client.py @@ -43,18 +43,96 @@ for _k in ("NO_PROXY", "no_proxy"): 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 "" + 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") if not 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 -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 get_client(user_token=None): + """构建指向 new-api 的裸 OpenAI 兼容客户端(WU2 §3.7:凭据优先 per-user token、缺失回落 env)。""" + 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): @@ -87,14 +165,18 @@ def _cache_hit_from_usage(usage): 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)重试。 返回 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 折缓存价。 + WU2 §3.7:user_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 for attempt in range(tries): try: @@ -130,6 +212,11 @@ def chat(model, system, user, max_tokens=16000, temperature=0.0, tries=3, retry_ "id": getattr(resp, "id", None), } 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 if attempt < tries - 1: time.sleep(retry_delay) diff --git a/wg1/gen-worker/worker/service.py b/wg1/gen-worker/worker/service.py index d4600b4e..18c367aa 100644 --- a/wg1/gen-worker/worker/service.py +++ b/wg1/gen-worker/worker/service.py @@ -54,6 +54,7 @@ if str(WORKER_DIR) not in sys.path: sys.path.insert(0, str(WORKER_DIR)) import run # noqa: E402 worker/run.py(GEN_DIR/REPO_ROOT 常量 + scaffold/build/play) +import _client # noqa: E402 裸 new-api 客户端(WU2 §3.9:NewapiQuotaExhaustedError 分辨额度耗尽) from agent_loop.studio import run_studio # noqa: E402 L2 agentic 闭环(design 展一句话→九门) 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}") else: 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: failure_reason = f"{type(e).__name__}: {e}" log(f"❌ 真生成异常 trace_id={trace_id}: {failure_reason}")