games-development-ai/cheap-worker/cheap_service_app.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

714 lines
45 KiB
Python
Raw Permalink 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_app.py — 便宜档 · 独立 Agent Service 服务壳(镜像 tier2 service/app.py,cheap 各起进程;决策②)。
把便宜档从「裸 HTTP /generate + 进程内 for-resume」归并到 AgentScope Service /chat:七工具/软预算/trace
经工厂 per-turn 注入。v3 模式不装旧 RepairMiddlewarewriter 完成后由 driver 跑机械门与 v3只有 final verified
reject 才允许同一 session 修一次。RepairMiddleware 仅保留给显式 v1/v2 历史回放。
与 tier2 Service 各起独立进程:cheap credential 走 OpenAI 兼容路(base+/v1)、预算 soft+¥10 软目标(成本上界靠
max_repairs 优雅终止、三闸=150 失控兜底)、六工具面、并发按 session 从进程内端口池派生(避撞固定 4320/9222)。
两工厂经会话注册表 sidecar 把框架分配的 session_id 解析回后端 gameId(C2:产物目录/评门/回调统一后端 gameId)。
【惰性 import 红线】顶层绝不 import agentscope/redis/fastapi;重依赖只在 build_cheap_app / 工厂体内 import。
【代理旁路必保】build_cheap_app 启动时 client.resolve_base_url() 装 NO_PROXY,否则 M3 请求被 clash fake-ip 转走 → 502。
"""
from __future__ import annotations
import hashlib
import json
import os
import queue as _queue
import stat
import sys
import threading
from pathlib import Path
from typing import TYPE_CHECKING
# cheap-worker/ 加进 sys.path(_bootstrap 再把 tier2/gen-worker 加进去,worker.*/service.*/observability.* 可 import)。
sys.path.insert(0, str(Path(__file__).resolve().parent))
import _bootstrap # noqa: E402,F401 仅加 sys.path + key(import 时不触网)
if TYPE_CHECKING: # pragma: no cover
from fastapi import FastAPI
# 服务标题(OpenAPI docs 显示)。
SERVICE_TITLE = "cheap-littlejs-agent-service"
_SESSION_CFG_MAX_BYTES = 1 * 1024 * 1024
_REFERENCE_RECEIPTS_FILE = ".reference-receipts.json"
def _reference_snapshot_dir(session_cfg_path: Path) -> Path:
"""由固定 sidecar 路径派生 session 快照目录,不读取配置中的路径。"""
return session_cfg_path.with_name(f"{session_cfg_path.stem}.reference-assets")
# ── 并发端口池(决策③):便宜档并发≤15,play 用固定 4320/9222 会撞;按 session 从池派生唯一端口对 ──
# 与 tier2(4330/9322)、旧 cheap 串行默认(4320/9222)错开,避免同机多 Service 撞端口。池大小 16 > 上限 15。
# 续修 check 每次 acquire 一对、finally release;单局各 check 串行(一 reply 内),并发多局各持一对,池 16 足够。
_CHEAP_PORT_BASE = int(os.environ.get("CHEAP_GATE_PORT_BASE", "4400"))
_CHEAP_CDP_BASE = int(os.environ.get("CHEAP_GATE_CDP_BASE", "9400"))
_CHEAP_PORT_POOL_SIZE = int(os.environ.get("CHEAP_GATE_POOL_SIZE", "16"))
_cheap_port_pool: "_queue.Queue | None" = None
_cheap_port_pool_lock = threading.Lock()
def _port_pool() -> "_queue.Queue":
"""惰性建端口池(线程安全):16 对 (port, cdp) 从 base 递增派生。
用 LifoQueue(栈式)而非 FIFO Queue:release 归还的端口对下一次 acquire 优先复用(最近释放先取),
这样「释放→再取」拿回的就是刚归还的那对,便于验证不泄漏(池验收契约 test_port_pool_derives_distinct_pairs);
FIFO 会先派发池底未用过的对、掩盖泄漏。LifoQueue is-a Queue,对外 get/put 接口不变、并发正确性不受影响。
"""
global _cheap_port_pool
with _cheap_port_pool_lock:
if _cheap_port_pool is None:
pool: "_queue.Queue" = _queue.LifoQueue()
for i in range(_CHEAP_PORT_POOL_SIZE):
pool.put((_CHEAP_PORT_BASE + i, _CHEAP_CDP_BASE + i))
_cheap_port_pool = pool
return _cheap_port_pool
def acquire_ports(timeout: float = 900.0) -> tuple:
"""取一对空闲端口(并发满时阻塞至有空,timeout 兜底;上游并发≤15 < 池 16,正常不阻塞)。"""
return _port_pool().get(timeout=timeout)
def release_ports(pair) -> None:
"""归还端口对(best-effort,绝不抛)。"""
try:
_port_pool().put_nowait(pair)
except Exception: # noqa: BLE001 —— 归还失败只丢一对端口,不连累生成
pass
def _read_session_cfg(session_id: str) -> dict:
"""读本 session 的会话注册表 sidecar(driver 在 /chat 前按 session_id 写;C2 收口)。
因 AgentScope 工厂签名固定 (user_id, agent_id, session_id)、拿不到后端 gameId 也拿不到 session 记录,
且 AgentData/SessionConfig 无自由字段,改用 worker↔Service 同机共享 FS sidecar
game-runtime/games/_cheap-sessions/<session_id>.json 把「session_id → 后端 gameId + 本 session 配置」传给工厂。
缺失 sidecar → {},保持无 sidecar 的旧路径兼容;已存在的文件若 JSON 损坏、顶层不是对象,或派生
snapshot 存在但缺少冻结 policy均直接拒绝避免在工具工厂边界回落到活目录。
存在的文件必须是固定上限内的普通文件symlink、特殊文件、超限或读取竞态直接拒绝。
"""
import cheap_run # noqa: PLC0415
p = cheap_run.session_cfg_path(session_id)
try:
before = os.lstat(p)
except FileNotFoundError:
# 只有 cfg 与派生 snapshot 都不存在才是真正缺 sidecar孤立冻结 snapshot 必须拒绝 live 回落。
snapshot_dir = _reference_snapshot_dir(p)
try:
os.lstat(snapshot_dir)
except FileNotFoundError:
return {}
except OSError as e:
raise ValueError("reference asset session snapshot 状态不可读") from e
raise ValueError("reference asset session snapshot 存在但 session-cfg 缺失")
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_MAX_BYTES:
raise ValueError("session-cfg 超过固定读取上限")
try:
flags = os.O_RDONLY | getattr(os, "O_NOFOLLOW", 0)
fd = os.open(p, flags)
try:
opened = os.fstat(fd)
if (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_MAX_BYTES + 1 - total))
if not chunk:
break
total += len(chunk)
if total > _SESSION_CFG_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 读取期间发生漂移")
finally:
os.close(fd)
raw = b"".join(chunks)
except OSError as e:
raise ValueError(f"session-cfg 普通文件边界拒绝:{type(e).__name__}") from e
except ValueError:
raise
try:
obj = json.loads(raw.decode("utf-8"))
except (UnicodeDecodeError, json.JSONDecodeError) as e:
print(f"[cheap-service] session-cfg JSON 非法,拒绝回落活目录:{type(e).__name__}: {e}", flush=True)
raise ValueError("session-cfg JSON 非法") from e
if not isinstance(obj, dict):
raise ValueError("session-cfg 顶层必须是对象")
# 冻结 snapshot 与 policy 必须成对存在;只剩 snapshot 时拒绝把工具绑定回原资产活目录。
if obj.get("reference_asset_policy") is None:
snapshot_dir = _reference_snapshot_dir(p)
try:
os.lstat(snapshot_dir)
except FileNotFoundError:
pass
except OSError as e:
raise ValueError("reference asset session snapshot 状态不可读") from e
else:
raise ValueError("reference asset session snapshot 存在但 sidecar 缺 policy")
return obj
def _is_sha256(value) -> bool:
"""校验 sidecar 中稳定 SHA-256 文本,拒绝宽松大小写和非字符串。"""
return isinstance(value, str) and len(value) == 64 and all(char in "0123456789abcdef" for char in value)
def _load_reference_snapshot(session_id: str, cfg: dict):
"""按 sidecar 索引重验 session snapshot返回 Toolkit 使用的只读 files/roots。
只有 cfg 显式声明冻结 policy 才进入本路径任何目录、文件、索引、receipt 或 canonical hash 漂移
都在 ``CheapSession`` 与工具创建前抛错,绝不读取原资产活目录或 best-effort 回落。
"""
policy = cfg.get("reference_asset_policy")
if policy is None:
return None, None
expected_keys = {"policy_id", "mode", "snapshot_hash", "receipt_hash", "roots", "files"}
if not isinstance(policy, dict) or set(policy) != expected_keys:
raise ValueError("reference asset policy sidecar 结构非法")
if policy["policy_id"] != "survivor-gold-v1" or policy["mode"] != "frozen_preflight":
raise ValueError("reference asset policy/mode 不受信")
if not _is_sha256(policy["snapshot_hash"]) or not _is_sha256(policy["receipt_hash"]):
raise ValueError("reference asset sidecar hash 非法")
if not isinstance(policy["roots"], dict) or not policy["roots"]:
raise ValueError("reference asset roots 索引非法")
entries = policy["files"]
if not isinstance(entries, list) or not entries or len(entries) > 512:
raise ValueError("reference asset 文件索引非法")
import artifact_snapshot # noqa: PLC0415 复用 fd/O_NOFOLLOW 可信读取边界
import cheap_run # noqa: PLC0415
import reference_asset_gate # noqa: PLC0415
snapshot_dir = _reference_snapshot_dir(cheap_run.session_cfg_path(session_id))
try:
root_stat = os.lstat(snapshot_dir)
except FileNotFoundError as exc:
raise ValueError("reference asset session snapshot 目录缺失") from exc
if stat.S_ISLNK(root_stat.st_mode) or not stat.S_ISDIR(root_stat.st_mode):
raise ValueError("reference asset session snapshot 根不是普通目录")
paths = []
by_path = {}
total_declared = 0
for entry in entries:
if not isinstance(entry, dict) or set(entry) != {"path", "size", "sha256"}:
raise ValueError("reference asset 文件索引条目非法")
path, size, digest = entry["path"], entry["size"], entry["sha256"]
if (not isinstance(path, str) or not path or path in by_path
or not isinstance(size, int) or isinstance(size, bool) or size < 0
or not _is_sha256(digest)):
raise ValueError("reference asset 文件索引漂移")
paths.append(path)
by_path[path] = entry
total_declared += size
if total_declared > reference_asset_gate.MAX_TOTAL_BYTES:
raise ValueError("reference asset snapshot 超过 128 MiB 硬帽")
if paths != sorted(paths, key=lambda item: item.encode("utf-8")):
raise ValueError("reference asset 文件索引顺序漂移")
files = {}
total_observed = 0
for path in paths:
try:
captured = artifact_snapshot.capture_selected_files(
snapshot_dir,
[path],
limits={
"max_files": 1,
"max_file_bytes": reference_asset_gate.MAX_FILE_BYTES,
"max_record_bytes": reference_asset_gate.MAX_FILE_BYTES,
"max_total_bytes": reference_asset_gate.MAX_FILE_BYTES,
},
)
except artifact_snapshot.ArtifactSnapshotError as exc:
raise ValueError(f"reference asset snapshot 文件拒绝:{exc.code}") from exc
content = captured.files[path]
total_observed += len(content)
expected = by_path[path]
if len(content) != expected["size"] or hashlib.sha256(content).hexdigest() != expected["sha256"]:
raise ValueError("reference asset 文件索引漂移")
if total_observed > reference_asset_gate.MAX_TOTAL_BYTES:
raise ValueError("reference asset snapshot 超过 128 MiB 硬帽")
files[path] = content
if total_observed != total_declared:
raise ValueError("reference asset 文件索引总量漂移")
if reference_asset_gate.snapshot_hash(files) != policy["snapshot_hash"]:
raise ValueError("reference asset snapshot hash 漂移")
try:
receipt_capture = artifact_snapshot.capture_selected_files(
snapshot_dir,
[_REFERENCE_RECEIPTS_FILE],
limits={"max_files": 1, "max_file_bytes": 1 * 1024 * 1024,
"max_record_bytes": 1 * 1024 * 1024, "max_total_bytes": 1 * 1024 * 1024},
)
except artifact_snapshot.ArtifactSnapshotError as exc:
raise ValueError(f"reference asset receipt 拒绝:{exc.code}") from exc
receipt_bytes = receipt_capture.files[_REFERENCE_RECEIPTS_FILE]
if hashlib.sha256(receipt_bytes).hexdigest() != policy["receipt_hash"]:
raise ValueError("reference asset receipt hash 漂移")
try:
receipt_payload = json.loads(receipt_bytes.decode("utf-8"))
except Exception as exc: # noqa: BLE001 receipt 必须是 canonical JSON 对象
raise ValueError("reference asset receipt 内容非法") from exc
receipts = receipt_payload.get("receipts") if isinstance(receipt_payload, dict) else None
if (not isinstance(receipts, list) or not receipts
or any(not isinstance(item, dict)
or item.get("finalSnapshotHash") != policy["snapshot_hash"] for item in receipts)):
raise ValueError("reference asset receipt 与 snapshot 不一致")
if reference_asset_gate.canonical_json_bytes(receipt_payload) != receipt_bytes:
raise ValueError("reference asset receipt canonical 字节漂移")
return files, policy["roots"]
def _resolve_external_game_id(session_id: str, cfg: dict) -> str:
"""把框架分配的 session_id 解析回后端 gameId(C2:产物目录/评门/回调统一后端 gameId)。
cfg.external_game_id 存在即用它;缺失(sidecar 未写/坏)→ 回落 session_id —— 此时产物目录与 driver scaffold
的后端 gameId 目录不一致,生成会尽早失败(agent read 不到起点 / write 越界),是响亮失败、不静默用错目录。
"""
ext = cfg.get("external_game_id")
if ext:
return str(ext)
print(f"[cheap-service] ⚠ 会话注册表缺 external_game_id(session={session_id}),回落 game_id=session_id;"
"产物目录可能与 scaffold 不一致、生成将失败——正常路 driver 恒在 /chat 前写映射。", flush=True)
return session_id
def _resolve_write_whitelist(cfg: dict):
"""据会话注册表算写边界白名单(I1 fail-closed)。
create 路(非 restricted)→ None(不收窄,全 L3 自由写);reskin/modify 路(restricted:true)→ 收窄成
write_whitelist 的 set;**restricted 为真但白名单缺失/空 → 空集**(禁写 L3,绝不回落 create 的不收窄)——
受限会话的 sidecar 若被写坏/漏字段,宁可禁写也不放开写边界(Codex I1:受限模式缺失 fail-closed)。
"""
if not cfg.get("restricted"):
return None # create 路:不收窄
wl = cfg.get("write_whitelist")
return set(wl) if isinstance(wl, list) and wl else set() # restricted 缺白名单 → 空集(禁写)
async def _cheap_tools_factory(user_id: str, agent_id: str, session_id: str) -> list:
"""extra_agent_tools 工厂:每回合产便宜档六工具(game_id=后端 gameId,经 sidecar 解析;C2)。
C2:六工具的写/评门目录必须与 driver scaffold 的目录、system prompt 里的 ⟦G⟧ 路径一致——都用**后端 gameId**
(driver scaffold 在 amgen-<后端 gameId>、prompt 写 amgen-<后端 gameId>)。工厂只拿到 session_id,故经会话注册表
sidecar 把 session_id 解析回 external_game_id(后端 gameId)再绑 CheapSession,避免「prompt 指 A、工具锁 B」的写越界。
write_whitelist 走 I1 fail-closed helper。复用 tier2 service.app._extract_function_tools 从 build_toolkit 的 Toolkit
抠 FunctionTool(丢 group.mcps,九门不做 MCP)。
"""
from service.app import _extract_function_tools # noqa: PLC0415 —— 复用 tier2 抠取(纯函数,6c6g 安全)
from cheap_toolkit import CheapSession, build_toolkit # noqa: PLC0415
cfg = _read_session_cfg(session_id)
game_id = _resolve_external_game_id(session_id, cfg) # C2:session_id → 后端 gameId
write_whitelist = _resolve_write_whitelist(cfg) # I1:create None / restricted 缺白名单则空集(禁写)
reference_files, reference_roots = _load_reference_snapshot(session_id, cfg)
session = CheapSession(
game_id=game_id,
reference_files=reference_files,
reference_roots=reference_roots,
) # 六工具据后端 gameId 管产物目录;显式 policy 同时绑定已复核的 session 只读快照。
toolkit = build_toolkit(session, write_whitelist=write_whitelist)
return _extract_function_tools(toolkit)
def _get_collector_cls():
"""惰性定义 collector middleware 类(subclass MiddlewareBase,缓存一次;顶层不 import agentscope 保 6c6g)。"""
global _COLLECTOR_CLS
if _COLLECTOR_CLS is not None:
return _COLLECTOR_CLS
from agentscope.middleware import MiddlewareBase # noqa: PLC0415
class _CheapRunCollector(MiddlewareBase):
"""便宜档收口采集(cheap-only,不碰 tier2 共用类):把成本/trace/续修数写 evidence sidecar 供 driver 跨进程读。
归并后 breaker/tracer/repair 都在 Service 进程内(per-turn),worker 侧 driver 拿不到它们的内存态;
本 middleware 把 costRmb/rmbGate/repairs/budgetSoftTripped/traceSummary 写 game_dir(后端 gameId)/evidence/
service-run-summary.json,driver 读它 + 九门 verdict 组 result-out(result_out 的 trace 需 costRmb/repairs 承重,
D11 就绪分据此算)。**C1 竞态修复**:在 ReplyEndEvent 流经时**同步 flush**(在 yield 它给下游 SSE publisher 前),
关闭「driver 读到 REPLY_END 时 Service 侧 finally 还没 flush → costRmb=0」竞态;finally 仍兜底(降级三路见 on_reply)。
纯旁路 best-effort:采集/写盘任何异常只告警、不影响生成;放注入序最外层,finally 在内层全跑完后执行。
"""
def __init__(self, *, game_id, breaker, repair, tracer) -> None:
self._game_id = game_id # C2:后端 gameId(经 sidecar 解析),与 driver 读 evidence 的目录一致
self._breaker = breaker
self._repair = repair
self._tracer = tracer
self._flush_logged = False # 工单 f:幂等打印标记——多次 _flush 只打一行「收口采集落盘」(落盘覆盖语义不变)
# 面二观测:本 run 累计 token(从 ModelCallEndEvent 只读嗅探)。生产 Service 路 in/out 可得;cached 拿不到
# (ModelCallEndEvent 不带 cached,唯一经手 cached 的是 tier2 共用 breaker,不动它)→ 发 metric 时 cached 传 None 降级。
self._tok_in = 0
self._tok_out = 0
self._metric_emitted = False # 面二 metric 只发一次(_flush 三条降级路会调多次,防重复计数)
# WU2 §3.9:模型调用抛错(new-api 非 2xx)时按 one-api 惯例分辨的失败因;仅 quota_exhausted 才写进
# 收口 sidecar 供 driver→result_out 落成 quota_exhausted(凭据失效/其他维持 llm_error,不透传)。默认 None。
self._failure_reason = None
async def on_reply(self, agent, input_kwargs, next_handler):
# C1:框架 _agent.py:615 先 yield ReplyEndEvent 再 yield finish Msg 再收尾;collector 是最外层 on_reply,
# 在把 ReplyEndEvent yield 给 SSE publisher 前先落盘(flush 同步无 await)→ SSE 侧收到 REPLY_END 时盘上必有,
# 关闭 driver 的 read-before-flush 竞态。此时最后一次模型调用已 accumulate、breaker.spent_rmb 为终值。
try:
async for evt in next_handler(**input_kwargs):
# 面二观测:只读嗅探每次模型调用的 token 用量(in/out);best-effort,绝不改事件、绝不抛。
# 在既有事件透传循环里嗅探,不新挂 on_model_call(不碰模型调用流式关键路径,零咬主链风险)。
if type(evt).__name__ == "ModelCallEndEvent":
try:
self._tok_in += int(getattr(evt, "input_tokens", 0) or 0)
self._tok_out += int(getattr(evt, "output_tokens", 0) or 0)
except Exception: # noqa: BLE001 —— 嗅探纯旁路,解析异常忽略,绝不咬生成
pass
if type(evt).__name__ == "ReplyEndEvent":
self._flush() # 正常路:REPLY_END 前落盘
yield evt # 纯旁路透传,不拦截、不修改
except Exception as e: # noqa: BLE001 —— T5 合成终结事件:run 崩快失败出口(600s 慢失败封口)
# run 崩(压缩崩 / 熔断 / 框架异常):框架 ChatService.run 只 log 吞异常、**不向 bus 发任何
# 终结事件**(app/_service/_chat.py:166-176)→ SSE 消费方等不到 REPLY_END、只能 600s 总超时
# 慢失败。这里在异常穿出前先 yield 一枚【合成 ReplyEndEvent】——消费侧 _run_impl 的 async-for
# 先拿到并 publish 到 bus(异常在下一次推进时才抛),driver 立刻收到回合终结、按盘面组 failed
# summary 快速失败;异常原样重抛,不掩盖失败(框架日志照记 exception)。合成失败(极端:事件类
# 构造变化)只告警、退回慢失败路,绝不遮原异常。
# WU2 §3.9:run 崩多因 new-api 非 2xx(额度耗尽 402/403 是其一)。best-effort 按 one-api 惯例分辨
# ——仅额度耗尽写 self._failure_reason,随 _flush 落进收口 sidecar,让 driver→result_out 落成
# quota_exhausted(而非含糊 llm_error);凭据失效/其他维持 llm_error。纯旁路,任何异常绝不遮原异常。
try:
import result_out as _ro # noqa: PLC0415 —— 复用单一分辨口径(one-api 惯例,阶段2 精确 message)
_reason = _ro.classify_newapi_failure_reason(*_ro.newapi_error_status_text(e))
if _reason == "quota_exhausted":
self._failure_reason = _reason
print(f"[cheap-service] 模型调用非 2xx 按 §3.9 分辨为额度耗尽 game={self._game_id}"
f"(quota_exhausted):{type(e).__name__}: {str(e)[:160]}", flush=True)
except Exception: # noqa: BLE001 —— 分辨纯旁路,失败不影响崩溃收口与原异常上抛
pass
self._flush() # 崩溃路也先落盘(部分成本/repairs 可回收)
try:
from agentscope.event import ReplyEndEvent # noqa: PLC0415
st = getattr(agent, "state", None)
synth = ReplyEndEvent(
session_id=str(getattr(st, "session_id", "") or ""),
reply_id=str(getattr(st, "reply_id", "") or ""),
)
print(f"[cheap-service] run 崩({type(e).__name__}: {str(e)[:200]})→ 发合成 ReplyEndEvent"
f"(session={synth.session_id})快终结,异常继续上抛。", flush=True)
yield synth
except Exception as synth_err: # noqa: BLE001 —— 合成兜底自身失败:退回慢失败,不遮原异常
print(f"[cheap-service] 合成终结事件失败(退回慢失败路):{type(synth_err).__name__}: "
f"{synth_err}", flush=True)
raise
finally:
# 降级三路:① 正常 REPLY_END 已 flush(此处幂等覆盖、不怕重复);② breaker 硬熔断(execute_chain 异常
# 沿 async-for 上抛)→ 没走到 REPLY_END,except 支已 flush+合成终结,finally 幂等再兜;③ 进程中途崩
# (未走 finally)→ 盘上无 sidecar,driver 有界轮询后诚实降级 costRmb=0/degraded(不伪造)。
self._flush()
def _flush(self) -> None:
try:
import json # noqa: PLC0415
import cheap_run # noqa: PLC0415
b = self._breaker
summary = {
"costRmb": round(float(getattr(b, "spent_rmb", 0.0) or 0.0), 4),
"rmbGate": "active" if getattr(b, "_rmb_gate_active", False) else "degraded",
"repairs": getattr(self._repair, "repairs", 0),
"budgetSoftTripped": bool(getattr(b, "budget_soft_tripped", False)),
"traceSummary": self._tracer.summary(),
}
# WU2 §3.9:仅当模型调用非 2xx 被分辨为额度耗尽时才带 failureReason(driver 据它落 quota_exhausted);
# 正常/其他失败不带此键 → driver 不透传 → result_out 走既有 llm_error 兜底,不误标。
if self._failure_reason:
summary["failureReason"] = self._failure_reason
ev = cheap_run.game_dir(self._game_id) / "evidence"
ev.mkdir(parents=True, exist_ok=True)
(ev / "service-run-summary.json").write_text(
json.dumps(summary, ensure_ascii=False, indent=2), encoding="utf-8")
# 工单 f:_flush 在 on_reply 的 REPLY_END / except / finally 三条降级路会被调 ≥2 次,落盘每次覆盖
# 是幂等的(采集语义不变、driver 恒读到最新),但成功日志只需一行——重复打「收口采集落盘」纯观感噪声。
# 用实例标记只在首次成功落盘时打,后续静默覆盖(失败行不受此门、每次都报,便于诊断)。
if not self._flush_logged:
self._flush_logged = True
print(f"[cheap-service] 收口采集落盘 game={self._game_id} costRmb={summary['costRmb']} "
f"repairs={summary['repairs']} softTripped={summary['budgetSoftTripped']}", flush=True)
except Exception as e: # noqa: BLE001 —— 采集 best-effort,失败绝不影响生成
print(f"[cheap-service] 收口采集失败(忽略,不影响生成):{type(e).__name__}: {e}", flush=True)
# 面二观测:发 token 计数 metric(独立于落盘成败、只发一次)。
self._emit_metrics()
def _emit_metrics(self) -> None:
"""面二:发 llm_tokens{kind} 计数(默认关 / best-effort / 只发一次,防 _flush 多调重复计数)。
生产 Service 路 cached 拿不到(见 on_reply 嗅探注)→ cached_tokens 传 None,record_llm_token_metrics
据此只发 input/output 计数、不发 cached 计数与 cache_hit_rate(如实降级,不伪造命中率)。
取不到 CHEAP_OTLP_ENDPOINT 即 no-op(与 trace sink 默认关同源);面五不发 token 计数,不双计;
发送任何异常吞掉,绝不影响生成。
"""
if self._metric_emitted:
return
self._metric_emitted = True
try:
import cheap_otlp_sink # noqa: PLC0415 惰性 import(顶层不牵 otel)
cheap_otlp_sink.record_llm_token_metrics(self._tok_in, self._tok_out, None)
except Exception as e: # noqa: BLE001 —— 面二 metric best-effort,绝不影响生成主链
print(f"[cheap-service] 面二 metric 发送异常(忽略,不影响生成):{type(e).__name__}: {e}",
flush=True)
_COLLECTOR_CLS = _CheapRunCollector
return _COLLECTOR_CLS
_COLLECTOR_CLS = None
async def _cheap_middlewares_factory(user_id: str, agent_id: str, session_id: str) -> list:
"""按验收模式装配中间件v3 为 collector/trace/breaker历史模式才额外装旧 repair。
历史续修 check = finish 点独立跑便宜档九门(cheap_gates.run_cheap_gates 经 to_thread,端口按 session 从池派生)→
judge_cheap_verdictv3 不调用此闭包,机械门改由 driver 在回合结束后只跑一次。
breaker soft_budget=Truev3 只有一次协议修复max_repairs=6 只约束历史回放;
三道次数/轮数闸(max_tool_calls/max_model_calls/max_iters)抬到 150 当纯失控兜底(Codex C4:不设 max_tool_calls
会吃默认 60 先撞);超时按最坏九门 ~390s 放宽(C3)。评门 game_id = 后端 gameId(经 sidecar 解析,C2)。
注入序(外→内):collector / tracer / [legacy repair] / breakerturns 观测件最后追加且不改控制流。
"""
import asyncio # noqa: PLC0415
import cheap_budget # noqa: PLC0415 —— 预算两段式同源工厂(与 CLI cheap_studio 读同一配置源)
import cheap_gates # noqa: PLC0415
import cheap_run # noqa: PLC0415
from worker import genconfig # noqa: PLC0415
from worker.gate_judge import GateJudgment, judge_cheap_verdict # noqa: PLC0415
from worker.middleware import ( # noqa: PLC0415
RepairMiddleware,
Tier2TraceMiddleware,
)
# C2:评门/collector/trace 全绑后端 gameId(经会话注册表 sidecar 把 session_id 解析回后端 gameId;与 driver
# scaffold、system prompt 的 ⟦G⟧、七工具写目录一致)。缺映射则回落 session_id(响亮失败,见 _resolve_external_game_id)。
_cfg = _read_session_cfg(session_id)
game_id = _resolve_external_game_id(session_id, _cfg)
# 面四断点③:从 sidecar 取 worker 这跳写下的 W3C traceparent(driver 恒在 /chat 前写),组入站 carrier
# 传给 build_trace_sink → 让 OTLP 生成 span 挂到 Java→worker 这条入站 trace 下当子 span(不再私有 root)。
# reply 在框架后台任务里跑、与收 /chat 的 HTTP 请求解耦,读不到出站 header 的实时 context,故靠 driver 写的
# sidecar 跨进程桥 traceparent(与 external_game_id 同一 C2 桥)。无 traceparent(默认关)→ None,退回私有 root。
_tp = _cfg.get("traceparent")
parent_carrier = None
if _tp:
parent_carrier = {"traceparent": _tp}
_ts = _cfg.get("tracestate")
if _ts:
parent_carrier["tracestate"] = _ts
# trace sink:落 game_dir/trace.jsonl(与旧 cheap_studio 同路径);import/构造失败回落无 sink(非阻塞)。
trace_sink = None
try:
from observability.trace import make_jsonl_sink # noqa: PLC0415
trace_sink = make_jsonl_sink(cheap_run.game_dir(game_id) / "trace.jsonl")
except Exception as e: # noqa: BLE001 —— trace 非阻塞:接线失败只告警、回落无 sink
print(f"[cheap-service] trace sink 接线失败(回落无 sink,不阻断):{type(e).__name__}: {e}", flush=True)
# 阶段四观测 波③ T3-1:jsonl 落盘之外并接一个 OTLP span sink(发 mini-infra Collector,service.name=cheap-gen-worker)。
# **默认关**(env CHEAP_OTLP_ENDPOINT 未设)→ build_trace_sink 原样返回上面的 jsonl sink,现有行为字节不变;
# 开则包一层 fan-out(jsonl 先落盘照旧、OTLP 额外发一份,两侧各 best-effort 互不连累)。观测挂掉绝不咬生成(设计 §8)。
# 惰性 import:cheap_otlp_sink 顶层零重依赖(otel 只在其函数体内 import),不违 build_cheap_app 惰性红线。
try:
import cheap_otlp_sink # noqa: PLC0415
trace_sink = cheap_otlp_sink.build_trace_sink(trace_sink, trace_id=game_id,
parent_carrier=parent_carrier)
except Exception as e: # noqa: BLE001 —— OTLP 并接非阻塞:失败回落原 jsonl sink、只告警
print(f"[cheap-service] OTLP sink 并接失败(回落原 jsonl sink,不阻断):{type(e).__name__}: {e}", flush=True)
tracer = Tier2TraceMiddleware(trace_id=game_id, sink=trace_sink) # C2:trace_id 用后端 gameId(与 trace 文件目录/trace.gameId 一致)
max_repairs = genconfig.get("iteration", "max_resumes", 6) # 成本上界=优雅终止(决策①)
# C3:单次九门收口最坏 ~390s(cheap_run 超时上界 stage≤90 + smoke≤120 + play≤180);门跑期间 to_thread 阻塞、
# reply 不产 AgentEvent、SSE 只剩心跳。故 gate 基准 ≥390(取 420),单步静默/墙钟按它放宽,绝不用旧 200(< 390 会误熔断)。
gate_timeout_s = genconfig.get("budget", "cheap_gate_timeout_s", 420) # ≥ 最坏单次门跑,别再取 200
repair_step_s = genconfig.get("budget", "cheap_repair_step_timeout_s", 600) # > 单次门跑 + 推理余量(治慢门误熔断);driver idle 同源派生
repair_wall_s = genconfig.get("budget", "cheap_repair_wall_timeout_s",
(max_repairs + 1) * (gate_timeout_s + 180) + 300) # (轮数)×(门+推理)+余量,容 N 次串行真门
# 预算两段式(W-ARCH②,创始人 2026-07-03 裁决;取代旧「soft 无地板」):软停线 ¥10(越线只许收尾类
# 动作,on_acting 拦新增生成调用)+ RMB 硬地板 ¥15(×1.5 fail-closed,数学封死上界——软停不等于无界)。
# 经 cheap_budget.build_cheap_breaker 同源工厂,与 CLI(cheap_studio)读同一 genconfig 配置源。
# Codex C4:三道次数/轮数闸抬到 150 当纯失控兜底(非成本尺;成本上界=硬地板+max_repairs 优雅终止)。
# max_tool_calls 不显式设会吃默认 60(middleware.py)→ 7 轮 ReAct 工具步先撞,故必须显式给 150。
cheap_max_tool_calls = genconfig.get("budget", "cheap_max_tool_calls", 150)
cheap_max_model_calls = genconfig.get("budget", "cheap_max_model_calls", 150)
breaker = cheap_budget.build_cheap_breaker(
max_tool_calls=cheap_max_tool_calls, # C4 失控兜底:显式 150(不设吃默认 60、会先撞)
max_model_calls=cheap_max_model_calls, # 失控兜底 150(非成本尺;成本上界=两段式硬地板)
wall_timeout_s=repair_wall_s, # 墙钟容纳 (max_repairs+1) 次串行真门(C3)
step_timeout_s=repair_step_s, # 单步静默阈值 > 单次门跑最坏(门跑期间无 evt 不误杀,C3)
)
async def _cheap_check(agent):
# 独立跑九门:端口按 session 从池派生(决策③避撞),subprocess 经 to_thread 防阻塞事件循环;
# 异常兜底(T2-b 工厂闭包 try,镜像 tier2 app.py:210-216):任何异常判未过续修,绝不穿透 on_reasoning。
# #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])
except Exception as e: # noqa: BLE001 —— 门执行任何异常都判未过续修
print(f"[cheap-repair] 九门执行异常(判未过、续修):{type(e).__name__}: {e}", flush=True)
return GateJudgment(passed=False, failed_gates=["cheap_gates_error"],
feedback=f"九门执行异常,请确认工程可 check/build/运行后再交付:{type(e).__name__}: {e}")
finally:
release_ports(pair)
# C6 接线(W-S1 单③):传 game_id + staged 目录,反馈带 phaseNow/driverType/game-log 摘要,
# 并把结构化反馈落 staged evidence/verdict-feedback.json(证据留痕,gate_judge 内 best-effort)。
return judge_cheap_verdict(verdict, game_id=game_id,
staged_dir=cheap_run.wg1_game_dir(game_id))
acceptance_mode = str(genconfig.get("acceptance", "mode", "v3_shadow"))
repair = None
if acceptance_mode not in ("v3", "v3_shadow"):
# 仅历史 v1/v2 回放保留九门 gameplay resumev3 的唯一修复权归 final verified reject。
repair = RepairMiddleware(
check=_cheap_check,
max_repairs=max_repairs,
budget_exhausted=lambda: (
breaker._rmb_gate_active and breaker.spent_rmb >= breaker.rmb_hard_limit),
)
# 收口采集(cost/trace/repairs 跨进程回收):reply 收尾写 evidence/service-run-summary.json,供 T3 driver 组 result-out。
collector = _get_collector_cls()(game_id=game_id, breaker=breaker, repair=repair, tracer=tracer)
# W-AXIS 波1 真相层:逐 turn 全文落 amgen-<gameId>/turns.jsonl(模型文本/工具全参/工具返回/门续修反馈)。
# observe-only、best-effort,默认开(CHEAP_TURNS_ENABLED=0 关);接线失败/关 → None,不加入列表(现有行为字节不变)。
# 历史模式下 RepairMiddleware 会 mid-reply 注入 name=gatev3 没有该事件。turns 只读 context 尾部,
# 两种模式都可放最内层,不改变控制流。
mws = [collector, tracer]
if repair is not None:
mws.append(repair)
mws.append(breaker)
try:
import cheap_turns_sink # noqa: PLC0415 —— 顶层零 agentscope,工厂内惰性建中间件
turns_mw = cheap_turns_sink.build_turns_middleware(game_id, trace_id=game_id)
if turns_mw is not None:
mws.append(turns_mw) # 观测旁路,置于最内层(只读事件、原样 yield,顺序不影响拦截语义)
except Exception as e: # noqa: BLE001 —— 真相层接线失败绝不阻断生成
print(f"[cheap-service] turns 真相层中间件接线失败(降级不落 turns,不阻断生成):{type(e).__name__}: {e}", flush=True)
return mws
def build_cheap_app(*, title: str = SERVICE_TITLE) -> "FastAPI":
"""组装便宜档独立 Agent Service(create_app + Redis storage/message_bus + 本地 workspace;决策②:与 tier2 各起进程)。
部署入口:`from cheap_service_app import build_cheap_app; app = build_cheap_app()` 再 uvicorn(见 main)。
重依赖(agentscope.app / redis)在此惰性 import;代理旁路必保(I2):**先 client.resolve_base_url() 装 NO_PROXY,
再 import agentscope.app**(agentscope 牵出 httpx,trust_env=True 默认吃 HTTP_PROXY,装旁路必须先于 import,client.py:51)。
"""
from worker import client # noqa: PLC0415 —— 先导 client(不牵 httpx),装旁路
# I2 代理旁路必保:装 NO_PROXY(内网 M3 请求经 clash fake-ip 会 502),**必须在 import agentscope/httpx 前**(红线)。
client.resolve_base_url()
# B浅(cutover 接缝):Service 守护进程注入 NEWAPI_KEY 兜底 —— 成本子系统 / 活取价读 env 的地方兜底。
# 生成本身走 per-POST 凭据(工厂 per-turn 装 model)、不依赖它;但 CircuitBreaker 的 ¥ 累进硬闸 + newapi_pricing
# 活取价经 client.get_api_key() 读 env NEWAPI_KEY,进程无 key 会降级(costRmb=0/rmbGate=degraded、¥ 硬闸失效)。
# B深层已让成本子系统优先用 per-POST 凭据 key(middleware._extract_credential_key),env 仅作回落兜底,故此处
# best-effort:ensure_api_key_env 从 docs/内网凭据与端点.md 解析注入(内网单一事实源),失败也不阻断 Service 启动。
try:
_bootstrap.ensure_api_key_env()
except Exception as e: # noqa: BLE001 —— key 兜底失败不阻断启动(per-POST 凭据仍可取价;成本子系统至多降级)
print(f"[cheap-service] NEWAPI_KEY 进程注入兜底失败(忽略,per-POST 凭据仍可取价):{type(e).__name__}: {e}", flush=True)
from agentscope.app import create_app # noqa: PLC0415 —— 旁路已装,此后 import 才安全
from agentscope.app.message_bus import RedisMessageBus # noqa: PLC0415
from agentscope.app.storage import RedisStorage # noqa: PLC0415
from agentscope.app.workspace_manager import LocalWorkspaceManager # noqa: PLC0415
from service import infra_config # noqa: PLC0415 —— 复用 tier2 Redis 参数(mini-infra)
# fix400 纵深兜底:M3 经 new-api 的并行 tool call 不区分流式 chunk index,agentscope 2.0.2 聚合会把第二个
# call 的 arguments 拼进首桶(非法 JSON → 服务端丢 call → 孤儿 tool result → 400 code 2013 → run 崩、
# SSE 无终结、driver 慢失败)。主防线在 driver(session parameters parallel_tool_calls=False 源头禁并行);
# 此补丁兜「网关/模型不尊重该参数仍并行」的残余路径:同 index 异 id 按 id 分桶。
# 【钉 agentscope==2.0.3,升级必须复核】运行时 monkeypatch、绝不改 venv 文件本体(细则见 m3_stream_patch)。
from m3_stream_patch import apply_m3_stream_patch # noqa: PLC0415
apply_m3_stream_patch()
# T5 压缩崩封口:M3 对压缩摘要 schema 遵从不稳('important_discoveries' required 4 次实证)→ 压缩
# ValidationError 穿透 reply → run 崩 → SSE 无终结 → driver 600s 慢失败。fail-open 补丁让压缩失败
# 只跳过不杀主链(连续 2 次断路防白烧);快失败出口另由 collector 合成终结事件兜。
# 【钉 agentscope==2.0.3,升级必须复核】纪律同 m3_stream_patch(不改 venv 本体)。
from ctx_compress_patch import apply_ctx_compress_patch # noqa: PLC0415
apply_ctx_compress_patch()
# 工具面封口:2.0.2 Service 路把 workspace 内建(Bash/Edit/Glob/Grep/Read/Write)无条件并入 agent
# 工具面 → 便宜档六工具写白名单被内建 Write/Bash 旁路 + M3 拿内建 Bash 相对路径死圈(80009/80011
# 生产实证)。本进程只跑便宜档 agent,进程内关内建作用域恰好;tier2 服务独立进程不受影响。
# 【钉 agentscope==2.0.3,升级必须复核】纪律同 m3_stream_patch(不改 venv 本体)。
from ws_builtin_tools_patch import apply_ws_builtin_tools_patch # noqa: PLC0415
apply_ws_builtin_tools_patch()
# 工具面再封口(工单 h):2.0.2 get_toolkit 除内建外还无条件并入 Planning(Task*)/Team(Team*·AgentCreate),
# 并按 session model 挂 Schedule(schedule_tools 组)——便宜档单机单游戏用不到,纯工具面噪声(80011 flail:
# 12+ 工具语义重叠是根因之一)。本进程只跑便宜档 agent,进程内收窄作用域恰好;tier2 独立进程不受影响。
# 【钉 agentscope==2.0.3,升级必须复核】纪律同 ws_builtin_tools_patch(不改 venv 本体)。
from planning_team_tools_patch import apply_planning_team_tools_patch # noqa: PLC0415
apply_planning_team_tools_patch()
storage_params = infra_config.redis_params(for_message_bus=False)
bus_params = infra_config.redis_params(for_message_bus=True)
# cheap 独立 workspace 根(与 tier2 的 _service-workspaces 分开;本线生成产物另落 game-runtime/games/amgen-<id>)。
ws_basedir = os.environ.get(
"CHEAP_SERVICE_WORKSPACES",
str(Path(__file__).resolve().parent / "_cheap-service-workspaces"),
)
print(f"[cheap-service] build_cheap_app: redis storage db{storage_params['db']} bus db{bus_params['db']} "
f"auth={'on' if storage_params['password'] else 'off'} ws_basedir={ws_basedir}", flush=True)
return create_app(
storage=RedisStorage(**storage_params),
message_bus=RedisMessageBus(**bus_params),
workspace_manager=LocalWorkspaceManager(basedir=ws_basedir),
extra_agent_tools=_cheap_tools_factory,
extra_agent_middlewares=_cheap_middlewares_factory,
custom_subagent_templates=[],
custom_agent_cls=None,
title=title,
)
def main() -> None:
"""部署入口:起 uvicorn 跑便宜档 Agent Service(默认 0.0.0.0:8300,与 tier2 8200 错开)。只在 mini-desktop 跑。"""
import uvicorn # noqa: PLC0415
host = os.environ.get("CHEAP_SERVICE_HOST", "0.0.0.0")
port = int(os.environ.get("CHEAP_SERVICE_PORT", "8300"))
print(f"[cheap-service] 启动 Agent Service:http://{host}:{port}", flush=True)
# genconfig 热源接线(配置控制面阶段二·步骤4,与 tier2 service/app.py:main 对称):TIER2_GENCONFIG_NACOS=1
# 才启用,默认关 = 行为字节不变。工厂自带启用旗判断 + best-effort(失败只告警不抛、不 attach),订阅
# cheap 档 dataId(gen-hot-params-cheap)——与 tier2 各进程各热源,天然不串档(worker.genconfig 注释详解)。
from worker import genconfig_nacos # noqa: PLC0415 —— 惰性 import,默认关时不牵 nacos 依赖
genconfig_nacos.build_and_attach_hot_source(tier="cheap")
# 知识包激活生效链接线(W-CFG-KB K3,与 genconfig 热源对称):TIER2_KB_ROOT 未配 → no-op(默认关、read_file
# 字节不变);已配 → 启动期按激活版物化知识根 +(Nacos 启用时)订阅知识包 dataId,激活推送即重新物化。best-effort:
# 任何异常只告警、不阻断 Service 启动 / 生成。
from worker import kb_nacos # noqa: PLC0415 —— 惰性 import,默认关时不牵 nacos 依赖
kb_nacos.setup_knowledge_activation()
uvicorn.run(build_cheap_app(), host=host, port=port)
if __name__ == "__main__":
main()