feat(tier2): C1 轻量观测——控制面把 SSE 事件记成每局可读 trace

C1 观测的轻量兑现(创始人:没有完整的 C,B 迭代很慢)。不起 Node Studio(它在 6c6g 非 mini-desktop、
是 Node monorepo,重)——而是在控制面客户端侧顺手把逐帧消费的 SSE 结构性事件(跳过 *_DELTA 逐 token 噪声)
记成 per-run JSONL 时间线:模型调用/工具调用(带工具名)/工具结果/推理/reply 边界/卡死信号
(REQUIRE_USER_CONFIRM/REQUIRE_EXTERNAL_EXECUTION)。落工程 workdir/control-plane-trace.jsonl,跨 attempt
追加=整局多轮时间线。零服务改动、零 Studio 依赖、best-effort(写失败只告警不中断)。

迭代时一眼看清「agent 调了哪些工具、按什么序、卡在哪」——这正是这轮 debug 反复要的可见性。
observability/studio_sink.py 的 OTLP/Studio 全量 UI 路仍在、是将来可选(待 mini-desktop 起 Studio)。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
zizi 2026-06-24 11:01:06 +00:00
parent 34336d3204
commit b1217301c8

View File

@ -142,7 +142,8 @@ def _collect_on_disk_project(game_id: str) -> dict | None:
async def _wait_for_turn_end(base_url: str, agent_id: str, session_id: str, *,
user_id: str, timeout_s: float, idle_timeout_s: float) -> dict:
user_id: str, timeout_s: float, idle_timeout_s: float,
trace_path: Any = None, attempt: int = 0) -> dict:
"""消费 SSE GET /sessions/{session_id}/stream,等本回合 chat run 结束(读到 REPLY_END / EXCEED_MAX_ITERS)。
源码核验(agentscope/app/_router/_session.py:425-534):
@ -211,6 +212,9 @@ async def _wait_for_turn_end(base_url: str, agent_id: str, session_id: str, *,
print(f"[tier2-control] ⚠ SSE 帧解析失败(已跳过): {type(e).__name__}: {e}", flush=True)
continue
etype = event.get("type")
# C1 观测:把结构性事件(跳过 *_DELTA 逐 token 噪声)记进 per-run trace 时间线。
if etype and not etype.endswith("_DELTA"):
_trace_sse_event(trace_path, event, now - t_start, attempt)
# 只认本回合的结束事件(REPLY_END / EXCEED_MAX_ITERS);其余事件(模型流 / 工具调用)略过。
if etype in _TURN_END_EVENTS:
print(f"[tier2-control] SSE 读到本回合结束事件 type={etype}"
@ -224,6 +228,41 @@ async def _wait_for_turn_end(base_url: str, agent_id: str, session_id: str, *,
return {"ended": False, "reason": f"sse_error:{type(e).__name__}", "endEvent": None}
# ── C1 轻量观测:把 SSE 事件流记成每局可读 trace(客户端侧,零服务改动 / 零 Studio 依赖)──
# 迭代痛点 = 看不见 run 在哪卡住。控制面本就逐帧消费 SSE,顺手把【结构性事件】(跳过 *_DELTA 逐 token 噪声)
# 记成一份 per-run JSONL 时间线:模型调用 / 工具调用(带工具名)/ 工具结果 / 推理 / reply 边界 / 卡死信号
# (REQUIRE_USER_CONFIRM / REQUIRE_EXTERNAL_EXECUTION)。落在工程 workdir/control-plane-trace.jsonl,
# 跨 attempt 追加 = 整局多轮时间线。best-effort:写失败只告警、绝不中断主链。这是「完整的 C」里 C1 观测的
# 轻量兑现(无需起 Node Studio;Studio/OTLP 那条 observability.studio_sink 仍在、是将来全量 UI 的可选路)。
def _compact_sse_event(event: dict) -> dict:
"""把一个 SSE 事件压成 trace 行的精简形(只留有意义字段,长文本截断;绝不抛)。"""
etype = event.get("type")
out: dict = {"type": etype, "reply_id": event.get("reply_id")}
# 工具调用:记工具名(看 agent 调了哪些工具、按什么顺序——卡在哪一目了然)。
for k in ("tool_name", "name", "id"):
if event.get(k):
out[k] = event.get(k)
# 截断可能的长文本字段(text/delta/output/content 的字符串形),只留前 200 字。
for k in ("text", "output", "content", "reason", "error"):
v = event.get(k)
if isinstance(v, str) and v:
out[k] = v[:200]
return out
def _trace_sse_event(trace_path: Any, event: dict, elapsed_s: float, attempt: int) -> None:
"""把一个结构性 SSE 事件追加进 per-run trace JSONL(best-effort,绝不抛)。"""
if not trace_path:
return
try:
import json as _json # noqa: PLC0415
rec = {"t": round(elapsed_s, 1), "attempt": attempt, **_compact_sse_event(event)}
with open(trace_path, "a", encoding="utf-8") as f:
f.write(_json.dumps(rec, ensure_ascii=False) + "\n")
except Exception as e: # noqa: BLE001 —— 观测落盘失败绝不中断主链(C1 observe-only)
print(f"[tier2-control][trace] ⚠ trace 落盘失败(忽略):{type(e).__name__}: {e}", flush=True)
# ── 续跑指令(等价 CLI run_studio:454-458 的「停了但门没绿,按失败门继续修、门绿再 finish」)──
def _resume_feedback_text(verdict_feedback: str) -> str:
"""把失败门反馈包成续跑指令(口径对齐 studio.py:454-458)。"""
@ -307,6 +346,13 @@ async def drive_generation(
gate_game_id = session_id
print(f"[tier2-control] game={game_id} 启动: agent={agent_id} session={session_id} "
f"(评门 game_id=session_id;max_resumes={max_resumes})", flush=True)
# C1 观测:本局 SSE 事件 trace 落工程 workdir/control-plane-trace.jsonl(跨 attempt 追加,可读时间线)。
trace_path = None
try:
trace_path = str(_run._workdir(gate_game_id) / "control-plane-trace.jsonl")
print(f"[tier2-control] 观测 trace → {trace_path}", flush=True)
except Exception: # noqa: BLE001 —— 取 trace 路径失败不影响主链(观测 best-effort)
trace_path = None
if not session_id or not agent_id:
result["stopped_reason"] = "start_missing_ids(start_new_game 未回 agent_id/session_id)"
result["wall_s"] = round(time.perf_counter() - t0, 1)
@ -319,7 +365,8 @@ async def drive_generation(
# a. 等本回合 chat run 结束(SSE;超时/断流据 on-disk 继续,绝不卡死)。
await _wait_for_turn_end(
base_url, agent_id, session_id, user_id=user_id,
timeout_s=sse_turn_timeout_s, idle_timeout_s=sse_idle_timeout_s)
timeout_s=sse_turn_timeout_s, idle_timeout_s=sse_idle_timeout_s,
trace_path=trace_path, attempt=attempt)
# b. 独立评门:run.run_gates 纯代码判 on-disk 源工程(零 LLM、确定性、绝不自评翻绿)。
try: