feat(obs): 阶段四观测面二 Python 侧 token metric——llm_tokens{kind} + llm_cache_hit_rate
Some checks failed
contract-gates / contract-gates (push) Has been cancelled
docs-gate / docs-gate (push) Has been cancelled

补面五(game-cloud GenMetrics)拿不到的两个纯 Python 侧 token 维度 metric:
llm_tokens(Counter,kind=input/output/cached)与 llm_cache_hit_rate(Histogram,
每 run 命中率 [0,1])。复用波③ 既有 OTLP 基础设施(开关 CHEAP_OTLP_ENDPOINT、
端点解析、代理旁路、service.name=cheap-gen-worker、best-effort 铁律),换 OTLP
metrics 通道;默认关字节不变、绝不咬生成、与面五 5 个 metric 零重叠不双计。

关键约束(主控亲验坐实):生产 Service 路拿不到 cached——trace.py 组 cost.tokens
只填 in/out,agentscope ModelCallEndEvent 只有 input_tokens/output_tokens 无
cached 字段。故 Service 路(on_reply 只读嗅探 in/out、_metric_emitted 幂等)传
cached=None 如实降级:只发 input/output 计数,不发 cached 计数与 cache_hit_rate,
绝不塞 0 伪造命中率。CLI/批跑路(RecordingOpenAIChatModel.records[2])cached 可得、
两 metric 全发。生产 Service 路发 cached 需 per-POST RecordingModel 或 tier2
breaker 存 cached token,均超本次最小改动边界,列 follow-up。

主控终审:亲跑 tests/test_cheap_otlp_sink.py + test_cheap_service_app.py 29 passed、
全量回归 392 passed/1 skipped;逐条读承重 diff(默认关守卫/best-effort/不双计/
cached 诚实降级/MeterProvider 进程单例防线程泄漏/cached 截 [0,input]/除零守卫)。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
lili 2026-07-05 23:57:40 -07:00
parent 8b91f34b13
commit 958f78f523
5 changed files with 400 additions and 0 deletions

View File

