JsonlFileSink 把五字段 trace 落 workdir/trace.jsonl(非阻塞·写失败不阻断生成);studio.py sink 接线;SAA 扩展段 schema(接口对称五核心+内容不对称 ext);管理面 3 只读端点 roles/traces/cost(惰性 import 保住「仅 import service.app 不牵 fastapi」红线)。测试 test_jsonl_sink 3 / test_admin_routes 6 全绿。 Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
520 lines
26 KiB
Python
520 lines
26 KiB
Python
"""observability/trace.py —— tier2 富游戏自治线 · 统一 trace adapter(H1/H2 观测管道)。
|
|
|
|
职责(对 docs/architecture/架构/生成引擎/tier2细节图说-H-观测与成本.md 图 H1/H2):
|
|
订阅 AgentScope 官方 typed Event System(reply_stream() 吐的 AgentEvent 流)→ 映射成统一 trace 形状。
|
|
tier2 只「订阅 + 映射」,不「埋点 + 采集」—— 事件流是框架现成件(源码 agentscope/event/_event.py 已核)。
|
|
|
|
统一 trace 形状(H1 接口对称、内容不对称):
|
|
- 公共核心子集(两条生成线都必填、字段同名同义):
|
|
traceId —— 一次生成的轨迹主键(贯穿,对账/成本关联键,对接 verdict.evidence.traceId)。
|
|
step —— 第几步(tier2=ReAct 第几轮 reply / 第几次工具调用;SAA=第几节点)。
|
|
cost —— 这一步折成¥的成本(本 adapter 仅在 ModelCallEnd 处带 token,折¥由 cost.py 事后做;
|
|
trace step 先填 token 量,cost_rmb 由编排器关联 cost.py 回填,保持 best-effort)。
|
|
verdict —— 这一步/这道门的裁决(tier2 在 ToolResultEnd / ReplyEnd 处可带 state;门裁由 judge 回填)。
|
|
timestamp —— 发生时刻(ISO8601,取事件 created_at)。
|
|
- tier2 扩展段(扩展段是一个 JSON 列,tier2 塞自己独有的「推理/动作/观察」三段):
|
|
reasoning —— ThinkingBlock* / TextBlock*(reason 段)。
|
|
action —— ToolCall*(act 段,带工具名/参数片段)。
|
|
observation —— ToolResult*(observe 段,带工具结果状态)。
|
|
raw —— 原始事件类型 + 关键字段(调试/反查用)。
|
|
|
|
best-effort 铁律(H1/H2 黄带):轨迹写失败默认 **不阻塞主生成流程,但落一条告警** —— 不能让一次落库抖动
|
|
把整局已经跑出来的生成废掉,也不能让它无声丢失。本 adapter 的 ingest() / 落 sink 全包 try,异常只告警不抛。
|
|
|
|
contracts/trace 现行边界:统一 trace 契约的正式 schema 落处 = contracts/trace/(additive 立位,**当前还没建**,
|
|
随 spike 或控制面 phase-1 才新立)。本 adapter 产出的就是「按那份契约该写进去的形状」;契约 schema 一旦落地,
|
|
本 adapter 的字段即对它对齐(字段名已按 H1 公共核心子集 + tier2 扩展段定,避免日后改名)。
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import json
|
|
from pathlib import Path
|
|
from typing import Any, Callable, Optional
|
|
|
|
|
|
# ── 公共核心子集字段名(常量,与 H1 钉死;契约 schema 落地后以这套名对齐)──
|
|
CORE_FIELDS = ("traceId", "step", "cost", "verdict", "timestamp")
|
|
|
|
|
|
class TraceStep(dict):
|
|
"""一条统一 trace 记录(dict 子类:既是结构化对象、又天然 JSON 可序列化落 sink)。
|
|
|
|
形状 = 公共核心子集五字段 + 扩展段 ext(tier2 推理/动作/观察)。用 dict 子类而非 dataclass:
|
|
落 sink(MySQL/对象存储/JSONL)时直接 json.dumps,零转换;也便于 best-effort 容错(缺字段不崩)。
|
|
"""
|
|
|
|
def __init__(
|
|
self,
|
|
trace_id: str,
|
|
step: int,
|
|
timestamp: str,
|
|
*,
|
|
cost: Optional[dict] = None,
|
|
verdict: Optional[str] = None,
|
|
ext: Optional[dict] = None,
|
|
) -> None:
|
|
super().__init__(
|
|
traceId=trace_id,
|
|
step=step,
|
|
cost=cost, # {tokens?: {in,out,cached}, cost_rmb?: float};token 在此填,¥由 cost.py 回填
|
|
verdict=verdict, # 这一步/门的裁决(可空;门裁由 judge 回填)
|
|
timestamp=timestamp,
|
|
ext=ext or {}, # tier2 扩展段:reasoning/action/observation/raw
|
|
)
|
|
|
|
|
|
# ── 事件类型名 → ReAct 三段的映射(以类名字符串判,不强依赖 import 具体事件类,降级友好)──
|
|
# reason 段:思考/文本块的边界与增量。
|
|
_REASONING_EVENTS = {
|
|
"ThinkingBlockStartEvent",
|
|
"ThinkingBlockDeltaEvent",
|
|
"ThinkingBlockEndEvent",
|
|
"TextBlockStartEvent",
|
|
"TextBlockDeltaEvent",
|
|
"TextBlockEndEvent",
|
|
}
|
|
# act 段:工具调用。
|
|
_ACTION_EVENTS = {
|
|
"ToolCallStartEvent",
|
|
"ToolCallDeltaEvent",
|
|
"ToolCallEndEvent",
|
|
}
|
|
# observe 段:工具结果。
|
|
_OBSERVATION_EVENTS = {
|
|
"ToolResultStartEvent",
|
|
"ToolResultTextDeltaEvent",
|
|
"ToolResultDataDeltaEvent",
|
|
"ToolResultEndEvent",
|
|
}
|
|
|
|
# ── tier2 自合成「钩子相位标记」事件类型(非 AgentScope 原生事件)──
|
|
# 由 U3 给 Tier2TraceMiddleware 补的三道 onion 钩子(on_reasoning/on_acting/on_model_call)产出。
|
|
# 为什么要它:on_reply 看到的是逐个 block 级事件(ThinkingBlockDelta/ToolCallStart…),粒度细但「不分相」;
|
|
# 而 on_reasoning/on_acting/on_model_call 各自只裹住一个**相位**的进出边界(图说 A4 同心环),
|
|
# 能产出比逐事件视图更结构化的「相位段」——一次推理相位、一次单工具 I/O(带工具名/耗时)、一次裸模型调用
|
|
# (带实测 token 用量)。这三类标记走与原生事件**同一条 ingest 路**(不另起 adapter 方法),由下方
|
|
# to_trace_step 识别后落进 ext 的对应段。形状是轻量 dict(含 type + 关键字段),best-effort、零侵入事件流。
|
|
PHASE_REASONING = "Tier2ReasoningPhaseEvent" # on_reasoning 相位进/出(推理段边界)
|
|
PHASE_ACTING = "Tier2ActingPhaseEvent" # on_acting 单工具 I/O 进/出(动手段,带工具名/耗时)
|
|
PHASE_MODEL_CALL = "Tier2ModelCallPhaseEvent" # on_model_call 裸模型调用(模型调用段,带实测 token / ¥ 闸状态)
|
|
_PHASE_MARKER_EVENTS = {PHASE_REASONING, PHASE_ACTING, PHASE_MODEL_CALL}
|
|
|
|
|
|
def make_phase_marker(
|
|
phase_type: str,
|
|
boundary: str,
|
|
*,
|
|
timestamp: str = "",
|
|
fields: Optional[dict] = None,
|
|
) -> dict:
|
|
"""构造一条 tier2 相位标记事件(轻量 dict,喂给 TraceAdapter.ingest,走原生同一条路)。
|
|
|
|
:param phase_type: PHASE_REASONING / PHASE_ACTING / PHASE_MODEL_CALL 之一。
|
|
:param boundary: 'enter' | 'exit' | 'error'(相位进入 / 正常退出 / 异常退出)。
|
|
:param timestamp: 发生时刻(ISO8601;调用方一般传 datetime.now().isoformat(),缺则映射时记空串)。
|
|
:param fields: 该相位的关键字段(如 acting 的 tool/elapsedMs、model_call 的 tokens/¥闸状态),原样并入 ext 段。
|
|
:return: dict(含 type=phase_type + boundary + created_at + 余字段),由 to_trace_step 识别落 ext。
|
|
"""
|
|
evt: dict = {"type": phase_type, "boundary": boundary, "created_at": timestamp}
|
|
if fields:
|
|
evt.update(fields)
|
|
return evt
|
|
|
|
|
|
def _event_type_name(event: Any) -> str:
|
|
"""取事件类型名:优先 type(event).__name__,兜底事件 .type 字段(EventType 枚举值)。"""
|
|
name = type(event).__name__
|
|
if name in ("dict",) or name == "object":
|
|
# 兜底:若传进来的是已序列化的 dict(如从 SSE 流解出),按其 type 字段。
|
|
t = event.get("type") if isinstance(event, dict) else None
|
|
return str(t) if t else name
|
|
return name
|
|
|
|
|
|
def _get(event: Any, name: str, default: Any = None) -> Any:
|
|
"""兼容对象属性 / dict 键取值(事件可能是 EventBase 对象,也可能是已序列化 dict)。"""
|
|
if isinstance(event, dict):
|
|
return event.get(name, default)
|
|
return getattr(event, name, default)
|
|
|
|
|
|
def to_trace_step(event: Any, trace_id: str, step: int) -> TraceStep:
|
|
"""把单个 AgentScope AgentEvent 映射成一条统一 trace 形状(纯函数,无副作用)。
|
|
|
|
映射规则(对 H2 七类强类型事件):
|
|
- ModelCallEndEvent:End 带 token 用量(input_tokens/output_tokens)→ 填 cost.tokens(¥由 cost.py 回填)。
|
|
- ThinkingBlock*/TextBlock* → ext.reasoning(reason 段)。
|
|
- ToolCall*(name) → ext.action(act 段)。
|
|
- ToolResult*(state) → ext.observation(observe 段);ToolResultEnd 的 state 同时填 verdict。
|
|
- ReplyStart/End → 一轮回复边界(框住一步完整 reason→act→observe)。
|
|
|
|
:param event: 一个 AgentEvent(对象或已序列化 dict)。
|
|
:param trace_id: 本次生成的 traceId(贯穿)。
|
|
:param step: 当前步序(由 TraceAdapter 维护并递增)。
|
|
:return: TraceStep(公共核心子集 + tier2 扩展段)。
|
|
"""
|
|
etype = _event_type_name(event)
|
|
ts = _get(event, "created_at") or ""
|
|
|
|
cost: Optional[dict] = None
|
|
verdict: Optional[str] = None
|
|
ext: dict = {"raw": {"event": etype}}
|
|
|
|
# ── ModelCallEnd:带 token 用量(H3 成本台账从这里抓 token)──
|
|
if etype == "ModelCallEndEvent":
|
|
cost = {
|
|
"tokens": {
|
|
"in": int(_get(event, "input_tokens", 0) or 0),
|
|
"out": int(_get(event, "output_tokens", 0) or 0),
|
|
}
|
|
# cost_rmb 不在此填:trace 只记 token,折¥由 cost.py 按 new-api quota 事后回填(权威口径)。
|
|
}
|
|
ext["raw"]["model_call_end"] = True
|
|
|
|
elif etype == "ModelCallStartEvent":
|
|
ext["raw"]["model_name"] = _get(event, "model_name")
|
|
|
|
# ── reason 段:思考 / 文本块 ──
|
|
elif etype in _REASONING_EVENTS:
|
|
seg: dict = {"phase": "reasoning", "kind": etype}
|
|
delta = _get(event, "delta")
|
|
if delta is not None:
|
|
seg["delta"] = delta # 增量文本/思考片段(落 sink 时可截断;此处保真)
|
|
ext["reasoning"] = seg
|
|
|
|
# ── act 段:工具调用 ──
|
|
elif etype in _ACTION_EVENTS:
|
|
seg = {"phase": "action", "kind": etype}
|
|
tool_name = _get(event, "tool_call_name")
|
|
if tool_name is not None:
|
|
seg["tool"] = tool_name
|
|
tool_id = _get(event, "tool_call_id")
|
|
if tool_id is not None:
|
|
seg["toolCallId"] = tool_id
|
|
delta = _get(event, "delta")
|
|
if delta is not None:
|
|
seg["argsDelta"] = delta # 工具参数 JSON 片段(增量)
|
|
ext["action"] = seg
|
|
|
|
# ── observe 段:工具结果 ──
|
|
elif etype in _OBSERVATION_EVENTS:
|
|
seg = {"phase": "observation", "kind": etype}
|
|
tool_id = _get(event, "tool_call_id")
|
|
if tool_id is not None:
|
|
seg["toolCallId"] = tool_id
|
|
tool_name = _get(event, "tool_call_name")
|
|
if tool_name is not None:
|
|
seg["tool"] = tool_name
|
|
# ToolResultEnd 带最终执行状态(ToolResultState;use_enum_values → 多为字符串值)。
|
|
state = _get(event, "state")
|
|
if state is not None:
|
|
state_val = getattr(state, "value", state)
|
|
seg["state"] = str(state_val)
|
|
# 工具结果状态作为这一步的 verdict(成功/失败/拒绝/中断),门裁另由 judge 回填。
|
|
verdict = str(state_val)
|
|
ext["observation"] = seg
|
|
|
|
# ── 回复边界:框住一步完整 reason→act→observe ──
|
|
elif etype in ("ReplyStartEvent", "ReplyEndEvent"):
|
|
ext["raw"]["reply_boundary"] = etype
|
|
ext["raw"]["reply_id"] = _get(event, "reply_id")
|
|
if etype == "ReplyStartEvent":
|
|
ext["raw"]["agent_name"] = _get(event, "name")
|
|
|
|
# ── tier2 相位标记(U3 三道钩子产):比逐事件更结构化的「相位段」边界 ──
|
|
elif etype in _PHASE_MARKER_EVENTS:
|
|
boundary = _get(event, "boundary", "enter")
|
|
ext["raw"]["phaseBoundary"] = boundary
|
|
if etype == PHASE_REASONING:
|
|
# 推理相位进/出(on_reasoning 裹的推理+模型决策那一层的边界)→ 归 reasoning 段。
|
|
ext["reasoning"] = {"phase": "reasoning", "kind": etype, "boundary": boundary}
|
|
elif etype == PHASE_ACTING:
|
|
# 单次工具 I/O 进/出(on_acting 只裹 toolkit.call_tool 那层)→ 归 action 段,带工具名/耗时/状态。
|
|
seg = {"phase": "action", "kind": etype, "boundary": boundary}
|
|
for k in ("tool", "toolCallId", "elapsedMs", "state"):
|
|
v = _get(event, k)
|
|
if v is not None:
|
|
seg[k] = v
|
|
ext["action"] = seg
|
|
# acting 退出时若带最终工具状态,顺带填这一步 verdict(与 ToolResultEnd 同义,便于按相位反查)。
|
|
st = _get(event, "state")
|
|
if boundary != "enter" and st is not None:
|
|
verdict = str(st)
|
|
elif etype == PHASE_MODEL_CALL:
|
|
# 裸模型调用相位(on_model_call)——「钱算在哪的地方」:退出时带实测 token 用量 + ¥ 累进闸状态。
|
|
ext["raw"]["model_call_phase"] = True
|
|
in_tok = _get(event, "input_tokens")
|
|
out_tok = _get(event, "output_tokens")
|
|
if in_tok is not None or out_tok is not None:
|
|
# 与 ModelCallEndEvent 同口径填 cost.tokens(¥由 cost.py 事后按 new-api quota 回填权威值);
|
|
# 注意:这是相位标记带的「实测」token,与 on_reply 看到的 ModelCallEndEvent 同源,
|
|
# 收口聚合成本时应择一计、勿重复累加(本标记主供细粒度审计/¥闸状态留痕)。
|
|
cost = {"tokens": {"in": int(in_tok or 0), "out": int(out_tok or 0)}}
|
|
# ¥ 累进硬熔断的当次留痕(花了多少、闸限多少、是否回落次数闸),供对账「¥闸真按金额拦」。
|
|
for k in ("model_name", "estRmb", "spentRmb", "limitRmb", "gateMode"):
|
|
v = _get(event, k)
|
|
if v is not None:
|
|
ext["raw"][k] = v
|
|
|
|
# ── 熔断 / HITL 等其余事件:原样记类型(反查用,不强解)──
|
|
else:
|
|
# ExceedMaxIters / RequireUserConfirm / Custom 等:落原始类型即可,供反查。
|
|
pass
|
|
|
|
return TraceStep(
|
|
trace_id=trace_id,
|
|
step=step,
|
|
timestamp=ts,
|
|
cost=cost,
|
|
verdict=verdict,
|
|
ext=ext,
|
|
)
|
|
|
|
|
|
class TraceAdapter:
|
|
"""统一 trace adapter:订阅 AgentScope Event System,把每个 AgentEvent 映射成统一 trace 形状并落 sink。
|
|
|
|
设计(H2「订阅 + 映射」非「埋点 + 采集」):
|
|
- tier2 这条线走官方 Event System(reply_stream() 吐 AgentEvent),adapter 只做映射,不自埋点。
|
|
- 公共核心子集对称、扩展段不对称:本 adapter 产出 tier2 那一轨;SAA 那一轨由 Java 线另一个 adapter 产。
|
|
- best-effort:ingest / 落 sink 全包 try,任何异常只告警不抛 —— 轨迹写失败不阻塞主生成流程。
|
|
|
|
用法(在 ReAct 循环外包一层):
|
|
adapter = TraceAdapter(trace_id="<gen-task-trace>", sink=my_jsonl_sink)
|
|
async for event in agent.reply_stream(user_msg):
|
|
adapter.ingest(event) # 映射 + 落 sink,best-effort
|
|
# ...(其余消费,如熔断 middleware 已在 agent 内处理)...
|
|
steps = adapter.steps # 内存里的 trace 列表(也可只走 sink、不留内存)
|
|
|
|
sink 契约:sink(trace_step: dict) -> None。可为写 JSONL / 入库 / 推 message_bus 等;
|
|
sink 内部异常由 adapter best-effort 兜住(告警,不阻断)。sink=None 时只留内存 steps。
|
|
"""
|
|
|
|
def __init__(
|
|
self,
|
|
trace_id: str,
|
|
*,
|
|
sink: Optional[Callable[[dict], None]] = None,
|
|
keep_in_memory: bool = True,
|
|
) -> None:
|
|
"""
|
|
Args:
|
|
trace_id: 本次生成的 traceId(贯穿;对接 verdict.evidence.traceId / 成本关联键)。
|
|
sink: 落库回调 sink(trace_step_dict);None → 只留内存。
|
|
keep_in_memory: 是否在 self.steps 留一份内存副本(长 run 可置 False 只走 sink 省内存)。
|
|
"""
|
|
self.trace_id = trace_id
|
|
self._sink = sink
|
|
self._keep = keep_in_memory
|
|
self.steps: list[TraceStep] = []
|
|
self._step = 0 # 步序计数器,每 ingest 一条递增
|
|
self._dropped = 0 # best-effort 丢弃计数(供收口告警/审计)
|
|
|
|
def ingest(self, event: Any) -> Optional[TraceStep]:
|
|
"""吃一个 AgentEvent → 映射成 TraceStep → 落 sink(全程 best-effort,不抛)。
|
|
|
|
:return: 映射出的 TraceStep(成功);best-effort 失败时返回 None。
|
|
"""
|
|
try:
|
|
step = self._step
|
|
self._step += 1
|
|
trace_step = to_trace_step(event, self.trace_id, step)
|
|
if self._keep:
|
|
self.steps.append(trace_step)
|
|
if self._sink is not None:
|
|
self._emit_to_sink(trace_step)
|
|
return trace_step
|
|
except Exception as exc: # best-effort:映射/记录任何异常只告警不抛
|
|
self._dropped += 1
|
|
print(
|
|
f"[tier2-trace] ingest 失败(best-effort 丢弃,traceId={self.trace_id}):{exc}",
|
|
flush=True,
|
|
)
|
|
return None
|
|
|
|
def _emit_to_sink(self, trace_step: TraceStep) -> None:
|
|
"""落 sink,sink 内部异常单独兜住(一条写失败不连累后续,只计告警)。"""
|
|
try:
|
|
self._sink(dict(trace_step)) # 传纯 dict 给 sink,解耦 TraceStep 类型
|
|
except Exception as exc:
|
|
self._dropped += 1
|
|
print(
|
|
f"[tier2-trace] sink 写失败(best-effort 计告警,traceId={self.trace_id}):{exc}",
|
|
flush=True,
|
|
)
|
|
|
|
async def consume_stream(self, event_stream: Any) -> list[TraceStep]:
|
|
"""便利件:直接消费一个异步事件流(如 agent.reply_stream(...))→ 全程 ingest。
|
|
|
|
注意:这会把流里所有 AgentEvent 都 ingest 进来。若调用方还要自己消费事件(如取最终 Msg),
|
|
请改用「在自己的 async for 里逐条 adapter.ingest(event)」而非本便利件(本件会吃光流)。
|
|
|
|
:return: 本次消费产出的 TraceStep 列表(keep_in_memory=True 时即 self.steps)。
|
|
"""
|
|
produced: list[TraceStep] = []
|
|
try:
|
|
async for event in event_stream:
|
|
ts = self.ingest(event)
|
|
if ts is not None:
|
|
produced.append(ts)
|
|
except Exception as exc: # best-effort:流本身异常也不连累(轨迹尽力收,主流程另由编排器兜)
|
|
print(
|
|
f"[tier2-trace] 事件流消费中断(best-effort,traceId={self.trace_id}):{exc}",
|
|
flush=True,
|
|
)
|
|
return produced
|
|
|
|
@property
|
|
def dropped(self) -> int:
|
|
"""best-effort 期间丢弃(映射或落 sink 失败)的条数,供收口告警/审计。"""
|
|
return self._dropped
|
|
|
|
def summary(self) -> dict:
|
|
"""收口摘要:本次 trace 的步数、丢弃数、traceId(供编排器记一行收口日志/对账)。"""
|
|
return {
|
|
"traceId": self.trace_id,
|
|
"steps": len(self.steps) if self._keep else self._step,
|
|
"dropped": self._dropped,
|
|
}
|
|
|
|
|
|
class JsonlFileSink:
|
|
"""把统一 trace 事件(TraceStep dict)追加写到 JSONL 文件的落盘 sink。
|
|
|
|
用法:
|
|
sink = JsonlFileSink(workdir / "trace.jsonl")
|
|
adapter = TraceAdapter(trace_id="...", sink=sink)
|
|
|
|
为什么要它:
|
|
- 五字段 schema(traceId / step / cost / verdict / timestamp)原样落盘,供成本对账、replay 用;
|
|
- 每条事件是独立 JSON 行(JSONL),追加写、无需整文件加锁,适合长跑生成;
|
|
- 落盘路径 = 生成产物 workdir / "trace.jsonl",与 game_id 一一对应,随产物一起留存。
|
|
|
|
best-effort 铁律(与 TraceAdapter 一致):
|
|
- 目录不存在则自动创建;
|
|
- 写失败只 print 告警,**绝不抛异常**、绝不阻断主生成流程(trace 是旁路观测)。
|
|
"""
|
|
|
|
def __init__(self, output_path: "Path | str") -> None:
|
|
"""
|
|
Args:
|
|
output_path: trace.jsonl 落盘路径(全路径);父目录不存在时在首次写前自动创建。
|
|
"""
|
|
self._path = Path(output_path)
|
|
# 预建父目录(best-effort;目录已存在不报错)。
|
|
try:
|
|
self._path.parent.mkdir(parents=True, exist_ok=True)
|
|
except Exception as exc:
|
|
# 预建失败(如路径不可写)只告警,每次 __call__ 时再尝试一次。
|
|
print(
|
|
f"[tier2-trace] 预建 trace.jsonl 目录失败(best-effort,写时再试):"
|
|
f"{self._path.parent} — {exc}",
|
|
flush=True,
|
|
)
|
|
|
|
def __call__(self, trace_step: dict) -> None:
|
|
"""把一条 TraceStep dict 序列化为 JSON 追加写一行到 trace.jsonl。
|
|
|
|
五字段(traceId / step / cost / verdict / timestamp)由 TraceAdapter.ingest 已填进
|
|
trace_step(TraceStep 是 dict 子类),此处原样序列化落盘,保留 ext 扩展段。
|
|
|
|
写失败(磁盘满 / 权限错等)只 print 告警,绝不抛 —— 一次落库失败不得废掉已成功的生成产物。
|
|
"""
|
|
try:
|
|
# 防御性再建一次目录(首次 __init__ 预建可能因时序失败)。
|
|
self._path.parent.mkdir(parents=True, exist_ok=True)
|
|
line = json.dumps(trace_step, ensure_ascii=False)
|
|
with self._path.open("a", encoding="utf-8") as f:
|
|
f.write(line + "\n")
|
|
except Exception as exc:
|
|
print(
|
|
f"[tier2-trace] trace.jsonl 写失败(best-effort,不阻断生成,path={self._path}):{exc}",
|
|
flush=True,
|
|
)
|
|
|
|
|
|
def make_jsonl_sink(output_path: "Path | str") -> "JsonlFileSink":
|
|
"""构造 JSONL 落盘 sink,目标路径 output_path;可直接塞进 TraceAdapter(sink=...)。
|
|
|
|
用法示例(在编排器 run_studio 里):
|
|
sink = make_jsonl_sink(run._workdir(game_id) / "trace.jsonl")
|
|
tracer = Tier2TraceMiddleware(trace_id=game_id, sink=sink)
|
|
|
|
:param output_path: trace.jsonl 落盘路径(全路径)。
|
|
:return: JsonlFileSink 实例(Callable[[dict], None])。
|
|
"""
|
|
return JsonlFileSink(output_path)
|
|
|
|
|
|
def with_studio(
|
|
trace_id: str,
|
|
*,
|
|
sink: Optional[Callable[[dict], None]] = None,
|
|
keep_in_memory: bool = True,
|
|
studio_url: Optional[str] = None,
|
|
) -> "TraceAdapter":
|
|
"""便利件:建一个把每步既落原 sink、又旁路推 AgentScope Studio 的 TraceAdapter(C1 观测接线)。
|
|
|
|
为什么是「便利件」而不改 TraceAdapter 构造:TraceAdapter 的 sink 槽是单 sink 契约(sink(dict)->None),
|
|
本件不破坏它 —— 而是把「原 sink」与「Studio sink」用一个 _fanout 串成一个 sink 再塞进去(对调用方透明)。
|
|
Studio sink 默认关闭(取不到 Studio 地址即 no-op,见 studio_sink.make_studio_sink),所以:
|
|
- 没开 Studio(无 TIER2_STUDIO_URL / infra_config[studio]):行为 == 直接 TraceAdapter(sink=sink),零影响;
|
|
- 开了 Studio:每步在原 sink 之外多推一份 OTLP span 给 Studio(observe-only,best-effort,不阻塞主链)。
|
|
|
|
studio_sink 在函数内**惰性 import**:6c6g 无 opentelemetry 依赖时本模块顶层仍可裸 import(只有真用到本件、
|
|
且 Studio 地址可达时才触发 otel 依赖)。Studio sink 的 .shutdown() 由编排器在 run 收口处调
|
|
(经返回 adapter 的 _studio_sink 属性拿到,或调用方自己持有 sink 调 shutdown)。
|
|
|
|
:param trace_id: 本次生成 traceId(贯穿)。
|
|
:param sink: 原有落库 sink(JSONL/入库/message_bus 等);None → 不另落,仅(可能)推 Studio + 内存。
|
|
:param keep_in_memory: 是否在 adapter.steps 留内存副本(透传 TraceAdapter)。
|
|
:param studio_url: 显式 Studio HTTP base(一般留空,走 env / infra_config)。
|
|
:return: 配好 fanout sink 的 TraceAdapter;其 ._studio_sink 属性 = Studio sink(供收口调 shutdown)。
|
|
"""
|
|
# 惰性 import:把 otel 依赖关进本函数,保证 trace.py 顶层 6c6g 仍可裸 import。
|
|
from .studio_sink import make_studio_sink # noqa: PLC0415
|
|
|
|
studio_sink = make_studio_sink(trace_id, studio_url=studio_url)
|
|
|
|
def _fanout(trace_step: dict) -> None:
|
|
"""把一条 trace_step 同时喂给原 sink 与 Studio sink;**两者各自 best-effort,互不连累**。
|
|
|
|
注意:本 fanout 仍被 TraceAdapter._emit_to_sink 的外层 try 兜一次;但为「原 sink 失败不连累推 Studio」、
|
|
「推 Studio 失败不连累原 sink」,这里再各包一层 —— 任一侧异常只告警、不抛、不阻断另一侧与主链(observe-only)。
|
|
"""
|
|
if sink is not None:
|
|
try:
|
|
sink(trace_step)
|
|
except Exception as exc: # 原 sink 失败:不连累 Studio 推送,也不抛(交外层计 dropped)
|
|
print(f"[tier2-trace] 原 sink 写失败(fanout 兜,traceId={trace_id}):{exc}", flush=True)
|
|
# Studio sink 自身已 best-effort(取不到地址 = no-op;发送失败只告警),这里再兜一层防御性 try。
|
|
try:
|
|
studio_sink(trace_step)
|
|
except Exception as exc: # 理论上 studio_sink 自己已兜;防御性再兜,绝不让观测旁路波及主链
|
|
print(f"[tier2-trace] Studio sink 旁路异常(忽略,traceId={trace_id}):{exc}", flush=True)
|
|
|
|
adapter = TraceAdapter(trace_id, sink=_fanout, keep_in_memory=keep_in_memory)
|
|
# 把 Studio sink 挂到 adapter 上,供编排器在 run 收口处 adapter._studio_sink.shutdown()(flush 残留 span)。
|
|
adapter._studio_sink = studio_sink # type: ignore[attr-defined]
|
|
return adapter
|
|
|
|
|
|
if __name__ == "__main__": # pragma: no cover —— 本地自检:用桩事件(dict 形态)跑映射链路
|
|
import json
|
|
|
|
# 桩:模拟已序列化的 AgentEvent dict(无 agentscope 也能跑映射纯逻辑)。
|
|
fake_events = [
|
|
{"type": "ReplyStartEvent", "created_at": "2026-06-23T10:00:00", "reply_id": "r1", "name": "writer"},
|
|
{"type": "ThinkingBlockDeltaEvent", "created_at": "2026-06-23T10:00:01", "delta": "先建 scene 树"},
|
|
{"type": "ToolCallStartEvent", "created_at": "2026-06-23T10:00:02", "tool_call_id": "t1", "tool_call_name": "write_source"},
|
|
{"type": "ToolResultEndEvent", "created_at": "2026-06-23T10:00:03", "tool_call_id": "t1", "tool_call_name": "write_source", "state": "success"},
|
|
{"type": "ModelCallEndEvent", "created_at": "2026-06-23T10:00:04", "input_tokens": 1200, "output_tokens": 800},
|
|
{"type": "ReplyEndEvent", "created_at": "2026-06-23T10:00:05", "reply_id": "r1"},
|
|
]
|
|
adapter = TraceAdapter(trace_id="demo-trace-001")
|
|
for e in fake_events:
|
|
adapter.ingest(e)
|
|
print(json.dumps([dict(s) for s in adapter.steps], ensure_ascii=False, indent=2))
|
|
print("summary:", json.dumps(adapter.summary(), ensure_ascii=False))
|