zizi bd8b2c80a6 feat(tier2): B门重机器并行落地——控制面有界resume + Studio观测 + MCP工具面 + 基线client/Prompt门
B 门(AgentScope 重机器)四单元并行建设,文件不相交,均过 6c6g 静态校验(py_compile +
惰性 import 安全 + 机器门)。真跑收敛/Studio/门验证属 mini-desktop,作为下一步 e2e。

B1 控制面有界 resume(service/control_plane.py · 皇冠):
  补上 service/app.py 与 bootstrap.py 文件头标的 followup——CLI run_studio:404-458 的外层有界
  resume(M3 看一次 verdict 就停 → 带反馈踹回去续修)此前在服务态(create_app :8200)缺失。
  drive_generation() 复用 bootstrap REST 原语 + run 纯代码门:每回合消费 SSE(REPLY_END/
  EXCEED_MAX_ITERS)等 chat run 结束 → 独立跑 run.run_gates 机器判门 → 门绿据 on-disk 源工程
  重建七要素 + persist 落库 / 否则 verdict_feedback 续修,直到门绿/预算耗尽/熔断。不靠内存
  session 判收敛(服务态九工具每回合工厂新建 Tier2Session),贴合「门机器判、控制面绝不自评翻绿」。

C1 AgentScope Studio 观测(observability/studio_sink.py + trace.py):
  把统一 trace 事件按 OTLP(经核实 Studio = OpenTelemetry,非私有 REST;HTTP 3000 /v1/traces)
  映射成 span 推 Studio,traceId→gen_ai.conversation.id 聚成一条 run。私有 TracerProvider 避免
  抢 otel 全局单例。默认关闭(无 TIER2_STUDIO_URL 即 no-op)、observe-only、绝不改 decision。

C2 MCP 工具面(worker/mcp_tools.py + toolkit.py 加性挂点):
  图说 E 族——build_toolkit 末尾加配置门控挂点,默认关闭时 build_mcp_clients() 返回 []、
  Toolkit(tools, mcps=[]) 与原构造等价,九工具(spike 已验证主链)行为零改。Toolkit 接受 mcps=
  形参已对源码核实(tool/_toolkit.py:93)。

C3 基线 client + Prompt 一致性机器门(config.py + genconfig.py + generation.yaml + ci/):
  build_baseline_model(Opus 经 new-api Anthropic 原生)供 A 门「证路」,build_model M3 路零回归;
  check-prompt-registry.sh 文件级三向一致 + import 级命中 registry,6c6g 实跑 exit 0。

修复 design_team 配置外置兑现(design_team.py):
  leader_max_iters/per_expert_cap/timeout_s 改 None→genconfig.get('design_team',...),此前硬编码
  默认 240 导致改 generation.yaml 的 timeout 不生效;已验 genconfig 真读到 YAML 值。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-24 07:46:59 +00:00

