games-development-ai/cheap-worker/cheap_turns_sink.py
lili dc89701fcc feat(cheap-gen): W-AXIS 波1 真相层——turns 全文落盘+per-run 归档+evidence 清单化+failureLayer
治「真相层缺失」病根(历史上模型文本/工具返回/门反馈零落盘、失败 run 被同 gid 重跑
覆盖,归因只能建在聚合数字上):

- cheap_turns_sink.py:逐 turn 全文落 amgen-<gid>/turns.jsonl(模型文本/工具全参/
  工具返回/门续修反馈);observe-only middleware、best-effort,接线失败绝不阻断生成;
  CLI 与 Service 双路接线(cheap_run/cheap_service_app)
- per-run 归档:同 gid 重跑前工程整体归档 _amgen-archive/<gid>-<ts>/,失败现场不再被覆盖
- evidence 清单化清理:按白名单保留,跨 run 残留(旧截图/旧 verdict)开跑即清,防证据串局
- failureLayer 三层归因:run-summary 落 generation/gate/gameplay 失败层字段,
  聚合报表可分「生成没跑完/机械门挂/玩法层拒」

主会话亲验:两局真跑归档 179KB、turns 含门反馈全文;pytest 新增三个测试档全绿。

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-10 04:27:17 -07:00

538 lines
26 KiB
Python

