C1 此前建了 observability.with_studio + studio_sink(OTLP→Studio),但没 wire 进真正跑的 Tier2TraceMiddleware(它在 worker/middleware.py、C1 没动)。本提交把中间件新建 adapter 的路径从 TraceAdapter(...) 换成 with_studio(...):drop-in——TIER2_STUDIO_URL/infra_config[studio] 取不到时 studio sink 自动 no-op,行为与原 TraceAdapter 完全一致(CLI/服务两线不开 Studio 即无感);取到地址则 每步在原 sink 外旁路推一份 OTLP span 给 AgentScope Studio(observe-only,best-effort,不阻塞主链)。 存 _studio_sink 句柄供收口 flush。配合 mini-desktop 起 Studio(创始人:now)即得真跑可视化。 Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
679 lines
40 KiB
Python
679 lines
40 KiB
Python
"""middleware.py —— tier2 单写 ReAct agent 的四道熔断 + ¥ 累进硬闸(C3/A4 图说)+ 软刹 + 全钩子 trace(H1/H2)。
|
||
|
||
为什么要熔断:放开 max_iters 的自治 ReAct 循环不能无限烧 —— 模型可能卡在「改→构建失败→再改」的
|
||
死圈、或 thinking 跑飞、或单步外部调用挂死。四道熔断是循环的硬地板,任一触发即停,落 verdict.breakerKind
|
||
(对接 tier2-verdict.schema.json breakerKind enum:step_cap/budget/stuck/timeout)。
|
||
|
||
接法(AgentScope v2.0.2 源码核验):熔断做成 `MiddlewareBase` 子类(middleware/_base.py)。四道熔断挂
|
||
`on_reply` 洋葱钩子(最外层,_agent.py:505),包裹整个 reply 生成器,每个 AgentEvent 流经时巡检四道闸:
|
||
① step_cap —— 步数硬顶:数 ToolCallStartEvent(每次工具调用 = 一步),超 max_tool_calls 即断。
|
||
② budget —— 预算闸:数 ModelCallStartEvent(每次模型推理 = 一次计费),超 max_model_calls 即断
|
||
(这是**廉价代理闸**;真金额硬闸见下方 ¥ 累进硬闸,二者同 kind=budget、语义递进)。
|
||
③ stuck —— 卡死探测:连续 N 次工具调用结果都是同一个失败签名(如反复 build 失败同错)→ 判死圈。
|
||
④ timeout —— 双层超时:整 reply 墙钟超 wall_timeout_s,或单步(两个相邻事件间)静默超 step_timeout_s。
|
||
|
||
¥ 累进硬闸(U3,图说 A4「token×new-api 单价换算成 ¥ 累进台账,越硬上限即 fail-closed」+ C3「强制硬闸
|
||
规模化前必须落地」):挂 `on_model_call` 洋葱钩子(最内层,裸模型 API 调用层,_agent.py:2121)。每次裸模型
|
||
调用**前**按「累计已花 ¥ + 本次预估 ¥」判,越 rmb_hard_limit 即 fail-closed 抛 Tier2CircuitBreak(budget);
|
||
调用后按实测 usage 经 cost.compute(new-api quota 同口径)折 ¥ 累加进 spent_rmb(权威)。活价经
|
||
newapi_pricing.fetch_pricing_params 首调惰性取(best-effort);**取不到则降级为「模型调用次数闸」**
|
||
(不假装拦得住金额,degrade)。这道把原「次数代理闸」升级成「金额硬地板」,补上 C3 标注的那处设计缝。
|
||
|
||
触发动作:记下 tripped(breakerKind + 原因),抛 Tier2CircuitBreak。编排器 catch 它 → 据 breakerKind 落
|
||
verdict(decision=kill / reasons 含 circuit_break)。
|
||
|
||
注:handoff brief 提到「官方 ReplyBudgetControlMiddleware 软刹」——经 v2.0.2 源码核验,该类**不存在**
|
||
(middleware/__init__.py 仅导出 TracingMiddleware/TTSMiddleware)。故软刹也在本模块自实现:
|
||
软刹 = 临近硬顶时(达 soft_ratio×硬顶)往 agent 注入一条 system-reminder,提醒它「快收敛、尽早 finish」,
|
||
给模型一个体面收尾的机会,而非到硬顶才粗暴 kill。软刹经 on_system_prompt 钩子(transformer 管线)实现。
|
||
|
||
trace 接线(H1/H2,Tier2TraceMiddleware,本模块第二个 middleware;U3 升级为「洋葱全钩子」):
|
||
H 族图说要求 tier2 这条线「订阅 + 映射」官方 typed Event System 产统一 trace 形状。本 middleware 现挂**四道
|
||
onion 钩子全套**(图说 A4 同心环):
|
||
- `on_reply`(最外层洋葱,_agent.py:505):看见一次 reply 流经的全部逐 block 级 AgentEvent ——
|
||
ReplyStart/End、ModelCallStart/End(End 带 token)、ThinkingBlock*/TextBlock*(reason)、ToolCall*(act)、
|
||
ToolResult*(observe)——逐事件 ingest 进 observability.trace.TraceAdapter,再原样 yield。
|
||
- `on_reasoning` / `on_acting` / `on_model_call`(U3 补齐,同心环更内三层):各只框住**一个相位**的进/出边界,
|
||
产更细粒度的「相位段」typed event(推理段 / 单工具 I/O 段带工具名耗时 / 模型调用段带实测 token),
|
||
走 adapter 同一条 ingest 路(make_phase_marker / PHASE_*),与 on_reply 那条轨**零重复**(标记是另一类 type)。
|
||
**不并进 CircuitBreakerMiddleware** 是为了单一职责:trace 只观测、不拦截、不抛;熔断只拦截。两者都挂多道钩子时
|
||
由框架按 middlewares 列表序串成洋葱链(_agent.py execute_chain 递归),互不耦合。
|
||
best-effort 铁律(H1/H2 黄带):trace ingest 全包 try(TraceAdapter.ingest 内部已 best-effort),本 middleware
|
||
再兜一层(_safe_ingest)——trace 失败绝不中断生成,只告警;相位钩子里下层真异常照常上抛(不吞)。
|
||
"""
|
||
|
||
import time
|
||
from datetime import datetime
|
||
from typing import Optional
|
||
|
||
from agentscope.middleware import MiddlewareBase
|
||
|
||
# 生成配置层(运行时读外部 generation.yaml + env 覆盖 + 缺省回落内置默认;读不到/脏 → 内置默认,绝不抛)。
|
||
# 本模块的熔断/软刹/¥ 闸旋钮默认值改为从这里读、default=现值,故「调成本上限/失控保护 = 改 YAML 重跑」。
|
||
# 包内/直跑兼容导入(与下方 observability 同款兜底:直跑时把 gen-worker/ 加进 sys.path)。
|
||
try:
|
||
from . import genconfig # type: ignore
|
||
except Exception: # pragma: no cover —— 直跑兜底:worker 包未就位时把 gen-worker/ 加进 sys.path
|
||
import sys
|
||
from pathlib import Path
|
||
|
||
sys.path.insert(0, str(Path(__file__).resolve().parents[1]))
|
||
from worker import genconfig # type: ignore
|
||
|
||
# trace adapter + 相位标记构造器在 observability 模块(H1/H2 产出件);包内/直跑兼容导入。
|
||
# make_phase_marker / PHASE_* 是 U3 给三道新钩子产「相位段」typed event 用的(走 adapter 同一条 ingest 路)。
|
||
try:
|
||
from observability.trace import ( # type: ignore
|
||
TraceAdapter, with_studio, make_phase_marker, PHASE_REASONING, PHASE_ACTING, PHASE_MODEL_CALL)
|
||
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 observability.trace import ( # type: ignore
|
||
TraceAdapter, with_studio, make_phase_marker, PHASE_REASONING, PHASE_ACTING, PHASE_MODEL_CALL)
|
||
|
||
# 成本折算(¥ 累进硬熔断用):compute 折单条 usage → ¥(纯函数,显式传参);活价取不到时回落次数闸。
|
||
# newapi_pricing.fetch_pricing_params 一步取齐 new-api 计费三件套(失败 best-effort 返回 None)。
|
||
try:
|
||
from observability.cost import compute as _cost_compute # type: ignore
|
||
from observability.newapi_pricing import ( # type: ignore
|
||
fetch_pricing_params as _fetch_pricing_params)
|
||
except Exception: # pragma: no cover —— 直跑/路径未就位兜底
|
||
import sys
|
||
from pathlib import Path
|
||
|
||
sys.path.insert(0, str(Path(__file__).resolve().parents[1]))
|
||
from observability.cost import compute as _cost_compute # type: ignore
|
||
from observability.newapi_pricing import ( # type: ignore
|
||
fetch_pricing_params as _fetch_pricing_params)
|
||
|
||
|
||
# ── 模块级小工具(供 trace 三道新钩子 + ¥ 熔断钩子复用,全 best-effort、绝不抛)──
|
||
def _now_iso() -> str:
|
||
"""当前时刻 ISO8601(相位标记 created_at;失败回空串,映射侧容忍空 timestamp)。"""
|
||
try:
|
||
return datetime.now().isoformat()
|
||
except Exception: # pragma: no cover —— 极端环境取时失败也不连累生成
|
||
return ""
|
||
|
||
|
||
def _is_async_gen(obj) -> bool:
|
||
"""判断对象是否异步生成器(on_model_call 返回可能是 ChatResponse 或 AsyncGenerator[ChatResponse])。"""
|
||
import inspect as _inspect # 函数内 import,避免模块级再加一行顶层依赖
|
||
|
||
return _inspect.isasyncgen(obj)
|
||
|
||
|
||
def _usage_tokens(usage) -> tuple[int, int]:
|
||
"""从 ChatUsage(或 None)安全取 (input_tokens, output_tokens);脏值/None → (0,0),绝不抛。"""
|
||
if usage is None:
|
||
return (0, 0)
|
||
try:
|
||
def _si(v) -> int:
|
||
try:
|
||
n = int(v)
|
||
return n if n >= 0 else 0
|
||
except (TypeError, ValueError):
|
||
return 0
|
||
|
||
return (
|
||
_si(getattr(usage, "input_tokens", 0)),
|
||
_si(getattr(usage, "output_tokens", 0)),
|
||
)
|
||
except Exception: # pragma: no cover —— usage 解析任何异常按 0 计,不阻断
|
||
return (0, 0)
|
||
|
||
|
||
def _usage_cached(usage) -> int:
|
||
"""从 ChatUsage 取命中缓存的 prompt token(¥ 折算缓存折扣用);缺/脏 → 0。
|
||
|
||
口径同 observability.cost.cache_hit_tokens:Anthropic 命中字段被 AgentScope 归一为 cache_input_tokens。
|
||
"""
|
||
if usage is None:
|
||
return 0
|
||
try:
|
||
cands = []
|
||
for name in ("cache_input_tokens", "cache_read_input_tokens"):
|
||
v = getattr(usage, name, None)
|
||
if v is not None:
|
||
try:
|
||
cands.append(int(v))
|
||
except (TypeError, ValueError):
|
||
pass
|
||
return max(cands) if cands else 0
|
||
except Exception: # pragma: no cover
|
||
return 0
|
||
|
||
|
||
class Tier2CircuitBreak(Exception):
|
||
"""熔断触发异常:编排器 catch 它,据 kind 落 verdict.breakerKind。
|
||
|
||
kind ∈ {step_cap, budget, stuck, timeout}(对接 tier2-verdict.schema.json breakerKind enum)。
|
||
"""
|
||
|
||
def __init__(self, kind: str, reason: str) -> None:
|
||
self.kind = kind
|
||
self.reason = reason
|
||
super().__init__(f"[熔断:{kind}] {reason}")
|
||
|
||
|
||
# ── ¥ 累进硬熔断上限(directional;真值待 mini-desktop 跑 n≥30 标定成本分布后定)──
|
||
# 单次 tier2 富游戏生成(自治多轮)的 ¥ 硬上限。当「已花 ¥ + 本次预估 ¥」越过它即 fail-closed 抛熔断。
|
||
# 口径依据:图说 A4「token × new-api 单价换算成 ¥ 累进台账,越硬上限即 fail-closed 抛错终止」+ C3
|
||
# 「强制硬闸规模化前必须落地」。当前值是方向性占位(¥3.0/局):0号 spike 富游戏单局观测约 ¥0.x~¥1 量级,
|
||
# 留 ~3×头寸防失控发散;待 batch n≥30 的 cost_rmb 分布出来后,按 P95 + 安全裕度回填真值。
|
||
# 值改为从生成配置层读(default=现值 3.0):标定后改 generation.yaml 的 budget.rmb_hard_limit 即可、不必改码。
|
||
# 模块级常量在 import 时取一次(作对外可见的「当前默认上限」快照 + genconfig 兜底 default);
|
||
# __init__ 每次实例化会再运行时读一遍(sentinel=None),保证改 YAML 后新建的 middleware 立即用新值。
|
||
DEFAULT_RMB_HARD_LIMIT = genconfig.get("budget", "rmb_hard_limit", 3.0) # directional —— 待标定后改 YAML
|
||
|
||
|
||
class CircuitBreakerMiddleware(MiddlewareBase):
|
||
"""四道熔断 + 软刹 + ¥ 累进硬闸(on_reply / on_model_call onion + on_system_prompt transformer)。
|
||
|
||
四道熔断:step_cap(步数硬顶)/ budget(预算闸)/ stuck(卡死)/ timeout(双层超时),挂 on_reply。
|
||
¥ 累进硬闸(U3 新增,图说 A4/C3 的「强制硬闸」):挂 on_model_call —— 每次裸模型调用**前**按
|
||
「累计已花 ¥ + 本次预估 ¥」判,越 rmb_hard_limit 即 fail-closed 抛 Tier2CircuitBreak(kind=budget)。
|
||
它**升级**了 on_reply 里那道「模型调用次数代理闸」:有活价时按真金额拦(更准),取不到活价时回落次数闸
|
||
(degrade,不假装拦得住金额)。两道同 kind=budget,语义递进、不冲突:次数闸是廉价兜底,¥ 闸是金额硬地板。
|
||
|
||
单实例对应单 agent 单次 reply 生命周期(状态计数器在实例上累加;复用前 reset() 或新建)。
|
||
"""
|
||
|
||
def __init__(
|
||
self,
|
||
*,
|
||
max_tool_calls: int | None = None,
|
||
max_model_calls: int | None = None,
|
||
wall_timeout_s: float | None = None,
|
||
step_timeout_s: float | None = None,
|
||
stuck_repeat_threshold: int | None = None,
|
||
soft_ratio: float | None = None,
|
||
rmb_hard_limit: float | None = None,
|
||
pricing_params: Optional[dict] = None,
|
||
enable_rmb_gate: bool = True,
|
||
group_ratio: float | None = None,
|
||
) -> None:
|
||
"""
|
||
Args:
|
||
max_tool_calls: 步数硬顶(工具调用总次数;C3 建议单系统≤8 轮/整局≤40 轮,这里给工具粒度的总闸)。
|
||
max_model_calls: 预算闸(模型推理总次数的廉价代理;¥ 活价取不到时由它兜底拦飞车)。
|
||
wall_timeout_s: 整 reply 墙钟超时(秒)。
|
||
step_timeout_s: 单步静默超时(相邻事件间隔超此判单步卡死;秒)。
|
||
stuck_repeat_threshold: 连续同一失败签名达此次数 → 判死圈。
|
||
soft_ratio: 软刹触发比例(达 soft_ratio×硬顶时注入收敛提醒)。
|
||
rmb_hard_limit: ¥ 累进硬上限(单次生成;越过即 fail-closed)。
|
||
pricing_params: new-api 计费三件套 {pricing, qpu, usd_rate};None → 首次模型调用时经
|
||
fetch_pricing_params 惰性活读取(best-effort,取不到则 ¥ 闸降级为次数闸)。显式传入便于测试/复用。
|
||
enable_rmb_gate: 是否启用 ¥ 累进硬闸(False → 只走原四道;留个总开关便于排障/对照)。
|
||
group_ratio: new-api 分组倍率(成本折算用)。
|
||
|
||
旋钮外置(运行时读):上述各熔断/软刹/¥ 闸/倍率旋钮未显式传入(None)→ 实例化时从生成配置层
|
||
(generation.yaml 的 budget 区)读、default=现值;故「调成本上限/失控保护强度 = 改 YAML 重跑」。
|
||
sentinel=None 而非签名写死,保证每次新建 middleware 都读当前 YAML(改 YAML 后下一次 run 生效)。
|
||
显式传入的参数(如测试/编排器点名)优先,不被配置覆盖。
|
||
"""
|
||
# 未显式传入(None)→ 运行时读 generation.yaml 的 budget 区(default=现值;读不到/脏由 genconfig 兜底为现值)。
|
||
if max_tool_calls is None:
|
||
max_tool_calls = genconfig.get("budget", "max_tool_calls", 60)
|
||
if max_model_calls is None:
|
||
max_model_calls = genconfig.get("budget", "max_model_calls", 80)
|
||
if wall_timeout_s is None:
|
||
wall_timeout_s = genconfig.get("budget", "wall_timeout_s", 1800.0)
|
||
if step_timeout_s is None:
|
||
step_timeout_s = genconfig.get("budget", "step_timeout_s", 420.0)
|
||
if stuck_repeat_threshold is None:
|
||
stuck_repeat_threshold = genconfig.get("budget", "stuck_repeat_threshold", 4)
|
||
if soft_ratio is None:
|
||
soft_ratio = genconfig.get("budget", "soft_ratio", 0.8)
|
||
if rmb_hard_limit is None:
|
||
rmb_hard_limit = genconfig.get("budget", "rmb_hard_limit", 3.0)
|
||
if group_ratio is None:
|
||
group_ratio = genconfig.get("budget", "group_ratio", 1.0)
|
||
|
||
self.max_tool_calls = max_tool_calls
|
||
self.max_model_calls = max_model_calls
|
||
self.wall_timeout_s = wall_timeout_s
|
||
self.step_timeout_s = step_timeout_s
|
||
self.stuck_repeat_threshold = stuck_repeat_threshold
|
||
self.soft_ratio = soft_ratio
|
||
# ¥ 累进硬闸配置(directional 上限 + 计费参数;参数为 None 时 on_model_call 首调惰性取价)。
|
||
self.rmb_hard_limit = rmb_hard_limit
|
||
self.enable_rmb_gate = enable_rmb_gate
|
||
self.group_ratio = group_ratio
|
||
self._pricing_params = pricing_params # {pricing, qpu, usd_rate} 或 None(惰性取价后填)
|
||
self._pricing_fetch_tried = pricing_params is not None # 已有则不再去取
|
||
self.reset()
|
||
|
||
def reset(self) -> None:
|
||
"""复用前重置计数器(单实例跨多次 reply 复用时调用)。"""
|
||
self.tool_calls = 0
|
||
self.model_calls = 0
|
||
self.start_ts: float | None = None
|
||
self.last_event_ts: float | None = None
|
||
# 卡死探测:记最近一次工具失败签名 + 连续重复计数。
|
||
self._last_fail_sig: str | None = None
|
||
self._fail_repeat = 0
|
||
# ¥ 累进硬闸:累计已花 ¥(每次模型调用后按实测 usage 折算累加)+ 当前是否走金额闸(取到活价)。
|
||
self.spent_rmb = 0.0
|
||
self._rmb_gate_active = self._pricing_params is not None # 有计费参数才按金额拦,否则降级次数闸
|
||
# 触发记录(供编排器/调试读)。
|
||
self.tripped: dict | None = None
|
||
|
||
# ── 软刹:在硬顶临近时往 system prompt 追加收敛提醒(transformer 钩子,顺序管线)──
|
||
async def on_system_prompt(self, agent, current_prompt: str) -> str:
|
||
"""达 soft_ratio×硬顶时,在 system prompt 末尾追加一条收敛提醒(软刹,非强制停)。
|
||
|
||
给模型一个体面收尾的机会:提醒它接近步数/预算上限、该尽快产出可交付源工程并调 finish。
|
||
"""
|
||
near_step = self.tool_calls >= self.soft_ratio * self.max_tool_calls
|
||
near_budget = self.model_calls >= self.soft_ratio * self.max_model_calls
|
||
if near_step or near_budget:
|
||
return (
|
||
current_prompt
|
||
+ "\n\n<system-reminder>你已接近本次会话的步数/预算上限。"
|
||
"请停止发散探索,基于当前最好的工程状态尽快收敛:补齐能过 L1 门的最小可玩闭环,"
|
||
"然后调用 finish 交付源工程。再不收敛将触发硬熔断、本次作废。</system-reminder>"
|
||
)
|
||
return current_prompt
|
||
|
||
# ── 四道硬熔断:包裹整个 reply 生成器,逐事件巡检(洋葱钩子)──
|
||
async def on_reply(self, agent, input_kwargs, next_handler):
|
||
"""巡检四道熔断:每个 AgentEvent 流经时检查,任一触发即抛 Tier2CircuitBreak。"""
|
||
self.start_ts = time.monotonic()
|
||
self.last_event_ts = self.start_ts
|
||
|
||
async for evt in next_handler(**input_kwargs):
|
||
now = time.monotonic()
|
||
|
||
# ── ④ timeout(双层):整 reply 墙钟 + 单步静默 ──
|
||
if now - self.start_ts > self.wall_timeout_s:
|
||
self._trip("timeout", f"整 reply 墙钟超时(>{self.wall_timeout_s}s)")
|
||
if self.last_event_ts is not None and now - self.last_event_ts > self.step_timeout_s:
|
||
self._trip(
|
||
"timeout",
|
||
f"单步静默超时(相邻事件间隔 >{self.step_timeout_s}s,疑单步卡死)",
|
||
)
|
||
self.last_event_ts = now
|
||
|
||
evt_type = type(evt).__name__
|
||
|
||
# ── ② budget:每次模型推理计一次 ──
|
||
if evt_type == "ModelCallStartEvent":
|
||
self.model_calls += 1
|
||
if self.model_calls > self.max_model_calls:
|
||
self._trip("budget", f"模型推理次数超预算闸(>{self.max_model_calls})")
|
||
|
||
# ── ① step_cap:每次工具调用计一步 ──
|
||
if evt_type == "ToolCallStartEvent":
|
||
self.tool_calls += 1
|
||
if self.tool_calls > self.max_tool_calls:
|
||
self._trip("step_cap", f"工具调用步数超硬顶(>{self.max_tool_calls})")
|
||
|
||
# ── ③ stuck:连续同一失败签名 → 死圈 ──
|
||
if evt_type == "ToolResultEndEvent":
|
||
# state 是 ToolResultState(use_enum_values=True → 多为字符串值);失败态计入签名。
|
||
state = getattr(evt, "state", None)
|
||
state_val = getattr(state, "value", state)
|
||
tool_call_id = getattr(evt, "tool_call_id", "")
|
||
if str(state_val).lower() in ("error", "denied", "interrupted"):
|
||
# 失败签名 = 状态值(tool_call_id 每次不同,不入签名以便识别「同类失败反复」)。
|
||
sig = str(state_val).lower()
|
||
if sig == self._last_fail_sig:
|
||
self._fail_repeat += 1
|
||
else:
|
||
self._last_fail_sig = sig
|
||
self._fail_repeat = 1
|
||
if self._fail_repeat >= self.stuck_repeat_threshold:
|
||
self._trip(
|
||
"stuck",
|
||
f"连续 {self._fail_repeat} 次工具失败(签名={sig},疑死圈)"
|
||
f" 最后 tool_call_id={tool_call_id}",
|
||
)
|
||
else:
|
||
# 一次成功就清零卡死计数(死圈判定要求「连续」失败)。
|
||
self._last_fail_sig = None
|
||
self._fail_repeat = 0
|
||
|
||
# 未触发任何熔断 → 正常向下游转发事件。
|
||
yield evt
|
||
|
||
# ── ¥ 累进硬闸:挂 on_model_call(裸模型调用层),调用前 fail-closed 判金额 ──
|
||
async def on_model_call(self, agent, input_kwargs, next_handler):
|
||
"""¥ 累进硬熔断(图说 A4/C3 强制硬闸):每次裸模型调用**前**判「已花 ¥ + 本次预估 ¥」越限即 fail-closed。
|
||
|
||
非生成器(源码 _agent.py:2121 同形:return next_handler 的结果,可为 ChatResponse 或 AsyncGenerator)。
|
||
流程:
|
||
① 取价:首次调用惰性活读 new-api 计费参数(best-effort);取到 → 走金额闸,取不到 → 降级次数闸(on_reply
|
||
那道 max_model_calls 已在兜底,这里不重复拦次数,只把闸状态记成 degrade,绝不假装拦得住金额)。
|
||
② fail-closed:调用**前**按 spent_rmb + 本次预估 ¥ 判;越 rmb_hard_limit 即抛 Tier2CircuitBreak(budget)。
|
||
预估用「上一次同模型的均价」或保守 token 估算(调用前拿不到本次真 token,故用预估;调用后按实测回填累计)。
|
||
③ 调真模型 → 调用后读实测 usage,按 cost.compute 折 ¥ 累加进 spent_rmb(权威口径,供下次判 + 收口对账)。
|
||
best-effort 边界:取价 / 折算 / 估算的任何异常都不冤杀生成(降级或按 0 计);唯有「金额确凿越限」才硬熔断。
|
||
"""
|
||
if not self.enable_rmb_gate:
|
||
return await next_handler(**input_kwargs) # 总开关关:不挂 ¥ 闸,直接放行(次数闸仍在 on_reply 兜)
|
||
|
||
# 取当前模型名(折价按模型倍率;取不到则后续折算 compute 用保守默认倍率 0/1/1,折出 0 不误拦)。
|
||
cur = input_kwargs.get("current_model") if isinstance(input_kwargs, dict) else None
|
||
model_name = getattr(cur, "model", None) or ""
|
||
|
||
# ① 惰性取价(仅首次;best-effort,失败则本局降级次数闸,不再重试以免每次调用都打网关)。
|
||
self._ensure_pricing()
|
||
|
||
# ② 调用前 fail-closed 金额判(只有取到活价、即金额闸生效时才硬拦;否则降级、由次数闸兜)。
|
||
if self._rmb_gate_active:
|
||
est_rmb = self._estimate_call_rmb(model_name)
|
||
if self.spent_rmb + est_rmb > self.rmb_hard_limit:
|
||
self._trip(
|
||
"budget",
|
||
f"¥ 累进硬闸越限 fail-closed:已花 ¥{self.spent_rmb:.4f} + 本次预估 ¥{est_rmb:.4f}"
|
||
f" > 上限 ¥{self.rmb_hard_limit:.2f}(模型={model_name or '未知'})",
|
||
)
|
||
|
||
# ③ 调真模型,调用后按实测 usage 折 ¥ 累加(非流式同步读;流式包一层从末块读)。
|
||
result = await next_handler(**input_kwargs)
|
||
if not _is_async_gen(result):
|
||
self._accumulate_call_rmb(model_name, getattr(result, "usage", None))
|
||
return result
|
||
|
||
# 流式:包一层,在末块拿到 usage 时再折算累加(不提前吃光流;消费完才更新 spent_rmb)。
|
||
accumulate = self._accumulate_call_rmb
|
||
|
||
async def _wrap():
|
||
last = None
|
||
try:
|
||
async for chunk in result:
|
||
last = chunk
|
||
yield chunk
|
||
finally:
|
||
accumulate(model_name, getattr(last, "usage", None))
|
||
|
||
return _wrap()
|
||
|
||
def _ensure_pricing(self) -> None:
|
||
"""首次模型调用时惰性活读 new-api 计费参数(best-effort);取到则激活金额闸,取不到则本局降级次数闸。
|
||
|
||
只取一次(_pricing_fetch_tried 守门):取价要打网关,不能每次模型调用都打;取不到就认本局按次数闸兜底。
|
||
"""
|
||
if self._pricing_fetch_tried:
|
||
return
|
||
self._pricing_fetch_tried = True
|
||
try:
|
||
params = _fetch_pricing_params() # best-effort,内部失败返回 None(不抛)
|
||
except Exception as exc: # noqa: BLE001 —— 取价任何异常都降级,绝不中断生成
|
||
params = None
|
||
print(
|
||
f"[tier2-circuit] ¥ 闸取价异常(降级为次数闸):{type(exc).__name__}: {exc}",
|
||
flush=True,
|
||
)
|
||
if params:
|
||
self._pricing_params = params
|
||
self._rmb_gate_active = True
|
||
print(
|
||
f"[tier2-circuit] ¥ 累进硬闸已激活:上限 ¥{self.rmb_hard_limit:.2f}/局"
|
||
f"(qpu={params.get('qpu')} usd_rate={params.get('usd_rate')} "
|
||
f"models={len(params.get('pricing') or {})})",
|
||
flush=True,
|
||
)
|
||
else:
|
||
# 取价失败:¥ 闸降级为次数闸(on_reply 的 max_model_calls 兜底)。明示 degrade,不假装拦得住金额。
|
||
self._rmb_gate_active = False
|
||
print(
|
||
f"[tier2-circuit] ¥ 累进硬闸取价失败,本局降级为「模型调用次数闸」"
|
||
f"(max_model_calls={self.max_model_calls});真金额拦截不可用。",
|
||
flush=True,
|
||
)
|
||
|
||
def _estimate_call_rmb(self, model_name: str) -> float:
|
||
"""估算「本次模型调用」的 ¥(调用前拿不到真 token,故估;best-effort,任何异常 → 0 不误拦)。
|
||
|
||
策略:有历史调用 → 用「已花 ¥ / 已调用次数」的均价作为本次预估(自适应,贴合本局实际单价);
|
||
无历史(首次调用)→ 保守按一笔典型富游戏调用 token 估(prompt 偏大 + 一定补全),用 compute 折真价。
|
||
预估只为「调用前 fail-closed」服务;调用后会用实测 usage 覆盖式累加,故预估不必精确,够拦飞车即可。
|
||
"""
|
||
try:
|
||
if self.model_calls > 0 and self.spent_rmb > 0:
|
||
# 自适应均价:本局已花 ¥ 摊到已调用次数(model_calls 在 on_reply 处按 ModelCallStart 累加)。
|
||
return self.spent_rmb / max(1, self.model_calls)
|
||
# 首次调用:保守典型估(富游戏单写轮 prompt 偏大)。用真计费参数折,取不到价则后面不会走到这(降级了)。
|
||
# 估算 token 量级外置(运行时读 generation.yaml budget.est_*;default=现值 12000/2000)。
|
||
pp = self._pricing_params or {}
|
||
return _cost_compute(
|
||
model_name,
|
||
prompt_tokens=genconfig.get("budget", "est_prompt_tokens", 12000), # 典型富游戏单写轮 prompt 量级(系统提示 + 历史 + 设计稿)
|
||
completion_tokens=genconfig.get("budget", "est_completion_tokens", 2000), # 一轮补全 + tool_use 量级
|
||
pricing=pp.get("pricing", {}),
|
||
qpu=pp.get("qpu", 0) or 1,
|
||
usd_rate=pp.get("usd_rate", 0) or 0,
|
||
group_ratio=self.group_ratio,
|
||
).get("rmb", 0.0)
|
||
except Exception as exc: # noqa: BLE001 —— 估算异常按 0,绝不因估算 bug 冤杀生成
|
||
print(f"[tier2-circuit] ¥ 本次预估失败(按 0 计,不误拦):{exc}", flush=True)
|
||
return 0.0
|
||
|
||
def _accumulate_call_rmb(self, model_name: str, usage) -> None:
|
||
"""调用后按实测 usage 折 ¥ 累加进 spent_rmb(权威口径,best-effort:折算失败不累加、不抛)。"""
|
||
if not self._rmb_gate_active or usage is None:
|
||
return
|
||
try:
|
||
in_tok, out_tok = _usage_tokens(usage)
|
||
cached = _usage_cached(usage)
|
||
pp = self._pricing_params or {}
|
||
detail = _cost_compute(
|
||
model_name,
|
||
prompt_tokens=in_tok,
|
||
completion_tokens=out_tok,
|
||
pricing=pp.get("pricing", {}),
|
||
qpu=pp.get("qpu", 0) or 1,
|
||
usd_rate=pp.get("usd_rate", 0) or 0,
|
||
group_ratio=self.group_ratio,
|
||
cached_tokens=cached,
|
||
)
|
||
self.spent_rmb += float(detail.get("rmb", 0.0) or 0.0)
|
||
except Exception as exc: # noqa: BLE001 —— 折算异常不累加、不抛(成本少记一笔好过冤杀生成)
|
||
print(
|
||
f"[tier2-circuit] ¥ 实测折算失败(本次未计入累计,不中断):{exc}",
|
||
flush=True,
|
||
)
|
||
|
||
def _trip(self, kind: str, reason: str) -> None:
|
||
"""记下触发并抛出 Tier2CircuitBreak(供编排器 catch → 落 verdict.breakerKind)。"""
|
||
self.tripped = {"kind": kind, "reason": reason}
|
||
# 可追溯日志(创始人铁律:错误路径要有可追溯日志)。
|
||
print(f"[tier2-circuit] 熔断触发 kind={kind} reason={reason}", flush=True)
|
||
raise Tier2CircuitBreak(kind, reason)
|
||
|
||
|
||
class Tier2TraceMiddleware(MiddlewareBase):
|
||
"""trace 接线(H1/H2):把官方 typed Event System 经**四道 onion 钩子全挂**映射进 TraceAdapter。
|
||
|
||
职责(对 tier2细节图说-H §H1/H2「订阅 + 映射」+ 图说 A4「洋葱全钩子」,非「埋点 + 采集」):
|
||
- on_reply(最外层洋葱,_agent.py:505/532):看得见一次 reply 流经的**全部**逐 block 级 AgentEvent
|
||
(ReplyStart/End、ModelCallStart/End、ThinkingBlock*/TextBlock*、ToolCall*/ToolResult*),逐事件
|
||
喂 TraceAdapter.ingest(内部经纯函数 to_trace_step 映射成统一 trace 形状并落 sink),再原样 yield。
|
||
- on_reasoning / on_acting / on_model_call(U3 补齐,图说 A4 同心环更内三层):各只框住**一个相位**的
|
||
进/出边界,产「相位段」typed event(推理段 / 单工具 I/O 段带工具名耗时 / 模型调用段带实测 token)——
|
||
粒度比 on_reply 逐事件视图更「分相」,且与 on_reply 那条轨**零重复**(相位标记是另一类 type,见
|
||
observability.trace.make_phase_marker / PHASE_*)。这三道不重复 ingest 逐 block 事件(on_reply 已收)。
|
||
- on_system_prompt **不在本类**(它是 transformer 不是 onion;软刹的 system prompt 注入归 CircuitBreaker)。
|
||
|
||
与 CircuitBreakerMiddleware 的关系:两者都挂 on_reply,各自一个 middleware(单一职责)。框架把多个
|
||
on_reply middleware 串成洋葱链(_agent.py:511 execute_chain 递归);本 middleware 与熔断的相对顺序
|
||
由 Agent(middlewares=[...]) 列表序决定,二者互不依赖(trace 只读事件、不拦截、不抛)。
|
||
|
||
best-effort 铁律(H1/H2 黄带):trace ingest 失败绝不中断生成。TraceAdapter.ingest 内部已 best-effort
|
||
(映射/落 sink 异常只告警不抛),本 middleware 再兜一层 try——哪怕 adapter 本身抛了,也只告警、照样
|
||
把事件 yield 下去,保证生成主链不被一次轨迹抖动废掉。
|
||
|
||
单实例对应单 agent 单次 reply;trace_id 在构造时钉定,本次 reply 所有 step 共用(对接 verdict.evidence.traceId)。
|
||
"""
|
||
|
||
def __init__(self, trace_id: str, *, adapter: Optional[TraceAdapter] = None,
|
||
sink=None, keep_in_memory: bool = True) -> None:
|
||
"""
|
||
Args:
|
||
trace_id: 本次生成的 traceId(贯穿;对接 tier2-verdict evidence.traceId / 成本关联键)。
|
||
adapter: 复用外部已建的 TraceAdapter(便于编排器收口读 steps/summary);None → 本类新建一个。
|
||
sink: 落库回调 sink(trace_step_dict);None → adapter 只留内存 steps(收口由编排器读)。
|
||
keep_in_memory: 是否在 adapter.steps 留内存副本(长 run 可置 False 只走 sink 省内存)。
|
||
"""
|
||
self.trace_id = trace_id
|
||
# 复用传入 adapter 或新建;adapter 是 trace 的产出件(H1/H2),本 middleware 只负责把事件喂进去。
|
||
# 新建走 with_studio(C1 接线):若 TIER2_STUDIO_URL / infra_config[studio] 可解析,则每步在原 sink
|
||
# 之外旁路推一份 OTLP span 给 AgentScope Studio(observe-only);取不到地址 → studio sink 自动 no-op,
|
||
# 行为 == 原 TraceAdapter(sink=...),零影响(故 CLI/服务两条线都白拿 Studio 可观测、不开即无感)。
|
||
self.adapter = adapter or with_studio(
|
||
trace_id=trace_id, sink=sink, keep_in_memory=keep_in_memory)
|
||
# studio sink 句柄(收口 flush 残留 span 用;with_studio 把它挂在 adapter._studio_sink,无则 None)。
|
||
self._studio_sink = getattr(self.adapter, "_studio_sink", None)
|
||
# 兜底层告警计数(adapter 自身抛异常的极端情况;正常应为 0,因 adapter 内部已 best-effort)。
|
||
self._middleware_dropped = 0
|
||
|
||
async def on_reply(self, agent, input_kwargs, next_handler):
|
||
"""最外层洋葱:逐事件 ingest 进 TraceAdapter 后原样 yield(旁路观测,best-effort 不中断)。"""
|
||
async for evt in next_handler(**input_kwargs):
|
||
# trace ingest 全包 try:adapter.ingest 内部已 best-effort,这里再兜一层防御
|
||
# (哪怕 adapter 本身抛了也不连累生成流;trace 失败只告警,事件照常向下游转发)。
|
||
self._safe_ingest(evt)
|
||
# 原样向下游转发事件(不修改、不拦截 —— trace 是纯旁路)。
|
||
yield evt
|
||
|
||
# ── U3:补齐另外三道 onion 钩子,产更细粒度的「相位段」typed event(图说 A4 同心环)──
|
||
# 设计要点(为什么不在这三道里重新 ingest 逐 block 事件):on_reply 已是最外层洋葱,
|
||
# 已经看到并 ingest 了 ThinkingBlock*/ToolCall*/ModelCall* 全部逐事件;若这三道再把同样的事件
|
||
# ingest 一遍,会在 adapter 里重复计步、污染 on_reply 那条轨。所以这三道只产「相位边界标记」——
|
||
# on_reasoning 框一次推理相位、on_acting 框单次工具 I/O(带工具名/耗时)、on_model_call 框一次裸
|
||
# 模型调用(带实测 token):粒度比逐事件视图更「分相」,且与 on_reply 轨零重复(标记是另一类 type)。
|
||
# 全 best-effort:相位标记产失败绝不中断生成(失败只 yield 原值,不影响熔断/生成)。
|
||
|
||
async def on_reasoning(self, agent, input_kwargs, next_handler):
|
||
"""推理相位钩子:在推理(模型决策)相位进/出各产一条标记,中间原样透传子层事件。
|
||
|
||
input_kwargs = {'tool_choice': ...}(源码 _agent.py:739);next_handler 吐的是该相位的逐事件
|
||
(ModelCallStart/Thinking*/Text*/ToolCall*/ModelCallEnd),这里不重复 ingest(on_reply 已收),
|
||
只在边界打「推理相位 enter/exit」标记,给轨迹一条「这是第几段推理」的结构化锚。
|
||
"""
|
||
self._safe_ingest(make_phase_marker(
|
||
PHASE_REASONING, "enter", timestamp=_now_iso()))
|
||
boundary = "exit"
|
||
try:
|
||
async for item in next_handler(**input_kwargs):
|
||
yield item # 原样透传(纯旁路,不拦截推理事件流)
|
||
except Exception:
|
||
boundary = "error" # 推理相位异常退出也打标(便于反查),异常照常向上抛(不吞)
|
||
raise
|
||
finally:
|
||
self._safe_ingest(make_phase_marker(
|
||
PHASE_REASONING, boundary, timestamp=_now_iso()))
|
||
|
||
async def on_acting(self, agent, input_kwargs, next_handler):
|
||
"""单工具 I/O 相位钩子:框住一次 toolkit.call_tool 的进/出,带工具名 + 耗时 + 最终状态。
|
||
|
||
input_kwargs = {'tool_call': ToolCallBlock}(源码 _agent.py:1579);next_handler 吐
|
||
ToolChunk | ToolResponse(非 AgentEvent)。这一层是「动手段」最细的真 I/O 边界——on_reply 看不到
|
||
单工具耗时,这里能算(enter→exit 墙钟),并从末个 ToolResponse 取最终 state(success/error/…)。
|
||
"""
|
||
tool_call = input_kwargs.get("tool_call") if isinstance(input_kwargs, dict) else None
|
||
tool_name = getattr(tool_call, "name", None)
|
||
tool_id = getattr(tool_call, "id", None)
|
||
t_enter = time.monotonic()
|
||
self._safe_ingest(make_phase_marker(
|
||
PHASE_ACTING, "enter", timestamp=_now_iso(),
|
||
fields={"tool": tool_name, "toolCallId": tool_id}))
|
||
boundary = "exit"
|
||
last_state = None
|
||
try:
|
||
async for item in next_handler(**input_kwargs):
|
||
# 顺手记最后一个产物的 state(ToolResponse/ToolChunk 都有 .state;use_enum_values→多为字符串值)。
|
||
st = getattr(item, "state", None)
|
||
if st is not None:
|
||
last_state = getattr(st, "value", st)
|
||
yield item # 原样透传(不修改工具产物)
|
||
except Exception:
|
||
boundary = "error"
|
||
raise
|
||
finally:
|
||
elapsed_ms = int((time.monotonic() - t_enter) * 1000)
|
||
self._safe_ingest(make_phase_marker(
|
||
PHASE_ACTING, boundary, timestamp=_now_iso(),
|
||
fields={
|
||
"tool": tool_name,
|
||
"toolCallId": tool_id,
|
||
"elapsedMs": elapsed_ms,
|
||
"state": (str(last_state) if last_state is not None else None),
|
||
}))
|
||
|
||
async def on_model_call(self, agent, input_kwargs, next_handler):
|
||
"""裸模型调用相位钩子(非生成器:return next_handler 的结果 —— 源码 _agent.py:2121 同形)。
|
||
|
||
「钱算在哪的地方」(图说 A4):本钩子拿到 current_model + messages,调真模型后从 ChatResponse.usage
|
||
读实测 token,产一条带 token 的模型调用相位标记(收口聚合成本时与 on_reply 的 ModelCallEnd 同源、择一计)。
|
||
非流式(tier2 writer 默认 stream=False)直接读 resp.usage;流式则包一层异步生成器,从末块读 usage —— 两路都
|
||
best-effort,trace 失败绝不动返回值/不中断生成。注意:¥ 累进硬熔断**不在本类**,在 CircuitBreakerMiddleware
|
||
的同名钩子里(单一职责:trace 只观测、不拦截)。
|
||
"""
|
||
model_name = None
|
||
try:
|
||
cur = input_kwargs.get("current_model") if isinstance(input_kwargs, dict) else None
|
||
model_name = getattr(cur, "model", None)
|
||
except Exception: # noqa: BLE001 —— 取模型名失败不影响调用,只是标记少一个字段
|
||
model_name = None
|
||
|
||
# 真调下层(模型 API);best-effort 包裹仅用于「调用失败也补一条 error 相位标记」,异常照常上抛(不吞)。
|
||
try:
|
||
result = await next_handler(**input_kwargs)
|
||
except Exception:
|
||
self._safe_ingest(make_phase_marker(
|
||
PHASE_MODEL_CALL, "error", timestamp=_now_iso(),
|
||
fields={"model_name": model_name}))
|
||
raise
|
||
|
||
# 非流式:result 是 ChatResponse,可同步读 usage → 立即落带 token 的相位标记。
|
||
if not _is_async_gen(result):
|
||
in_tok, out_tok = _usage_tokens(getattr(result, "usage", None))
|
||
self._safe_ingest(make_phase_marker(
|
||
PHASE_MODEL_CALL, "exit", timestamp=_now_iso(),
|
||
fields={"model_name": model_name,
|
||
"input_tokens": in_tok, "output_tokens": out_tok}))
|
||
return result
|
||
|
||
# 流式:result 是 AsyncGenerator[ChatResponse];包一层在消费末块时读 usage 落标记(不提前吃光流)。
|
||
adapter = self.adapter
|
||
safe_ingest = self._safe_ingest
|
||
|
||
async def _wrap():
|
||
last = None
|
||
try:
|
||
async for chunk in result:
|
||
last = chunk
|
||
yield chunk
|
||
finally:
|
||
in_tok, out_tok = _usage_tokens(getattr(last, "usage", None))
|
||
safe_ingest(make_phase_marker(
|
||
PHASE_MODEL_CALL, "exit", timestamp=_now_iso(),
|
||
fields={"model_name": model_name,
|
||
"input_tokens": in_tok, "output_tokens": out_tok}))
|
||
_ = adapter # 显式留引用,表明 wrap 与本 adapter 同源(收口反查一致)
|
||
|
||
return _wrap()
|
||
|
||
def _safe_ingest(self, evt) -> None:
|
||
"""把一个事件/相位标记喂进 adapter,全包 try(best-effort:trace 任何异常都不中断生成主链)。"""
|
||
try:
|
||
self.adapter.ingest(evt)
|
||
except Exception as exc: # noqa: BLE001 —— best-effort:trace 任何异常都不中断生成主链
|
||
self._middleware_dropped += 1
|
||
print(
|
||
f"[tier2-trace] middleware ingest 兜底失败(best-effort,traceId={self.trace_id}):"
|
||
f"{type(exc).__name__}: {exc}",
|
||
flush=True,
|
||
)
|
||
|
||
def summary(self) -> dict:
|
||
"""收口摘要:本次 trace 的 traceId、步数、丢弃数(透传 adapter.summary,另带 middleware 兜底丢弃数)。"""
|
||
s = self.adapter.summary()
|
||
s["middlewareDropped"] = self._middleware_dropped
|
||
return s
|