@ -562,6 +562,206 @@ def build_trace_sink(
return _FanoutSink([jsonl_sink, otlp]) return _FanoutSink([jsonl_sink, otlp])
# ══════════════════════════════════════════════════════════════════════════════
# 面二 · Python 侧业务 metric(OTLP metrics)—— 只补面五(game-cloud GenMetrics)拿不到的两个
# ──────────────────────────────────────────────────────────────────────────────
# 面五已从回调 trace 发 gen_task_total{status} / gen_queue_depth / gen_gate_fail_total{gate} /
# gen_duration_seconds / llm_cost —— 成本、耗时、成功率、门失败、队列深度都归它;面二绝不重发这些(重发=双计)。
# 面二只发面五缺的两个纯 Python 侧 token 维度 metric:
# · llm_tokens(Counter{kind=input|output|cached})—— 一次生成累计的输入/输出/缓存 token 计数;
# Prometheus 侧规范化后呈现为 `llm_tokens_total{kind}`(counter 由 exporter 追加 `_total` 后缀,
# 是 OTel→Prometheus 的既定约定,故 instrument 本身不自带 `_total`)。
# · llm_cache_hit_rate(Histogram)—— 一次生成的缓存命中率 cached/input([0,1]);Prometheus 侧
# 展开为 llm_cache_hit_rate_{bucket,sum,count},既能看单次命中率分布、也能按 run 聚合平均。
# 与上面 span sink 复用同一套:开关(env CHEAP_OTLP_ENDPOINT / infra_config,默认关)、端点解析、
# 代理旁路(install_proxy_bypass)、service.name(cheap-gen-worker)、best-effort 铁律;只是换 OTLP metrics 通道。
# 注:token 加权的整体命中率(Σcached/Σinput)由 Prometheus 侧用 llm_tokens_total{kind} 两条 counter 相除算,
# 这里的 histogram 记的是「每次 run 的命中率」——两者互补,一个看整体、一个看单次分布与告警。
# ══════════════════════════════════════════════════════════════════════════════
# OTel instrument 名(InMemoryMetricReader / SDK 层按这个名读;Prometheus 呈现名见上注)。
METRIC_TOKENS = "llm_tokens" # Counter → Prometheus llm_tokens_total{kind}
METRIC_CACHE_HIT_RATE = "llm_cache_hit_rate" # Histogram([0,1] 每 run 命中率)
# 命中率 histogram 的桶边界:cache_hit_rate 是 [0,1] 比率,默认延迟型桶(0/5/10/…)不适用,
# 显式给一组 [0,1] 细分,让「命中率落在哪一档」在 Grafana 里可分辨。
_CACHE_HIT_RATE_BUCKETS = (0.0, 0.1, 0.25, 0.5, 0.75, 0.9, 0.99, 1.0)
# 进程级共享 MeterProvider 单例(差异④同 span:长驻 Service 逐局复用一条 PeriodicExportingMetricReader
# 后台上报线程,不 per-run 建 provider 泄漏线程)。counter/histogram instrument 一并缓存(建一次)。
_shared_meter: dict = {
"provider": None, "counter": None, "histogram": None,
"endpoint": None, "atexit": False, "failed": False,
}
_shared_meter_lock = threading.Lock()
def _build_meter_provider(endpoint: str, metric_reader: Any = None) -> Any:
"""惰性建私有 MeterProvider(装 OTLP metric reader);失败抛给上层 best-effort 兜。仅函数体内 import otel(惰性红线)。
:param endpoint: Collector OTLP/HTTP base( http://100.64.0.8:4318;真正发时拼 /v1/metrics)
:param metric_reader: 仅测试注入给一个 InMemoryMetricReader(record .get_metrics_data() 立即可断言,
无网络无导出周期);生产恒 None 建真 OTLPMetricExporter + PeriodicExportingMetricReader 异步后台发
:return: 配好 reader + service.name resource 的私有 MeterProvider(不碰 otel 全局单例)
"""
from opentelemetry.sdk.metrics import MeterProvider # noqa: PLC0415
from opentelemetry.sdk.resources import Resource # noqa: PLC0415
# resource.service.name = cheap-gen-worker:与 span 同服务名,Tempo/Prometheus 按它把便宜档这条线归一类。
resource = Resource.create({"service.name": SERVICE_NAME})
if metric_reader is not None:
# 测试路:注入 InMemoryMetricReader,同步聚合、无导出周期,record 后立即 get_metrics_data() 断言。
return MeterProvider(metric_readers=[metric_reader], resource=resource)
# 生产路:代理旁路(与 span 同源 T3-4)必须在建 OTLPMetricExporter(内部起 requests session)之前装——把
# Collector host 并入 NO_PROXY,否则 trust_env 的 http 客户端把发往 100.64.0.8 的上报经 clash fake-ip 转走被拦。
try:
from worker import client # noqa: PLC0415
client.install_proxy_bypass(endpoint)
except Exception as exc: # best-effort:装旁路失败不阻断建 provider(host 级 NO_PROXY 可能已由启动/span 侧装过)
_warn_once("proxy-bypass-metric", f"装 OTLP metric 代理旁路失败(忽略,可能已装过):{exc}")
from opentelemetry.exporter.otlp.proto.http.metric_exporter import ( # noqa: PLC0415
OTLPMetricExporter,
)
from opentelemetry.sdk.metrics.export import ( # noqa: PLC0415
PeriodicExportingMetricReader,
)
# OTLP/HTTP metric exporter:显式给 metrics 全路径 base + /v1/metrics(与 span 的 /v1/traces 同策略,
# 避免不同 SDK 版本对 base/自动拼路径行为不一致;Collector 的 OTLP/HTTP receiver 收 /v1/metrics)。
exporter = OTLPMetricExporter(endpoint=endpoint.rstrip("/") + "/v1/metrics")
# 周期性导出 reader:后台攒批定期发,不卡生成主链(observe-only 不阻塞)。
reader = PeriodicExportingMetricReader(exporter)
return MeterProvider(metric_readers=[reader], resource=resource)
def _get_shared_meter_instruments(endpoint: str, metric_reader: Any = None):
"""取进程级共享 (counter, histogram);惰性建一次、并发安全;建失败 best-effort 返回 (None, None)、永久降级(不每条重试)。
差异④:便宜档 Service 长驻,provider per-run 会逐局泄漏后台上报线程; MeterProvider/instrument 做进程级单例
"""
# 快路:已建好直接返回;已判失败直接 (None, None)(永久降级)。
if _shared_meter["counter"] is not None:
return _shared_meter["counter"], _shared_meter["histogram"]
if _shared_meter["failed"]:
return None, None
with _shared_meter_lock:
# 双检:抢锁期间别的线程可能已建好 / 已判失败。
if _shared_meter["counter"] is not None:
return _shared_meter["counter"], _shared_meter["histogram"]
if _shared_meter["failed"]:
return None, None
try:
provider = _build_meter_provider(endpoint, metric_reader=metric_reader)
meter = provider.get_meter("cheap-gen-worker", "1.0.0")
counter = meter.create_counter(
METRIC_TOKENS,
unit="{token}",
description="便宜档生成每次 run 的 LLM token 用量(kind=input/output/cached);"
"面五不发 token 计数,面二补(Prometheus 侧 llm_tokens_total{kind})",
)
histogram = meter.create_histogram(
METRIC_CACHE_HIT_RATE,
unit="1",
description="便宜档生成每次 run 的缓存命中率 cached/input([0,1]);面五不发,面二补",
explicit_bucket_boundaries_advisory=list(_CACHE_HIT_RATE_BUCKETS),
)
_shared_meter.update({
"provider": provider, "counter": counter,
"histogram": histogram, "endpoint": endpoint,
})
# 进程退出兜底 flush:注册一次(避免 per-run 注册堆积)。
if not _shared_meter["atexit"]:
import atexit # noqa: PLC0415
atexit.register(flush_shared_meter)
_shared_meter["atexit"] = True
_warn_once("meter-init-ok",
f"便宜档 OTLP metric 已开(service.name={SERVICE_NAME}{endpoint}/v1/metrics)")
return counter, histogram
except Exception as exc: # best-effort:缺依赖 / 建 provider 失败 → 永久降级不发,不抛、不阻塞生成
_shared_meter["failed"] = True
_warn_once("meter-init-fail",
f"建 OTLP metric→Collector 通道失败,便宜档降级为不发 metric:{exc}")
return None, None
def flush_shared_meter() -> None:
"""进程退出 / 收口把 PeriodicExportingMetricReader 攒批里残留的 metric 尽力发出(best-effort,绝不抛)。"""
provider = _shared_meter.get("provider")
if provider is None:
return
try:
provider.force_flush(timeout_millis=5000) # 给上限超时,避免收口/退出卡住
except Exception as exc: # best-effort:flush 失败只告警
_warn_once("meter-flush-fail", f"force_flush 残留 metric 失败(忽略):{exc}")
def reset_shared_meter_for_test() -> None:
"""仅测试用:清进程级共享 meter 单例 + 告警去重集,让各用例互不串(生产路绝不调)。"""
provider = _shared_meter.get("provider")
if provider is not None:
try:
provider.shutdown()
except Exception: # noqa: BLE001 —— 测试清理 best-effort
pass
_shared_meter.update({
"provider": None, "counter": None, "histogram": None,
"endpoint": None, "failed": False,
})
_WARNED.clear()
def record_llm_token_metrics(
input_tokens: int,
output_tokens: int,
cached_tokens: Optional[int] = None,
*,
endpoint: Optional[str] = None,
metric_reader: Any = None,
) -> None:
"""一次生成 run 收口时发面二两个 token 维度 metric(默认关 / best-effort / 绝不抛、绝不阻断生成)。
调用契约:生成 run 收尾拿到本次累计 token 后调一次;取不到 endpoint(默认关)直接返回字节不变
Args:
input_tokens: 本次 run 累计输入(prompt)token CLI Σ RecordingModel.records[0](=usage_sum()[0]);
Service Σ ModelCallEndEvent.input_tokens
output_tokens: 本次 run 累计输出(completion)token(同上, records[1] / ModelCallEndEvent.output_tokens)
cached_tokens: 本次 run 命中缓存的 prompt token ** CLI 路可得(Σ RecordingOpenAIChatModel.records[2],
usage.cache_input_tokens 归一)**;生产 Service 路的框架默认 model recording wrapper
ModelCallEndEvent 不带 cached,故传 **None** 表示该路拿不到 不发 cached 计数不发
cache_hit_rate(如实降级,绝不塞 0 伪造成命中率 0)
endpoint: 显式 Collector base(一般留空, env CHEAP_OTLP_ENDPOINT / infra_config)
metric_reader: 仅测试注入 InMemoryMetricReader;生产恒 None 走真 OTLP metric exporter
"""
url = resolve_otlp_endpoint(endpoint)
if not url:
# 默认关:取不到 endpoint 即不发(与 span sink 默认关同源,现有便宜档真跑字节不变)。
return
try:
counter, histogram = _get_shared_meter_instruments(url, metric_reader=metric_reader)
if counter is None:
return # 建通道失败已永久降级(告警已落),observe-only 不影响主链
in_tok = int(input_tokens or 0)
out_tok = int(output_tokens or 0)
# token 计数:input / output 恒发(kind 属性区分)。
counter.add(in_tok, {"kind": "input"})
counter.add(out_tok, {"kind": "output"})
if cached_tokens is not None:
# cached 可得(CLI 路):发 cached 计数 + 命中率;cached 截到 [0, input](防脏数据命中数超 prompt)。
cached = max(0, min(int(cached_tokens or 0), in_tok)) if in_tok > 0 else max(0, int(cached_tokens or 0))
counter.add(cached, {"kind": "cached"})
# cache_hit_rate 仅 input>0 有意义(除零无值);命中率截到 [0,1]。
if in_tok > 0:
histogram.record(min(1.0, cached / in_tok))
except Exception as exc: # best-effort:发 metric 任何异常只告警一次,绝不抛、绝不连累生成主链
_warn_once("metric-emit-fail", f"发 OTLP metric 失败(best-effort 计告警):{exc}")
if __name__ == "__main__": # pragma: no cover —— 本地自检:无 endpoint 时应拿到 no-op,__call__/shutdown 不炸 if __name__ == "__main__": # pragma: no cover —— 本地自检:无 endpoint 时应拿到 no-op,__call__/shutdown 不炸
s = make_cheap_otlp_sink(trace_id="demo-trace-001") s = make_cheap_otlp_sink(trace_id="demo-trace-001")
print("sink type:", type(s).__name__) # 期望 _NoopSink(未设 CHEAP_OTLP_ENDPOINT) print("sink type:", type(s).__name__) # 期望 _NoopSink(未设 CHEAP_OTLP_ENDPOINT)

View File

@ -164,6 +164,11 @@ def _get_collector_cls():
self._repair = repair self._repair = repair
self._tracer = tracer self._tracer = tracer
self._flush_logged = False # 工单 f:幂等打印标记——多次 _flush 只打一行「收口采集落盘」(落盘覆盖语义不变) self._flush_logged = False # 工单 f:幂等打印标记——多次 _flush 只打一行「收口采集落盘」(落盘覆盖语义不变)
# 面二观测:本 run 累计 token(从 ModelCallEndEvent 只读嗅探)。生产 Service 路 in/out 可得;cached 拿不到
# (ModelCallEndEvent 不带 cached,唯一经手 cached 的是 tier2 共用 breaker,不动它)→ 发 metric 时 cached 传 None 降级。
self._tok_in = 0
self._tok_out = 0
self._metric_emitted = False # 面二 metric 只发一次(_flush 三条降级路会调多次,防重复计数)
async def on_reply(self, agent, input_kwargs, next_handler): async def on_reply(self, agent, input_kwargs, next_handler):
# C1:框架 _agent.py:615 先 yield ReplyEndEvent 再 yield finish Msg 再收尾;collector 是最外层 on_reply, # C1:框架 _agent.py:615 先 yield ReplyEndEvent 再 yield finish Msg 再收尾;collector 是最外层 on_reply,
@ -171,6 +176,14 @@ def _get_collector_cls():
# 关闭 driver 的 read-before-flush 竞态。此时最后一次模型调用已 accumulate、breaker.spent_rmb 为终值。 # 关闭 driver 的 read-before-flush 竞态。此时最后一次模型调用已 accumulate、breaker.spent_rmb 为终值。
try: try:
async for evt in next_handler(**input_kwargs): async for evt in next_handler(**input_kwargs):
# 面二观测:只读嗅探每次模型调用的 token 用量(in/out);best-effort,绝不改事件、绝不抛。
# 在既有事件透传循环里嗅探,不新挂 on_model_call(不碰模型调用流式关键路径,零咬主链风险)。
if type(evt).__name__ == "ModelCallEndEvent":
try:
self._tok_in += int(getattr(evt, "input_tokens", 0) or 0)
self._tok_out += int(getattr(evt, "output_tokens", 0) or 0)
except Exception: # noqa: BLE001 —— 嗅探纯旁路,解析异常忽略,绝不咬生成
pass
if type(evt).__name__ == "ReplyEndEvent": if type(evt).__name__ == "ReplyEndEvent":
self._flush() # 正常路:REPLY_END 前落盘 self._flush() # 正常路:REPLY_END 前落盘
yield evt # 纯旁路透传,不拦截、不修改 yield evt # 纯旁路透传,不拦截、不修改
@ -230,6 +243,27 @@ def _get_collector_cls():
f"repairs={summary['repairs']} softTripped={summary['budgetSoftTripped']}", flush=True) f"repairs={summary['repairs']} softTripped={summary['budgetSoftTripped']}", flush=True)
except Exception as e: # noqa: BLE001 —— 采集 best-effort,失败绝不影响生成 except Exception as e: # noqa: BLE001 —— 采集 best-effort,失败绝不影响生成
print(f"[cheap-service] 收口采集失败(忽略,不影响生成):{type(e).__name__}: {e}", flush=True) print(f"[cheap-service] 收口采集失败(忽略,不影响生成):{type(e).__name__}: {e}", flush=True)
# 面二观测:发 token 计数 metric(独立于落盘成败、只发一次)。
self._emit_metrics()
def _emit_metrics(self) -> None:
"""面二:发 llm_tokens{kind} 计数(默认关 / best-effort / 只发一次,防 _flush 多调重复计数)。
生产 Service cached 拿不到( on_reply 嗅探注) cached_tokens None,record_llm_token_metrics
据此只发 input/output 计数不发 cached 计数与 cache_hit_rate(如实降级,不伪造命中率)
取不到 CHEAP_OTLP_ENDPOINT no-op( trace sink 默认关同源);面五不发 token 计数,不双计;
发送任何异常吞掉,绝不影响生成
"""
if self._metric_emitted:
return
self._metric_emitted = True
try:
import cheap_otlp_sink # noqa: PLC0415 惰性 import(顶层不牵 otel)
cheap_otlp_sink.record_llm_token_metrics(self._tok_in, self._tok_out, None)
except Exception as e: # noqa: BLE001 —— 面二 metric best-effort,绝不影响生成主链
print(f"[cheap-service] 面二 metric 发送异常(忽略,不影响生成):{type(e).__name__}: {e}",
flush=True)
_COLLECTOR_CLS = _CheapRunCollector _COLLECTOR_CLS = _CheapRunCollector
return _COLLECTOR_CLS return _COLLECTOR_CLS