338 lines
20 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""observability/studio_sink.py —— tier2 富游戏自治线 · 把统一 trace 事件推到 AgentScope Studio(C1 观测接线)。
【这份解决什么问题】
tier2 一次生成 run 经 Tier2TraceMiddleware / TraceAdapter.ingest 把每步事件映射成统一 trace 形状
(observability/trace.py 的 TraceStep:traceId/step/cost/verdict/timestamp + 推理/动作/观察扩展段),
现状只落内存 adapter.steps、收口由编排器读 summary。本模块把这同一条事件流再接一个旁路 sink —— 把每步
推成 OTLP span 发给 AgentScope Studio,让 tier2 生成 run 在 Studio 的 Trace 页可视化(服务于创始人
「迭代过程拿证据/日志做优化」)。
【Studio 真实接入方式(源码/reference 核实,绝不臆造)】
AgentScope Studio 的可观测是 **OpenTelemetry / OTLP 协议** 的(不是某私有 REST 形状):
- studio README:「Tracing: OpenTelemetry-based trace visualization for LLM calls, token usage,
and agent invocations」(/root/oss/agentscope-studio/README.md)。
- studio tracing 教程(/root/oss/agentscope-studio/docs/tutorial/en/develop/tracing.md):Studio 默认暴露
· OTLP/Trace/gRPC:localhost:4317(env OTEL_GRPC_PORT 调)
· OTLP/Trace/HTTP:localhost:3000(env PORT 调)
「Studio currently only supports receiving Trace type data.」
- 该教程「Advanced Integration: Import Custom Trace Data」段给了任意来源接 Studio 的官方配方 ——
用 OpenTelemetry Python SDK 的 OTLPSpanExporter(endpoint=Studio HTTP)+ BatchSpanProcessor +
TracerProvider,自己 start_as_current_span 组 span(原文示例 endpoint="http://localhost:3000")。
- 接收侧实现在 studio server(/root/oss/agentscope-studio/packages/server/src/index.ts:136 提到
「traces will be received via HTTP endpoint /v1/traces」;OTLP HTTP exporter 自动在 base 末尾拼 /v1/traces)。
AgentScope 框架原生那条(agentscope.init(studio_url=...) + middleware.TracingMiddleware)同样落到这个
OTLP 端点(middleware/_tracing/_trace.py:用 opentelemetry 全局 tracer 发 span;setup 由 init 装 SDK provider)。
【本模块接的是哪条 / 为什么】
tier2 已经有自己一套结构化 trace 事件(TraceStep),且 C1 的任务是「把 TraceAdapter ingest 的事件接到
Studio」—— 所以走 **Advanced Integration 自组 span 这条**(把每条 TraceStep 映射成一个 gen-ai 语义约定 span),
而不是去替换/抢装 agentscope.init 的全局 TracerProvider(那是编排器/服务启动层的事,会与框架原生 middleware
争 set_tracer_provider 全局单例)。本 sink 自带一个 **独立私有 TracerProvider**(不碰 otel 全局),
span 经独立 OTLPSpanExporter 直发 Studio HTTP,与框架原生 trace 互不干扰、可单独开关。
【observe-only 铁律(C1 红线)】
trace / Studio 是纯旁路:绝不改 verdict / decision / L1 / 任何门判定,绝不中断生成主链。
- 取不到 Studio 地址 → sink 直接降级为 no-op(__call__ 立即 return),只留内存 adapter,绝不报错;
- 推送失败(Studio 没起 / 网络错 / 缺 otel 依赖)→ best-effort 记一次告警日志后继续,绝不抛、绝不阻塞;
- 本模块只「读」TraceStep 往外发,从不回写任何决策字段。
【惰性 import(6c6g 红线)】
opentelemetry SDK / exporter 在函数体内 import(6c6g 静态校验环境无这些 wheel),本模块顶层零重依赖 ——
保证 6c6g 能 py_compile 且裸 import。真正建 provider / 发 span 只在装了依赖的 mini-desktop / Mac 上发生。
【开关 / 地址来源】
Studio HTTP 地址按以下顺序取(默认关闭 = 取不到即 no-op):
① env TIER2_STUDIO_URL(单点覆盖,临时开一次实验最省事);
② service/infra_config.py 的 [studio].http_url(若该分区/项存在;infra.yaml 内网入仓口径);
③ 显式传给 make_studio_sink(studio_url=...) 的参数。
三者都没有 → resolve_studio_url() 返回 None → make_studio_sink 返回 no-op sink(现有 CLI 真跑不受影响)。
"""
from __future__ import annotations
import os
from typing import Any, Callable, Optional
# ── env 开关键(与 genconfig/infra_config 的 env 风格一致:全大写、TIER2_ 前缀)──
# 取到非空即视为「开 Studio 观测」并用作 Studio HTTP base(如 http://100.64.0.7:3000)。
ENV_STUDIO_URL = "TIER2_STUDIO_URL"
# ── tier2 这条 trace 在 Studio 里的服务名(OTLP resource service.name;Studio 按它聚合归类)──
SERVICE_NAME = "tier2-gen-worker"
# ── 一次性告警去重(同一 tag 只 print 一次,避免长 run 刷屏;对齐 genconfig._warn_once 风格)──
_WARNED: set[str] = set()
def _warn_once(tag: str, msg: str) -> None:
"""同一 tag 的告警只落一次(best-effort 告警去重,绝不抛)。"""
if tag in _WARNED:
return
_WARNED.add(tag)
print(f"[tier2-studio-sink] {msg}", flush=True)
def resolve_studio_url(explicit: Optional[str] = None) -> Optional[str]:
"""解析 Studio HTTP base 地址:env > infra_config[studio].http_url > 显式参数;都没有 → None(=关闭)。
:param explicit: 调用方显式传入的 Studio HTTP 地址(最低优先级,作兜底)。
:return: 归一化后的 Studio HTTP base(去尾斜杠);取不到任何来源 → None(make_studio_sink 据此降级 no-op)。
"""
# ① env 单点覆盖(最高优先级)。
env_val = os.environ.get(ENV_STUDIO_URL)
if env_val and env_val.strip():
return env_val.strip().rstrip("/")
# ② infra_config.py 的 [studio].http_url(惰性 import:infra_config 顶层无重依赖,但仍 best-effort 兜)。
try:
# 包内/直跑兼容导入(直跑 observability 下脚本时 service 包仍可解析)。
try:
from service import infra_config # type: ignore
except Exception: # pragma: no cover —— 直跑兜底:把 gen-worker/ 加进 sys.path 再取
import sys
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parents[1]))
from service import infra_config # type: ignore
cfg_val = infra_config.get("studio", "http_url", None)
if cfg_val and str(cfg_val).strip():
return str(cfg_val).strip().rstrip("/")
except Exception as exc: # best-effort:读配置任何异常都不阻断,落一次告警后继续
_warn_once("infra_config", f"读 infra_config[studio].http_url 失败(忽略):{exc}")
# ③ 显式参数兜底。
if explicit and explicit.strip():
return explicit.strip().rstrip("/")
return None
# ── TraceStep 工具结果状态 → OTLP span 状态码的映射(observe 段 verdict 决定 span 是否标 ERROR)──
# 注意:这只是「把工具/门结果状态如实反映到 span 颜色」,不改任何生成侧判定(observe-only)。
_ERROR_VERDICTS = {"error", "interrupted", "denied"}
class StudioSink:
"""把统一 trace 事件(TraceStep dict)推成 OTLP span 发给 AgentScope Studio 的旁路 sink。
用法(插进现有 TraceAdapter 的 sink 槽,零改 ingest / 事件结构):
sink = make_studio_sink(trace_id="<gen-task-trace>") # 取不到 Studio 地址 → 返回 no-op
adapter = TraceAdapter(trace_id="<gen-task-trace>", sink=sink)
async for event in agent.reply_stream(user_msg):
adapter.ingest(event) # 映射 + 落内存 + 经本 sink 推 Studio,全程 best-effort
sink.shutdown() # 收口:flush 残留 span(best-effort,不抛)
sink 契约 = Callable[[dict], None](与 TraceAdapter 的 sink 一致):吃一条 TraceStep dict → 发一个 span。
任何环节失败(缺 otel 依赖 / 建 provider 失败 / 发送失败)一律 best-effort:落一次告警、自动降级 no-op、绝不抛。
"""
def __init__(self, trace_id: str, studio_url: str) -> None:
"""
Args:
trace_id: 本次生成的 traceId(贯穿;作为 span 的 gen_ai.conversation.id,Studio 据它把同一 run 的步归一条 trace)。
studio_url: Studio HTTP base(已归一化;OTLP HTTP exporter 会自动在末尾拼 /v1/traces)。
"""
self.trace_id = trace_id
self.studio_url = studio_url
self._enabled = True # best-effort 自降级开关:首次建 provider 失败即永久置 False
self._provider: Any = None # 私有 TracerProvider(惰性建;不碰 otel 全局单例)
self._tracer: Any = None # 从私有 provider 取的 tracer
self._initialized = False # 是否已尝试过初始化(成功或失败都置 True,避免每条都重试)
def _ensure_tracer(self) -> bool:
"""惰性建私有 TracerProvider + OTLPSpanExporter(HTTP→Studio);成功返回 True,失败 best-effort 降级返回 False。
为什么用私有 provider 而非 otel 全局 set_tracer_provider:本 sink 是「把 tier2 自己的 TraceStep 另发一份给
Studio」的旁路,不应抢占 otel 全局单例(那是 agentscope.init / 框架原生 TracingMiddleware 的地盘);
私有 provider 经自带 exporter 直发,互不干扰、可独立开关。
"""
if self._initialized:
return self._enabled
self._initialized = True
try:
# 惰性 import:6c6g 无 opentelemetry SDK / exporter wheel,本模块仍可被裸 import(此处才触发依赖)。
from opentelemetry.sdk.trace import TracerProvider # noqa: PLC0415
from opentelemetry.sdk.trace.export import ( # noqa: PLC0415
BatchSpanProcessor,
)
from opentelemetry.exporter.otlp.proto.http.trace_exporter import ( # noqa: PLC0415
OTLPSpanExporter,
)
from opentelemetry.sdk.resources import Resource # noqa: PLC0415
# resource.service.name = tier2-gen-worker:Studio 按服务名聚合归类这条线的 trace。
resource = Resource.create({"service.name": SERVICE_NAME})
provider = TracerProvider(resource=resource)
# OTLP/HTTP exporter:endpoint 传 Studio HTTP base,SDK 默认在末尾拼 /v1/traces(对齐 studio /v1/traces 接收)。
# 显式给 traces 全路径,避免不同 SDK 版本对 base/自动拼路径行为不一致(reference 示例用 base,
# 这里给全路径更稳;两者 studio 都收)。
exporter = OTLPSpanExporter(endpoint=self.studio_url + "/v1/traces")
# 批处理 processor:异步攒批后台发,推送不卡 ingest 主路(observe-only 不阻塞生成)。
provider.add_span_processor(BatchSpanProcessor(exporter))
# 从私有 provider 取 tracer(不经 otel.trace.get_tracer 全局,避免取到框架/默认 provider)。
self._provider = provider
self._tracer = provider.get_tracer("tier2-gen-worker", "1.0.0")
_warn_once(
"init-ok",
f"Studio 观测已开(traceId={self.trace_id}{self.studio_url}/v1/traces)",
)
return True
except Exception as exc: # best-effort:缺依赖 / 建 provider 失败 → 永久降级 no-op,不抛、不阻塞
self._enabled = False
_warn_once(
"init-fail",
f"建 OTLP→Studio 通道失败,本 run 降级为不推 Studio(只留内存 trace):{exc}",
)
return False
def __call__(self, trace_step: dict) -> None:
"""吃一条 TraceStep dict → 发一个 OTLP span 给 Studio(全程 best-effort,绝不抛、绝不阻塞主链)。
映射(把 TraceStep 公共核心子集 + 扩展段铺成 gen-ai 语义约定 + agentscope 扩展约定的 span 属性):
- span.name = 取 ext.action.tool / 相位 kind / raw.event,组成可读步名(如 "step 3 · write_source")。
- gen_ai.conversation.id = traceId(Studio 据它把同一 run 的步归一条 trace —— 与框架 TracingMiddleware 同键)。
- agentscope.function.* / gen_ai.usage.* = 把 ext 三段(推理/动作/观察)+ cost.tokens 铺成属性,JSON 序列化。
- span 状态:observe 段 verdict ∈ {error/interrupted/denied} → 标 ERROR(如实反映工具/门结果,不改判定)。
"""
if not self._enabled:
return # 已降级 no-op(取不到地址 / 初始化失败):直接返回,observe-only 不影响主链
if not self._ensure_tracer():
return # 首次初始化失败已落告警并永久降级
try:
self._emit_span(trace_step)
except Exception as exc: # best-effort:单条 span 发送失败只告警一次,不连累后续、不抛
_warn_once(
"emit-fail",
f"推 Studio span 失败(best-effort 计告警,traceId={self.trace_id}):{exc}",
)
def _emit_span(self, trace_step: dict) -> None:
"""把单条 TraceStep 组成一个 span 并即时收尾(同步 start→set→end,交 BatchSpanProcessor 异步发出)。"""
import json # noqa: PLC0415 —— 仅序列化属性用,函数内 import 无妨
from opentelemetry.trace import StatusCode # noqa: PLC0415
step = trace_step.get("step")
ext = trace_step.get("ext") or {}
raw = ext.get("raw") or {}
# ── span 名:优先工具名 / 相位 kind / 原始事件类型,拼上步序,便于在 Studio 列表里一眼认出这一步在干嘛 ──
label = (
(ext.get("action") or {}).get("tool")
or (ext.get("action") or {}).get("kind")
or (ext.get("observation") or {}).get("kind")
or (ext.get("reasoning") or {}).get("kind")
or raw.get("event")
or "step"
)
span_name = f"step {step} · {label}"
# 用私有 tracer 起 span(end_on_exit 默认 True:with 退出即 end;本步是「一条已成事实的 trace 记录」,
# 无需跨步嵌套,即起即收最简单、也不会因异常漏 end)。
with self._tracer.start_as_current_span(span_name) as span:
# ── 公共核心子集 → 标准/约定属性 ──
# traceId 作 gen_ai.conversation.id(与框架原生 TracingMiddleware 同键:Studio 据此把同一 run 的步聚一条 trace)。
span.set_attribute("gen_ai.conversation.id", str(self.trace_id))
if step is not None:
span.set_attribute("agentscope.tier2.step", int(step))
ts = trace_step.get("timestamp")
if ts:
span.set_attribute("agentscope.tier2.timestamp", str(ts))
# ── cost.tokens → gen_ai.usage.*(Studio 的 token 用量视图按这两个标准键展示)──
cost = trace_step.get("cost") or {}
tokens = (cost or {}).get("tokens") or {}
if "in" in tokens:
span.set_attribute("gen_ai.usage.input_tokens", int(tokens.get("in") or 0))
if "out" in tokens:
span.set_attribute("gen_ai.usage.output_tokens", int(tokens.get("out") or 0))
# ── tier2 扩展三段 → agentscope.function.* + tier2 自有属性(JSON 序列化,Studio 在 Metadata 区展示)──
# 用 agentscope.function.input/output 这两个 reference 明确列出的「Opt-In」扩展键承载推理/动作/观察,
# 让 Studio 能按 AgentScope 扩展约定展示;原始三段也各存一份 tier2 私有键便于反查。
for seg_key in ("reasoning", "action", "observation"):
seg = ext.get(seg_key)
if seg:
span.set_attribute(
f"agentscope.tier2.{seg_key}",
json.dumps(seg, ensure_ascii=False),
)
# 动作段的工具名也铺到标准键(gen_ai.tool.name),Studio 工具调用视图按它高亮。
tool = (ext.get("action") or {}).get("tool")
if tool:
span.set_attribute("gen_ai.tool.name", str(tool))
# 原始事件类型(反查用)。
if raw.get("event"):
span.set_attribute("agentscope.tier2.raw_event", str(raw.get("event")))
# ── verdict → span 状态:如实反映工具/门结果(observe-only,不改任何生成判定)──
verdict = trace_step.get("verdict")
if verdict is not None:
span.set_attribute("agentscope.tier2.verdict", str(verdict))
if str(verdict).lower() in _ERROR_VERDICTS:
# 工具失败 / 被拒 / 中断 → span 标 ERROR(只是把已发生的结果反映到 trace 颜色,不影响主链)。
span.set_status(StatusCode.ERROR, str(verdict))
else:
span.set_status(StatusCode.OK)
def shutdown(self) -> None:
"""收口:flush + 关私有 provider,把 BatchSpanProcessor 攒批里残留的 span 发出(best-effort,绝不抛)。
编排器在一次 run 收口处调一次(对齐 TraceAdapter.summary 那一步);不调也不会崩(进程退出时 batch 可能丢尾批,
但 observe-only 容忍 —— trace 丢尾批不影响生成产物)。
"""
if self._provider is None:
return
try:
# force_flush:把攒批里没发的 span 立刻发出(给个上限超时,避免收口卡住)。
self._provider.force_flush(timeout_millis=5000)
except Exception as exc: # best-effort:flush 失败只告警
_warn_once("flush-fail", f"force_flush 残留 span 失败(忽略):{exc}")
try:
self._provider.shutdown()
except Exception as exc:
_warn_once("shutdown-fail", f"关 Studio provider 失败(忽略):{exc}")
# ── no-op sink:取不到 Studio 地址时返回它,插进 TraceAdapter.sink 槽等于「没接 Studio」(只留内存 adapter)──
class _NoopStudioSink:
"""关闭态占位 sink:__call__ 什么都不做、shutdown 什么都不做。让调用方无需判 None、统一拿 sink + 调 shutdown。"""
enabled = False
def __call__(self, trace_step: dict) -> None: # noqa: D401 —— 故意 no-op
return
def shutdown(self) -> None: # noqa: D401 —— 故意 no-op
return
def make_studio_sink(
trace_id: str,
*,
studio_url: Optional[str] = None,
) -> Callable[[dict], None]:
"""构造一个 Studio 旁路 sink:解析到 Studio 地址 → StudioSink;取不到 → no-op sink(默认关闭)。
这是给编排器/run 的唯一入口:无论开没开 Studio,都拿到一个可直接塞进 TraceAdapter(sink=...) 的可调用对象,
且都带 .shutdown()(no-op 时也安全可调)—— 调用方零分支判断,observe-only 不影响主链。
:param trace_id: 本次生成 traceId(贯穿)。
:param studio_url: 显式 Studio HTTP base(最低优先级;一般留空,走 env TIER2_STUDIO_URL / infra_config)。
:return: 一个 Callable[[dict], None] 且带 .shutdown() 的 sink 对象(StudioSink 或 _NoopStudioSink)。
"""
url = resolve_studio_url(studio_url)
if not url:
# 取不到地址 = 默认关闭:返回 no-op(现有 CLI 真跑不受任何影响)。
return _NoopStudioSink()
return StudioSink(trace_id=trace_id, studio_url=url)
if __name__ == "__main__": # pragma: no cover —— 本地自检:无 Studio 地址时应拿到 no-op,且 __call__/shutdown 不炸
s = make_studio_sink(trace_id="demo-trace-001")
print("sink type:", type(s).__name__) # 期望 _NoopStudioSink(未设 TIER2_STUDIO_URL)
s({"traceId": "demo-trace-001", "step": 0, "ext": {"raw": {"event": "ReplyStartEvent"}}})
s.shutdown()
print("noop sink okobserve-only 关闭态:吃事件不报错、shutdown 不报错)")