games-development-ai/cheap-worker/tests/test_cheap_otlp_sink.py
lili daa49e5864
Some checks failed
contract-gates / contract-gates (push) Has been cancelled
docs-gate / docs-gate (push) Has been cancelled
feat(obs): 便宜档主生成路 OTLP span sink + 双入口 fan-out 并接(阶段四观测 波③ T3-1)
便宜档 Service(cheap_service_app)+ CLI(cheap_studio)主生成路此前只经 make_jsonl_sink
写 trace.jsonl、不产 OTLP span,观测栈 Tempo 看不到便宜档生成。新接 cheap_otlp_sink:与
make_jsonl_sink 同签名,把每步 TraceStep 推成 OTLP span 发 mini-infra Collector
(service.name=cheap-gen-worker),span 属性映射对齐 tier2 studio_sink 的标准 gen-ai 键。

- 默认关(env CHEAP_OTLP_ENDPOINT / infra_config 取不到即 no-op),build_trace_sink 原样返回
  jsonl sink、现有行为字节不变;开则 fan-out(jsonl 先落盘照旧、OTLP 额外发一份,各 best-effort)。
- best-effort 铁律:缺 otel / Collector 不可达 / 发送失败 → 记一次告警、永久降级 no-op、绝不抛、
  绝不咬生成主链(设计 §8);observe-only 从不回写生成/门判定字段。
- 惰性 import otel(顶层零重依赖,对齐 build_cheap_app 惰性红线)+ 代理旁路(T3-4,建 exporter 前
  把 Collector host 并入 NO_PROXY,绕 clash fake-ip)。
- provider 做进程级共享单例(非 per-run):长驻 Service 中间件工厂 per-turn 建 sink,per-run 会逐局
  泄漏后台上报线程;trace_id 只作 per-span 属性,BatchSpanProcessor 一条后台线程跨局复用,atexit 兜底 flush。
- 双入口都接(Service :273 + CLI :214),fan-out 不破坏既有 jsonl 落盘。

终审(亲跑非照搬):cheap-worker/.venv 全套 350 passed/1 skipped(含 span sink 11 例)、零回归、
py_compile 三文件过、docs-gate 七检绿。跨机 OTLP 命门已验(mini-desktop→mini-infra:4318=200)。
真发 Collector→Tempo 见 span 排 mini-desktop 起服务后验。

T3-3 metrics 空口按红线#7(不留孤儿)剔除:本波无生产调用方,真值填充随波⑤ T5-2 带埋点调用方一起落。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-05 09:16:16 -07:00

186 lines
7.9 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()
if __name__ == "__main__":
import pytest # noqa: PLC0415
sys.exit(pytest.main([__file__, "-v"]))