fix(cheap): 阶段一② 终审 merge-before(红线③ unlink×2 防陈旧 verdict 假绿 + setup fail-fast + acquire_ports to_thread + 死码/ReplyEndEvent canary/restricted fail-closed) (切片一 阶段一②/终审fix)

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
lili 2026-07-02 13:34:08 -07:00
parent 8d9bf1e577
commit bf7f1aebfc
6 changed files with 167 additions and 10 deletions

View File

@ -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 {}

View File

@ -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])

View File

@ -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,

View File

@ -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(否则读回上一局绿=假绿放行)"

View File

@ -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"

View File

@ -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()