"""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-/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..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 。 _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-/turns.jsonl(单文件上限 + 轮转) # ══════════════════════════════════════════════════════════════════════════════ class CheapTurnsWriter: """把逐 turn 记录(dict)追加写到 turns.jsonl 的落盘件(单文件上限 → 轮转 turns..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..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-/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-/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