From b1217301c8b380f7af588ceabb375ae69de45fba Mon Sep 17 00:00:00 2001 From: zizi Date: Wed, 24 Jun 2026 11:01:06 +0000 Subject: [PATCH] =?UTF-8?q?feat(tier2):=20C1=20=E8=BD=BB=E9=87=8F=E8=A7=82?= =?UTF-8?q?=E6=B5=8B=E2=80=94=E2=80=94=E6=8E=A7=E5=88=B6=E9=9D=A2=E6=8A=8A?= =?UTF-8?q?=20SSE=20=E4=BA=8B=E4=BB=B6=E8=AE=B0=E6=88=90=E6=AF=8F=E5=B1=80?= =?UTF-8?q?=E5=8F=AF=E8=AF=BB=20trace?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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) --- tier2/gen-worker/service/control_plane.py | 51 ++++++++++++++++++++++- 1 file changed, 49 insertions(+), 2 deletions(-) diff --git a/tier2/gen-worker/service/control_plane.py b/tier2/gen-worker/service/control_plane.py index 7a8d8354..197eb4d2 100644 --- a/tier2/gen-worker/service/control_plane.py +++ b/tier2/gen-worker/service/control_plane.py @@ -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: