diff --git a/cheap-worker/cheap_otlp_sink.py b/cheap-worker/cheap_otlp_sink.py new file mode 100644 index 00000000..266518e8 --- /dev/null +++ b/cheap-worker/cheap_otlp_sink.py @@ -0,0 +1,436 @@ +"""cheap_otlp_sink.py — 便宜档主生成路 · OTLP span sink(阶段四观测 波③ · 工单 T3-1)。 + +【这份解决什么问题】 + 便宜档 MVP 主生成路(Service `cheap_service_app` + CLI `cheap_studio`)现在只经 + observability.trace.make_jsonl_sink 把每步 trace 写进产物目录的 trace.jsonl,**不产 OTLP span**, + 观测栈(Tempo)看不到便宜档生成。本模块新接一个 OTLP span sink:与 make_jsonl_sink 同签名 + (吃一条 TraceStep dict),把每步推成一个 OTLP span 发给 mini-infra 的 OTel Collector + (100.64.0.8:4318 · OTLP/HTTP),让便宜档生成在 Tempo 里自成一条 trace(端到端挂到 Java trace 下 + 归波④,本波先让它自产 span)。jsonl 落盘保留不动,OTLP 只是额外发一份(见 build_trace_sink 并接)。 + +【与 tier2 studio_sink 的关系(照它的形态写,三点不同 + 一处生命周期改进)】 + 形态照 tier2 observability/studio_sink.py(全仓唯一 OTLP exporter):私有 TracerProvider + + OTLPSpanExporter(HTTP) + BatchSpanProcessor + best-effort 降级 + 惰性 import otel + **私有 provider + 不抢 otel 全局单例**。四点差异: + ① 端点指 **Collector**(http://100.64.0.8:4318)而非 AgentScope Studio 私有端点; + ② service.name = **cheap-gen-worker**(Studio 那条是 tier2-gen-worker),Tempo 按它归类便宜档这条线; + ③ 端点来源默认关(env CHEAP_OTLP_ENDPOINT / infra_config,取不到即 no-op),开关默认态零影响; + ④ **provider 做进程级单例、不做 per-run**:便宜档 Service 是长驻进程,中间件工厂 per-turn 建 sink; + 若照 studio_sink 每 run 建一个私有 TracerProvider + BatchSpanProcessor 后台线程,长驻进程会 + 逐局泄漏 provider/线程。故这里把 provider 建成进程级共享单例(trace_id 只作 per-span 属性), + BatchSpanProcessor 一条后台线程跨局复用,进程退出时经单一 atexit 兜底 flush。 + 标准 gen-ai 语义键(gen_ai.conversation.id / gen_ai.usage.* / gen_ai.tool.name)与 studio_sink **保持一致**, + 便于将来在 Tempo 里按 conversation.id 跨线归组;service.name 有意不同(按线区分)。 + +【best-effort 铁律(设计 §8 安全红线:观测栈挂掉绝不能咬生成主链)】 + - 取不到 endpoint → make_cheap_otlp_sink 返回 no-op(现有便宜档生成一字节不变); + - 建 provider / 发 span 失败(缺 otel wheel / Collector 不可达 / 网络错)→ best-effort 记一次告警、 + 永久降级 no-op、绝不抛、绝不阻断生成; + - 本 sink 只「读」TraceStep 往外发,从不回写任何生成 / 门判定字段(observe-only)。 + +【惰性 import(6c6g / cheap_service_app 顶层惰性红线)】 + opentelemetry SDK / exporter 只在函数体内 import,本模块顶层零重依赖 —— 对齐 cheap_service_app + 「顶层绝不 import agentscope / 重依赖」。真正建 provider / 发 span 只在装了 otel wheel 的 mini-desktop / Mac 发生。 + +【代理旁路(工单 T3-4)】 + OTLP 打 100.64.0.8:4318 走 Tailscale,必须绕本机 clash fake-ip 代理(198.18.x),否则 exporter 被拦。 + 建 exporter 前复用 worker.client.install_proxy_bypass 把 Collector host 并入 NO_PROXY —— 与 new-api / + 本地 Service 旁路同源(便宜档 new-api 网关恰好也在 100.64.0.8,host 级 NO_PROXY 已覆盖,这里再显式装一次幂等兜底)。 + +【开关 / 地址来源】(默认关闭 = 取不到即 no-op) + ① env CHEAP_OTLP_ENDPOINT(单点覆盖,临时开一次最省事;内网恢复后设为 http://100.64.0.8:4318); + ② service/infra_config.py 的 [observability].otlp_endpoint(若该分区/项存在;infra.yaml 内网入仓口径); + ③ 显式传给 make_cheap_otlp_sink(endpoint=...) 的参数。 + 三者都没有 → resolve_otlp_endpoint() 返回 None → make_cheap_otlp_sink 返回 no-op sink(现有真跑不受影响)。 +""" + +from __future__ import annotations + +import os +import threading +from typing import Any, Callable, Optional + + +# ── env 开关键(与 genconfig / infra_config 的 env 风格一致:全大写、CHEAP_ 前缀)── +# 取到非空即视为「开便宜档 OTLP 观测」并用作 Collector OTLP/HTTP base。 +ENV_OTLP_ENDPOINT = "CHEAP_OTLP_ENDPOINT" + +# ── 便宜档这条 trace 在 Collector/Tempo 里的服务名(OTLP resource service.name;Tempo 按它聚合归类)── +SERVICE_NAME = "cheap-gen-worker" + +# ── Collector OTLP/HTTP 规范端点(内网恢复后把 env CHEAP_OTLP_ENDPOINT 设成它即开;默认不自动开)── +# 仅作文档常量:resolve_otlp_endpoint 不拿它当兜底(拿它兜底=默认自动开,违反「默认关」),开关必须显式给。 +DEFAULT_COLLECTOR_ENDPOINT = "http://100.64.0.8:4318" + +# ── 一次性告警去重(同一 tag 只 print 一次,避免长 run 刷屏;对齐 studio_sink._warn_once 风格)── +_WARNED: set[str] = set() + +# ── 进程级共享 TracerProvider 单例(差异④:长驻 Service 逐局复用一条后台上报线程,不 per-run 泄漏)── +# _shared 记 provider/tracer/endpoint/atexit 标记;_shared_lock 保并发建单例线程安全(Service 并发多局)。 +_shared: dict = {"provider": None, "tracer": None, "endpoint": None, "atexit": False, "failed": False} +_shared_lock = threading.Lock() + +# 工具结果状态 → span ERROR 的判据(observe 段 verdict 决定 span 颜色;只如实反映、不改任何生成判定)。 +_ERROR_VERDICTS = {"error", "interrupted", "denied"} + + +def _warn_once(tag: str, msg: str) -> None: + """同一 tag 的告警只落一次(best-effort 告警去重,绝不抛)。""" + if tag in _WARNED: + return + _WARNED.add(tag) + print(f"[cheap-otlp-sink] {msg}", flush=True) + + +def resolve_otlp_endpoint(explicit: Optional[str] = None) -> Optional[str]: + """解析 Collector OTLP/HTTP base:env > infra_config[observability].otlp_endpoint > 显式参数;都没有 → None(=关闭)。 + + :param explicit: 调用方显式传入的 endpoint(最低优先级,作兜底)。 + :return: 归一化后的 OTLP base(去尾斜杠);取不到任何来源 → None(make_cheap_otlp_sink 据此降级 no-op)。 + """ + # ① env 单点覆盖(最高优先级)。 + env_val = os.environ.get(ENV_OTLP_ENDPOINT) + if env_val and env_val.strip(): + return env_val.strip().rstrip("/") + + # ② infra_config.py 的 [observability].otlp_endpoint(best-effort:读配置任何异常都不阻断)。 + try: + # 包内/直跑兼容导入(直跑时 service 包仍可解析;_bootstrap 已把 tier2/gen-worker 加进 sys.path)。 + try: + from service import infra_config # type: ignore # noqa: PLC0415 + except Exception: # pragma: no cover —— 直跑兜底:把 gen-worker/ 加进 sys.path 再取 + import sys + from pathlib import Path + + sys.path.insert(0, str(Path(__file__).resolve().parent.parent / "tier2" / "gen-worker")) + from service import infra_config # type: ignore # noqa: PLC0415 + + cfg_val = infra_config.get("observability", "otlp_endpoint", None) + if cfg_val and str(cfg_val).strip(): + return str(cfg_val).strip().rstrip("/") + except Exception as exc: # best-effort:读配置任何异常都不阻断,落一次告警后继续 + _warn_once("infra_config", f"读 infra_config[observability].otlp_endpoint 失败(忽略):{exc}") + + # ③ 显式参数兜底。 + if explicit and explicit.strip(): + return explicit.strip().rstrip("/") + + return None + + +def _build_provider(endpoint: str, span_exporter: Any = None) -> Any: + """惰性建一个私有 TracerProvider(装 exporter);失败抛给上层 best-effort 兜。仅在函数体内 import otel(惰性红线)。 + + :param endpoint: Collector OTLP/HTTP base(如 http://100.64.0.8:4318;真正发时拼 /v1/traces)。 + :param span_exporter: 仅测试注入用——给一个内存 exporter(走 SimpleSpanProcessor 同步导出,断言 span 属性用); + 生产恒 None → 建真 OTLPSpanExporter(HTTP)+ BatchSpanProcessor 异步后台发。 + :return: 配好 processor 的私有 TracerProvider(不碰 otel 全局单例)。 + """ + # 惰性 import:6c6g 无 opentelemetry SDK / exporter wheel,本模块顶层仍可裸 import(此处才触发依赖)。 + from opentelemetry.sdk.resources import Resource # noqa: PLC0415 + from opentelemetry.sdk.trace import TracerProvider # noqa: PLC0415 + from opentelemetry.sdk.trace.export import ( # noqa: PLC0415 + BatchSpanProcessor, + SimpleSpanProcessor, + ) + + # resource.service.name = cheap-gen-worker:Tempo 按服务名把便宜档这条线的 span 归一类。 + resource = Resource.create({"service.name": SERVICE_NAME}) + provider = TracerProvider(resource=resource) + + if span_exporter is not None: + # 测试路:同步 processor + 内存 exporter,__call__ 后立即可断言 span 属性(无需网络 / flush 时序)。 + provider.add_span_processor(SimpleSpanProcessor(span_exporter)) + return provider + + # 生产路:代理旁路(T3-4)必须在建 OTLPSpanExporter(内部起 httpx / requests session)之前装 —— 把 + # Collector host 并入 NO_PROXY,否则 trust_env 的 http 客户端把发往 100.64.0.8 的上报经 clash fake-ip 转走被拦。 + # 复用 worker.client.install_proxy_bypass(与 new-api / 本地 Service 旁路同源),幂等、绝不抛。 + try: + from worker import client # noqa: PLC0415 + + client.install_proxy_bypass(endpoint) + except Exception as exc: # best-effort:装旁路失败不阻断建 provider(host 级 NO_PROXY 可能已由启动装过) + _warn_once("proxy-bypass", f"装 OTLP 代理旁路失败(忽略,可能已由启动装过):{exc}") + + from opentelemetry.exporter.otlp.proto.http.trace_exporter import ( # noqa: PLC0415 + OTLPSpanExporter, + ) + + # OTLP/HTTP exporter:显式给 traces 全路径 base + /v1/traces(与 studio_sink 同策略,避免不同 SDK 版本对 + # base/自动拼路径行为不一致;Collector 的 OTLP/HTTP receiver 收 /v1/traces)。 + exporter = OTLPSpanExporter(endpoint=endpoint.rstrip("/") + "/v1/traces") + # 批处理 processor:异步攒批后台发,推送不卡生成主链(observe-only 不阻塞)。 + provider.add_span_processor(BatchSpanProcessor(exporter)) + return provider + + +def _get_shared_tracer(endpoint: str, span_exporter: Any = None) -> Any: + """取进程级共享 tracer(惰性建一次、并发安全);建失败 best-effort 返回 None、永久降级(不每条重试)。 + + 差异④:便宜档 Service 长驻、中间件工厂 per-turn 建 sink,provider 若 per-run 会逐局泄漏后台线程; + 故 provider/tracer 做进程级单例,trace_id 只作 per-span 属性(见 CheapOtlpSink._emit_span)。 + """ + # 快路:已建好直接返回;已判失败直接 None(永久降级,不每条重试)。 + if _shared["tracer"] is not None: + return _shared["tracer"] + if _shared["failed"]: + return None + with _shared_lock: + # 双检:抢锁期间别的线程可能已建好 / 已判失败。 + if _shared["tracer"] is not None: + return _shared["tracer"] + if _shared["failed"]: + return None + try: + provider = _build_provider(endpoint, span_exporter=span_exporter) + tracer = provider.get_tracer("cheap-gen-worker", "1.0.0") + _shared["provider"] = provider + _shared["tracer"] = tracer + _shared["endpoint"] = endpoint + # 进程退出兜底 flush:注册一次(避免 per-run 注册堆积)。BatchSpanProcessor 平时按调度发, + # 进程退出时靠这条把攒批残留 span 尽力发出(best-effort、丢尾批可容忍,observe-only)。 + if not _shared["atexit"]: + import atexit # noqa: PLC0415 + + atexit.register(flush_shared_tracer) + _shared["atexit"] = True + _warn_once("init-ok", f"便宜档 OTLP 观测已开(service.name={SERVICE_NAME} → {endpoint}/v1/traces)") + return tracer + except Exception as exc: # best-effort:缺依赖 / 建 provider 失败 → 永久降级 no-op,不抛、不阻塞生成 + _shared["failed"] = True + _warn_once("init-fail", f"建 OTLP→Collector 通道失败,便宜档降级为不发 span(只留 jsonl):{exc}") + return None + + +def flush_shared_tracer() -> None: + """进程退出 / 收口时把 BatchSpanProcessor 攒批里残留的 span 尽力发出(best-effort,绝不抛)。 + + atexit 注册一次;CLI 收口也可显式调(波④ 若要精确 flush)。Service 长驻期间靠 BatchSpanProcessor 自身调度, + 无需 per-turn 调本函数(shutdown 单例 provider 会杀掉后台上报线程,长驻进程不该 per-turn 关)。 + """ + provider = _shared.get("provider") + if provider is None: + return + try: + provider.force_flush(timeout_millis=5000) # 给上限超时,避免收口/退出卡住 + except Exception as exc: # best-effort:flush 失败只告警 + _warn_once("flush-fail", f"force_flush 残留 span 失败(忽略):{exc}") + + +def reset_shared_tracer_for_test() -> None: + """仅测试用:清进程级共享单例 + 告警去重集,让各用例互不串(生产路绝不调)。""" + provider = _shared.get("provider") + if provider is not None: + try: + provider.shutdown() + except Exception: # noqa: BLE001 —— 测试清理 best-effort + pass + _shared.update({"provider": None, "tracer": None, "endpoint": None, "failed": False}) + _WARNED.clear() + + +class CheapOtlpSink: + """把统一 trace 事件(TraceStep dict)推成 OTLP span 发给 Collector 的旁路 sink(便宜档主路)。 + + 契约 = Callable[[dict], None](与 make_jsonl_sink 的 sink 一致):吃一条 TraceStep dict → 发一个 span。 + 任何环节失败(缺 otel 依赖 / 建 provider 失败 / 发送失败)一律 best-effort:落一次告警、自动降级 no-op、绝不抛。 + span 属性映射与 tier2 studio_sink 的标准 gen-ai 键保持一致(便于将来 Tempo 跨线按 conversation.id 归组)。 + """ + + enabled = True + + def __init__(self, trace_id: str, endpoint: str, *, span_exporter: Any = None) -> None: + """ + Args: + trace_id: 本次生成 traceId(贯穿;作 span 的 gen_ai.conversation.id,Tempo 据它把同一 run 的步归一条 trace)。 + endpoint: Collector OTLP/HTTP base(已归一化;发时拼 /v1/traces)。 + span_exporter: 仅测试注入(内存 exporter);生产恒 None → 走真 OTLP exporter。 + """ + self.trace_id = trace_id + self.endpoint = endpoint + self._span_exporter = span_exporter + + def __call__(self, trace_step: dict) -> None: + """吃一条 TraceStep dict → 发一个 OTLP span(全程 best-effort,绝不抛、绝不阻塞生成主链)。""" + tracer = _get_shared_tracer(self.endpoint, span_exporter=self._span_exporter) + if tracer is None: + return # 建通道失败已永久降级 no-op(告警已落),observe-only 不影响主链 + try: + self._emit_span(tracer, trace_step) + except Exception as exc: # best-effort:单条 span 发送失败只告警一次,不连累后续、不抛 + _warn_once("emit-fail", f"推 OTLP span 失败(best-effort 计告警,traceId={self.trace_id}):{exc}") + + def _emit_span(self, tracer: Any, trace_step: dict) -> None: + """把单条 TraceStep 组成一个 span 并即时收尾(同步 start→set→end,交 processor 发出)。 + + 映射(与 studio_sink._emit_span 结构对齐;标准 gen-ai 键一致,私有键用 agentscope.cheap.* 区分线): + - span.name = 取 ext.action.tool / 相位 kind / raw.event,组成可读步名(如 "step 3 · write_source")。 + - gen_ai.conversation.id = traceId(与 studio_sink 同键:Tempo 据它把同一 run 的步归一条 trace)。 + - gen_ai.usage.* = cost.tokens 的 in/out(标准 token 用量键)。 + - gen_ai.tool.name = 动作段工具名。 + - agentscope.cheap.* = tier2 三段(推理/动作/观察)+ step/timestamp/verdict(JSON 序列化,便于反查)。 + - span 状态:observe 段 verdict ∈ {error/interrupted/denied} → 标 ERROR(如实反映,不改任何生成判定)。 + """ + import json # noqa: PLC0415 —— 仅序列化属性用 + from opentelemetry.trace import StatusCode # noqa: PLC0415 + + step = trace_step.get("step") + ext = trace_step.get("ext") or {} + raw = ext.get("raw") or {} + + # ── span 名:优先工具名 / 相位 kind / 原始事件类型,拼上步序,便于在 Tempo 列表一眼认出这步在干嘛 ── + label = ( + (ext.get("action") or {}).get("tool") + or (ext.get("action") or {}).get("kind") + or (ext.get("observation") or {}).get("kind") + or (ext.get("reasoning") or {}).get("kind") + or raw.get("event") + or "step" + ) + span_name = f"step {step} · {label}" + + # 即起即收(with 退出即 end):本步是「一条已成事实的 trace 记录」,无需跨步嵌套,也不会因异常漏 end。 + with tracer.start_as_current_span(span_name) as span: + # traceId → gen_ai.conversation.id(与 studio_sink 同键,跨线归组用)。 + span.set_attribute("gen_ai.conversation.id", str(self.trace_id)) + if step is not None: + span.set_attribute("agentscope.cheap.step", int(step)) + ts = trace_step.get("timestamp") + if ts: + span.set_attribute("agentscope.cheap.timestamp", str(ts)) + + # cost.tokens → gen_ai.usage.*(标准 token 用量键)。 + cost = trace_step.get("cost") or {} + tokens = (cost or {}).get("tokens") or {} + if "in" in tokens: + span.set_attribute("gen_ai.usage.input_tokens", int(tokens.get("in") or 0)) + if "out" in tokens: + span.set_attribute("gen_ai.usage.output_tokens", int(tokens.get("out") or 0)) + + # tier2 扩展三段 → agentscope.cheap.*(JSON 序列化,反查用)。 + for seg_key in ("reasoning", "action", "observation"): + seg = ext.get(seg_key) + if seg: + span.set_attribute(f"agentscope.cheap.{seg_key}", json.dumps(seg, ensure_ascii=False)) + tool = (ext.get("action") or {}).get("tool") + if tool: + span.set_attribute("gen_ai.tool.name", str(tool)) + if raw.get("event"): + span.set_attribute("agentscope.cheap.raw_event", str(raw.get("event"))) + + # verdict → span 状态(observe-only,不改任何生成判定)。 + verdict = trace_step.get("verdict") + if verdict is not None: + span.set_attribute("agentscope.cheap.verdict", str(verdict)) + if str(verdict).lower() in _ERROR_VERDICTS: + span.set_status(StatusCode.ERROR, str(verdict)) + else: + span.set_status(StatusCode.OK) + + def shutdown(self) -> None: + """收口:flush 进程级共享 provider 的攒批残留 span(best-effort,绝不抛)。 + + 注意:只 flush 不 shutdown provider —— provider 是进程级单例(长驻 Service 逐局复用),不能 per-turn 关掉 + 它的后台上报线程;真正 shutdown 交进程退出的 atexit(flush_shared_tracer)一次性做。 + """ + flush_shared_tracer() + + +class _NoopSink: + """关闭态占位 sink:__call__ / shutdown 都不做。让调用方无需判 None、统一拿 sink。""" + + enabled = False + + def __call__(self, trace_step: dict) -> None: # noqa: D401 —— 故意 no-op + return + + def shutdown(self) -> None: # noqa: D401 —— 故意 no-op + return + + +def make_cheap_otlp_sink( + trace_id: str, + *, + endpoint: Optional[str] = None, + span_exporter: Any = None, +) -> Callable[[dict], None]: + """构造便宜档 OTLP 旁路 sink:解析到 endpoint → CheapOtlpSink;取不到 → no-op sink(默认关闭)。 + + :param trace_id: 本次生成 traceId(贯穿)。 + :param endpoint: 显式 Collector base(最低优先级;一般留空,走 env CHEAP_OTLP_ENDPOINT / infra_config)。 + :param span_exporter: 仅测试注入(内存 exporter);生产恒 None。 + :return: Callable[[dict], None] 且带 .shutdown() 的 sink(CheapOtlpSink 或 _NoopSink)。 + """ + url = resolve_otlp_endpoint(endpoint) + if not url: + # 取不到 endpoint = 默认关闭:返回 no-op(现有便宜档真跑不受任何影响)。 + return _NoopSink() + return CheapOtlpSink(trace_id=trace_id, endpoint=url, span_exporter=span_exporter) + + +class _FanoutSink: + """把一条 trace_step 同时喂给多个 sink 的 fan-out 壳:每侧各 best-effort、互不连累(observe-only 不咬主链)。 + + 用于把既有 jsonl 落盘 sink 与新 OTLP sink 并接成「一个 sink」交给 Tier2TraceMiddleware(它只吃单 sink 契约)。 + 任一侧抛异常只告警、不抛、不连累另一侧与生成主链(镜像 observability.trace.with_studio 的 _fanout 语义)。 + """ + + def __init__(self, sinks: list) -> None: + self._sinks = [s for s in sinks if s is not None] + + def __call__(self, trace_step: dict) -> None: + for s in self._sinks: + try: + s(trace_step) + except Exception as exc: # noqa: BLE001 —— 一侧失败只告警,不连累其余 sink 与主链 + _warn_once("fanout", f"fan-out 某 sink 写失败(忽略,不连累其余):{type(exc).__name__}: {exc}") + + def shutdown(self) -> None: + """把带 .shutdown() 的子 sink 依次收口(best-effort;jsonl sink 无 shutdown 自然跳过)。""" + for s in self._sinks: + fn = getattr(s, "shutdown", None) + if callable(fn): + try: + fn() + except Exception: # noqa: BLE001 —— 收口 best-effort + pass + + +def build_trace_sink( + jsonl_sink: Optional[Callable[[dict], None]], + *, + trace_id: str, + endpoint: Optional[str] = None, + span_exporter: Any = None, +) -> Optional[Callable[[dict], None]]: + """把既有 jsonl sink 与(可选)OTLP sink 并接成一个交给 Tier2TraceMiddleware 的 sink。 + + **默认关(取不到 endpoint):原样返回传入的 jsonl sink,零包装 —— 现有行为字节不变。** + 开(env / config / 显式给 endpoint):返回 fan-out sink(jsonl 先落盘照旧、OTLP 额外发一份,两侧各 best-effort)。 + jsonl 落盘永远保留不动;OTLP 只是旁路多发一份,挂掉绝不咬生成(设计 §8)。 + + :param jsonl_sink: 既有 make_jsonl_sink 产的落盘 sink(可能为 None:上游 jsonl 接线已失败降级)。 + :param trace_id: 本次生成 traceId(贯穿)。 + :param endpoint: 显式 Collector base(一般留空,走 env / infra_config)。 + :param span_exporter: 仅测试注入。 + :return: 交给 Tier2TraceMiddleware(sink=...) 的单 sink(jsonl 原件 / OTLP 单件 / fan-out / None)。 + """ + otlp = make_cheap_otlp_sink(trace_id, endpoint=endpoint, span_exporter=span_exporter) + otlp_on = not isinstance(otlp, _NoopSink) + if not otlp_on: + # 默认关:零包装,原样返回 jsonl(现有行为字节不变;jsonl 为 None 时也原样返回 None)。 + return jsonl_sink + if jsonl_sink is None: + # jsonl 上游已降级无 sink:只发 OTLP(仍不阻断)。 + return otlp + # 两条都在:fan-out(jsonl 先落盘、OTLP 后发,各 best-effort)。 + return _FanoutSink([jsonl_sink, otlp]) + + +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) + s({"traceId": "demo-trace-001", "step": 0, "ext": {"raw": {"event": "ReplyStartEvent"}}}) + s.shutdown() + print("noop sink ok(observe-only 关闭态:吃事件不报错、shutdown 不报错)") diff --git a/cheap-worker/cheap_service_app.py b/cheap-worker/cheap_service_app.py index fd3e886b..3ac67220 100644 --- a/cheap-worker/cheap_service_app.py +++ b/cheap-worker/cheap_service_app.py @@ -273,6 +273,16 @@ async def _cheap_middlewares_factory(user_id: str, agent_id: str, session_id: st trace_sink = make_jsonl_sink(cheap_run.game_dir(game_id) / "trace.jsonl") except Exception as e: # noqa: BLE001 —— trace 非阻塞:接线失败只告警、回落无 sink print(f"[cheap-service] trace sink 接线失败(回落无 sink,不阻断):{type(e).__name__}: {e}", flush=True) + # 阶段四观测 波③ T3-1:jsonl 落盘之外并接一个 OTLP span sink(发 mini-infra Collector,service.name=cheap-gen-worker)。 + # **默认关**(env CHEAP_OTLP_ENDPOINT 未设)→ build_trace_sink 原样返回上面的 jsonl sink,现有行为字节不变; + # 开则包一层 fan-out(jsonl 先落盘照旧、OTLP 额外发一份,两侧各 best-effort 互不连累)。观测挂掉绝不咬生成(设计 §8)。 + # 惰性 import:cheap_otlp_sink 顶层零重依赖(otel 只在其函数体内 import),不违 build_cheap_app 惰性红线。 + try: + import cheap_otlp_sink # noqa: PLC0415 + + trace_sink = cheap_otlp_sink.build_trace_sink(trace_sink, trace_id=game_id) + except Exception as e: # noqa: BLE001 —— OTLP 并接非阻塞:失败回落原 jsonl sink、只告警 + print(f"[cheap-service] OTLP sink 并接失败(回落原 jsonl sink,不阻断):{type(e).__name__}: {e}", flush=True) tracer = Tier2TraceMiddleware(trace_id=game_id, sink=trace_sink) # C2:trace_id 用后端 gameId(与 trace 文件目录/trace.gameId 一致) max_repairs = genconfig.get("iteration", "max_resumes", 6) # 成本上界=优雅终止(决策①) diff --git a/cheap-worker/cheap_studio.py b/cheap-worker/cheap_studio.py index 2b233038..6f94e7d8 100644 --- a/cheap-worker/cheap_studio.py +++ b/cheap-worker/cheap_studio.py @@ -214,6 +214,14 @@ async def run_studio(game_id, brief, *, max_iters=40, max_resumes=6, max_tokens= _rec(f"trace sink 接线 → {_trace_path}") except Exception as _sink_err: _rec(f"trace sink 构造异常(只告警,回落无 sink,不阻断生成):{_sink_err}") + # 阶段四观测 波③ T3-1:CLI 路同款并接 OTLP span sink(发 Collector,service.name=cheap-gen-worker)。 + # 默认关(env CHEAP_OTLP_ENDPOINT 未设)→ 原样返回上面的 jsonl sink,CLI 现有行为字节不变; + # 开则包 fan-out(jsonl 落盘照旧、OTLP 额外发一份)。best-effort:观测挂掉绝不咬生成(设计 §8)。 + try: + import cheap_otlp_sink # noqa: PLC0415 函数体内惰性 import,顶层不牵 otel + _trace_sink = cheap_otlp_sink.build_trace_sink(_trace_sink, trace_id=game_id) + except Exception as _otlp_err: + _rec(f"OTLP sink 并接异常(回落原 jsonl sink,不阻断生成):{_otlp_err}") tracer = Tier2TraceMiddleware(trace_id=game_id, sink=_trace_sink) writer = Agent( diff --git a/cheap-worker/tests/test_cheap_otlp_sink.py b/cheap-worker/tests/test_cheap_otlp_sink.py new file mode 100644 index 00000000..5b486cd7 --- /dev/null +++ b/cheap-worker/tests/test_cheap_otlp_sink.py @@ -0,0 +1,185 @@ +"""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"]))