"""cheap_service_driver.py — 便宜档 · 单 POST 消费方(镜像 tier2 bootstrap + control_plane.single_post,cheap 收口)。 worker_service.py 收 §6.1 job 后,不再进程内 run_studio,而是驱动 cheap Service /chat: 注册 OpenAI 兼容凭据 → 建 agent(system prompt 的 ⟦G⟧=后端 gameId)/session → scaffold(后端 gameId,起点落 amgen-<后端 gameId>)+ 写 session→gameId 映射注册表 sidecar → 设 BYPASS → 发 kick → SSE 等这一次(内部续修多轮) 回合真结束 → 读九门 verdict + Service 收口采集 sidecar 组 run-summary → reply 外非阻塞丰富度评分。续修/门判/软预算 v3 下旧 gameplay resume 已摘,消费方首回合后跑 v3;仅 verified reject 可向同一 session 再发一次修复 POST。**C2**: 产物目录/评门/回调 trace 全用后端 gameId, Service 两工厂经会话注册表把框架分配的 session_id 解析回后端 gameId(工厂拿不到后端 gameId,故靠 sidecar 桥接)。 【惰性 import 红线】顶层零重依赖;httpx/agentscope 牵出的件在函数体内 import。 【代理旁路必保(I2)】**先 client.resolve_base_url()+install_proxy_bypass(base_url) 再 import httpx**(装 NO_PROXY 含 new-api host 与本地 Service 127.0.0.1);回调侧旁路由 worker_service._NO_PROXY_OPENER 兜。 """ from __future__ import annotations import hashlib import json import os import shutil import stat import tempfile import time from pathlib import PurePosixPath _REFERENCE_RECEIPTS_FILE = ".reference-receipts.json" _SESSION_CFG_WRITE_MAX_BYTES = 1 * 1024 * 1024 def _reference_snapshot_dir(session_cfg_path): """从服务端固定 sidecar 路径派生 session 独立快照目录,不接受 job 自报路径。""" return session_cfg_path.with_name(f"{session_cfg_path.stem}.reference-assets") def _validated_snapshot_path(value: str) -> str: """防御性复核验证器返回路径,确保物化永远留在临时 snapshot 根内。""" if not isinstance(value, str) or not value or "\\" in value or "\x00" in value: raise ValueError("reference snapshot 路径非法") path = PurePosixPath(value) if path.is_absolute() or any(part in ("", ".", "..") for part in path.parts): raise ValueError("reference snapshot 路径越界") return value def _read_existing_session_cfg_for_write(path): """写 session-cfg 前读取既有文件,无法确认时统一拒绝覆盖。 返回 ``None`` 表示文件不存在;返回字典表示已确认是普通 JSON 对象。读取使用固定上限、 ``O_NOFOLLOW`` 以及打开前后文件身份/元数据复核,避免把损坏文件、特殊文件或竞态中的文件 当成可安全修复的旧 sidecar。调用方据此区分「普通无 policy 可兼容覆盖」与「冻结/不可确认必须止损」。 """ try: before = os.lstat(path) except FileNotFoundError: return None except OSError as e: raise ValueError("已有 session-cfg 状态不可确认") from e # 既有路径必须是稳定的普通文件;symlink 或特殊文件都不能作为覆盖前的安全依据。 if stat.S_ISLNK(before.st_mode) or not stat.S_ISREG(before.st_mode): raise ValueError("已有 session-cfg 必须是普通文件且不得为 symlink") if before.st_size > _SESSION_CFG_WRITE_MAX_BYTES: raise ValueError("已有 session-cfg 超过固定读取上限") nofollow = getattr(os, "O_NOFOLLOW", None) if nofollow is None: raise ValueError("当前平台无法确认 session-cfg 非 symlink") fd = None try: fd = os.open(path, os.O_RDONLY | nofollow) opened = os.fstat(fd) if (not stat.S_ISREG(opened.st_mode) or (opened.st_dev, opened.st_ino) != (before.st_dev, before.st_ino)): raise ValueError("已有 session-cfg 读取期间发生替换") chunks = [] total = 0 while True: chunk = os.read( fd, min(64 * 1024, _SESSION_CFG_WRITE_MAX_BYTES + 1 - total)) if not chunk: break total += len(chunk) if total > _SESSION_CFG_WRITE_MAX_BYTES: raise ValueError("已有 session-cfg 超过固定读取上限") chunks.append(chunk) after = os.fstat(fd) if any(getattr(opened, field) != getattr(after, field) for field in ("st_dev", "st_ino", "st_size", "st_mtime_ns", "st_ctime_ns")): raise ValueError("已有 session-cfg 读取期间发生漂移") raw = b"".join(chunks) except OSError as e: raise ValueError(f"已有 session-cfg 读取边界拒绝:{type(e).__name__}") from e finally: if fd is not None: try: os.close(fd) except OSError: pass try: cfg = json.loads(raw.decode("utf-8")) except Exception as e: # noqa: BLE001 —— 任何无法确认的 JSON 都不得被覆盖修复 raise ValueError("已有 session-cfg JSON 不可确认") from e if not isinstance(cfg, dict): raise ValueError("已有 session-cfg 顶层必须是对象") return cfg def _resolve_base() -> str: """new-api host 根(已装代理旁路);cheap 凭据在其后补 /v1(OpenAI 兼容路)。""" from worker import client # noqa: PLC0415 return client.resolve_base_url() 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 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(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(user_token), "base_url": base}} def _write_session_cfg(session_id: str, *, external_game_id: str, write_whitelist=None, scaffold_template=None, traceparent=None, tracestate=None, reference_assets=None) -> bool: """写本 session 的会话注册表 sidecar(C2:Service 两工厂读它把 session_id 解析回后端 gameId + 取 write_whitelist)。 按 session_id 键写 game-runtime/games/_cheap-sessions/.json(worker↔Service 同机共享 FS); **external_game_id = 后端 gameId 恒写**(工厂据它绑六工具/评门/collector 目录,与 driver scaffold 一致)。 restricted = 是否受限写(create 路 False + write_whitelist=None;reskin/modify 路 True + 白名单;I1 fail-closed 依赖 restricted 标记)。traceparent/tracestate = 面四断点③ 桥:worker 这跳的 W3C context,让工厂据它把 cheap-service 生成 span 挂到入站 trace 下(reply 在后台任务跑、读不到出站 header 的实时 context,故走 sidecar 桥)。 无策略路径写失败仍按旧行为只告警;显式策略调用方必须检查 False 并在 /chat 前终止。 sidecar 始终通过同目录临时文件 + os.replace 原子发布;显式策略额外先原子发布独立 snapshot 目录。 """ import cheap_run # noqa: PLC0415 import reference_asset_gate # noqa: PLC0415 snapshot_dir = None snapshot_published = False temp_snapshot = None temp_cfg = None try: # 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) existing_cfg = _read_existing_session_cfg_for_write(p) if existing_cfg is not None and "reference_asset_policy" in existing_cfg: # 只要既有 cfg 留下过冻结声明,就算 snapshot 消失也不得用默认写入抹掉证据。 raise FileExistsError("已有 session-cfg 含 reference asset policy") snapshot_dir = _reference_snapshot_dir(p) try: os.lstat(snapshot_dir) except FileNotFoundError: pass except OSError as e: raise ValueError("session reference snapshot 状态不可读") from e else: # 同一 session 一旦有冻结 snapshot,默认写入也不得覆盖其 policy 声明。 raise FileExistsError("session reference snapshot 已存在") cfg = { "external_game_id": str(external_game_id), "write_whitelist": wl, "scaffold_template": scaffold_template, "restricted": restricted, } # 面四断点③ 桥:仅当有 traceparent 时才写(默认关/无传播时不写、sidecar 字节不变)。 if traceparent: cfg["traceparent"] = traceparent if tracestate: cfg["tracestate"] = tracestate if reference_assets is not None: if not isinstance(reference_assets, reference_asset_gate.VerifiedReferenceAssets): raise TypeError("reference_assets 必须是 VerifiedReferenceAssets") files = dict(reference_assets.reference_files) total_bytes = sum(len(content) for content in files.values()) if total_bytes > reference_asset_gate.MAX_TOTAL_BYTES: raise ValueError("reference snapshot 超过 128 MiB 硬帽") if reference_asset_gate.snapshot_hash(files) != reference_assets.snapshot_hash: raise ValueError("reference snapshot hash 与验证结果不一致") temp_snapshot = tempfile.mkdtemp( prefix=f".{p.stem}.reference-assets.", dir=str(p.parent)) temp_snapshot_path = p.parent / os.path.basename(temp_snapshot) file_index = [] for logical_path in sorted(files, key=lambda item: item.encode("utf-8")): safe_path = _validated_snapshot_path(logical_path) content = files[logical_path] if not isinstance(content, bytes): raise TypeError("reference snapshot 文件必须是 bytes") target = temp_snapshot_path.joinpath(*PurePosixPath(safe_path).parts) target.parent.mkdir(parents=True, exist_ok=True) with target.open("xb") as handle: handle.write(content) handle.flush() os.fsync(handle.fileno()) file_index.append({ "path": safe_path, "size": len(content), "sha256": hashlib.sha256(content).hexdigest(), }) receipt_payload = { "receipts": reference_asset_gate.to_json_value(reference_assets.receipts), } receipt_bytes = reference_asset_gate.canonical_json_bytes(receipt_payload) receipt_hash = hashlib.sha256(receipt_bytes).hexdigest() receipt_path = temp_snapshot_path / _REFERENCE_RECEIPTS_FILE with receipt_path.open("xb") as handle: handle.write(receipt_bytes) handle.flush() os.fsync(handle.fileno()) cfg["reference_asset_policy"] = { "policy_id": "survivor-gold-v1", "mode": "frozen_preflight", "snapshot_hash": reference_assets.snapshot_hash, "receipt_hash": receipt_hash, "roots": reference_asset_gate.to_json_value(reference_assets.reference_roots), "files": file_index, } os.replace(temp_snapshot_path, snapshot_dir) snapshot_published = True temp_snapshot = None cfg_bytes = json.dumps(cfg, ensure_ascii=False, separators=(",", ":")).encode("utf-8") fd, temp_cfg = tempfile.mkstemp(prefix=f".{p.name}.", dir=str(p.parent)) with os.fdopen(fd, "wb") as handle: handle.write(cfg_bytes) handle.flush() os.fsync(handle.fileno()) os.replace(temp_cfg, p) temp_cfg = None return True except Exception as e: # noqa: BLE001 —— sidecar best-effort,写失败 Service 侧回落 session_id print(f"[cheap-driver] session-cfg 写失败(Service 侧将回落 game_id=session_id、生成响亮失败):" f"{type(e).__name__}: {e}", flush=True) if temp_cfg: try: os.unlink(temp_cfg) except OSError: pass if temp_snapshot: shutil.rmtree(temp_snapshot, ignore_errors=True) if snapshot_published and snapshot_dir is not None: shutil.rmtree(snapshot_dir, ignore_errors=True) return False def _read_last_cheap_verdict(game_id: str): """读服务端最后一次九门 play 已落的 verdict(_wg1-gen//evidence/verdict.json)。缺失/坏 → None(判未过,不伪造)。""" import json # noqa: PLC0415 import cheap_run # noqa: PLC0415 vp = cheap_run.wg1_game_dir(game_id) / "evidence" / "verdict.json" if not vp.exists(): return None try: return json.loads(vp.read_text(encoding="utf-8")) except Exception as e: # noqa: BLE001 print(f"[cheap-driver] verdict 读失败(按未产出):{type(e).__name__}: {e}", flush=True) return None def _read_service_run_summary(game_id: str, *, wait_s: float = 5.0) -> dict: """读 Service 收口采集 sidecar(game_dir/evidence/service-run-summary.json:costRmb/rmbGate/repairs/…)。 C1 竞态兜底:Service 侧 collector 已在 yield REPLY_END 前**同步 flush**,正常 driver 读 REPLY_END 时盘上必有; 但为防极端时序(SSE 到达早于文件系统可见 / collector 走 finally 兜底路),这里再加**有界轮询**——最多等 wait_s (每 100ms 探一次「存在且可解析」),仍读不到才诚实降级 {}(driver 落 costRmb=0/degraded,不伪造)。竞态窗口是 ms 级,短轮询即闭合;这是 M3b「result_out 只产 trace.cost → D11 三维恒中性」坑的归并版防线。 """ import json # noqa: PLC0415 import time as _time # noqa: PLC0415 import cheap_run # noqa: PLC0415 p = cheap_run.game_dir(game_id) / "evidence" / "service-run-summary.json" deadline = _time.time() + max(0.0, wait_s) while True: if p.exists(): try: obj = json.loads(p.read_text(encoding="utf-8")) if isinstance(obj, dict): return obj except Exception: # noqa: BLE001 —— 半成品/坏 JSON:轮询期内再等,过期后按空降级 pass if _time.time() >= deadline: print(f"[cheap-driver] service-run-summary 有界轮询({wait_s:.1f}s)未读到,诚实降级 costRmb=0/degraded。", flush=True) return {} _time.sleep(0.1) def _read_driver_type(game_id: str): """W-AXIS-V2 波2:tap-targets/key-cycle 驱动器随取证契约退役,验收改由测试 agent 视觉引导真玩驱动。 停读已退役的 play-spec.json.driver.type,改读测试员真相层(evidence/playtest/playtest.json):存在则回其 模型标识(playtest:),缺失(_build_summary 在 run_acceptance 跑测试员之前先组 summary,此刻尚无)则回 常量 'playtest-agent'——两路都如实反映「v2 的驱动主体是测试 agent」。trace.driverType 现为纯观测维,D11 首局 信号已迁 trace.playtest(见 plan §5)。 """ import json # noqa: PLC0415 import cheap_run # noqa: PLC0415 p = cheap_run.wg1_game_dir(game_id) / "evidence" / "playtest" / "playtest.json" if not p.exists(): return "playtest-agent" try: model = (json.loads(p.read_text(encoding="utf-8")) or {}).get("model") return f"playtest:{model}" if model else "playtest-agent" except Exception: # noqa: BLE001 读不出只回通用标记,绝不抛 return "playtest-agent" def _build_summary(game_id: str, brief: str, verdict, svc: dict, turn: dict, t0: float) -> dict: """据九门 verdict + Service 收口采集 sidecar 组一份与 run_studio 同形的 run-summary(供 result_out.build_result_out)。 复用 cheap_studio 的 _verdict_brief(verdict 简报)+ build_trace_source(9d-trace 源维度 verdictFull/driverType/models/stage); costRmb/rmbGate/repairs 来自 Service 收口采集(worker 进程拿不到 Service 内存态,靠 sidecar 跨进程)。 attempts = repairs+1(result_out 的 trace.repairs = max(0, attempts-1) 会反推回 repairs)。 """ import cheap_run # noqa: PLC0415 import cheap_studio # noqa: PLC0415 from worker.gate_judge import classify_cheap_failure_layer, judge_cheap_verdict # noqa: PLC0415 # C6 接线:带 game_id/staged 目录判(驱动器族/日志摘要口径与 Service 续修一致;此处只为 summary 归一,不再落 sidecar 也无妨——staged_dir 传入使口径同源)。 j = judge_cheap_verdict(verdict or {}, game_id=game_id, staged_dir=cheap_run.wg1_game_dir(game_id)) # F0-归因分层(W-AXIS 波1):把九门失败归到三层(driver_contract/mechanical/gameplay)进 run-summary, # 供 result_out/批次账消费与回喂措辞对齐真实层次;放行(passed)/未失败时 layer=none。 failure_layer = classify_cheap_failure_layer(verdict or {}, staged_dir=cheap_run.wg1_game_dir(game_id)) vb = cheap_studio._verdict_brief(verdict) if verdict else None svc = svc or {} repairs = svc.get("repairs") attempts = (repairs + 1) if isinstance(repairs, int) else 1 driver_type = _read_driver_type(game_id) stage = "play" if verdict is not None else "code" # 简化:有 verdict=跑到 play;无=最远到写码(细分阶段列 follow-up) model = cheap_studio._bootstrap.SPIKE_MODEL summary = { "ok": j.passed, "gameId": game_id, "brief": brief, "model": model, "finished": j.passed, "attempts": attempts, "verdict": vb, "costRmb": round(float(svc.get("costRmb", 0.0) or 0.0), 4), "rmbGate": svc.get("rmbGate", "degraded"), "budgetSoftTripped": bool(svc.get("budgetSoftTripped", False)), "wallSec": round(time.time() - t0, 1), "stoppedReason": (turn.get("reason") if not turn.get("ended") else "turn_ended"), } # 9d-trace 源维度(verdictFull/driverType/models/stage);对照 None 时各键自然降级(result_out 省略)。 summary.update(cheap_studio.build_trace_source(verdict, driver_type, attempts, stage, model)) # F0-归因分层:失败三层归因随 summary 透传(additive;成功/未失败为 layer=none,不误标坏死)。 summary["failureLayer"] = failure_layer # 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 def _failed_summary(game_id: str, reason: str) -> dict: """早失败(scaffold 失败等)时的最小 failed summary(verdict None → result_out 走 failed 路)。""" import cheap_studio # noqa: PLC0415 print(f"[cheap-driver] game={game_id} 早失败:{reason}", flush=True) return { "ok": False, "gameId": game_id, "model": cheap_studio._bootstrap.SPIKE_MODEL, "finished": False, "attempts": 1, "verdict": None, "costRmb": 0.0, "rmbGate": "degraded", "stage": "scaffold", "stoppedReason": reason, } async def _run_v3_floor_gates(game_id: str) -> dict: """Service v3 回合结束后只跑一次机械收口;不在九门结果上做 gameplay resume。""" import asyncio # noqa: PLC0415 import cheap_gates # noqa: PLC0415 import cheap_verify # noqa: PLC0415 port, cdp_port = cheap_verify._derive_playtest_ports(f"{game_id}:service-floor") try: return await asyncio.to_thread(cheap_gates.run_cheap_gates, game_id, port, cdp_port) except Exception as e: # noqa: BLE001 —— v3 会把缺失四门证据诚实判 tester_error print(f"[cheap-driver] game={game_id} v3 机械门异常:{type(e).__name__}: {e}", flush=True) return {} def _merge_service_turns(first: dict, second: dict) -> dict: """合并两次 Service writer 回合成本;第二回合就是 v3 唯一修复。""" out = dict(second or {}) out["costRmb"] = round(float((first or {}).get("costRmb") or 0.0) + float((second or {}).get("costRmb") or 0.0), 4) out["repairs"] = 1 out["budgetSoftTripped"] = bool((first or {}).get("budgetSoftTripped") or (second or {}).get("budgetSoftTripped")) return out async def drive_cheap_generation(job: dict, *, base_url: str | None = None, user_id: str = "cheap"): """驱动 cheap Service /chat 跑一局便宜档生成;v3 最多首轮 + 一次 verified-reject 修复。返回摘要与目录。 与旧 worker_service._default_run_fn 同契约(process_job 零改动消费)。create 路默认:通用 scaffold + 通用系统提示 + write_whitelist=None(对齐现 _default_run_fn 的 run_studio(game_id, brief) 无 scaffold/whitelist)。 """ import asyncio # noqa: PLC0415 from worker import client, config, genconfig # noqa: PLC0415 —— client 先导(不牵 httpx),装旁路 base_url = base_url or os.environ.get("CHEAP_SERVICE_URL", "http://127.0.0.1:8300") # I2 代理旁路必保(**必须先于 import httpx**):resolve_base_url 把 new-api host 并入 NO_PROXY(装凭据/M3 侧); # driver→Service 打的是 127.0.0.1:8300、resolve_base_url 不含它,故再 install_proxy_bypass(base_url) 把本地 # Service host 也并入 NO_PROXY,避免 trust_env=True 的 httpx 被 clash fake-ip 拦本地环回(Opus/Codex I2)。 client.resolve_base_url() client.install_proxy_bypass(base_url) import httpx # noqa: PLC0415 —— 旁路已装,此后 import 才安全 import _bootstrap # noqa: PLC0415 import cheap_otlp_sink # noqa: PLC0415 —— 面四三跳传播:取当前 context 的 W3C carrier(顶层零重依赖) import cheap_run # noqa: PLC0415 import cheap_verify # noqa: PLC0415 import cheap_studio # noqa: PLC0415 from cheap_roles import build_system_prompt, SCAFFOLD_DESC_BY_TEMPLATE # noqa: PLC0415 from service.control_plane import _wait_for_turn_end # noqa: PLC0415 —— 复用 tier2 SSE 等回合 game_id = job.get("gameId") or job.get("job_id") # 真后端 gameId 是 Java Long → JSON 数字 → int;cheap_run 的 subprocess/路径要求 str(同 worker_service 边界归一)。 game_id = str(game_id) if game_id is not None else None brief = job.get("brief") or "" t0 = time.time() acceptance_mode = cheap_studio.acceptance_v3_mode() v3_enabled = acceptance_mode in ("v3", "v3_shadow") reference_asset_policy_id = job.get("referenceAssetPolicyId") frozen_reference_assets = None frozen_reference_constraint_block = None reference_asset_generation_receipts = None if reference_asset_policy_id is not None: try: # 只把服务端 job 的 policyId 作为选择信号;路径、release 和 hash 全由统一 helper 固定。 frozen_reference_assets = cheap_verify.preflight_reference_asset_policy( reference_asset_policy_id, acceptance_mode) frozen_reference_constraint_block = cheap_verify.build_frozen_reference_asset_constraint_block( frozen_reference_assets) import reference_asset_gate # noqa: PLC0415 reference_asset_generation_receipts = reference_asset_gate.to_json_value( frozen_reference_assets.receipts) except Exception as exc: # noqa: BLE001 可信预检失败不得创建 Writer 请求 reason = f"reference asset frozen preflight 失败:{type(exc).__name__}: {exc}" return _failed_summary(game_id, reason), cheap_run.game_dir(game_id) task_binding_hash = None if v3_enabled: try: task_binding_hash = cheap_verify.task_binding_hash_v3(job.get("traceId") or job.get("job_id")) except ValueError as exc: failed = cheap_studio.apply_v3_entry_failure( _failed_summary(game_id, str(exc)), str(exc), mode=acceptance_mode) return failed, cheap_run.game_dir(game_id) interaction_profile_id = job.get("interactionProfileId") # 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(与旧行为一致,路由只加不减)。 # write_whitelist=None(create 不收窄,对齐现 _default_run_fn)。 from cheap_genre_route import route_genre # noqa: PLC0415 scaffold_template, genre = route_genre(brief) scaffold_desc = None write_whitelist = None if scaffold_template: # 2026-07-08 创始人纠偏:规范 = 文件位置(工程规范) + 绝对通用 plumbing(已由 cheap_run._L1_FIXED 锁 # host-config/game/index/entry/main 的 boot 链与装配),**不把模型限死到只改几个文件**——只改几个文件游戏 # 做不好玩、也切不中 brief 主题(浏览器验收实测:写锁 core.js 致「榫卯」brief 出「陶艺」游戏,主题漂)。 # 故 create 路给 per-genre 完整可玩起点(scaffold_template)+ 一句话描述引导(scaffold_desc),但**不加写锁**: # game-logic/core/render/assets 全放开,模型自由把它做成好玩且切题的游戏。render.js 整写截断类失败靠 harness # 门(check 的 ESM 导出对账 + stripCode 假阴性已修)兜住并回喂续修,不靠锁模型的手。 scaffold_desc = SCAFFOLD_DESC_BY_TEMPLATE.get(scaffold_template) print(f"[cheap-driver] game={game_id} 品类路由命中:genre={genre} template={scaffold_template} " f"(per-genre 起点+引导,不写锁)", flush=True) if interaction_profile_id is not None: try: # job 只能显式选择已冻结的 shadow 校准 profile;必须先于凭据、网络、归档和 scaffold。 cheap_verify.preflight_interaction_profile_v3( interaction_profile_id, acceptance_mode, genre=str(genre or ""), template_route=str(scaffold_template or "")) except Exception as exc: # noqa: BLE001 不可信 profile 不得触发 Writer 或产物写入 reason = f"interaction profile 预检失败:{type(exc).__name__}: {exc}" failed = cheap_studio.apply_v3_entry_failure( _failed_summary(game_id, reason), reason, mode=acceptance_mode) return failed, cheap_run.game_dir(game_id) # 失控兜底轮数(与三闸 150 对齐;成本上界=max_repairs=6 优雅终止,而非轮数闸;放开防单 POST 多续修被 EXCEED_MAX_ITERS 提前切断)。 max_iters = genconfig.get("budget", "cheap_max_iters", 150) max_tokens = genconfig.get("model", "max_tokens", 16000) _bootstrap.ensure_api_key_env() # 保证 NEWAPI_KEY(装凭据要用) # 面四断点②:把「worker 这跳的当前 context」(worker_span_scope 在 worker 线程 attach、经 asyncio.run 透传进本 # 协程)注入成 W3C carrier,随所有出站到 :8300 的 header 带上,让 cheap-service 端能接住入站 context。同一份 # carrier 也写进会话 sidecar(下方 _write_session_cfg)——因 /chat 的 reply 在框架后台任务里跑、生成 span 产在 # 那条与 HTTP 请求解耦的任务里,读不到出站 header 的实时 context,故靠 driver 恒在 /chat 前写的 sidecar 把 # traceparent 跨进程桥给工厂(与 external_game_id 同一 C2 桥)。纯增量 best-effort:无 context / otel 不可用 → # carrier 为空,header 少带一项、生成 span 退回私有 root,绝不咬生成(设计 §8)。 otel_carrier = cheap_otlp_sink.current_traceparent_carrier() headers = {"X-User-Id": user_id} 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(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 往下走—— # 否则后续 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 = { "name": "cheap-writer", # create 路统一用通用 create prompt + per-genre scaffold_desc 引导(模型自由改全部游戏文件,不锁)。 "system_prompt": build_system_prompt(game_id, scaffold_desc), "context_config": ctx_cfg.model_dump(mode="json"), "react_config": {"max_iters": max_iters}, } 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 分离)。 # fix400 主防线:parallel_tool_calls=False 源头禁并行 tool call——M3 经 new-api 的并行 call 不区分 # 流式 chunk index,agentscope 2.0.2 聚合会把第二个 call 的 arguments 拼进首桶(非法 JSON)→ 下一请求 # 服务端丢弃该 call → tool result 孤儿 → 400 code 2013 → run 崩、SSE 无终结、driver 慢失败 600s。 # 经 get_model 的 Parameters(**parameters) 透传,_call_api 对 API 带 parallel_tool_calls=false # (venv agentscope/model/_openai_chat/_model.py:264-265)。纵深兜底见 m3_stream_patch(Service 进程内按 id 分桶)。 session_body = { "agent_id": agent_id, "chat_model_config": { "type": "openai_credential", "credential_id": credential_id, "model": _bootstrap.SPIKE_MODEL, "parameters": {"max_tokens": max_tokens, "parallel_tool_calls": False}, }, } 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)) # ④a per-run 归档(W-AXIS 波1 F0-b):scaffold 会 rmSync 整个 amgen-/、重跑同 gid 直接抹掉上一 run 的 # trace.jsonl/turns.jsonl/证据。故 scaffold 前先把上一 run 的 amgen-/ 与 _wg1-gen// 整体移进 # 带时间戳归档位(永不同 gid 覆盖),返回归档指针记进 run-summary/批次账。best-effort:关/失败退回旧覆盖行为。 archived_to = await asyncio.to_thread(cheap_run.archive_prior_run, game_id) # ④b scaffold(clone 模板 + SAA 信封)——必须在 /chat 前,**game_id=后端 gameId**(与 system prompt 的 ⟦G⟧ 一致), # agent 起点就位在 amgen-<后端 gameId>(subprocess 经 to_thread)。C2:产物目录统一后端 gameId、不用 session_id。 if interaction_profile_id is None: # 普通 job 保持历史调用形状,既不选择 profile,也不向 scaffold 传空占位参数。 sc = await asyncio.to_thread(cheap_run.scaffold, game_id, scaffold_template) else: sc = await asyncio.to_thread( cheap_run.scaffold, game_id, scaffold_template, interaction_profile_id=interaction_profile_id) if not sc["ok"]: _fs = _failed_summary(game_id, f"scaffold 失败:{sc['output'][:300]}") if archived_to: _fs["archivedPriorRunTo"] = archived_to return _fs, cheap_run.game_dir(game_id) # ④c evidence 清理清单化(W-AXIS 波1 F0-c,取代旧红线③只清 verdict.json 一个文件):清 _wg1-gen//evidence/ # 残留(verdict/续修反馈/日志/截图全列)+ 兜底清 amgen-/evidence/,封「读回上一 run 反馈/绿 verdict → 误归因/假绿」。 await asyncio.to_thread(cheap_run.clean_stale_evidence, game_id) acceptance_identity = None if v3_enabled: try: acceptance_identity = cheap_verify.build_acceptance_v3_identity( game_id, brief, genre=str(genre or ""), template_route=str(scaffold_template or ""), repair_ordinal=0, task_binding_hash=task_binding_hash, interaction_profile_id=interaction_profile_id, reference_asset_record_ids=( [receipt["recordId"] for receipt in reference_asset_generation_receipts] if reference_asset_generation_receipts is not None else None), consumer_ref=(reference_asset_gate.POLICY_CONSUMER_REF if reference_asset_generation_receipts is not None else None)) except Exception as exc: # noqa: BLE001 —— 无可信 route/profile 时禁止向 Writer 发首条消息 failed = _failed_summary( game_id, f"Writer 前无法冻结 v3 acceptance identity:{type(exc).__name__}: {exc}") failed = cheap_studio.apply_v3_entry_failure( failed, failed.get("fail") or "acceptance identity 非法", mode=acceptance_mode) if archived_to: failed["archivedPriorRunTo"] = archived_to return failed, cheap_run.game_dir(game_id) # ⑤ 写会话注册表 sidecar(C2:Service 两工厂据 session_id 读它解析回后端 gameId、绑六工具/评门/collector; # create 路 write_whitelist=None → restricted=False)。external_game_id 恒写、正常路工厂必读到。 sidecar_ok = _write_session_cfg( session_id, external_game_id=game_id, write_whitelist=write_whitelist, scaffold_template=scaffold_template, traceparent=otel_carrier.get("traceparent"), tracestate=otel_carrier.get("tracestate"), reference_assets=frozen_reference_assets, ) if not sidecar_ok: if frozen_reference_assets is not None: failed = _failed_summary(game_id, "reference asset session snapshot 原子发布失败") if archived_to: failed["archivedPriorRunTo"] = archived_to return failed, cheap_run.game_dir(game_id) # 旧默认 sidecar 写失败仍保持 best-effort;但发现已有冻结 snapshot 时必须停在 /chat 前。 try: os.lstat(_reference_snapshot_dir(cheap_run.session_cfg_path(session_id))) except FileNotFoundError: pass except OSError: failed = _failed_summary(game_id, "已有 reference asset session snapshot 状态不可确认") if archived_to: failed["archivedPriorRunTo"] = archived_to return failed, cheap_run.game_dir(game_id) else: failed = _failed_summary(game_id, "默认 sidecar 不得覆盖已有 reference asset session snapshot") if archived_to: failed["archivedPriorRunTo"] = archived_to return failed, cheap_run.game_dir(game_id) try: existing_cfg = _read_existing_session_cfg_for_write( cheap_run.session_cfg_path(session_id)) except ValueError: failed = _failed_summary(game_id, "已有 reference asset session-cfg 状态不可确认") if archived_to: failed["archivedPriorRunTo"] = archived_to return failed, cheap_run.game_dir(game_id) if existing_cfg is not None and "reference_asset_policy" in existing_cfg: failed = _failed_summary(game_id, "默认 sidecar 不得覆盖已有 reference asset policy cfg") if archived_to: failed["archivedPriorRunTo"] = archived_to return failed, cheap_run.game_dir(game_id) # ⑥ 设 BYPASS 权限(六工具默认 ASK,服务态无人确认,不设首个工具调用即卡死;PATCH 需 query agent_id,同 tier2)。 await http.patch(f"{base_url}/sessions/{session_id}", json={"permission_mode": "bypass"}, headers=headers, params={"agent_id": agent_id}, timeout=30.0) # ⑦ 发 kick(引导 read skill → 写 game-logic.js → check/build → finish;文本对齐旧 run_studio 的 create kick)。 kick = (f"请按这个 brief 造一款游戏:「{brief}」。先 read_file 读手册(.agents/skills/littlejs-game-dev.md)" "和你的起点 game-logic.js 再动手;核心玩法实现完、check 与 build 都绿了就立即 finish。" + (f"\n\n{frozen_reference_constraint_block}" if frozen_reference_constraint_block else "")) await http.post(f"{base_url}/chat/", json={ "agent_id": agent_id, "session_id": session_id, "input": {"name": "user", "role": "user", "content": [{"type": "text", "text": kick}]}, }, headers=headers, timeout=30.0) # ⑧ SSE 等这一次回合真结束(内部续修多轮)。C3:idle 从 genconfig 派生(不硬编码 300;慢门 ~390s 静默不误弃单), # 总超时按 (max_repairs+1)×(门+推理)放宽、且 > breaker 墙钟(让 breaker 先优雅收、driver 别抢先弃单)。 sse_idle_s = genconfig.get("budget", "cheap_repair_step_timeout_s", 600) # = breaker step_timeout,> 单次门跑最坏 sse_timeout = genconfig.get("budget", "cheap_sse_turn_timeout_s", (genconfig.get("iteration", "max_resumes", 6) + 1) * (genconfig.get("budget", "cheap_gate_timeout_s", 420) + 180) + 600) 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_s) if not turn.get("ended"): # 回合未真结束(SSE 总超时/断流):Service 可能仍在续修写盘,读中间态会误判 → 据现有产物 best-effort 组 summary、诚实标 reason。 print(f"[cheap-driver] game={game_id} 回合未真结束(reason={turn.get('reason')}),据现有产物组 summary。", flush=True) # ⑨ 读九门 verdict + Service 收口采集 sidecar → 组 run-summary(供 result_out;C2:全用后端 gameId)。 # _read_service_run_summary 带有界轮询(C1:关闭 collector flush 与 driver 读 REPLY_END 的竞态)。 # SSE 未真结束时 writer 可能仍在改盘,禁止并发起机械门;交 v3 以缺失证据诚实判 tester_error。 verdict = (await _run_v3_floor_gates(game_id) if v3_enabled and turn.get("ended") else ({} if v3_enabled else _read_last_cheap_verdict(game_id))) svc = _read_service_run_summary(game_id) summary = _build_summary(game_id, brief, verdict, svc, turn, t0) # F0-b:上一 run 归档指针随 run-summary 透传(批次账每 run 记一笔;首跑/关/失败为 None,不带该键)。 if archived_to: summary["archivedPriorRunTo"] = archived_to if v3_enabled and genre not in cheap_verify._V3_GENRES: summary = cheap_studio.apply_v3_entry_failure( summary, f"v3 无法从 brief/template 元数据确定 canonical genre:{genre or '-'}", mode=acceptance_mode) print(f"[cheap-driver] game={game_id} v3 genre 缺失,冻结发布且不回落旧验收。", flush=True) elif v3_enabled: acceptance = await cheap_verify.run_acceptance_v3(cheap_studio.build_acceptance_v3_request( game_id, brief, verdict, acceptance_identity=acceptance_identity, idempotency_key=f"{job.get('traceId') or game_id}:service-v3", writer_cost_rmb=float(svc.get("costRmb") or 0.0), reference_asset_policy_id=reference_asset_policy_id, reference_asset_generation_receipts=reference_asset_generation_receipts)) decision = acceptance.get("decision") or {} acceptance_first_pass = None print(f"[cheap-driver] game={game_id} 验收 v3 mode={acceptance_mode} outcome={decision.get('outcome')} " f"accepted={decision.get('accepted')} publishFrozen={decision.get('publishFrozen')} " f"repairEligible={decision.get('repairEligible')}", flush=True) if cheap_verify.is_v3_repair_authorized(acceptance): # Service v3 不装 RepairMiddleware;只有 final verified reject 才额外 POST 同一会话一次。 acceptance_first_pass = acceptance acceptance_identity = cheap_verify.build_acceptance_v3_identity( game_id, brief, genre=acceptance_identity["genre"], template_route=acceptance_identity["templateRoute"], source_artifact_hash=acceptance["artifactHash"], parent_acceptance_request_hash=acceptance_identity["acceptanceRequestHash"], repair_ordinal=1, proof_profile_id=acceptance_identity["proofProfileId"], proof_registry_version=acceptance_identity["proofRegistryVersion"], task_binding_hash=acceptance_identity["taskBindingHash"], interaction_binding=acceptance_identity.get("interactionBinding"), design_ref=acceptance_identity.get("designRef"), reference_asset_record_ids=acceptance_identity.get("referenceAssetRecordIds"), consumer_ref=acceptance_identity.get("consumerRef"), ) feedback = cheap_studio.build_v3_repair_prompt(decision.get("repairFeedback")) async with httpx.AsyncClient() as http: await http.post(f"{base_url}/chat/", json={ "agent_id": agent_id, "session_id": session_id, "input": {"name": "user", "role": "user", "content": [{"type": "text", "text": feedback}]}, }, headers=headers, timeout=30.0) repair_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_s) # 修复回合若未真结束,同样禁止在 writer 改盘中并发验收,也不把首回合 sidecar 误读成第二回合成本。 repaired_verdict = await _run_v3_floor_gates(game_id) if repair_turn.get("ended") else {} repaired_svc = _read_service_run_summary(game_id) if repair_turn.get("ended") else {} svc = _merge_service_turns(svc, repaired_svc) summary = _build_summary(game_id, brief, repaired_verdict, svc, repair_turn, t0) if archived_to: summary["archivedPriorRunTo"] = archived_to acceptance = await cheap_verify.run_acceptance_v3(cheap_studio.build_acceptance_v3_request( game_id, brief, repaired_verdict, acceptance_identity=acceptance_identity, idempotency_key=f"{job.get('traceId') or game_id}:service-v3", repair_count=1, parent_run_id=acceptance.get("runId"), # 首轮成本由 sealed parent decision 读取;请求只提交修复回合 writer 的新增成本。 writer_cost_rmb=float(repaired_svc.get("costRmb") or 0.0), reference_asset_policy_id=reference_asset_policy_id, reference_asset_generation_receipts=reference_asset_generation_receipts)) verdict = repaired_verdict summary = cheap_studio.apply_acceptance_v3(summary, acceptance, first_pass=acceptance_first_pass) # ⑩ reply 外收口:非阻塞丰富度 LLM 评分(additive;不进 verdict、不改 ok;token 不污染生成成本台账)。 # genre 透传(T4):品类路由命中时同一次评分 additive 追加品类扩展条目(与 run_studio 的 genre 参数同义)。 try: summary["richness"] = await cheap_verify.verify_richness(game_id, brief=brief, genre=genre) rv = summary["richness"] print(f"[cheap-driver] game={game_id} 丰富度 score={rv.get('score')}/{rv.get('max')} " f"genre={genre or '-'} groups={rv.get('groups')} degraded={rv.get('degraded')}", flush=True) except Exception as e: # noqa: BLE001 —— richness 非阻塞铁律:绝不影响回调 summary["richness"] = {"score": None, "degraded": True, "reason": f"richness 接线异常:{type(e).__name__}: {e}"} # ⑪ 统一验收编排器(W-AXIS-V2 波1 · 拆着杀):四门投影(A/B/C/D)作预筛权威,过筛者交测试 agent 视觉引导 # 真玩裁 broken/hollow/off-brief,阻断放行(accepted = floor ∧ 测试员;acceptance.mode 三态 v2/shadow/v1)。 # verdict 传入做 floor 投影;端口缺省按 game_id 派生(driver 收口可能并发,测试员起服避撞)。 # run_acceptance 内建 fail-closed 与顶层兜底、绝不抛,additive 写 floor/playtest/judge/acceptanceVersion, # 更新 ok/accepted——result_out 据 accepted 落 status。 if not v3_enabled: summary = await cheap_verify.run_acceptance(summary, game_id=game_id, brief=brief, verdict=verdict) _js = summary.get("judge") or {} _pt = summary.get("playtest") or {} print(f"[cheap-driver] game={game_id} 历史验收 mode={summary.get('acceptanceVersion')} " f"floor={(summary.get('floor') or {}).get('pass')} playtest.accepted={_pt.get('accepted')} " f"rolls={_pt.get('rollCount')} degraded={_pt.get('degraded')} costRmb={_pt.get('costRmb')} " f"→ accepted={summary.get('accepted')} ok={summary.get('ok')}", flush=True) print(f"[cheap-driver] game={game_id} Service 驱动结束: ok={summary['ok']} attempts={summary['attempts']} " f"costRmb={summary['costRmb']} wallSec={summary['wallSec']}", flush=True) return summary, cheap_run.game_dir(game_id)