"""cheap_turns_sink.py — 便宜档生成链路逐 turn 全文落盘(W-AXIS 波1 · 真相层 F0-a)。
【这份解决什么问题】
诊断档 §2.3 F0-a 坐实:trace.jsonl 只落 TraceStep 事件骨架——模型推理/回复文本被碎成逐 delta
分散在成千行里、工具返回文本(ToolResultTextDeltaEvent.delta)在 observability.trace.to_trace_step
的 observe 段被整段丢弃、门反馈(RepairMiddleware 注入 / CLI kick)根本不走 trace 事件。于是
「模型每轮写了什么、harness 每轮回了什么」磁盘上物理没有这份数据,历史结论只能建在聚合数字上。
【本模块补的真相层】
一个便宜档专属、observe-only 的中间件 CheapTurnsMiddleware,挂 on_reply 嗅探同一条 AgentScope
事件流,把逐 delta 事件【按一轮 ReAct(一次模型调用 + 其工具调用 + 工具返回)聚合成整段】,
连同【harness→模型的输入 / 门续修反馈】一起逐 turn 落 amgen-<gid>/turns.jsonl(人读全文、
可归档、可归因)。与 trace.jsonl / OTLP 三者并存互不影响:trace.jsonl 给成本对账机读、
OTLP 给 Tempo、turns.jsonl 给人读链路。
【两条 harness→模型输入都要收(两条生成路各有一条)】
· CLI 路(cheap_studio):外层 resume 循环每 attempt 重新 writer.reply(kick),门反馈就是那个 kick
(name="user")——在 on_reply 起点读 input_kwargs["inputs"] 收得到。
· Service 路(cheap_service_app + RepairMiddleware):单次 reply 内于 finish 点 mid-reply 注入
UserMsg(name="gate", ...),不在 inputs 里——每见新一轮 ModelCallStartEvent 就扫 agent.state.context
尾部新增的 name="gate" 消息收下(去重防重复落)。
【best-effort 铁律(与 collector / breaker / trace sink 同款)】
on_reply 只读事件、原样 yield;写盘 / 聚合 / 读 context 全程吞异常,只告警、绝不抛、绝不阻断生成主链。
默认开(真相层是本波地基);env CHEAP_TURNS_ENABLED=0 一键关(纯增量旁路,关开关即回,现有行为字节不变)。
【单文件上限 + 轮转 + 密钥脱敏】
长 run(如烧到 180K token 的死圈)turns.jsonl 会很大——超单文件上限即把当前文件轮转成
turns.<n>.jsonl 另起,防单文件撑爆。落盘前对所有文本做 new-api key / Bearer / token 类脱敏
(工具若回显凭据,不泄进盘)。
【惰性 import 红线】顶层零 agentscope 依赖:纯落盘/聚合/脱敏件(CheapTurnsWriter / _TurnAggregator /
_redact)在顶层,可脱离 agentscope 单测;真正 subclass MiddlewareBase 的中间件类只在
build_turns_middleware() 里惰性定义(镜像 cheap_service_app._get_collector_cls),对齐工厂惰性红线。
"""
from __future__ import annotations
import json
import os
import re
import threading
from datetime import datetime
from pathlib import Path
from typing import Any, Callable, Optional
# ── 开关 / 上限 env(全大写 CHEAP_ 前缀,与 genconfig / OTLP 风格一致)──
ENV_ENABLED = "CHEAP_TURNS_ENABLED" # 默认开;"0"/"false"/"no" 关(纯增量旁路)
ENV_MAX_BYTES = "CHEAP_TURNS_MAX_BYTES" # 单文件上限(字节),超即轮转;默认 32MB
# 单文件默认上限:32MB。单局 turns.jsonl 正常几百 KB~数 MB,死圈长 run 才可能逼近;超即轮转另起。
_DEFAULT_MAX_BYTES = 32 * 1024 * 1024
# 单字段文本上限(防单个工具返回把一行 JSON 撑成几 MB):工具返回按 read_file 200KB 截断上界给到 256KB,
# 模型文本/思考/参数各给 128KB——正常远够(保「全文」),只挡病态超大,超出留 …[截断 N 字节] 标记不静默丢。
_FIELD_TEXT_CAP = 256 * 1024
_MODEL_TEXT_CAP = 128 * 1024
# ── 事件类型分组(以类名字符串判,不强依赖 import 具体事件类;与 observability.trace 同口径)──
_REASONING_TEXT_EVENTS = {"TextBlockStartEvent", "TextBlockDeltaEvent", "TextBlockEndEvent"}
_THINKING_EVENTS = {"ThinkingBlockStartEvent", "ThinkingBlockDeltaEvent", "ThinkingBlockEndEvent"}
_TOOLCALL_EVENTS = {"ToolCallStartEvent", "ToolCallDeltaEvent", "ToolCallEndEvent"}
_TOOLRESULT_EVENTS = {
"ToolResultStartEvent", "ToolResultTextDeltaEvent", "ToolResultDataDeltaEvent", "ToolResultEndEvent",
}
# ── 一次性告警去重(同 tag 只 print 一次,避免长 run 刷屏)──
_WARNED: set[str] = set()
_WARN_LOCK = threading.Lock()
def _warn_once(tag: str, msg: str) -> None:
"""同一 tag 的告警只落一次(best-effort,绝不抛)。"""
with _WARN_LOCK:
if tag in _WARNED:
return
_WARNED.add(tag)
print(f"[cheap-turns] {msg}", flush=True)
def reset_warned_for_test() -> None:
"""仅测试用:清告警去重集(让各用例互不串)。"""
with _WARN_LOCK:
_WARNED.clear()
def turns_enabled() -> bool:
"""真相层落盘是否开(默认开;env CHEAP_TURNS_ENABLED ∈ {0,false,no,off} 关)。"""
v = os.environ.get(ENV_ENABLED)
if v is None:
return True
return v.strip().lower() not in ("0", "false", "no", "off", "")
def _max_bytes() -> int:
"""单文件上限(env CHEAP_TURNS_MAX_BYTES 覆盖;非法值回落默认 32MB)。"""
raw = os.environ.get(ENV_MAX_BYTES)
if raw and raw.strip():
try:
n = int(raw.strip())
if n > 0:
return n
except ValueError:
_warn_once("bad-max-bytes", f"CHEAP_TURNS_MAX_BYTES 非法({raw!r}),回落默认 {_DEFAULT_MAX_BYTES}")
return _DEFAULT_MAX_BYTES
def _now_iso() -> str:
"""当前时刻 ISO8601(落盘 ts;取不到返空串,best-effort)。"""
try:
return datetime.now().isoformat(timespec="milliseconds")
except Exception: # noqa: BLE001
return ""
# ══════════════════════════════════════════════════════════════════════════════
# 密钥脱敏 —— 落盘前把 new-api 凭据 / Bearer / token 类字面量抹掉(工具回显凭据不泄进盘)
# ══════════════════════════════════════════════════════════════════════════════
_REDACTED = "***REDACTED***"
# sk- 前缀密钥(OpenAI 兼容 / new-api 常见形态)。
_RE_SK = re.compile(r"sk-[A-Za-z0-9_\-]{6,}")
# Bearer <token>。
_RE_BEARER = re.compile(r"(?i)(bearer\s+)[A-Za-z0-9._\-]{8,}")
# 带键名的赋值:api_key/apikey/token/secret/password/authorization = "值" 或 : 值。
_RE_LABELED = re.compile(
r'(?i)("?(?:api[_-]?key|apikey|access[_-]?token|auth(?:orization)?|token|secret|password)"?\s*[:=]\s*"?)'
r"([^\s\"',}]{6,})"
)
def _redact(text: Optional[str]) -> str:
"""把文本里的凭据类字面量脱敏(best-effort:任何异常回退原文,绝不因脱敏丢内容/抛)。
覆盖三类:sk- 前缀密钥、Bearer token、带键名的赋值;另外若进程 env 有 NEWAPI_KEY,
连它的字面量一起抹(最强一层——即便工具原样回显了当前生成用的真实凭据也不落盘)。
"""
if not text:
return text or ""
try:
out = str(text)
# 进程真实凭据字面量(最强,先抹):env 里那把 key 若出现在文本中,整条替换。
real_key = os.environ.get("NEWAPI_KEY")
if real_key and len(real_key) >= 8 and real_key in out:
out = out.replace(real_key, _REDACTED)
out = _RE_SK.sub(_REDACTED, out)
out = _RE_BEARER.sub(lambda m: m.group(1) + _REDACTED, out)
out = _RE_LABELED.sub(lambda m: m.group(1) + _REDACTED, out)
return out
except Exception: # noqa: BLE001 —— 脱敏失败宁可回退原文也不丢内容/不抛
return str(text)
def _cap(text: str, cap: int) -> str:
"""单字段文本上限:超出截断并留 …[截断 N 字节] 标记(不静默丢,保可诊断)。"""
if text is None:
return ""
if len(text) <= cap:
return text
dropped = len(text) - cap
return text[:cap] + f"\n…[截断 {dropped} 字节]"
# ══════════════════════════════════════════════════════════════════════════════
# CheapTurnsWriter —— 逐 turn 记录落 amgen-<gid>/turns.jsonl(单文件上限 + 轮转)
# ══════════════════════════════════════════════════════════════════════════════
class CheapTurnsWriter:
"""把逐 turn 记录(dict)追加写到 turns.jsonl 的落盘件(单文件上限 → 轮转 turns.<n>.jsonl)。
契约:write_record(rec: dict) -> None。JSONL 追加写、每条一行,写失败只告警绝不抛(best-effort)。
轮转:每次写前若现文件已达上限,把它 rename 成 turns.1.jsonl / turns.2.jsonl … 再从空文件续写,
防单文件被病态长 run 撑爆(归档时整个 amgen 目录连同轮转分片一起归档)。
"""
def __init__(self, output_path: "Path | str", *, max_bytes: Optional[int] = None) -> None:
self._path = Path(output_path)
self._max_bytes = max_bytes if (max_bytes and max_bytes > 0) else _max_bytes()
self._rotated = 0
self._lock = threading.Lock() # Service 长驻并发多局各自 writer,单 writer 内写串行保序
try:
self._path.parent.mkdir(parents=True, exist_ok=True)
except Exception as exc: # noqa: BLE001 —— 预建失败写时再试
_warn_once("mkdir", f"预建 turns.jsonl 目录失败(写时再试):{self._path.parent}{exc}")
def _maybe_rotate(self) -> None:
"""现文件达上限 → rename 成 turns.<n>.jsonl 另起(best-effort,失败只告警继续写原文件)。"""
try:
if self._path.exists() and self._path.stat().st_size >= self._max_bytes:
self._rotated += 1
rotated = self._path.with_name(f"{self._path.stem}.{self._rotated}{self._path.suffix}")
self._path.rename(rotated)
print(f"[cheap-turns] turns.jsonl 达上限({self._max_bytes} 字节)→ 轮转 {rotated.name}", flush=True)
except Exception as exc: # noqa: BLE001 —— 轮转失败不阻断写(继续写原文件)
_warn_once("rotate", f"turns.jsonl 轮转失败(继续写原文件):{exc}")
def write_record(self, rec: dict) -> None:
"""追加写一条 turn 记录(JSONL);写失败只告警绝不抛(observe-only 不咬生成)。"""
try:
with self._lock:
self._maybe_rotate()
self._path.parent.mkdir(parents=True, exist_ok=True)
line = json.dumps(rec, ensure_ascii=False)
with self._path.open("a", encoding="utf-8") as f:
f.write(line + "\n")
except Exception as exc: # noqa: BLE001 —— 一条写失败不连累后续、不阻断生成
_warn_once("write", f"turns.jsonl 写失败(best-effort,不阻断生成,path={self._path}):{exc}")
# ══════════════════════════════════════════════════════════════════════════════
# _TurnAggregator —— 把逐 delta 事件聚合成整段 turn(纯逻辑,脱 agentscope 单测)
# ══════════════════════════════════════════════════════════════════════════════
def _evt_type(evt: Any) -> str:
"""事件类型名:优先 type().__name__,dict 形态兜底取 'type' 键(与 observability.trace 同口径)。"""
name = type(evt).__name__
if name == "dict":
t = evt.get("type") if isinstance(evt, dict) else None
return str(t) if t else name
return name
def _evt_get(evt: Any, key: str, default: Any = None) -> Any:
"""兼容对象属性 / dict 键取值(事件可能是 typed 对象,也可能是已序列化 dict)。"""
if isinstance(evt, dict):
return evt.get(key, default)
return getattr(evt, key, default)
def _msg_name(msg: Any) -> str:
"""取 Msg.name(区分 harness kick=user / 门续修=gate);取不到返空串。"""
return str(_evt_get(msg, "name", "") or "")
def _msg_text(msg: Any) -> str:
"""把一个 Msg 的 content 抽成纯文本:str 直取;list(块)拼各块 text;取不到返空串(best-effort)。"""
try:
content = _evt_get(msg, "content", None)
if content is None:
return ""
if isinstance(content, str):
return content
if isinstance(content, list):
parts = []
for blk in content:
if isinstance(blk, dict):
if blk.get("type") in (None, "text") and blk.get("text"):
parts.append(str(blk.get("text")))
else:
t = getattr(blk, "text", None)
if t:
parts.append(str(t))
return "\n".join(parts)
return str(content)
except Exception: # noqa: BLE001
return ""
def _as_msg_list(inputs: Any) -> list:
"""把 on_reply input_kwargs['inputs'] 归一成 Msg 列表(单条/列表/None 都吃)。"""
if inputs is None:
return []
if isinstance(inputs, (list, tuple)):
return list(inputs)
return [inputs]
class _TurnAggregator:
"""把一条 reply 的逐事件流聚合成「逐 turn 整段记录」并经 emit 落盘(纯逻辑,可脱 agentscope 单测)。
turn 边界 = 一次模型调用:ModelCallStartEvent 起新 turn(先冲上一 turn),其间累积模型思考/文本/
工具调用参数,ModelCallEnd 记 token;工具返回(在 ModelCallEnd 之后、下一 ModelCallStart 之前的
acting 段)累积进【当前仍未冲的】这一 turn;ReplyEndEvent 冲最后一 turn。
harness→模型输入(kick / 门续修反馈)另走 note_reply_inputs / scan_context_gate,按文件追加序即时落
(内容 hash 去重防重复)。所有记录带统一信封 {gid, traceId, seq, ts, kind, ...}。
"""
def __init__(self, gid: str, trace_id: str, emit: Callable[[dict], None]) -> None:
self._gid = gid
self._trace_id = trace_id
self._emit = emit
self._seq = 0
self._turn_idx = 0
self._cur: Optional[dict] = None
self._seen_input: set[int] = set() # 已落 harness/gate 输入的内容 hash(去重)
# ── 统一信封 + 落盘 ──
def _write(self, kind: str, body: dict) -> None:
rec = {"gid": self._gid, "traceId": self._trace_id, "seq": self._seq,
"ts": _now_iso(), "kind": kind}
rec.update(body)
self._seq += 1
try:
self._emit(rec)
except Exception as exc: # noqa: BLE001 —— emit(落盘)失败已在 writer 内兜;此处再兜一层防御
_warn_once("emit", f"turn 记录 emit 失败(best-effort):{exc}")
# ── harness→模型输入(两条路)──
def note_reply_inputs(self, inputs: Any) -> None:
"""on_reply 起点:把触发本次 reply 的输入(CLI kick / Service 首个 kick)落成 harness_input/gate_feedback。"""
for msg in _as_msg_list(inputs):
name = _msg_name(msg)
text = _msg_text(msg)
kind = "gate_feedback" if name == "gate" else "harness_input"
self._emit_input(kind, name, text)
def scan_context_gate(self, context: Any) -> None:
"""每见新一轮 ModelCallStart 扫 agent.state.context 尾部的 name='gate' 续修注入(Service mid-reply 路)。"""
try:
for msg in (context or []):
if _msg_name(msg) == "gate":
self._emit_input("gate_feedback", "gate", _msg_text(msg))
except Exception: # noqa: BLE001 —— 读 context best-effort,失败忽略
pass
def _emit_input(self, kind: str, name: str, text: str) -> None:
text = (text or "").strip()
if not text:
return
h = hash((kind, text))
if h in self._seen_input:
return
self._seen_input.add(h)
self._write(kind, {"name": name, "text": _cap(_redact(text), _FIELD_TEXT_CAP)})
# ── 事件聚合 ──
def ingest_event(self, evt: Any) -> None:
"""吃一个 AgentEvent(对象/序列化 dict)→ 累积进当前 turn / 冲 turn(全 best-effort)。"""
try:
et = _evt_type(evt)
if et == "ModelCallStartEvent":
self._flush()
self._start_turn(_evt_get(evt, "model_name"))
return
if et in _THINKING_EVENTS:
delta = _evt_get(evt, "delta")
if delta:
self._ensure_turn()["thinking"] += str(delta)
return
if et in _REASONING_TEXT_EVENTS:
delta = _evt_get(evt, "delta")
if delta:
self._ensure_turn()["text"] += str(delta)
return
if et in _TOOLCALL_EVENTS:
self._on_tool_call(et, evt)
return
if et in _TOOLRESULT_EVENTS:
self._on_tool_result(et, evt)
return
if et == "ModelCallEndEvent":
cur = self._cur
if cur is not None:
cur["tokensIn"] = int(_evt_get(evt, "input_tokens", 0) or 0)
cur["tokensOut"] = int(_evt_get(evt, "output_tokens", 0) or 0)
return
if et in ("ReplyEndEvent",):
self._flush()
return
except Exception as exc: # noqa: BLE001 —— 聚合任何异常只告警,绝不抛、绝不咬生成
_warn_once("ingest", f"turn 事件聚合失败(best-effort 丢该事件):{exc}")
def _start_turn(self, model_name: Any) -> None:
self._turn_idx += 1
self._cur = {
"turn": self._turn_idx,
"modelName": str(model_name) if model_name else None,
"thinking": "", "text": "",
"toolCalls": {}, # tool_call_id -> {"name","args"}
"toolResults": {}, # tool_call_id -> {"name","text","state"}
"tokensIn": None, "tokensOut": None,
}
def _ensure_turn(self) -> dict:
"""无当前 turn(极少数事件早于首个 ModelCallStart)时惰性起一个,保证不漏内容。"""
if self._cur is None:
self._start_turn(None)
return self._cur # type: ignore[return-value]
def _on_tool_call(self, et: str, evt: Any) -> None:
tcid = str(_evt_get(evt, "tool_call_id", "") or "")
if not tcid:
return
cur = self._ensure_turn()
slot = cur["toolCalls"].setdefault(tcid, {"name": None, "args": ""})
name = _evt_get(evt, "tool_call_name")
if name and not slot["name"]:
slot["name"] = str(name)
delta = _evt_get(evt, "delta")
if delta:
slot["args"] += str(delta)
def _on_tool_result(self, et: str, evt: Any) -> None:
tcid = str(_evt_get(evt, "tool_call_id", "") or "")
if not tcid:
return
cur = self._ensure_turn()
slot = cur["toolResults"].setdefault(tcid, {"name": None, "text": "", "state": None})
name = _evt_get(evt, "tool_call_name")
if name and not slot["name"]:
slot["name"] = str(name)
if et == "ToolResultTextDeltaEvent":
delta = _evt_get(evt, "delta")
if delta:
slot["text"] += str(delta) # ← 补上 to_trace_step 丢掉的工具返回文本
elif et == "ToolResultEndEvent":
state = _evt_get(evt, "state")
if state is not None:
slot["state"] = str(getattr(state, "value", state))
def _flush(self) -> None:
"""冲当前 turn:有实质内容才落(纯空 turn 不落),把 toolCalls/toolResults dict 转有序 list + 脱敏 + 截断。"""
cur = self._cur
self._cur = None
if cur is None:
return
has_content = bool(
cur["thinking"].strip() or cur["text"].strip()
or cur["toolCalls"] or cur["toolResults"]
or cur["tokensIn"] or cur["tokensOut"]
)
if not has_content:
return
body: dict = {"turn": cur["turn"]}
if cur["modelName"]:
body["modelName"] = cur["modelName"]
if cur["thinking"].strip():
body["thinking"] = _cap(_redact(cur["thinking"]), _MODEL_TEXT_CAP)
if cur["text"].strip():
body["text"] = _cap(_redact(cur["text"]), _MODEL_TEXT_CAP)
if cur["toolCalls"]:
body["toolCalls"] = [
{"id": tcid, "name": v["name"], "args": _cap(_redact(v["args"]), _FIELD_TEXT_CAP)}
for tcid, v in cur["toolCalls"].items()
]
if cur["toolResults"]:
body["toolResults"] = [
{"id": tcid, "name": v["name"], "state": v["state"],
"text": _cap(_redact(v["text"]), _FIELD_TEXT_CAP)}
for tcid, v in cur["toolResults"].items()
]
if cur["tokensIn"] is not None or cur["tokensOut"] is not None:
body["tokens"] = {"in": cur["tokensIn"] or 0, "out": cur["tokensOut"] or 0}
self._write("model_turn", body)
def flush(self) -> None:
"""对外冲当前 turn(中间件在 ModelCallStart 处先冲上一 turn、再落门反馈、再起新 turn,保时间序)。"""
self._flush()
def close(self) -> None:
"""收口:冲最后一 turn(ReplyEnd 缺失/异常退出时兜底)。"""
self._flush()
# ══════════════════════════════════════════════════════════════════════════════
# 中间件工厂 —— 惰性 subclass MiddlewareBase(顶层不 import agentscope,对齐工厂惰性红线)
# ══════════════════════════════════════════════════════════════════════════════
_TURNS_MW_CLS = None
def _get_turns_middleware_cls():
"""惰性定义并缓存 CheapTurnsMiddleware(subclass MiddlewareBase;顶层不 import agentscope 保惰性红线)。"""
global _TURNS_MW_CLS
if _TURNS_MW_CLS is not None:
return _TURNS_MW_CLS
from agentscope.middleware import MiddlewareBase # noqa: PLC0415
class CheapTurnsMiddleware(MiddlewareBase):
"""便宜档真相层中间件:只挂 on_reply,observe-only 嗅探事件流→逐 turn 全文落 turns.jsonl。
与 collector/breaker 同款「读事件、原样 yield」姿势,绝不拦截/修改事件、绝不抛。只 override
on_reply(框架按 override 检测只把本类接进 on_reply 洋葱,不参与 reasoning/acting/model_call)。
"""
def __init__(self, writer: CheapTurnsWriter, *, gid: str, trace_id: str) -> None:
self._agg = _TurnAggregator(gid, trace_id, writer.write_record)
self._seen_model_call = False # 见过首个 ModelCallStart(之后每轮扫 context 收 gate 注入)
async def on_reply(self, agent, input_kwargs, next_handler):
# ① reply 起点:落触发本次 reply 的输入(CLI 每 attempt 的 kick / Service 首个 kick)。
try:
self._agg.note_reply_inputs(
input_kwargs.get("inputs") if isinstance(input_kwargs, dict) else None)
except Exception: # noqa: BLE001 —— 输入落盘 best-effort,失败不影响事件透传
pass
try:
async for evt in next_handler(**input_kwargs):
# ② 每见新一轮 ModelCallStart:按时间序【先冲上一 turn → 再落门续修(它在时间上晚于上一 turn)
# → 再起新 turn】。故先 flush() 冲上一 turn,再扫 agent.state.context 收 Service 路 mid-reply
# 注入的 name=gate 续修(内容 hash 去重防重复),最后 ingest 起新 turn(其内部 flush 已 no-op)。
try:
if _evt_type(evt) == "ModelCallStartEvent":
self._agg.flush()
self._agg.scan_context_gate(getattr(getattr(agent, "state", None), "context", None))
self._agg.ingest_event(evt)
except Exception: # noqa: BLE001 —— 嗅探纯旁路,任何异常不连累事件透传
pass
yield evt # 原样透传(observe-only,不改不拦)
finally:
# 收口:冲最后一 turn(即便 ReplyEnd 缺失 / 异常退出也不漏最后一轮)。
try:
self._agg.close()
except Exception: # noqa: BLE001
pass
_TURNS_MW_CLS = CheapTurnsMiddleware
return _TURNS_MW_CLS
def build_turns_middleware(game_id: str, *, trace_id: Optional[str] = None,
writer: Optional[CheapTurnsWriter] = None):
"""构造便宜档真相层中间件(默认开;关或接线失败 → 返回 None,调用方据 None 不加入 middleware 列表)。
落盘路径 = amgen-<game_id>/turns.jsonl(与 trace.jsonl / evidence 同目录,随产物一起归档/定位)。
best-effort:关(CHEAP_TURNS_ENABLED=0)/ 建 writer 失败 / import 失败一律返回 None,绝不抛、绝不阻断生成。
:param game_id: 后端 gameId(产物目录键;与 scaffold/trace 一致)。
:param trace_id: 贯穿 traceId(缺省用 game_id,与 tracer.trace_id 同源)。
:param writer: 仅测试注入(内存/临时 writer);生产恒 None → 落 amgen-<gid>/turns.jsonl。
:return: CheapTurnsMiddleware 实例 或 None(关/失败)。
"""
if not turns_enabled():
return None
try:
if writer is None:
import cheap_run # noqa: PLC0415 —— 惰性 import(顶层不牵产物路径解析)
writer = CheapTurnsWriter(cheap_run.game_dir(game_id) / "turns.jsonl")
cls = _get_turns_middleware_cls()
return cls(writer, gid=str(game_id), trace_id=str(trace_id or game_id))
except Exception as exc: # noqa: BLE001 —— 接线失败降级 None(不落 turns,不阻断生成)
_warn_once("build", f"真相层 turns 中间件接线失败(降级不落 turns,不阻断生成):{type(exc).__name__}: {exc}")
return None