diff --git a/tier2/gen-worker/service/app.py b/tier2/gen-worker/service/app.py index a80f3ac2..8045081a 100644 --- a/tier2/gen-worker/service/app.py +++ b/tier2/gen-worker/service/app.py @@ -169,29 +169,60 @@ async def _tier2_tools_factory(user_id: str, agent_id: str, session_id: str) -> async def _tier2_middlewares_factory(user_id: str, agent_id: str, session_id: str) -> list: - """extra_agent_middlewares 工厂:每个 chat 回合产出 tier2 的四熔断 + trace 中间件。 + """extra_agent_middlewares 工厂:每回合产 [trace, 续修, 熔断] 三件。 - 签名严格对齐 AgentMiddlewareFactory((user_id, agent_id, session_id) -> Awaitable[list[MiddlewareBase]])。 - 产出的中间件接在框架中间件之后(_chat.py:273-280);trace 列在 breaker 前 → 它是更外层洋葱 - (先 ingest 事件再进熔断巡检,与 studio.py 同序)。trace_id 用 session_id 贯穿本款生成。 - - 重依赖(worker.middleware → 经包顶层连带 agentscope)在此惰性 import。 - - Returns: - list[MiddlewareBase]:[Tier2TraceMiddleware(trace_id=session_id), CircuitBreakerMiddleware()]。 + 注入序(列表序=外→内):tracer / 续修 / 熔断;框架在其外另前置 InboxMiddleware(挂 on_reasoning、每轮 drain + inbox 后透传全部 evt、不吞 finish Msg,故不阻断续修拦截)+ StateChange/ToolOffload(只挂 on_acting)。 + 实际 on_reasoning 链 = [InboxMiddleware, tracer, repair];repair 仍是最内层、能拦到 _reasoning 产的原始 finish。 + 续修 check = finish 点独立重跑 tier2 run_gates(同步 subprocess、to_thread 防阻塞)+ 异常兜底(评审 I-2: + run_gates 非 subprocess 段也可能抛,不兜会穿透 on_reasoning → Service 无 REPLY_END、只能靠超时兜)。game_id=session_id。 + breaker 传 soft_budget=True(软停、不误伤 cheap hard 档)+ 续修模式放大的超时(容纳 N 次串行真门,评审 C1)。 + budget_exhausted 读实测已花 ¥ 超上限(非预估标记;配 T2 的 repairs>0 保护,软预算不旁路续修)。 """ - # 惰性 import:worker.middleware 经 worker 包顶层会牵出 agentscope;故只在回合内 import。 - from worker.middleware import ( # noqa: PLC0415 —— 惰性 import 红线 + import asyncio # noqa: PLC0415 + from worker import genconfig, run as _run # noqa: PLC0415 —— 惰性 import 红线 + from worker.gate_judge import GateJudgment, judge_tier2_verdict # noqa: PLC0415 + from worker.middleware import ( # noqa: PLC0415 CircuitBreakerMiddleware, + RepairMiddleware, Tier2TraceMiddleware, ) - # trace 在外、breaker 在内(与 studio.py middlewares=[tracer, breaker] 同序);sink=None → 步留内存, - # 真落库 sink 随控制面 phase-1 接(同 studio 现状)。每回合一组新实例(熔断计数按回合,非跨回合累计—— - # 这是与 CLI 单局累计的差异点,列入 followup:跨回合累计 ¥ 硬闸需把计数挂到 session 维度)。 tracer = Tier2TraceMiddleware(trace_id=session_id) - breaker = CircuitBreakerMiddleware() - return [tracer, breaker] + max_repairs = genconfig.get("iteration", "max_resumes", 6) + gate_timeout_s = 300 # run_gates 单次真门上限(run.py:812 的 timeout_s 默认) + # 续修模式重标定 breaker 超时(评审 C1):整局在一个 reply 内串行跑 (max_repairs+1) 次真门,墙钟/单步静默 + # 都要覆盖它,否则续修被自己的门判撑爆 timeout 熔断。directional,待 mini-desktop 量真实门时延后收紧(见风险段)。 + repair_wall_s = genconfig.get("budget", "repair_wall_timeout_s", + (max_repairs + 1) * (gate_timeout_s + 120) + 300) # 门 + 单轮推理 + 余量 + repair_step_s = genconfig.get("budget", "repair_step_timeout_s", gate_timeout_s + 180) # > 单次门判上限,治静默误触 + breaker = CircuitBreakerMiddleware( + soft_budget=True, # tier2 单 POST 续修路走软停(cheap 仍默认 hard) + wall_timeout_s=repair_wall_s, # 放大墙钟容纳 N 次串行真门 + step_timeout_s=repair_step_s, # 放大单步静默阈值 > 单次门判上限(门判期间无 evt 不误杀) + ) + + game_id = session_id # 评门 game_id = session_id(九工具据此管工程目录) + + async def _check(agent): + # 独立重跑门:run_gates 同步 subprocess(~300s、起 chrome),to_thread 防阻塞;异常兜底为「门未过」续修。 + # play_spec 服务态先 None(run_gates 内按品类默认 business-sim 驱动;定制驱动随后续接口下发,列 followup)。 + try: + gate_result = await asyncio.to_thread(_run.run_gates, game_id, None) + except Exception as e: # noqa: BLE001 —— 门执行任何异常都判未过续修,绝不穿透 on_reasoning(评审 I-2) + print(f"[repair] run_gates 异常(判未过、续修): {type(e).__name__}: {e}", flush=True) + return GateJudgment(passed=False, failed_gates=["run_gates_error"], + feedback=f"门执行异常,请检查工程可构建/可运行后再交付:{type(e).__name__}: {e}") + return judge_tier2_verdict(gate_result) + + repair = RepairMiddleware( + check=_check, + max_repairs=max_repairs, + # 实测已花 ¥ 超上限(非预估软停标记):配 T2 的 repairs>0 保护,首个未绿 finish 必先修一次再因预算放行。 + budget_exhausted=lambda: ( + breaker._rmb_gate_active and breaker.spent_rmb >= breaker.rmb_hard_limit), + ) + return [tracer, repair, breaker] def build_app(*, title: str = SERVICE_TITLE) -> "FastAPI": diff --git a/tier2/gen-worker/service/control_plane.py b/tier2/gen-worker/service/control_plane.py index 3bbaf908..d556dbca 100644 --- a/tier2/gen-worker/service/control_plane.py +++ b/tier2/gen-worker/service/control_plane.py @@ -295,6 +295,7 @@ async def drive_generation( writer_max_iters: int | None = None, user_id: str = "tier2", do_design: bool = True, + single_post: bool = True, sse_turn_timeout_s: float = 1200.0, sse_idle_timeout_s: float = 300.0, ) -> dict: @@ -312,6 +313,9 @@ async def drive_generation( model_name: M3 模型名(默认 MiniMax-M3,经 new-api 走 Anthropic 原生)。 writer_max_iters: 单写 ReAct 放开轮数;None → genconfig.get('iteration','writer_max_iters',40)。 user_id: 多租户用户标识(经 X-User-Id 头)。 + single_post: True(默认,阶段一①)= 只发一次 kick,续修在 Service 端 on_reasoning middleware 内完成, + 消费方等这一次回合真结束后读服务端已落 verdict 判落库,不再外层多轮 resume;False = 走现有外层 + 有界 resume 循环(fallback,续修 middleware spike 不成时回落)。 sse_turn_timeout_s: 单回合 SSE 等结束的总超时(兜底,绝不卡死)。 sse_idle_timeout_s: 单回合 SSE 空闲(无 data 帧)断流判定超时。 @@ -319,8 +323,10 @@ async def drive_generation( {game_id, agent_id, session_id, finished(bool=门绿+落库), attempts, last_verdict, store_addressing, wall_s, stopped_reason}。 """ - # 惰性 import:bootstrap(REST 原语,牵出 httpx/worker)/ run(门 + 落库,牵出 agentscope)/ genconfig。 + # 惰性 import:bootstrap(REST 原语,牵出 httpx/worker)/ run(门 + 落库,牵出 agentscope)/ genconfig / + # gate_judge(single_post 读服务端已落 verdict 后经它归一判门绿,与门线口径一致、不放松)。 from worker import genconfig, run as _run # noqa: PLC0415 + from worker.gate_judge import judge_tier2_verdict # noqa: PLC0415 from . import bootstrap # noqa: PLC0415 t0 = time.perf_counter() @@ -369,8 +375,53 @@ async def drive_generation( result["wall_s"] = round(time.perf_counter() - t0, 1) return result + # ── ②' 单 POST 路(阶段一①):续修在 Service 端 on_reasoning middleware 内完成,消费方只发一次 kick、 + # 等这一次(内部续修多轮)回合真结束、读服务端已落 verdict 判落库;不再外层多轮 resume。老循环留 fallback。── + if single_post: + # SSE 总超时须覆盖续修 N 次串行真门(评审 C-2):调用方未显式放大时按 budget.repair_wall_timeout_s 兜底, + # 与工厂 breaker 的 repair_wall_s 同源对齐;绝不用默认 1200s(否则续修中的合法局被误判 not_ended、丢产物)。 + sse_timeout = max(sse_turn_timeout_s, + genconfig.get("budget", "repair_wall_timeout_s", + (max_resumes + 1) * (300 + 120) + 300)) + turn = await _wait_for_turn_end( + base_url, agent_id, session_id, user_id=user_id, + timeout_s=sse_timeout, idle_timeout_s=sse_idle_timeout_s, + trace_path=trace_path, attempt=0) + result["attempts"] = 1 + if not turn.get("ended"): + # 回合未真结束(SSE 总超时/断流):服务端可能仍在续修写盘,评中间态会误判 + 与写盘竞态 → 不评、不落库。 + print(f"[tier2-control] game={game_id} single_post 回合未真结束" + f"(reason={turn.get('reason')}),不评中间态、不落库。", flush=True) + result["stopped_reason"] = f"single_post_turn_not_ended:{turn.get('reason')}" + result["wall_s"] = round(time.perf_counter() - t0, 1) + return result + # 回合真结束(REPLY_END/EXCEED_MAX_ITERS):读服务端续修已落的 verdict(不重跑门,评审 I5)。 + v = _run.read_last_verdict(gate_game_id) + judgment = judge_tier2_verdict({"rc": 0, "verdict": v, "log": ""}) + result["last_verdict"] = v or None + if judgment.passed: + # 门绿:据 on-disk 源工程落库(逻辑同外层循环 c 分支)。 + built = _collect_on_disk_project(gate_game_id) + if built: + addr = _run.persist_source_project( + gate_game_id, built["source_project"], built["file_list"], now_ts=time.time()) + result["store_addressing"] = addr + result["finished"] = bool(addr) + result["stopped_reason"] = "gates_green_persisted" if addr else "gates_green_but_persist_failed" + else: + result["stopped_reason"] = "gates_green_but_no_on_disk_source" + else: + # 软停/续修耗尽放行的尽力产物:门没绿 → 只留 on-disk workdir、不入 store(创始人 2026-07-02:质量门守住)。 + print(f"[tier2-control] game={game_id} single_post 门未绿(未过门={judgment.failed_gates})," + "尽力产物留 on-disk workdir、不入库。", flush=True) + result["stopped_reason"] = "single_post_gates_not_green" + result["wall_s"] = round(time.perf_counter() - t0, 1) + print(f"[tier2-control] game={game_id} single_post 结束: finished={result['finished']} " + f"reason={result['stopped_reason']} wall={result['wall_s']}s", flush=True) + return result + last_verdict: dict | None = None - # ── ② 外层有界 resume 循环(对齐 run_studio:424-458)── + # ── ② 外层有界 resume 循环(fallback:single_post=False;续修 middleware spike 不成时回落这条,设计 §6)── for attempt in range(max_resumes + 1): result["attempts"] = attempt + 1 # a. 等本回合 chat run 结束(SSE;超时/断流据 on-disk 继续,绝不卡死)。 diff --git a/tier2/gen-worker/tests/test_repair_wiring.py b/tier2/gen-worker/tests/test_repair_wiring.py new file mode 100644 index 00000000..8fd29127 --- /dev/null +++ b/tier2/gen-worker/tests/test_repair_wiring.py @@ -0,0 +1,129 @@ +"""服务集成单测:middleware 工厂注入续修(注入序 + soft_budget + 续修超时)+ drive_generation single_post 分支 +(接收 ended + read_last_verdict 判落库,mock、无真 Service)。 +跑:PYTHONPATH=tier2/gen-worker cheap-worker/.venv/bin/python -m pytest tier2/gen-worker/tests/test_repair_wiring.py -v +""" +import asyncio + +from worker.middleware import CircuitBreakerMiddleware, RepairMiddleware, Tier2TraceMiddleware + + +def test_factory_injects_repair_with_soft_budget_and_repair_timeouts(): + # extra 注入序 [tracer, repair, breaker];breaker 走 soft_budget=True + 续修模式放大超时(评审 C1/C2)。 + from service.app import _tier2_middlewares_factory + mws = asyncio.run(_tier2_middlewares_factory("u1", "a1", "s1")) + assert [type(m).__name__ for m in mws] == [ + "Tier2TraceMiddleware", "RepairMiddleware", "CircuitBreakerMiddleware"] + tracer, repair, breaker = mws + assert isinstance(tracer, Tier2TraceMiddleware) + assert isinstance(repair, RepairMiddleware) + assert isinstance(breaker, CircuitBreakerMiddleware) + assert breaker.soft_budget is True, "tier2 续修路走 soft 档(cheap 仍 hard)" + assert breaker.wall_timeout_s > 1800, "续修模式墙钟须放大容纳 N 次串行真门" + assert breaker.step_timeout_s > 300, "续修模式单步静默阈值须 > 单次门判上限" + + +def test_repair_budget_exhausted_reads_measured_spend(): + # 续修 budget_exhausted 读同回合 breaker 实测已花 ¥ 是否超上限(非预估标记):未超 False、超 True。 + from service.app import _tier2_middlewares_factory + _, repair, breaker = asyncio.run(_tier2_middlewares_factory("u1", "a1", "s1")) + breaker._rmb_gate_active = True + breaker.rmb_hard_limit = 3.0 + breaker.spent_rmb = 0.0 + assert repair._budget_exhausted() is False + breaker.spent_rmb = 3.5 + assert repair._budget_exhausted() is True + + +def test_single_post_persists_on_green(monkeypatch): + # single_post + 回合真结束(ended=True)+ 读到门绿 verdict → 落库、attempts=1、消费端不重跑门(评审 I5)。 + from service import control_plane as cp + import service.bootstrap as _bs + import worker.run as _run + + calls = {"wait": 0, "read": 0, "gates": 0} + + async def _fake_start(*a, **k): + return {"agent_id": "a1", "session_id": "s1"} + + async def _fake_wait(*a, **k): + calls["wait"] += 1 + return {"ended": True, "reason": "REPLY_END", "endEvent": {}} + + def _fake_read_verdict(gid): + calls["read"] += 1 + return {"decision": "accept", "layerResults": {"L1": {"passed": True, "gateResults": []}}} + + def _fake_gates(*a, **k): + calls["gates"] += 1 # 消费端不应调它 + return {"rc": 0, "verdict": None, "log": ""} + + monkeypatch.setattr(_bs, "start_new_game", _fake_start) + monkeypatch.setattr(cp, "_wait_for_turn_end", _fake_wait) + monkeypatch.setattr(_run, "read_last_verdict", _fake_read_verdict) + monkeypatch.setattr(_run, "run_gates", _fake_gates) + monkeypatch.setattr(cp, "_collect_on_disk_project", + lambda gid: {"source_project": {"ok": 1}, "file_list": [{"path": "game.js", "content": "x"}]}) + monkeypatch.setattr(_run, "persist_source_project", + lambda gid, sp, fl, **k: {"id": gid, "versionId": "v1", "sourceHash": "abc"}) + + res = asyncio.run(cp.drive_generation("http://x", "g1", "做个游戏", single_post=True, max_resumes=6)) + assert res["finished"] is True + assert res["attempts"] == 1 + assert calls["wait"] == 1 and calls["read"] == 1 + assert calls["gates"] == 0, "消费端应读 verdict、不重跑门(评审 I5)" + assert res["store_addressing"]["versionId"] == "v1" + + +def test_single_post_not_green_stays_on_disk(monkeypatch): + # single_post 回合真结束但门没绿(软停/续修耗尽):不落库、reason=single_post_gates_not_green、attempts=1。 + from service import control_plane as cp + import service.bootstrap as _bs + import worker.run as _run + + async def _fake_start(*a, **k): + return {"agent_id": "a1", "session_id": "s1"} + + async def _fake_wait(*a, **k): + return {"ended": True, "reason": "EXCEED_MAX_ITERS", "endEvent": {}} + + def _fake_read_verdict(gid): + return {"decision": "accept", "layerResults": {"L1": {"passed": False, "gateResults": [ + {"gate": "E_live", "passed": False, "fatal": True, "detail": ""}]}}} + + monkeypatch.setattr(_bs, "start_new_game", _fake_start) + monkeypatch.setattr(cp, "_wait_for_turn_end", _fake_wait) + monkeypatch.setattr(_run, "read_last_verdict", _fake_read_verdict) + + res = asyncio.run(cp.drive_generation("http://x", "g1", "做个游戏", single_post=True)) + assert res["finished"] is False + assert res["attempts"] == 1 + assert res["stopped_reason"] == "single_post_gates_not_green" + + +def test_single_post_turn_not_ended_does_not_evaluate(monkeypatch): + # SSE 超时未真结束(ended=False):不评中间态、不读 verdict、不落库、reason 标 not_ended(评审 I2)。 + from service import control_plane as cp + import service.bootstrap as _bs + import worker.run as _run + + calls = {"read": 0} + + async def _fake_start(*a, **k): + return {"agent_id": "a1", "session_id": "s1"} + + async def _fake_wait(*a, **k): + return {"ended": False, "reason": "total_timeout", "endEvent": None} + + def _fake_read_verdict(gid): + calls["read"] += 1 + return None + + monkeypatch.setattr(_bs, "start_new_game", _fake_start) + monkeypatch.setattr(cp, "_wait_for_turn_end", _fake_wait) + monkeypatch.setattr(_run, "read_last_verdict", _fake_read_verdict) + + res = asyncio.run(cp.drive_generation("http://x", "g1", "做个游戏", single_post=True)) + assert res["finished"] is False + assert res["attempts"] == 1 + assert res["stopped_reason"].startswith("single_post_turn_not_ended") + assert calls["read"] == 0, "回合未真结束不应读 verdict、不评中间态" diff --git a/tier2/gen-worker/worker/run.py b/tier2/gen-worker/worker/run.py index afec2339..1193c5f8 100644 --- a/tier2/gen-worker/worker/run.py +++ b/tier2/gen-worker/worker/run.py @@ -878,6 +878,22 @@ def run_gates(game_id: str, play_spec: dict | None = None, *, return {"rc": rc, "verdict": verdict, "log": log} +def read_last_verdict(game_id: str) -> dict | None: + """读工程 evidence/verdict.json(服务端续修 middleware 最后一次 run_gates 已落),不重跑门。 + + single_post 消费端据此判落库(评审 I5:避免消费端再起一次 ~300s 真门 + 与服务端 chrome 端口 4330/9322 竞态)。 + 文件不存在 / 解析失败 → None(判未过,不伪造门绿)。 + """ + vpath = _workdir(game_id) / "evidence" / "verdict.json" + if not vpath.exists(): + return None + try: + return json.loads(vpath.read_text(encoding="utf-8")) + except Exception as e: # noqa: BLE001 —— 解析失败按 None(判未过),不中断 + print(f"[tier2-run] read_last_verdict 解析失败(按未产出): {type(e).__name__}: {e}", flush=True) + return None + + def persist_source_project(game_id: str, source_project: dict | None, file_list: list[dict] | None, *, now_ts: float, store=None) -> dict | None: