补面五(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>
304 lines
14 KiB
Python
304 lines
14 KiB
Python
"""test_cheap_otlp_sink.py — 便宜档 OTLP span sink 单测(阶段四观测 波③ T3-1)。
|
|
|
|
只本地跑、零网络:default-off 降 no-op / fan-out 不吞 jsonl / best-effort 不咬主链 / span 属性映射
|
|
(经注入内存 exporter 断言,不发真 Collector)。真发 Collector→Tempo 见 span 排内网恢复后 mini-desktop 验。
|
|
|
|
跑:cheap-worker/.venv/bin/python -m pytest cheap-worker/tests/test_cheap_otlp_sink.py -v
|
|
"""
|
|
import os
|
|
import sys
|
|
from pathlib import Path
|
|
|
|
# sys.path:cheap-worker/ + tier2/gen-worker/(observability / service / worker 包)。
|
|
_HERE = os.path.dirname(os.path.abspath(__file__))
|
|
_CW = os.path.dirname(_HERE)
|
|
_REPO_ROOT = os.path.dirname(_CW)
|
|
_GEN_WORKER = os.path.join(_REPO_ROOT, "tier2", "gen-worker")
|
|
for _p in (_CW, _GEN_WORKER):
|
|
if _p not in sys.path:
|
|
sys.path.insert(0, _p)
|
|
|
|
import cheap_otlp_sink as S # noqa: E402
|
|
from opentelemetry.sdk.trace.export.in_memory_span_exporter import ( # noqa: E402
|
|
InMemorySpanExporter,
|
|
)
|
|
|
|
|
|
# ── resolve_otlp_endpoint:默认关 / env 覆盖 ──────────────────────────────────────
|
|
def test_resolve_endpoint_default_off(monkeypatch):
|
|
"""无 env + infra 无 observability.otlp_endpoint → None(默认关)。"""
|
|
monkeypatch.delenv(S.ENV_OTLP_ENDPOINT, raising=False)
|
|
# 隔离真 infra.yaml:强制 infra_config[observability].otlp_endpoint 缺省。
|
|
from service import infra_config
|
|
monkeypatch.setattr(infra_config, "get", lambda section, key, default=None: default)
|
|
assert S.resolve_otlp_endpoint() is None
|
|
|
|
|
|
def test_resolve_endpoint_env_override(monkeypatch):
|
|
"""env CHEAP_OTLP_ENDPOINT 设了 → 归一化(去尾斜杠)返回。"""
|
|
monkeypatch.setenv(S.ENV_OTLP_ENDPOINT, "http://100.64.0.8:4318/")
|
|
assert S.resolve_otlp_endpoint() == "http://100.64.0.8:4318"
|
|
|
|
|
|
# ── make_cheap_otlp_sink:默认关 → no-op,__call__/shutdown 不炸 ──────────────────
|
|
def test_make_sink_off_is_noop(monkeypatch):
|
|
monkeypatch.setattr(S, "resolve_otlp_endpoint", lambda explicit=None: None)
|
|
sink = S.make_cheap_otlp_sink(trace_id="t1")
|
|
assert isinstance(sink, S._NoopSink)
|
|
assert sink.enabled is False
|
|
sink({"traceId": "t1", "step": 0, "ext": {"raw": {"event": "X"}}}) # 不抛
|
|
sink.shutdown() # 不抛
|
|
|
|
|
|
# ── build_trace_sink:默认关 = 零包装,原样返回 jsonl(现有行为字节不变) ───────────
|
|
def test_build_trace_sink_off_returns_jsonl_identity(monkeypatch):
|
|
monkeypatch.setattr(S, "resolve_otlp_endpoint", lambda explicit=None: None)
|
|
|
|
def _jsonl(step): # 桩 jsonl sink
|
|
return None
|
|
|
|
got = S.build_trace_sink(_jsonl, trace_id="g1")
|
|
assert got is _jsonl, "默认关时必须原样返回同一个 jsonl sink 对象(零包装、字节不变)"
|
|
|
|
|
|
def test_build_trace_sink_off_jsonl_none_stays_none(monkeypatch):
|
|
monkeypatch.setattr(S, "resolve_otlp_endpoint", lambda explicit=None: None)
|
|
assert S.build_trace_sink(None, trace_id="g1") is None
|
|
|
|
|
|
# ── build_trace_sink:开 → fan-out,jsonl 照落 + OTLP 额外发一份 ───────────────────
|
|
def test_build_trace_sink_on_fans_out_to_both(monkeypatch):
|
|
S.reset_shared_tracer_for_test()
|
|
mem = InMemorySpanExporter()
|
|
seen = []
|
|
|
|
def _jsonl(step):
|
|
seen.append(step) # jsonl 侧照常收到(落盘保留不动的等价验证)
|
|
|
|
fan = S.build_trace_sink(_jsonl, trace_id="gid-9", endpoint="http://collector.test:4318",
|
|
span_exporter=mem)
|
|
assert isinstance(fan, S._FanoutSink)
|
|
fan({"traceId": "gid-9", "step": 0, "cost": None, "verdict": None,
|
|
"timestamp": "2026-07-05T10:00:00", "ext": {"raw": {"event": "ReplyStartEvent"}}})
|
|
# jsonl 侧收到了(没被 OTLP 并接吞掉)
|
|
assert len(seen) == 1 and seen[0]["step"] == 0
|
|
# OTLP 侧也发出了一个 span
|
|
spans = mem.get_finished_spans()
|
|
assert len(spans) == 1
|
|
S.reset_shared_tracer_for_test()
|
|
|
|
|
|
# ── fan-out best-effort:一侧抛不连累另一侧、不向外抛 ─────────────────────────────
|
|
def test_fanout_one_side_raises_does_not_break_other_or_raise():
|
|
calls = {"good": 0}
|
|
|
|
def _bad(step):
|
|
raise IOError("模拟 jsonl 磁盘满")
|
|
|
|
def _good(step):
|
|
calls["good"] += 1
|
|
|
|
fan = S._FanoutSink([_bad, _good])
|
|
fan({"traceId": "t", "step": 0}) # 不抛
|
|
assert calls["good"] == 1, "一侧抛异常不得连累另一侧"
|
|
|
|
|
|
def test_fanout_shutdown_forwards_and_skips_plain(monkeypatch):
|
|
closed = {"n": 0}
|
|
|
|
class _WithShutdown:
|
|
def __call__(self, step):
|
|
return
|
|
|
|
def shutdown(self):
|
|
closed["n"] += 1
|
|
|
|
def _plain(step): # 无 shutdown,应被跳过、不炸
|
|
return
|
|
|
|
fan = S._FanoutSink([_plain, _WithShutdown()])
|
|
fan.shutdown()
|
|
assert closed["n"] == 1
|
|
|
|
|
|
# ── span 属性映射(注入内存 exporter,零网络):对齐 studio_sink 标准 gen-ai 键 ──────
|
|
def test_span_mapping_service_name_conversation_and_tokens():
|
|
S.reset_shared_tracer_for_test()
|
|
mem = InMemorySpanExporter()
|
|
sink = S.make_cheap_otlp_sink(trace_id="conv-42", endpoint="http://collector.test:4318",
|
|
span_exporter=mem)
|
|
# ModelCallEnd 形态(带 token)+ action 段(带工具名)
|
|
sink({
|
|
"traceId": "conv-42", "step": 3,
|
|
"cost": {"tokens": {"in": 200, "out": 80}},
|
|
"verdict": None, "timestamp": "2026-07-05T10:00:01",
|
|
"ext": {"action": {"tool": "write_source"}, "raw": {"event": "ModelCallEndEvent"}},
|
|
})
|
|
spans = mem.get_finished_spans()
|
|
assert len(spans) == 1
|
|
sp = spans[0]
|
|
# resource.service.name = cheap-gen-worker(Tempo 按它归类便宜档线)
|
|
assert sp.resource.attributes.get("service.name") == "cheap-gen-worker"
|
|
attrs = dict(sp.attributes)
|
|
assert attrs.get("gen_ai.conversation.id") == "conv-42" # 与 studio_sink 同键(跨线归组)
|
|
assert attrs.get("gen_ai.usage.input_tokens") == 200
|
|
assert attrs.get("gen_ai.usage.output_tokens") == 80
|
|
assert attrs.get("gen_ai.tool.name") == "write_source"
|
|
assert attrs.get("agentscope.cheap.step") == 3
|
|
assert sp.name == "step 3 · write_source"
|
|
S.reset_shared_tracer_for_test()
|
|
|
|
|
|
def test_span_mapping_error_verdict_sets_error_status():
|
|
S.reset_shared_tracer_for_test()
|
|
mem = InMemorySpanExporter()
|
|
sink = S.make_cheap_otlp_sink(trace_id="conv-err", endpoint="http://collector.test:4318",
|
|
span_exporter=mem)
|
|
sink({
|
|
"traceId": "conv-err", "step": 5, "cost": None,
|
|
"verdict": "error", "timestamp": "2026-07-05T10:00:02",
|
|
"ext": {"observation": {"kind": "ToolResultEndEvent"}, "raw": {"event": "ToolResultEndEvent"}},
|
|
})
|
|
from opentelemetry.trace import StatusCode
|
|
sp = mem.get_finished_spans()[0]
|
|
assert sp.status.status_code == StatusCode.ERROR, "error verdict 应把 span 标 ERROR(如实反映,不改判定)"
|
|
S.reset_shared_tracer_for_test()
|
|
|
|
|
|
# ── best-effort:endpoint 设了但 Collector 不可达 → __call__ 不抛(BatchSpanProcessor 异步、后台失败静默)──
|
|
def test_call_unreachable_collector_does_not_raise():
|
|
S.reset_shared_tracer_for_test()
|
|
# 真走 OTLP exporter + BatchSpanProcessor(不注入内存 exporter),端点连接会被拒;enqueue 不阻塞、发送在后台失败。
|
|
sink = S.make_cheap_otlp_sink(trace_id="t-unreach", endpoint="http://127.0.0.1:1")
|
|
try:
|
|
sink({"traceId": "t-unreach", "step": 0, "cost": None, "verdict": None,
|
|
"timestamp": "", "ext": {"raw": {"event": "X"}}})
|
|
sink.shutdown() # flush 也 best-effort、不抛
|
|
except Exception as e: # noqa: BLE001
|
|
raise AssertionError(f"Collector 不可达时 sink 绝不能抛(best-effort 铁律),但抛了:{e}") from e
|
|
finally:
|
|
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__":
|
|
import pytest # noqa: PLC0415
|
|
sys.exit(pytest.main([__file__, "-v"]))
|