From bf7f1aebfcce856b17a8fe04fbda844848d1d78f Mon Sep 17 00:00:00 2001 From: lili Date: Thu, 2 Jul 2026 13:34:08 -0700 Subject: [PATCH] =?UTF-8?q?fix(cheap):=20=E9=98=B6=E6=AE=B5=E4=B8=80?= =?UTF-8?q?=E2=91=A1=20=E7=BB=88=E5=AE=A1=20merge-before(=E7=BA=A2?= =?UTF-8?q?=E7=BA=BF=E2=91=A2=20unlink=C3=972=20=E9=98=B2=E9=99=88?= =?UTF-8?q?=E6=97=A7=20verdict=20=E5=81=87=E7=BB=BF=20+=20setup=20fail-fas?= =?UTF-8?q?t=20+=20acquire=5Fports=20to=5Fthread=20+=20=E6=AD=BB=E7=A0=81/?= =?UTF-8?q?ReplyEndEvent=20canary/restricted=20fail-closed)=20(=E5=88=87?= =?UTF-8?q?=E7=89=87=E4=B8=80=20=E9=98=B6=E6=AE=B5=E4=B8=80=E2=91=A1/?= =?UTF-8?q?=E7=BB=88=E5=AE=A1fix)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Opus 4.8 (1M context) --- cheap-worker/cheap_gates.py | 3 + cheap-worker/cheap_service_app.py | 4 +- cheap-worker/cheap_service_driver.py | 21 +++- cheap-worker/tests/test_cheap_gates.py | 24 ++++ cheap-worker/tests/test_cheap_service_app.py | 8 ++ .../tests/test_cheap_service_driver.py | 117 ++++++++++++++++-- 6 files changed, 167 insertions(+), 10 deletions(-) diff --git a/cheap-worker/cheap_gates.py b/cheap-worker/cheap_gates.py index 220b0faa..daddf306 100644 --- a/cheap-worker/cheap_gates.py +++ b/cheap-worker/cheap_gates.py @@ -23,5 +23,8 @@ def run_cheap_gates(game_id: str, port: int, cdp_port: int) -> dict: # ensure_play_spec 据 smoke 抓的 state 产 driver(已存在不覆盖),让九门 driven=true、E_live/H_progress 由 advisory 升致命, # 根治「裸 harness 无 driver → 假绿」。smoke 失败(state=None)时按保守 key-cycle 薄 spec,play 仍照跑。 cheap_run.ensure_play_spec(game_id, sm.get("state")) + # 红线③(封陈旧 verdict 假绿,镜像 tier2 run.py:866):起 play 前清残留 verdict,保证 judge 拿到的恒是本次真门产出。 + # serve-and-play.sh 早退路径(chrome/serve/cdp 未就绪)+ 超时都不写新 verdict 也不删旧的,不清则 play 读回上一局绿 verdict → 假绿放行。 + (cheap_run.wg1_game_dir(game_id) / "evidence" / "verdict.json").unlink(missing_ok=True) pr = cheap_run.play(game_id, port=port, cdp_port=cdp_port) return pr.get("verdict") or {} diff --git a/cheap-worker/cheap_service_app.py b/cheap-worker/cheap_service_app.py index fc190a0f..fe299fdc 100644 --- a/cheap-worker/cheap_service_app.py +++ b/cheap-worker/cheap_service_app.py @@ -270,7 +270,9 @@ async def _cheap_middlewares_factory(user_id: str, agent_id: str, session_id: st async def _cheap_check(agent): # 独立跑九门:端口按 session 从池派生(决策③避撞),subprocess 经 to_thread 防阻塞事件循环; # 异常兜底(T2-b 工厂闭包 try,镜像 tier2 app.py:210-216):任何异常判未过续修,绝不穿透 on_reasoning。 - pair = acquire_ports() + # #3:acquire_ports 池耗尽会阻塞(Queue.get),不能在事件循环里同步等——经 to_thread 卸到线程池, + # 避免冻死整 Service 事件循环(端口满时其他会话的 SSE/续修全被拖死)。下方 run_cheap_gates 已是 to_thread。 + pair = await asyncio.to_thread(acquire_ports) try: try: verdict = await asyncio.to_thread(cheap_gates.run_cheap_gates, game_id, pair[0], pair[1]) diff --git a/cheap-worker/cheap_service_driver.py b/cheap-worker/cheap_service_driver.py index 2e774042..475d51a0 100644 --- a/cheap-worker/cheap_service_driver.py +++ b/cheap-worker/cheap_service_driver.py @@ -16,7 +16,6 @@ from __future__ import annotations import os import time -from pathlib import Path def _resolve_base() -> str: @@ -59,7 +58,9 @@ def _write_session_cfg(session_id: str, *, external_game_id: str, import cheap_run # noqa: PLC0415 try: - restricted = bool(write_whitelist) + # I1 fail-closed:restricted 由「是否给定白名单(is not None)」判,不用 bool()——空集白名单 bool 为 False 会误成 + # create 的不收窄放开全写;is not None 让空集 → restricted True → T2 读侧收窄成空集禁写(受限会话缺白名单宁禁勿放)。 + restricted = write_whitelist is not None wl = sorted(write_whitelist) if write_whitelist else None p = cheap_run.session_cfg_path(session_id) p.parent.mkdir(parents=True, exist_ok=True) @@ -231,6 +232,11 @@ async def drive_cheap_generation(job: dict, *, base_url: str | None = None, user cred = (await http.post(f"{base_url}/credential/", json=_cheap_credential_payload(), 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 往下走—— + # 否则后续 POST/SSE 拿 None 当 id,Service 侧找不到资源、driver 在 _wait_for_turn_end 空等到 SSE 总超时(~600s)才弃单。 + if not credential_id: + return (_failed_summary(game_id, f"Service /credential 未返 id(setup 失败):{str(cred)[:200]}"), + cheap_run.game_dir(game_id)) # ② 建 agent(system_prompt=cheap build_system_prompt;context_config 历史压缩;react_config 放开轮数)。 ctx_cfg = config.build_context_config() agent_body = { @@ -241,6 +247,9 @@ async def drive_cheap_generation(job: dict, *, base_url: str | None = None, user } agent = (await http.post(f"{base_url}/agent/", json=agent_body, headers=headers, timeout=30.0)).json() agent_id = agent.get("agent_id") or agent.get("id") + if not agent_id: # #2 setup fail-fast:建 agent 失败,不带 None 往下走 + return (_failed_summary(game_id, f"Service /agent 未返 id(setup 失败):{str(agent)[:200]}"), + cheap_run.game_dir(game_id)) # ③ 建 session(chat_model_config = openai_credential + MiniMax-M3 + max_tokens,无 thinking 分离)。 session_body = { "agent_id": agent_id, @@ -253,12 +262,20 @@ async def drive_cheap_generation(job: dict, *, base_url: str | None = None, user } session = (await http.post(f"{base_url}/sessions/", json=session_body, headers=headers, timeout=30.0)).json() session_id = session.get("session_id") or session.get("id") + if not session_id: # #2 setup fail-fast:建 session 失败,不带 None 往下走(否则 PATCH/SSE 拿 None 空等超时) + return (_failed_summary(game_id, f"Service /sessions 未返 id(setup 失败):{str(session)[:200]}"), + cheap_run.game_dir(game_id)) # ④ scaffold(clone 模板 + SAA 信封)——必须在 /chat 前,**game_id=后端 gameId**(与 system prompt 的 ⟦G⟧ 一致), # agent 起点就位在 amgen-<后端 gameId>(subprocess 经 to_thread)。C2:产物目录统一后端 gameId、不用 session_id。 sc = await asyncio.to_thread(cheap_run.scaffold, game_id, scaffold_template) if not sc["ok"]: return _failed_summary(game_id, f"scaffold 失败:{sc['output'][:300]}"), cheap_run.game_dir(game_id) + # 红线③:清残留 verdict + service-run-summary,保证回合后 _read_last_cheap_verdict/收口读到的恒是本次 Service 真写。 + # 封零门跑路径(setup 静默失败 / SSE 超时未跑一次 check 时,读回上一局绿 verdict → 假 succeeded)。 + for _stale in (cheap_run.wg1_game_dir(game_id) / "evidence" / "verdict.json", + cheap_run.game_dir(game_id) / "evidence" / "service-run-summary.json"): + _stale.unlink(missing_ok=True) # ⑤ 写会话注册表 sidecar(C2:Service 两工厂据 session_id 读它解析回后端 gameId、绑六工具/评门/collector; # create 路 write_whitelist=None → restricted=False)。external_game_id 恒写、正常路工厂必读到。 _write_session_cfg(session_id, external_game_id=game_id, diff --git a/cheap-worker/tests/test_cheap_gates.py b/cheap-worker/tests/test_cheap_gates.py index f34b1fa3..3662cb40 100644 --- a/cheap-worker/tests/test_cheap_gates.py +++ b/cheap-worker/tests/test_cheap_gates.py @@ -31,3 +31,27 @@ def test_run_cheap_gates_stage_fail_returns_empty(monkeypatch): out = cheap_gates.run_cheap_gates("g1", 4407, 9407) assert out == {}, "stage 失败应短路返 {}(judge 判未过)" assert called["smoke"] == 0, "stage 失败不应再跑 smoke" + + +def test_run_cheap_gates_unlinks_stale_verdict_before_play(tmp_path, monkeypatch): + # 红线③(封陈旧 verdict 假绿,镜像 tier2 run.py:866):play 前必清残留 verdict —— 预铺上一局绿 verdict, + # 断言 play 被调用当刻它已不在盘上(证明 unlink 在 play 之前跑)。不清则 serve-and-play.sh 早退/超时不写新 + # verdict 时,play 会读回上一局绿 verdict → judge 只看 verdict 不看 rc → 未过门新代码被假绿放行。 + import json + monkeypatch.setattr(cheap_run, "wg1_game_dir", lambda gid: tmp_path / "_wg1-gen" / gid) + ev = tmp_path / "_wg1-gen" / "g1" / "evidence" + ev.mkdir(parents=True) + (ev / "verdict.json").write_text(json.dumps({"pass": True, "guards": {}}), encoding="utf-8") # 陈旧绿残留 + monkeypatch.setattr(cheap_run, "stage", lambda gid: {"ok": True, "output": ""}) + monkeypatch.setattr(cheap_run, "smoke", lambda gid, port, cdp_port: {"ok": True, "state": {"targets": []}}) + monkeypatch.setattr(cheap_run, "ensure_play_spec", lambda gid, state: {"ok": True}) + seen = {} + + def _stub_play(gid, port, cdp_port): + # 记录 play 被调用当刻,陈旧 verdict 是否还在盘上(封红线③:必须已被 run_cheap_gates 在 play 前清掉)。 + seen["verdict_exists_at_play"] = (ev / "verdict.json").exists() + return {"ok": True, "verdict": {"pass": False, "guards": {}}} + + monkeypatch.setattr(cheap_run, "play", _stub_play) + cheap_gates.run_cheap_gates("g1", 4407, 9407) + assert seen["verdict_exists_at_play"] is False, "play 前必须已清残留 verdict(否则读回上一局绿=假绿放行)" diff --git a/cheap-worker/tests/test_cheap_service_app.py b/cheap-worker/tests/test_cheap_service_app.py index 5986a537..433c30bd 100644 --- a/cheap-worker/tests/test_cheap_service_app.py +++ b/cheap-worker/tests/test_cheap_service_app.py @@ -137,3 +137,11 @@ def test_collector_flushes_on_reply_end(tmp_path, monkeypatch): assert asyncio.run(_drive()) is True, "REPLY_END 流经时 sidecar 应已落盘(同步 flush)" summary = json.loads((tmp_path / "amgen-70021" / "evidence" / "service-run-summary.json").read_text(encoding="utf-8")) assert summary["costRmb"] == 0.42 and summary["repairs"] == 2 + + +def test_reply_end_event_class_name_canary(): + # collector 用 type(evt).__name__ == "ReplyEndEvent" 字符串判(cheap_service_app.py on_reply)锁 REPLY_END 落盘时机; + # 此 canary 钉住 AgentScope 事件类真名:2.0.x 若改名 ReplyEndEvent,字符串判会静默退化成 finally-only flush + # (成功局 costRmb=0),本断言让改名即红、把静默退化变永久绊线。生产码保留字符串判(不引硬 import 依赖、保 6c6g 惰性)。 + from agentscope.event import ReplyEndEvent # 已实测可 import + assert ReplyEndEvent.__name__ == "ReplyEndEvent" diff --git a/cheap-worker/tests/test_cheap_service_driver.py b/cheap-worker/tests/test_cheap_service_driver.py index 1d1b2b94..6bd06a1c 100644 --- a/cheap-worker/tests/test_cheap_service_driver.py +++ b/cheap-worker/tests/test_cheap_service_driver.py @@ -72,6 +72,12 @@ def test_session_cfg_roundtrip_maps_session_to_backend_gameid(tmp_path, monkeypa cfg2 = A._read_session_cfg("sess-def") assert cfg2["restricted"] is True assert A._resolve_write_whitelist(cfg2) == {"game-logic.js"} + # 空集白名单(Minor#4 fail-closed):restricted 必须由「白名单是否给定(is not None)」决定、非 bool(空集会误成 False)。 + # 空集 → restricted True → T2 读侧 I1 收窄成空集禁写,不再退化成 restricted False 放开全写。 + D._write_session_cfg("sess-empty", external_game_id="70014", write_whitelist=set()) + cfg3 = A._read_session_cfg("sess-empty") + assert cfg3["restricted"] is True, "空集白名单 → restricted True(fail-closed,不因 bool(空集)退成 False)" + assert A._resolve_write_whitelist(cfg3) == set(), "T2 读侧空集禁写(不放开全写)" def test_build_summary_trace_contract_end_to_end(tmp_path, monkeypatch): @@ -149,16 +155,24 @@ def test_drive_cheap_generation_fake_sse(tmp_path, monkeypatch): import httpx monkeypatch.setattr(httpx, "AsyncClient", _FakeHttp) - # 造后端 gameId=70012 的 on-disk 产物(scaffold 后 Service 侧应写在这里,fake 直接铺好)。 + # 造后端 gameId=70012 的 on-disk 产物:bundle 由 scaffold 阶段就位(engineBundle 承重);verdict + 收口 sidecar + # 改为在 fake turn 内写(模拟 Service 回合中跑门写盘)——driver 现在回合前会清残留(红线③修 b),预铺在 drive + # 之前会被清掉,故必须让「本次回合」写入,才对齐真实时序(先清陈旧、回合中写本次、回合后读本次)。 (tmp_path / "amgen-70012").mkdir(parents=True) (tmp_path / "amgen-70012" / "bundle.iife.js").write_text("window.__GameBundle={};", encoding="utf-8") - vev = tmp_path / "_wg1-gen" / "70012" / "evidence"; vev.mkdir(parents=True) - (vev / "verdict.json").write_text(json.dumps({"pass": True, "guards": {"A_boot": {"pass": True}}}), encoding="utf-8") - sev = tmp_path / "amgen-70012" / "evidence"; sev.mkdir(parents=True) - (sev / "service-run-summary.json").write_text(json.dumps({"costRmb": 0.5, "rmbGate": "active", "repairs": 1}), encoding="utf-8") - # ① 回合真结束(REPLY_END)→ 读 verdict/sidecar 组 succeeded 形。 - async def _ended(*a, **k): return {"ended": True, "reason": "REPLY_END", "endEvent": {}} + # 有界轮询提速: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)) + + # ① 回合真结束(REPLY_END)→ fake turn 写本次 verdict/sidecar → driver 回合后读它组 succeeded 形。 + async def _ended(*a, **k): + vev = tmp_path / "_wg1-gen" / "70012" / "evidence"; vev.mkdir(parents=True, exist_ok=True) + (vev / "verdict.json").write_text(json.dumps({"pass": True, "guards": {"A_boot": {"pass": True}}}), encoding="utf-8") + sev = tmp_path / "amgen-70012" / "evidence"; sev.mkdir(parents=True, exist_ok=True) + (sev / "service-run-summary.json").write_text( + json.dumps({"costRmb": 0.5, "rmbGate": "active", "repairs": 1}), encoding="utf-8") + return {"ended": True, "reason": "REPLY_END", "endEvent": {}} monkeypatch.setattr(CP, "_wait_for_turn_end", _ended) summary, gdir = asyncio.run(D.drive_cheap_generation({"gameId": 70012, "traceId": "t9", "brief": "点球"})) assert scaffolded["gid"] == "70012", "C2:scaffold 用后端 gameId(str),非 session_id" @@ -174,6 +188,95 @@ def test_drive_cheap_generation_fake_sse(tmp_path, monkeypatch): assert s2["stoppedReason"] == "total_timeout" +def _install_fake_http(monkeypatch, cred=None): + """装 fake httpx.AsyncClient:setup 三 POST 返可控体(cred 缺省给全 id),patch 空体。供回合前 unlink / fail-fast 用例复用。""" + import httpx + + class _Resp: + def __init__(self, d): self._d = d + def json(self): return self._d + + class _FakeHttp: + def __init__(self, *a, **k): pass + async def __aenter__(self): return self + async def __aexit__(self, *a): return False + async def post(self, url, **k): + if url.endswith("/credential/"): + return _Resp(cred if cred is not None else {"credential_id": "c1"}) + return _Resp({"agent_id": "a1"} if url.endswith("/agent/") else + {"session_id": "sess-9"} if url.endswith("/sessions/") else {}) + async def patch(self, url, **k): return _Resp({}) + + monkeypatch.setattr(httpx, "AsyncClient", _FakeHttp) + + +def _install_common_stubs(tmp_path, monkeypatch): + """装 drive 公共桩:目录重定向到 tmp、scaffold ok、richness/key 旁路、凭据体桩(零网络/LLM)。""" + import _bootstrap + monkeypatch.setattr(cheap_run, "game_dir", lambda gid: tmp_path / f"amgen-{gid}") + monkeypatch.setattr(cheap_run, "wg1_game_dir", lambda gid: tmp_path / "_wg1-gen" / gid) + monkeypatch.setattr(cheap_run, "session_cfg_path", lambda sid: tmp_path / "_cheap-sessions" / f"{sid}.json") + monkeypatch.setattr(cheap_run, "scaffold", lambda gid, tpl=None: {"ok": True, "output": ""}) + + async def _fake_richness(gid, *, brief=""): return {"score": None, "degraded": True, "reason": "stub"} + 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"}}) + # 有界轮询提速: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)) + + +def test_drive_unlinks_stale_before_turn(tmp_path, monkeypatch): + # 红线③(修 b,driver 侧):回合前清残留 verdict + service-run-summary。预铺上一局绿 verdict + 旧成本, + # fake SSE not-ended(本回合 Service 零门跑、啥也没写)→ 断言 drive 后 ok=False、costRmb=0(没被上一局绿 verdict + # 顶成假 succeeded、没读回上一局成本)。这封「零门跑路径:setup 静默失败 / SSE 超时未跑一次 check」。 + import asyncio + import json + import service.control_plane as CP + + _install_common_stubs(tmp_path, monkeypatch) + _install_fake_http(monkeypatch) + + # 预铺上一局残留(模拟同 gameId 重跑局的陈旧盘态):绿 verdict + 旧 summary。 + (tmp_path / "amgen-70012").mkdir(parents=True) + (tmp_path / "amgen-70012" / "bundle.iife.js").write_text("window.__GameBundle={};", encoding="utf-8") + vev = tmp_path / "_wg1-gen" / "70012" / "evidence"; vev.mkdir(parents=True) + (vev / "verdict.json").write_text(json.dumps({"pass": True, "guards": {"A_boot": {"pass": True}}}), encoding="utf-8") + sev = tmp_path / "amgen-70012" / "evidence"; sev.mkdir(parents=True) + (sev / "service-run-summary.json").write_text( + json.dumps({"costRmb": 0.9, "rmbGate": "active", "repairs": 3}), encoding="utf-8") + + async def _not_ended(*a, **k): return {"ended": False, "reason": "total_timeout", "endEvent": None} + monkeypatch.setattr(CP, "_wait_for_turn_end", _not_ended) + summary, _ = asyncio.run(D.drive_cheap_generation({"gameId": 70012, "traceId": "tX", "brief": "点球"})) + assert summary["ok"] is False, "陈旧绿 verdict 必须已被回合前清掉,不得顶成假 succeeded" + assert summary["costRmb"] == 0.0, "陈旧 summary 也须清掉,成本不得读回上一局" + + +def test_drive_setup_failfast_on_missing_id(tmp_path, monkeypatch): + # #2:Service /credential 未返 id(setup 失败)→ driver fail-fast(ok=False + reason 反映 setup 失败), + # 且不进入 _wait_for_turn_end(计数坐实)——不再 id=None 硬挂 ~600s SSE 空等。 + import asyncio + import service.control_plane as CP + + _install_common_stubs(tmp_path, monkeypatch) + _install_fake_http(monkeypatch, cred={}) # /credential 返 {} 无 id → fail-fast + + called = {"wait": 0} + + async def _counting_wait(*a, **k): + called["wait"] += 1 + return {"ended": True, "reason": "REPLY_END"} + monkeypatch.setattr(CP, "_wait_for_turn_end", _counting_wait) + + summary, _ = asyncio.run(D.drive_cheap_generation({"gameId": 70012, "traceId": "tf", "brief": "b"})) + assert summary["ok"] is False + assert "setup" in (summary.get("stoppedReason") or ""), "reason 应反映 setup 失败" + assert called["wait"] == 0, "setup 失败必须 fail-fast,不进入 _wait_for_turn_end 空等" + + def test_worker_default_run_fn_is_service_driver(): # worker_service 默认 run_fn 已切到驱动 Service(不再进程内 run_studio)。 state = W.WorkerState()