From 5b34807e7320e80188102d40635850fa06cda185 Mon Sep 17 00:00:00 2001 From: lili Date: Sun, 5 Jul 2026 20:30:04 -0700 Subject: [PATCH] =?UTF-8?q?feat(cheap):=20=E9=9D=A2=E5=9B=9B=20Python=20?= =?UTF-8?q?=E4=B8=89=E8=B7=B3=20traceparent=20=E6=89=93=E9=80=9A=E2=80=94?= =?UTF-8?q?=E2=80=94=E4=BE=BF=E5=AE=9C=E6=A1=A3=E7=94=9F=E6=88=90=E4=B8=B2?= =?UTF-8?q?=E6=88=90=E5=8D=95=E6=9D=A1=20trace=20=E6=A0=91(=E8=A7=82?= =?UTF-8?q?=E6=B5=8B=E9=98=B6=E6=AE=B5=E5=9B=9B)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Java agent(面一b)已在派发出网自动注 traceparent;本片补 Python 侧三个断点,让一次便宜档生成 在 Tempo 里是 Java→worker(:9501)→cheap-service(:8300)生成 的单条 trace 树而非散落 root: - 断点①:worker_service.do_POST 抓入站 traceparent 进 job._otelCarrier(self.headers 仅此可得); process_job 用 worker_span_scope 在入站 context 下建 worker SERVER span、包住真生成调用(span 建在 真处理 job 的 worker 线程,非只入队即 202 的 do_POST)。 - 断点②:driver 读当前 context(worker span 经 asyncio.run 透传进协程)inject 成 W3C carrier,随出站 :8300 各 POST header 带上;同份 carrier 写会话 sidecar。 - 断点③:/chat reply 在 AgentScope 后台任务跑、与 HTTP 请求解耦读不到实时 context——故走既有 C2 sidecar 桥:cheap_service_app 工厂读 sidecar traceparent → parent_carrier → CheapOtlpSink 每步生成 span 挂入站 context 当子 span(不再私有 root)。 传播 helper 集中 cheap_otlp_sink(extract/inject/worker_span_scope);全程 best-effort(otel 不可用/carrier 坏/Collector 不可达都退化成少串一段 trace、生成异常原样穿出)、CHEAP_OTLP_ENDPOINT 默认关时零行为变化 (parent=None 退私有 root、sidecar 不写 traceparent 字节不变)、代理旁路红线未破、业务 traceId 并存(§10 决策5)。 test_trace_propagation.py 14 例(含端到端三跳单条 trace 树,InMemorySpanExporter 证父子链+同 trace_id) + test_worker_service +2;主控亲跑 40 passed(3 文件)+ 全量 364 passed 1 skipped 无回退。 真跨进程 Tempo trace 树验证 = dev 开 CHEAP_OTLP_ENDPOINT 重跑 follow-up。 Co-Authored-By: Claude Opus 4.8 (1M context) --- cheap-worker/cheap_otlp_sink.py | 142 +++++++++++- cheap-worker/cheap_service_app.py | 17 +- cheap-worker/cheap_service_driver.py | 30 ++- cheap-worker/tests/test_trace_propagation.py | 214 +++++++++++++++++++ cheap-worker/tests/test_worker_service.py | 52 +++++ cheap-worker/worker_service.py | 25 ++- 6 files changed, 468 insertions(+), 12 deletions(-) create mode 100644 cheap-worker/tests/test_trace_propagation.py diff --git a/cheap-worker/cheap_otlp_sink.py b/cheap-worker/cheap_otlp_sink.py index 266518e8..db3337a6 100644 --- a/cheap-worker/cheap_otlp_sink.py +++ b/cheap-worker/cheap_otlp_sink.py @@ -46,6 +46,7 @@ from __future__ import annotations +import contextlib import os import threading from typing import Any, Callable, Optional @@ -118,6 +119,114 @@ def resolve_otlp_endpoint(explicit: Optional[str] = None) -> Optional[str]: return None +# ══════════════════════════════════════════════════════════════════════════════ +# 跨进程 W3C 传播(阶段四观测 · 面四三跳)—— 把一次便宜档生成串成单条 trace 树 +# ────────────────────────────────────────────────────────────────────────────── +# 一次生成的真实链路是三跳:Java 派发(:9501)→ worker_service 收 job → worker 打本机 +# cheap-service(:8300)真生成。三个断点原本让 trace 断成散落 root:worker 收到的 +# traceparent 没人提取、worker 出站 :8300 的 header 只有 X-User-Id、cheap-service 的生成 +# span 是私有 root 不读入站 context。下面三件把 W3C traceparent 接起来(全 best-effort, +# otel/传播任何环节挂掉都退化成「少串一段 trace」,绝不咬生成——设计 §8 旁路铁律)。 +# ══════════════════════════════════════════════════════════════════════════════ + +def _extract_remote_context(carrier: Optional[dict]) -> Any: + """把入站 W3C carrier(traceparent[/tracestate])extract 成一个 OTel parent context。 + + :param carrier: {'traceparent': ..., 'tracestate'?: ...};None / 无 traceparent → 返回 None。 + :return: 供 start_span(context=...) 当父的 OTel Context;取不到 / otel 不可用 → None(退回私有 root)。 + """ + try: + if not carrier or not carrier.get("traceparent"): + return None + from opentelemetry import propagate # noqa: PLC0415 —— 惰性 import(顶层零重依赖) + + return propagate.extract(carrier) + except Exception: # noqa: BLE001 —— 提取是旁路,失败退回无父,绝不咬主链 + return None + + +def current_traceparent_carrier() -> dict: + """把「当前 OTel context」注入成一份 W3C carrier(dict:traceparent[/tracestate]),供出站带上(断点②)。 + + 便宜档 worker 打本机 cheap-service 的 httpx 请求、以及写给工厂的会话 sidecar,都用这份 carrier 把 + 「worker 这跳的 span context」透传下去,让 cheap-service 端的生成 span 挂到同一条 trace 树下。 + 纯增量、best-effort:otel 不可用或当前无 context → 返回 {}(出站少带一个 header,不影响生成)。 + """ + try: + from opentelemetry import propagate # noqa: PLC0415 —— 惰性 import(顶层零重依赖) + + carrier: dict = {} + propagate.inject(carrier) # 读当前 context → 写 traceparent(/tracestate)进 carrier + return carrier + except Exception: # noqa: BLE001 —— 注入是旁路,失败返回空 carrier,绝不咬主链 + return {} + + +@contextlib.contextmanager +def worker_span_scope(carrier: Optional[dict], *, conversation_id: Any = None, + business_trace_id: Any = None): + """便宜档 worker 这跳(:9501 收 Java 派发的 job)的 server span 作用域(断点①)。 + + 背景:worker_service 是裸标准库 http.server,OTel 不给它自动埋点;job 又经内部有界队列跨线程处理, + 故 do_POST 只把入站 traceparent 抓进 job(self.headers 仅那里可得),真正的 worker span 在 worker + 线程处理该 job 时(本作用域)才开——span 恰好覆盖「调 cheap-service 生成」这一跳、并成为下游的父。 + 行为: + · OTLP 已开(进程级共享 tracer 可用):在入站 context 下开一个 cheap.worker.generate 的 SERVER span、 + attach 为当前 context;下游(driver→:8300 的 httpx inject、写工厂 sidecar 的 traceparent)据此继承。 + · OTLP 未开 / 建 span 失败:只把入站 context attach 为当前(无 span),让下游 inject 仍能把 Java 的 + traceparent 原样透传出去(纯增量、不依赖 OTLP 开关);无入站 context 时退化成什么都不做。 + §10 决策5「业务 traceId 并存」:业务 traceId(承担产物归位/回调关联)作 span 属性保留、不改其格式; + W3C traceparent 只负责跨进程贯穿。 + best-effort 铁律(设计 §8):提取 / 建 span / attach 全程吞异常,绝不让传播处理影响生成成败或阻塞。 + """ + parent_ctx = _extract_remote_context(carrier) + + # 尝试拿进程级共享 tracer(仅 OTLP 开时可用;复用波③ 同一 provider,不新造 root)。 + span = None + try: + endpoint = resolve_otlp_endpoint() + if endpoint: + tracer = _get_shared_tracer(endpoint) + if tracer is not None: + from opentelemetry.trace import SpanKind # noqa: PLC0415 + + span = tracer.start_span( + "cheap.worker.generate", context=parent_ctx, kind=SpanKind.SERVER) + # 业务 traceId 并存(§10 决策5):gameId 作 conversation.id(与 cheap-service 生成 span 同键、 + # Tempo 按它跨线归组),业务 traceId 另存一键;都只作属性,不改其格式、不参与 W3C 贯穿。 + if conversation_id is not None: + span.set_attribute("gen_ai.conversation.id", str(conversation_id)) + if business_trace_id is not None: + span.set_attribute("agentscope.cheap.business_trace_id", str(business_trace_id)) + except Exception: # noqa: BLE001 —— 建 span 失败退回「只 attach 入站 context」路径,绝不咬生成 + span = None + + if span is not None: + from opentelemetry.trace import use_span # noqa: PLC0415 + + # use_span:把 worker span attach 为当前 + 退出时 end(异常也 end 并记录);下游 inject 据此继承。 + with use_span(span, end_on_exit=True): + yield span + return + + # OTLP 未开 / 建 span 失败:只 attach 入站 context(让下游 inject 透传 Java traceparent);无则什么都不做。 + token = None + try: + if parent_ctx is not None: + from opentelemetry import context as otel_context # noqa: PLC0415 + + token = otel_context.attach(parent_ctx) + yield None + finally: + if token is not None: + try: + from opentelemetry import context as otel_context # noqa: PLC0415 + + otel_context.detach(token) + except Exception: # noqa: BLE001 —— detach 失败只影响 contextvar 清理,不咬主链 + pass + + def _build_provider(endpoint: str, span_exporter: Any = None) -> Any: """惰性建一个私有 TracerProvider(装 exporter);失败抛给上层 best-effort 兜。仅在函数体内 import otel(惰性红线)。 @@ -240,16 +349,23 @@ class CheapOtlpSink: enabled = True - def __init__(self, trace_id: str, endpoint: str, *, span_exporter: Any = None) -> None: + def __init__(self, trace_id: str, endpoint: str, *, span_exporter: Any = None, + parent_carrier: Optional[dict] = 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。 + parent_carrier: 入站 W3C carrier(traceparent[/tracestate];面四断点③)。给定则每步生成 span 挂到 + 入站 context 下当子 span(不再私有 root),使 Java→worker→cheap-service 串成单条 trace 树; + 缺省 None → 保持波③ 私有 root 行为(默认关时字节不变)。 """ self.trace_id = trace_id self.endpoint = endpoint self._span_exporter = span_exporter + self._parent_carrier = parent_carrier or None + self._parent_ctx: Any = None # extract 出的 parent context(惰性缓存一次) + self._parent_ctx_resolved = False def __call__(self, trace_step: dict) -> None: """吃一条 TraceStep dict → 发一个 OTLP span(全程 best-effort,绝不抛、绝不阻塞生成主链)。""" @@ -261,6 +377,16 @@ class CheapOtlpSink: except Exception as exc: # best-effort:单条 span 发送失败只告警一次,不连累后续、不抛 _warn_once("emit-fail", f"推 OTLP span 失败(best-effort 计告警,traceId={self.trace_id}):{exc}") + def _parent_context(self) -> Any: + """把入站 W3C carrier extract 成 parent context(本 sink 实例缓存一次;面四断点③)。 + + None(未给 carrier / extract 失败)→ start_as_current_span 用 None,退回私有 root(波③ 行为)。 + """ + if not self._parent_ctx_resolved: + self._parent_ctx_resolved = True + self._parent_ctx = _extract_remote_context(self._parent_carrier) + return self._parent_ctx + def _emit_span(self, tracer: Any, trace_step: dict) -> None: """把单条 TraceStep 组成一个 span 并即时收尾(同步 start→set→end,交 processor 发出)。 @@ -291,7 +417,9 @@ class CheapOtlpSink: span_name = f"step {step} · {label}" # 即起即收(with 退出即 end):本步是「一条已成事实的 trace 记录」,无需跨步嵌套,也不会因异常漏 end。 - with tracer.start_as_current_span(span_name) as span: + # 面四断点③:context=入站 parent → 每步生成 span 挂到 Java→worker 这条入站 trace 下当子 span + # (不再私有 root);parent 为 None(未开传播)时等价于原 start_as_current_span()(波③ 行为不变)。 + with tracer.start_as_current_span(span_name, context=self._parent_context()) 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: @@ -354,19 +482,22 @@ def make_cheap_otlp_sink( *, endpoint: Optional[str] = None, span_exporter: Any = None, + parent_carrier: Optional[dict] = 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。 + :param parent_carrier: 入站 W3C carrier(面四断点③);给定则生成 span 挂到入站 trace 下当子 span。 :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) + return CheapOtlpSink(trace_id=trace_id, endpoint=url, span_exporter=span_exporter, + parent_carrier=parent_carrier) class _FanoutSink: @@ -403,6 +534,7 @@ def build_trace_sink( trace_id: str, endpoint: Optional[str] = None, span_exporter: Any = None, + parent_carrier: Optional[dict] = None, ) -> Optional[Callable[[dict], None]]: """把既有 jsonl sink 与(可选)OTLP sink 并接成一个交给 Tier2TraceMiddleware 的 sink。 @@ -414,9 +546,11 @@ def build_trace_sink( :param trace_id: 本次生成 traceId(贯穿)。 :param endpoint: 显式 Collector base(一般留空,走 env / infra_config)。 :param span_exporter: 仅测试注入。 + :param parent_carrier: 入站 W3C carrier(面四断点③);给定则 OTLP 生成 span 挂到入站 trace 下当子 span。 :return: 交给 Tier2TraceMiddleware(sink=...) 的单 sink(jsonl 原件 / OTLP 单件 / fan-out / None)。 """ - otlp = make_cheap_otlp_sink(trace_id, endpoint=endpoint, span_exporter=span_exporter) + otlp = make_cheap_otlp_sink(trace_id, endpoint=endpoint, span_exporter=span_exporter, + parent_carrier=parent_carrier) otlp_on = not isinstance(otlp, _NoopSink) if not otlp_on: # 默认关:零包装,原样返回 jsonl(现有行为字节不变;jsonl 为 None 时也原样返回 None)。 diff --git a/cheap-worker/cheap_service_app.py b/cheap-worker/cheap_service_app.py index 3ac67220..bbc9a996 100644 --- a/cheap-worker/cheap_service_app.py +++ b/cheap-worker/cheap_service_app.py @@ -263,7 +263,19 @@ async def _cheap_middlewares_factory(user_id: str, agent_id: str, session_id: st # C2:评门/collector/trace 全绑后端 gameId(经会话注册表 sidecar 把 session_id 解析回后端 gameId;与 driver # scaffold、system prompt 的 ⟦G⟧、六工具写目录一致)。缺映射则回落 session_id(响亮失败,见 _resolve_external_game_id)。 - game_id = _resolve_external_game_id(session_id, _read_session_cfg(session_id)) + _cfg = _read_session_cfg(session_id) + game_id = _resolve_external_game_id(session_id, _cfg) + # 面四断点③:从 sidecar 取 worker 这跳写下的 W3C traceparent(driver 恒在 /chat 前写),组入站 carrier + # 传给 build_trace_sink → 让 OTLP 生成 span 挂到 Java→worker 这条入站 trace 下当子 span(不再私有 root)。 + # reply 在框架后台任务里跑、与收 /chat 的 HTTP 请求解耦,读不到出站 header 的实时 context,故靠 driver 写的 + # sidecar 跨进程桥 traceparent(与 external_game_id 同一 C2 桥)。无 traceparent(默认关)→ None,退回私有 root。 + _tp = _cfg.get("traceparent") + parent_carrier = None + if _tp: + parent_carrier = {"traceparent": _tp} + _ts = _cfg.get("tracestate") + if _ts: + parent_carrier["tracestate"] = _ts # trace sink:落 game_dir/trace.jsonl(与旧 cheap_studio 同路径);import/构造失败回落无 sink(非阻塞)。 trace_sink = None @@ -280,7 +292,8 @@ async def _cheap_middlewares_factory(user_id: str, agent_id: str, session_id: st try: import cheap_otlp_sink # noqa: PLC0415 - trace_sink = cheap_otlp_sink.build_trace_sink(trace_sink, trace_id=game_id) + trace_sink = cheap_otlp_sink.build_trace_sink(trace_sink, trace_id=game_id, + parent_carrier=parent_carrier) 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 一致) diff --git a/cheap-worker/cheap_service_driver.py b/cheap-worker/cheap_service_driver.py index e711eb90..85486553 100644 --- a/cheap-worker/cheap_service_driver.py +++ b/cheap-worker/cheap_service_driver.py @@ -45,13 +45,16 @@ def _cheap_credential_payload() -> dict: def _write_session_cfg(session_id: str, *, external_game_id: str, - write_whitelist=None, scaffold_template=None) -> None: + write_whitelist=None, scaffold_template=None, + traceparent=None, tracestate=None) -> None: """写本 session 的会话注册表 sidecar(C2:Service 两工厂读它把 session_id 解析回后端 gameId + 取 write_whitelist)。 按 session_id 键写 game-runtime/games/_cheap-sessions/.json(worker↔Service 同机共享 FS); **external_game_id = 后端 gameId 恒写**(工厂据它绑六工具/评门/collector 目录,与 driver scaffold 一致)。 restricted = 是否受限写(create 路 False + write_whitelist=None;reskin/modify 路 True + 白名单;I1 fail-closed - 依赖 restricted 标记)。best-effort:写失败只告警(Service 侧读缺失 → 回落 session_id、生成响亮失败)。 + 依赖 restricted 标记)。traceparent/tracestate = 面四断点③ 桥:worker 这跳的 W3C context,让工厂据它把 + cheap-service 生成 span 挂到入站 trace 下(reply 在后台任务跑、读不到出站 header 的实时 context,故走 sidecar 桥)。 + best-effort:写失败只告警(Service 侧读缺失 → 回落 session_id、生成响亮失败)。 """ import json # noqa: PLC0415 @@ -64,12 +67,18 @@ def _write_session_cfg(session_id: str, *, external_game_id: str, wl = sorted(write_whitelist) if write_whitelist else None p = cheap_run.session_cfg_path(session_id) p.parent.mkdir(parents=True, exist_ok=True) - p.write_text(json.dumps({ + cfg = { "external_game_id": str(external_game_id), "write_whitelist": wl, "scaffold_template": scaffold_template, "restricted": restricted, - }, ensure_ascii=False), encoding="utf-8") + } + # 面四断点③ 桥:仅当有 traceparent 时才写(默认关/无传播时不写、sidecar 字节不变)。 + if traceparent: + cfg["traceparent"] = traceparent + if tracestate: + cfg["tracestate"] = tracestate + p.write_text(json.dumps(cfg, ensure_ascii=False), encoding="utf-8") except Exception as e: # noqa: BLE001 —— sidecar best-effort,写失败 Service 侧回落 session_id print(f"[cheap-driver] session-cfg 写失败(Service 侧将回落 game_id=session_id、生成响亮失败):" f"{type(e).__name__}: {e}", flush=True) @@ -208,6 +217,7 @@ async def drive_cheap_generation(job: dict, *, base_url: str | None = None, user import httpx # noqa: PLC0415 —— 旁路已装,此后 import 才安全 import _bootstrap # noqa: PLC0415 + import cheap_otlp_sink # noqa: PLC0415 —— 面四三跳传播:取当前 context 的 W3C carrier(顶层零重依赖) import cheap_run # noqa: PLC0415 import cheap_verify # noqa: PLC0415 from cheap_roles import build_system_prompt # noqa: PLC0415 @@ -236,7 +246,15 @@ async def drive_cheap_generation(job: dict, *, base_url: str | None = None, user _bootstrap.ensure_api_key_env() # 保证 NEWAPI_KEY(装凭据要用) + # 面四断点②:把「worker 这跳的当前 context」(worker_span_scope 在 worker 线程 attach、经 asyncio.run 透传进本 + # 协程)注入成 W3C carrier,随所有出站到 :8300 的 header 带上,让 cheap-service 端能接住入站 context。同一份 + # carrier 也写进会话 sidecar(下方 _write_session_cfg)——因 /chat 的 reply 在框架后台任务里跑、生成 span 产在 + # 那条与 HTTP 请求解耦的任务里,读不到出站 header 的实时 context,故靠 driver 恒在 /chat 前写的 sidecar 把 + # traceparent 跨进程桥给工厂(与 external_game_id 同一 C2 桥)。纯增量 best-effort:无 context / otel 不可用 → + # carrier 为空,header 少带一项、生成 span 退回私有 root,绝不咬生成(设计 §8)。 + otel_carrier = cheap_otlp_sink.current_traceparent_carrier() headers = {"X-User-Id": user_id} + headers.update(otel_carrier) # traceparent(/tracestate,若有);出站 :8300 各 POST 都带 async with httpx.AsyncClient() as http: # ① 注册 OpenAI 兼容凭据(cheap M3 走 base+/v1)。 cred = (await http.post(f"{base_url}/credential/", json=_cheap_credential_payload(), @@ -294,7 +312,9 @@ async def drive_cheap_generation(job: dict, *, base_url: str | None = None, user # ⑤ 写会话注册表 sidecar(C2:Service 两工厂据 session_id 读它解析回后端 gameId、绑六工具/评门/collector; # create 路 write_whitelist=None → restricted=False)。external_game_id 恒写、正常路工厂必读到。 _write_session_cfg(session_id, external_game_id=game_id, - write_whitelist=write_whitelist, scaffold_template=scaffold_template) + write_whitelist=write_whitelist, scaffold_template=scaffold_template, + traceparent=otel_carrier.get("traceparent"), + tracestate=otel_carrier.get("tracestate")) # ⑥ 设 BYPASS 权限(六工具默认 ASK,服务态无人确认,不设首个工具调用即卡死;PATCH 需 query agent_id,同 tier2)。 await http.patch(f"{base_url}/sessions/{session_id}", json={"permission_mode": "bypass"}, headers=headers, params={"agent_id": agent_id}, timeout=30.0) diff --git a/cheap-worker/tests/test_trace_propagation.py b/cheap-worker/tests/test_trace_propagation.py new file mode 100644 index 00000000..3161bf6c --- /dev/null +++ b/cheap-worker/tests/test_trace_propagation.py @@ -0,0 +1,214 @@ +"""test_trace_propagation.py — 便宜档面四三跳 W3C traceparent 传播单测(阶段四观测)。 + +守的不变量(把一次便宜档生成串成 Java→worker→cheap-service 单条 trace 树的三个断点): + · 断点①提取:propagator 从构造的 traceparent header 提取出正确的 trace_id/span_id(is_remote)。 + · 断点①建 span:OTLP 开时 worker_span_scope 在入站 context 下开 cheap.worker.generate SERVER span 当子, + 业务 traceId 并存(§10 决策5)作属性;OTLP 关时只透传入站 context(不依赖开关)。 + · 断点②注入:worker_span_scope 内 current_traceparent_carrier() 注入出的 traceparent 与入站同 trace_id(往返一致)。 + · 断点③挂靠:CheapOtlpSink(parent_carrier=...) 让生成 span 挂到入站 context 下当子 span(不再私有 root); + 缺省 parent_carrier=None 时保持波③ 私有 root 行为(默认关字节不变)。 + · best-effort(设计 §8):坏 carrier 不抛;生成异常原样穿出 worker_span_scope(worker loop 据此发兜底 failed)。 + +只本地跑、零网络:用 InMemorySpanExporter / 手工 traceparent 构造断言,不发真 Collector。 +跑:cheap-worker/.venv/bin/python -m pytest cheap-worker/tests/test_trace_propagation.py -q +""" +import os +import sys +from pathlib import Path + +# sys.path:cheap-worker/ + tier2/gen-worker/(observability / service / worker 包;与 test_cheap_otlp_sink 同范式)。 +_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 pytest # noqa: E402 +import cheap_otlp_sink as S # noqa: E402 +from opentelemetry import trace as _otel_trace # noqa: E402 +from opentelemetry.trace import SpanKind # noqa: E402 +from opentelemetry.sdk.trace.export.in_memory_span_exporter import ( # noqa: E402 + InMemorySpanExporter, +) + +# 固定构造的入站 W3C 值(模拟 Java OTel agent 出网注入的 traceparent)。 +_TRACE_HEX = "11112222333344445555666677778888" +_SPAN_HEX = "1111222233334444" +_TP = f"00-{_TRACE_HEX}-{_SPAN_HEX}-01" + + +def _span_ctx(ctx): + """从 extract 出的 Context 取当前 span 的 SpanContext。""" + return _otel_trace.get_current_span(ctx).get_span_context() + + +# ── 断点①提取:propagator 从 traceparent header 提取正确 trace_id/span_id ───────────── +def test_extract_from_traceparent_header_yields_correct_ids(): + ctx = S._extract_remote_context({"traceparent": _TP}) + assert ctx is not None + sc = _span_ctx(ctx) + assert format(sc.trace_id, "032x") == _TRACE_HEX + assert format(sc.span_id, "016x") == _SPAN_HEX + assert sc.is_remote is True # 入站(远端)context + + +def test_extract_none_or_missing_traceparent_returns_none(): + assert S._extract_remote_context(None) is None + assert S._extract_remote_context({}) is None + assert S._extract_remote_context({"tracestate": "x=1"}) is None # 无 traceparent 键 + + +# ── 断点②注入:worker_span_scope(OTLP 关)内注入的 traceparent 与入站同 trace_id(往返一致)── +def test_inject_roundtrip_trace_id_consistent_when_otlp_off(monkeypatch): + monkeypatch.setattr(S, "resolve_otlp_endpoint", lambda explicit=None: None) # 强制 OTLP 关 + S.reset_shared_tracer_for_test() + with S.worker_span_scope({"traceparent": _TP}, conversation_id="w1", business_trace_id="tp1"): + out = S.current_traceparent_carrier() + assert "traceparent" in out + parts = out["traceparent"].split("-") + assert parts[1] == _TRACE_HEX, "出站 trace_id 必与入站一致(往返一致)" + # OTLP 关:无 worker span,当前 = 入站 remote context → span_id 段透传入站 Java span(直挂 Java 下)。 + assert parts[2] == _SPAN_HEX + + +def test_current_carrier_empty_when_no_context(): + """无任何当前 context → 注入空 carrier(出站少带 header,不抛)。""" + S.reset_shared_tracer_for_test() + out = S.current_traceparent_carrier() + assert out == {} or "traceparent" not in out + + +# ── 断点①建 span:OTLP 开 → worker span 在入站 context 下当子 + 业务 traceId 并存 ──────── +def test_worker_span_scope_creates_child_server_span_when_otlp_on(monkeypatch): + S.reset_shared_tracer_for_test() + mem = InMemorySpanExporter() + monkeypatch.setattr(S, "resolve_otlp_endpoint", lambda explicit=None: "http://collector.test:4318") + # 预置进程级共享 tracer 用内存 exporter(worker_span_scope 复用同一单例,无需真网络)。 + S._get_shared_tracer("http://collector.test:4318", span_exporter=mem) + + with S.worker_span_scope({"traceparent": _TP}, conversation_id="w1", business_trace_id="tp1"): + inj = S.current_traceparent_carrier() # 作用域内当前 = worker span + + spans = mem.get_finished_spans() + assert len(spans) == 1 + ws = spans[0] + assert ws.name == "cheap.worker.generate" + assert ws.kind == SpanKind.SERVER + # worker span 在入站 trace 下(同 trace_id),父 = 入站 Java span(不新造 root)。 + assert format(ws.context.trace_id, "032x") == _TRACE_HEX + assert ws.parent is not None and format(ws.parent.span_id, "016x") == _SPAN_HEX + # 断点②:作用域内注入的 traceparent = worker span 自身(同 trace,span_id = worker span)。 + p = inj["traceparent"].split("-") + assert p[1] == _TRACE_HEX and p[2] == format(ws.context.span_id, "016x") + # §10 决策5 业务 traceId 并存:gameId 作 conversation.id、业务 traceId 另存属性(都不改格式)。 + attrs = dict(ws.attributes) + assert attrs.get("gen_ai.conversation.id") == "w1" + assert attrs.get("agentscope.cheap.business_trace_id") == "tp1" + S.reset_shared_tracer_for_test() + + +# ── 断点③挂靠:CheapOtlpSink(parent_carrier) 让生成 span 挂到入站 context 下当子 span ──── +def test_cheap_service_gen_span_parents_under_inbound_carrier(): + S.reset_shared_tracer_for_test() + mem = InMemorySpanExporter() + sink = S.make_cheap_otlp_sink(trace_id="conv-x", endpoint="http://collector.test:4318", + span_exporter=mem, parent_carrier={"traceparent": _TP}) + sink({"traceId": "conv-x", "step": 0, "cost": None, "verdict": None, + "timestamp": "2026-07-05T10:00:00", "ext": {"raw": {"event": "ReplyStartEvent"}}}) + spans = mem.get_finished_spans() + assert len(spans) == 1 + sp = spans[0] + # 生成 span 挂在入站 trace 下(同 trace_id、父 = 入站 span);不再私有 root。 + assert format(sp.context.trace_id, "032x") == _TRACE_HEX + assert sp.parent is not None and format(sp.parent.span_id, "016x") == _SPAN_HEX + # 业务 traceId 并存:conversation.id 仍是业务 traceId(承担产物归位/回调关联),不被 W3C 贯穿覆盖。 + assert dict(sp.attributes).get("gen_ai.conversation.id") == "conv-x" + S.reset_shared_tracer_for_test() + + +def test_cheap_service_gen_span_is_root_when_no_parent_carrier(): + """缺省 parent_carrier=None(未开传播)→ 生成 span 仍私有 root(波③ 行为字节不变)。""" + S.reset_shared_tracer_for_test() + mem = InMemorySpanExporter() + sink = S.make_cheap_otlp_sink(trace_id="conv-y", endpoint="http://collector.test:4318", + span_exporter=mem) # parent_carrier 缺省 None + sink({"traceId": "conv-y", "step": 0, "ext": {"raw": {"event": "ReplyStartEvent"}}}) + sp = mem.get_finished_spans()[0] + assert sp.parent is None, "未开传播时保持私有 root(默认关字节不变)" + S.reset_shared_tracer_for_test() + + +def test_build_trace_sink_threads_parent_carrier_to_child_span(): + """build_trace_sink(parent_carrier=...) 开时:fan-out 的 OTLP 侧生成 span 也挂到入站 context 下当子。""" + S.reset_shared_tracer_for_test() + mem = InMemorySpanExporter() + seen = [] + fan = S.build_trace_sink(lambda step: seen.append(step), trace_id="g-9", + endpoint="http://collector.test:4318", span_exporter=mem, + parent_carrier={"traceparent": _TP}) + fan({"traceId": "g-9", "step": 0, "ext": {"raw": {"event": "ReplyStartEvent"}}}) + assert len(seen) == 1 # jsonl 侧照落不被吞 + sp = mem.get_finished_spans()[0] + assert format(sp.context.trace_id, "032x") == _TRACE_HEX # OTLP 侧挂到入站 trace 下 + assert sp.parent is not None and format(sp.parent.span_id, "016x") == _SPAN_HEX + S.reset_shared_tracer_for_test() + + +# ── 三跳串成单条 trace 树(端到端:入站 Java → worker span → cheap-service 生成 span)────── +def test_three_hops_form_single_trace_tree(monkeypatch): + """把三个断点串起来:worker span 挂在入站 Java span 下,生成 span 又挂在 worker span 下 —— 同一 trace_id、 + 父子链完整(正是 dev 真机在 Tempo 里要看到的 Java→worker→cheap-service 单条 trace 树)。""" + S.reset_shared_tracer_for_test() + mem = InMemorySpanExporter() + monkeypatch.setattr(S, "resolve_otlp_endpoint", lambda explicit=None: "http://collector.test:4318") + S._get_shared_tracer("http://collector.test:4318", span_exporter=mem) + + # 跳①:worker 线程在入站 Java context 下开 worker span,取其出站 carrier(= driver 会注入 header/写 sidecar 的)。 + with S.worker_span_scope({"traceparent": _TP}, conversation_id="w1", business_trace_id="tp1"): + worker_carrier = S.current_traceparent_carrier() + # 跳③:cheap-service 端拿 worker 的 carrier 当 parent(经 sidecar 桥),生成 span 挂到 worker span 下。 + gen_sink = S.make_cheap_otlp_sink(trace_id="w1", endpoint="http://collector.test:4318", + span_exporter=mem, parent_carrier=worker_carrier) + gen_sink({"traceId": "w1", "step": 0, "ext": {"raw": {"event": "ReplyStartEvent"}}}) + + spans = {sp.name: sp for sp in mem.get_finished_spans()} + worker_sp = spans["cheap.worker.generate"] + gen_sp = spans["step 0 · ReplyStartEvent"] + # 三段同一 trace_id(单条 trace)。 + assert format(worker_sp.context.trace_id, "032x") == _TRACE_HEX + assert format(gen_sp.context.trace_id, "032x") == _TRACE_HEX + # 父子链:worker span 父 = 入站 Java span;生成 span 父 = worker span。 + assert format(worker_sp.parent.span_id, "016x") == _SPAN_HEX + assert gen_sp.parent.span_id == worker_sp.context.span_id + S.reset_shared_tracer_for_test() + + +# ── best-effort 铁律(设计 §8):坏 carrier 不抛;生成异常原样穿出 ──────────────────────── +def test_worker_span_scope_bad_carrier_does_not_raise(monkeypatch): + monkeypatch.setattr(S, "resolve_otlp_endpoint", lambda explicit=None: None) + with S.worker_span_scope({"traceparent": "garbage-not-w3c"}, conversation_id="w"): + pass # 坏 traceparent 只导致无父,不抛 + + +def test_worker_span_scope_does_not_swallow_body_exception(monkeypatch): + """生成异常必须原样穿出 worker_span_scope —— worker loop 据此发兜底 failed 回调,绝不被观测吞掉。""" + monkeypatch.setattr(S, "resolve_otlp_endpoint", lambda explicit=None: None) + with pytest.raises(ValueError): + with S.worker_span_scope({"traceparent": _TP}): + raise ValueError("boom") + + +def test_worker_span_scope_body_exception_ends_span_when_otlp_on(monkeypatch): + """OTLP 开时生成异常:worker span 仍被 end(不漏 span)、异常照常穿出。""" + S.reset_shared_tracer_for_test() + mem = InMemorySpanExporter() + monkeypatch.setattr(S, "resolve_otlp_endpoint", lambda explicit=None: "http://collector.test:4318") + S._get_shared_tracer("http://collector.test:4318", span_exporter=mem) + with pytest.raises(RuntimeError): + with S.worker_span_scope({"traceparent": _TP}, conversation_id="w1"): + raise RuntimeError("gen crashed") + spans = mem.get_finished_spans() + assert len(spans) == 1 and spans[0].name == "cheap.worker.generate" # 异常路 span 照样收尾导出 + S.reset_shared_tracer_for_test() diff --git a/cheap-worker/tests/test_worker_service.py b/cheap-worker/tests/test_worker_service.py index c04322db..38c18721 100644 --- a/cheap-worker/tests/test_worker_service.py +++ b/cheap-worker/tests/test_worker_service.py @@ -294,6 +294,58 @@ def test_worker_loop_run_failure_sends_fallback_failed(): assert captured["payload"]["failureReason"] == "llm_error" +# ---------- 面四断点①(capture):do_POST 抓入站 traceparent 进 job ---------- + +def test_do_post_captures_inbound_traceparent_into_job(): + """POST /generate 带 traceparent header → worker 线程拿到的 job 带 _otelCarrier(供建 worker span)。""" + gd = _make_game_dir(_tmp()) + captured = {} + state = W.WorkerState( + run_fn=lambda job: captured.update(job=job) or ({"verdict": {"pass": True}}, gd), + send_fn=lambda url, payload, secret: (200, ""), + profile_fn=lambda j, g: None, + ) + W.start_worker(state) + server, port = W.start_server(state, port=0) + try: + tp = "00-11112222333344445555666677778888-1111222233334444-01" + body = json.dumps({"job_id": "tp1", "traceId": "tp1", "templateId": "generic", "brief": "x", + "gameId": "w1", "callback": {"target": "http://cb"}}).encode("utf-8") + req = urllib.request.Request(f"http://127.0.0.1:{port}/generate", data=body, method="POST", + headers={"Content-Type": "application/json", "traceparent": tp}) + with urllib.request.urlopen(req, timeout=5) as resp: + assert resp.status == 202 + state.queue.join() + # 入站 traceparent 被抓进 job(do_POST 只这里能读 header),供 process_job 在 worker 线程建 worker span。 + assert captured["job"]["_otelCarrier"]["traceparent"] == tp + finally: + server.shutdown() + + +def test_do_post_without_traceparent_adds_no_carrier(): + """无 traceparent header → job 不带 _otelCarrier(纯增量、不无中生有;默认关行为不变)。""" + gd = _make_game_dir(_tmp()) + captured = {} + state = W.WorkerState( + run_fn=lambda job: captured.update(job=job) or ({"verdict": {"pass": True}}, gd), + send_fn=lambda url, payload, secret: (200, ""), + profile_fn=lambda j, g: None, + ) + W.start_worker(state) + server, port = W.start_server(state, port=0) + try: + body = json.dumps({"job_id": "n1", "traceId": "n1", "brief": "x", "gameId": "w1", + "callback": {"target": "http://cb"}}).encode("utf-8") + req = urllib.request.Request(f"http://127.0.0.1:{port}/generate", data=body, method="POST", + headers={"Content-Type": "application/json"}) + with urllib.request.urlopen(req, timeout=5) as resp: + assert resp.status == 202 + state.queue.join() + assert "_otelCarrier" not in captured["job"] + finally: + server.shutdown() + + if __name__ == "__main__": _fns = [v for k, v in sorted(globals().items()) if k.startswith("test_") and callable(v)] _failed = 0 diff --git a/cheap-worker/worker_service.py b/cheap-worker/worker_service.py index eb526ef6..9a31d36e 100644 --- a/cheap-worker/worker_service.py +++ b/cheap-worker/worker_service.py @@ -26,6 +26,7 @@ from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer from pathlib import Path import cheap_classify +import cheap_otlp_sink # 面四三跳传播:worker span 作用域 + W3C carrier 提取(顶层零重依赖,otel 惰性) import dedup import result_out @@ -217,7 +218,15 @@ def process_job(state: WorkerState, job: dict) -> dict: return _process_modify_job(state, job, modify) # ↓↓↓ 以下 create/生成路保持原样、零改动 ↓↓↓ - summary, game_dir = state.run_fn(job) + # 面四断点①(span):在入站 W3C context 下开 worker 这跳的 SERVER span,包住真生成调用(run_fn 内经 + # driver 打 :8300)。span attach 为当前 context 后,driver 的 httpx inject 与写工厂 sidecar 的 traceparent + # 都据它继承,使 cheap-service 生成 span 挂到本 span(→ Java span)下当子 span,串成单条 trace 树。 + # 业务 traceId 并存(§10 决策5):gameId 作 conversation.id 与 cheap-service 生成 span 同键、业务 traceId 另存属性。 + # worker_span_scope 全 best-effort:OTLP 未开 / otel 不可用时退化成透传入站 context 或什么都不做,绝不咬生成。 + _game_id = job.get("gameId") or job.get("job_id") + with cheap_otlp_sink.worker_span_scope(job.get("_otelCarrier"), + conversation_id=_game_id, business_trace_id=trace_id): + summary, game_dir = state.run_fn(job) # ── M3b U3:D9 反同质化(vendored dedup)。撞重只告警不阻断(status 仍按九门判); # check_similarity 内部已对 FS 异常返 dupHit=None,另包一层 try 兜底——D9 整体失败不得让回调崩。── @@ -516,6 +525,20 @@ def make_handler(state: WorkerState): if job is None: self._send_json(400, {"accepted": False, "error": "bad job json"}) return + # 面四断点①(capture):抓入站 W3C traceparent 存进 job,供 worker 线程建 worker span(见 process_job)。 + # do_POST 只入队即 202 投递握手,真生成在 worker 线程发生;self.headers 仅这里可得,故此处只抓 header、 + # 不建 span(span 建在真处理该 job 的 worker 线程,才能覆盖「调 cheap-service 生成」这一跳并当下游的父)。 + # 纯增量 best-effort:无 traceparent / 抓取异常都不影响入队与生成(设计 §8 旁路铁律)。 + try: + _tp = self.headers.get("traceparent") + if _tp: + _carrier = {"traceparent": _tp} + _ts = self.headers.get("tracestate") + if _ts: + _carrier["tracestate"] = _ts + job["_otelCarrier"] = _carrier + except Exception: # noqa: BLE001 —— 抓 traceparent 是旁路,失败绝不咬入队/生成 + pass trace_id = job.get("traceId") or job.get("job_id") accepted, code = try_enqueue(state, job) log(f"收到 job trace_id={trace_id}, templateId={job.get('templateId')}, "