From 958f78f523a2fd799920e339078a46311b925f4e Mon Sep 17 00:00:00 2001 From: lili Date: Sun, 5 Jul 2026 23:57:40 -0700 Subject: [PATCH] =?UTF-8?q?feat(obs):=20=E9=98=B6=E6=AE=B5=E5=9B=9B?= =?UTF-8?q?=E8=A7=82=E6=B5=8B=E9=9D=A2=E4=BA=8C=20Python=20=E4=BE=A7=20tok?= =?UTF-8?q?en=20metric=E2=80=94=E2=80=94llm=5Ftokens{kind}=20+=20llm=5Fcac?= =?UTF-8?q?he=5Fhit=5Frate?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 补面五(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) --- cheap-worker/cheap_otlp_sink.py | 200 +++++++++++++++++++ cheap-worker/cheap_service_app.py | 34 ++++ cheap-worker/cheap_studio.py | 13 ++ cheap-worker/tests/test_cheap_otlp_sink.py | 118 +++++++++++ cheap-worker/tests/test_cheap_service_app.py | 35 ++++ 5 files changed, 400 insertions(+) diff --git a/cheap-worker/cheap_otlp_sink.py b/cheap-worker/cheap_otlp_sink.py index db3337a6..fe519579 100644 --- a/cheap-worker/cheap_otlp_sink.py +++ b/cheap-worker/cheap_otlp_sink.py @@ -562,6 +562,206 @@ def build_trace_sink( 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 不炸 s = make_cheap_otlp_sink(trace_id="demo-trace-001") print("sink type:", type(s).__name__) # 期望 _NoopSink(未设 CHEAP_OTLP_ENDPOINT) diff --git a/cheap-worker/cheap_service_app.py b/cheap-worker/cheap_service_app.py index 522f7586..9352ad3e 100644 --- a/cheap-worker/cheap_service_app.py +++ b/cheap-worker/cheap_service_app.py @@ -164,6 +164,11 @@ def _get_collector_cls(): self._repair = repair self._tracer = tracer 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): # 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 为终值。 try: 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": self._flush() # 正常路:REPLY_END 前落盘 yield evt # 纯旁路透传,不拦截、不修改 @@ -230,6 +243,27 @@ def _get_collector_cls(): f"repairs={summary['repairs']} softTripped={summary['budgetSoftTripped']}", flush=True) except Exception as e: # noqa: BLE001 —— 采集 best-effort,失败绝不影响生成 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 return _COLLECTOR_CLS diff --git a/cheap-worker/cheap_studio.py b/cheap-worker/cheap_studio.py index 6f94e7d8..f4782bc5 100644 --- a/cheap-worker/cheap_studio.py +++ b/cheap-worker/cheap_studio.py @@ -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 / "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']}") + + # ── 面二观测:发 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 deferred;spike 单局 run-summary 已含等价字段。 return summary diff --git a/cheap-worker/tests/test_cheap_otlp_sink.py b/cheap-worker/tests/test_cheap_otlp_sink.py index 5b486cd7..074e9911 100644 --- a/cheap-worker/tests/test_cheap_otlp_sink.py +++ b/cheap-worker/tests/test_cheap_otlp_sink.py @@ -180,6 +180,124 @@ def test_call_unreachable_collector_does_not_raise(): 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"])) diff --git a/cheap-worker/tests/test_cheap_service_app.py b/cheap-worker/tests/test_cheap_service_app.py index bf192121..5fb76f6e 100644 --- a/cheap-worker/tests/test_cheap_service_app.py +++ b/cheap-worker/tests/test_cheap_service_app.py @@ -205,3 +205,38 @@ def test_flush_landing_line_once_on_crash_path(tmp_path, monkeypatch, capsys): asyncio.run(_drive()) out = capsys.readouterr().out 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 路降级、不伪造命中率)"