View File

@ -385,6 +385,19 @@ async def run_studio(game_id, brief, *, max_iters=40, max_resumes=6, max_tokens=
ev_dir.mkdir(parents=True, exist_ok=True) ev_dir.mkdir(parents=True, exist_ok=True)
(ev_dir / "run-summary.json").write_text(json.dumps(summary, ensure_ascii=False, indent=2), encoding="utf-8") (ev_dir / "run-summary.json").write_text(json.dumps(summary, ensure_ascii=False, indent=2), encoding="utf-8")
_rec(f"run-summary 落盘 → {ev_dir / 'run-summary.json'} wallSec={summary['wallSec']}") _rec(f"run-summary 落盘 → {ev_dir / 'run-summary.json'} wallSec={summary['wallSec']}")
# ── 面二观测:发 llm_tokens{kind} 计数 + llm_cache_hit_rate(默认关 / best-effort;绝不咬生成)──
# 本路(CLI/批跑)model 是 RecordingOpenAIChatModel,records 每条 (in,out,cached) 三元组,cached 由
# usage.cache_input_tokens 归一 → 三项全可得(生产 Service 路框架默认 model 无 records,cached 拿不到、
# 见 cheap_service_app collector)。取不到 CHEAP_OTLP_ENDPOINT 即 no-op,与 trace sink 默认关同源;
# 面五(game-cloud GenMetrics)不发 token 计数,故不双计。发送异常只报告、绝不阻断生成。
try:
import cheap_otlp_sink # noqa: PLC0415 惰性 import(顶层不牵 otel)
_cached_tok = sum(r[2] for r in model.records if len(r) > 2)
cheap_otlp_sink.record_llm_token_metrics(tok_in, tok_out, _cached_tok)
except Exception as _me: # noqa: BLE001 面二 metric best-effort,绝不影响生成主链
_rec(f"面二 metric 发送异常(已忽略,不影响生成):{type(_me).__name__}: {_me}")
# 注n=5 收敛环批跑RunRecord/batch_run是 plan deferredspike 单局 run-summary 已含等价字段。 # 注n=5 收敛环批跑RunRecord/batch_run是 plan deferredspike 单局 run-summary 已含等价字段。
return summary return summary

View File

@ -180,6 +180,124 @@ def test_call_unreachable_collector_does_not_raise():
S.reset_shared_tracer_for_test() S.reset_shared_tracer_for_test()
# ══════════════════════════════════════════════════════════════════════════════
# 面二 · Python 侧 metric(llm_tokens{kind} / llm_cache_hit_rate)——注入 InMemoryMetricReader 断言
# ──────────────────────────────────────────────────────────────────────────────
# 只本地跑、零网络:token 计数记入 / 缓存命中率记入 / cached 不可得(None)如实降级 / 默认关不发 /
# best-effort 不咬主链。真发 Collector→Prometheus 见 metric 排内网恢复后 mini-infra 验。
# ══════════════════════════════════════════════════════════════════════════════
from opentelemetry.sdk.metrics.export import InMemoryMetricReader # noqa: E402
def _collect_points(reader):
"""从 InMemoryMetricReader 抽 {metric_name: [(attrs_dict, value_or_hist)]}。
counter data point .value;histogram data point .count/.sum 分别取出便于断言
"""
data = reader.get_metrics_data()
out: dict = {}
if data is None:
return out
for rm in data.resource_metrics:
for sm in rm.scope_metrics:
for metric in sm.metrics:
pts = []
for dp in metric.data.data_points:
attrs = dict(dp.attributes or {})
if hasattr(dp, "value"):
pts.append((attrs, dp.value)) # counter NumberDataPoint
else:
pts.append((attrs, {"count": dp.count, "sum": dp.sum})) # histogram
out[metric.name] = pts
return out
def test_record_metrics_default_off_emits_nothing(monkeypatch):
"""默认关(resolve → None):record 直接返回、根本不建 provider,注入的 reader 收不到任何 metric。"""
S.reset_shared_meter_for_test()
monkeypatch.setattr(S, "resolve_otlp_endpoint", lambda explicit=None: None)
reader = InMemoryMetricReader()
S.record_llm_token_metrics(1000, 500, 200, metric_reader=reader) # 默认关:不发
got = _collect_points(reader)
assert S.METRIC_TOKENS not in got and S.METRIC_CACHE_HIT_RATE not in got, "默认关时绝不发 metric"
S.reset_shared_meter_for_test()
def test_record_token_counts_and_cache_hit_rate():
"""cached 可得(CLI 路):llm_tokens 三 kind 计数正确 + llm_cache_hit_rate 记入 cached/input。"""
S.reset_shared_meter_for_test()
reader = InMemoryMetricReader()
# input=1000 output=400 cached=250 → 命中率 0.25
S.record_llm_token_metrics(1000, 400, 250,
endpoint="http://collector.test:4318", metric_reader=reader)
got = _collect_points(reader)
# llm_tokens{kind} 三条 counter data point
tokens = {a.get("kind"): v for a, v in got.get(S.METRIC_TOKENS, [])}
assert tokens.get("input") == 1000
assert tokens.get("output") == 400
assert tokens.get("cached") == 250, "cached 可得时必须发 cached 计数"
# llm_cache_hit_rate:一次 record,count=1,sum≈0.25
hist = got.get(S.METRIC_CACHE_HIT_RATE, [])
assert len(hist) == 1
_, h = hist[0]
assert h["count"] == 1
assert abs(h["sum"] - 0.25) < 1e-6, f"命中率应为 cached/input=0.25,实为 {h['sum']}"
S.reset_shared_meter_for_test()
def test_record_cached_none_degrades_no_cached_no_rate():
"""cached 不可得(生产 Service 路传 None):只发 input/output 计数,不发 cached 计数、不发 cache_hit_rate(不伪造)。"""
S.reset_shared_meter_for_test()
reader = InMemoryMetricReader()
S.record_llm_token_metrics(800, 300, None,
endpoint="http://collector.test:4318", metric_reader=reader)
got = _collect_points(reader)
tokens = {a.get("kind"): v for a, v in got.get(S.METRIC_TOKENS, [])}
assert tokens.get("input") == 800 and tokens.get("output") == 300
assert "cached" not in tokens, "cached=None 时绝不发 cached 计数(如实降级,不塞 0 伪造)"
# cache_hit_rate 完全不 record(metric 不出现,或无 data point)
assert not got.get(S.METRIC_CACHE_HIT_RATE), "cached 不可得时绝不发 cache_hit_rate(不伪造命中率)"
S.reset_shared_meter_for_test()
def test_record_cache_hit_rate_clamped_on_dirty_cached():
"""脏数据 cached > input:命中率截到 1.0、cached 计数截到 input(防 >100% 命中率与负值)。"""
S.reset_shared_meter_for_test()
reader = InMemoryMetricReader()
S.record_llm_token_metrics(100, 50, 999, # cached 远超 input(脏)
endpoint="http://collector.test:4318", metric_reader=reader)
got = _collect_points(reader)
tokens = {a.get("kind"): v for a, v in got.get(S.METRIC_TOKENS, [])}
assert tokens.get("cached") == 100, "cached 应截到 input=100"
_, h = got.get(S.METRIC_CACHE_HIT_RATE, [])[0]
assert abs(h["sum"] - 1.0) < 1e-6, "命中率应截到 1.0"
S.reset_shared_meter_for_test()
def test_record_input_zero_no_rate_but_counts_ok():
"""input=0(除零无命中率):仍发 token 计数(可能都 0),但不 record cache_hit_rate(避免除零)。"""
S.reset_shared_meter_for_test()
reader = InMemoryMetricReader()
S.record_llm_token_metrics(0, 0, 0,
endpoint="http://collector.test:4318", metric_reader=reader)
got = _collect_points(reader)
assert S.METRIC_TOKENS in got, "input=0 也发 token 计数(计数可为 0)"
assert not got.get(S.METRIC_CACHE_HIT_RATE), "input=0 时不 record 命中率(除零无意义)"
S.reset_shared_meter_for_test()
def test_record_metrics_best_effort_unreachable_does_not_raise():
"""真走 OTLP metric exporter(不注入 reader),端点不可达 → record 绝不抛(best-effort 铁律)。"""
S.reset_shared_meter_for_test()
try:
S.record_llm_token_metrics(100, 50, 10, endpoint="http://127.0.0.1:1") # 连接会被拒,后台异步失败
S.flush_shared_meter() # flush 也 best-effort、不抛
except Exception as e: # noqa: BLE001
raise AssertionError(f"Collector 不可达时 record 绝不能抛(best-effort 铁律),但抛了:{e}") from e
finally:
S.reset_shared_meter_for_test()
if __name__ == "__main__": if __name__ == "__main__":
import pytest # noqa: PLC0415 import pytest # noqa: PLC0415
sys.exit(pytest.main([__file__, "-v"])) sys.exit(pytest.main([__file__, "-v"]))

View File

@ -205,3 +205,38 @@ def test_flush_landing_line_once_on_crash_path(tmp_path, monkeypatch, capsys):
asyncio.run(_drive()) asyncio.run(_drive())
out = capsys.readouterr().out out = capsys.readouterr().out
assert out.count("收口采集落盘") == 1, "崩溃路两次 _flush 也只打一行收口采集落盘(工单 f)" assert out.count("收口采集落盘") == 1, "崩溃路两次 _flush 也只打一行收口采集落盘(工单 f)"
def test_collector_sniffs_tokens_and_emits_metric_once_cached_none(tmp_path, monkeypatch):
"""面二 Service 路 wiring:collector 只读嗅探 ModelCallEndEvent 累加 in/out,REPLY_END 收口发一次 metric;
cached None(生产 Service 路框架默认 model records拿不到 cached,如实降级);_flush 多调 metric 只发一次"""
col = _make_collector(A, "70066", tmp_path, monkeypatch)
# monkeypatch recorder:只记被调参数(record 本体的默认关/best-effort/降级已在 test_cheap_otlp_sink 覆盖,这里只验 wiring)。
import cheap_otlp_sink
calls = []
monkeypatch.setattr(cheap_otlp_sink, "record_llm_token_metrics",
lambda i, o, c=None, **k: calls.append((i, o, c)))
class ModelCallEndEvent: # 类名即判据(collector 认 type(evt).__name__ == "ModelCallEndEvent")
def __init__(self, i, o):
self.input_tokens, self.output_tokens = i, o
class ReplyEndEvent:
pass
async def _fake_next(**_):
yield ModelCallEndEvent(1200, 800) # 第一次模型调用
yield ModelCallEndEvent(300, 100) # 第二次
yield ReplyEndEvent() # 收口(触发 _flush → _emit_metrics)
async def _drive():
async for _evt in col.on_reply(None, {}, _fake_next):
pass
asyncio.run(_drive())
assert col._tok_in == 1500 and col._tok_out == 900, "嗅探累加 in/out(1200+300 / 800+100)"
# on_reply 内 REPLY_END + finally 已两调 _flush;再手调两次(模拟 except/finally 最坏叠加),metric 仍只发一次。
col._flush()
col._flush()
assert calls == [(1500, 900, None)], "收口只发一次 metric;cached=None(生产 Service 路降级、不伪造命中率)"