games-development-ai/cheap-worker/cheap_service_driver.py
lili 686aaaa421
Some checks failed
contract-gates / contract-gates (push) Has been cancelled
docs-gate / docs-gate (push) Has been cancelled
feat(reference-assets): 签认《山海行纪》并闭合可信消费门
冻结地图1二十分钟纵切版的平衡、证据与金标登记。

新增 ReferenceAsset/2 清单、策略、release、受信快照及 CLI/Service/acceptance provenance /4 消费链;保持 survivor live、R1 签名与部署关闭。
2026-07-28 09:11:54 -07:00

796 lines
48 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""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 "<none>"
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/<session_id>.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/<id>/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:<model>),缺失(_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)
# 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)
# 失控兜底轮数(与三闸 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-<id>/、重跑同 gid 直接抹掉上一 run 的
# trace.jsonl/turns.jsonl/证据。故 scaffold 前先把上一 run 的 amgen-<id>/ 与 _wg1-gen/<id>/ 整体移进
# 带时间戳归档位(永不同 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。
sc = await asyncio.to_thread(cheap_run.scaffold, game_id, scaffold_template)
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/<id>/evidence/
# 残留(verdict/续修反馈/日志/截图全列)+ 兜底清 amgen-<id>/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,
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"],
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)