已合并
【msserviceprofiler】【需求】【Tracing 1/3】支持通过Hook机制实现自定义tracing埋点 #446
【msserviceprofiler】【需求】【Tracing 1/3】支持通过Hook机制实现自定义tracing埋点 #446
已合并
ChaseChe77创建于 8月21日
6 个文件变更+963-0
@@ -0,0 +1,242 @@
1+# vLLM Hook Tracing 特性详细设计
2+ 
3+## 1. 文档信息
4+ 
5+| 项目 | 内容 |
6+|---|---|
7+| 特性 | vLLM/vLLM-Ascend Hook Tracing |
8+| 仓库 | msserviceprofiler |
9+| 版本 | 26.0.0 |
10+| 日期 | 2026-08-15 |
11+| 状态 | 开发验证 |
12+ 
13+## 2. 背景与设计结论
14+ 
15+Metrics 只能反映聚合结果,无法解释单个推理请求在请求接入、调度、模型执行和输出处理阶段的耗时关系。本特性通过现有 Hook/YAML 机制补充业务 Span,并支持 Jaeger 和 Perfetto 展示。
16+ 
17+vLLM 已提供基于 OpenTelemetry 的原生 Tracing。本方案不再实现第二套 Provider、OTLP exporter 或跨进程协议,而是将 vLLM 原生 Tracing 作为必选基础能力:
18+ 
19+1. vLLM 必须通过 `--otlp-traces-endpoint` 初始化全局 `TracerProvider`
20+2. Hook Span 复用该 Provider、采样策略、上下文和 OTLP exporter。
21+3. Jaeger 是原生 OTLP 输出;Perfetto 是 Hook Span 的可选附加输出。
22+4. 不支持“vLLM 未开启原生 Tracing,但 msServiceProfiler 单独创建 Trace”的兼容模式。
23+5. 工具架构硬约束:当前及后续版本均不得修改 vLLM/vLLM-Ascend 源码或 IPC,也不导入 `vllm.tracing` 私有接口。
24+ 
25+## 3. 目标与非目标
26+ 
27+### 3.1 目标
28+ 
29+- 通过同一份 profiling YAML 的 `trace` 字段配置 request、scheduler、model、output 埋点。
30+- 复用 vLLM 全局 OpenTelemetry Provider,将 Hook Span 随原生 Span 一起导出到 Jaeger。
31+- 可选将 msServiceProfiler Hook Span 旁路写成 Perfetto/Chrome Trace JSON。
32+- tracing 失败必须 fail-open,不能改变推理结果。
33+- profiling、metrics、MindIE C++ Trace 和原 OTLP Forwarder 保持原行为。
34+- tracing 与 profiling 可同时开启,同一业务函数只执行一次。
35+ 
36+### 3.2 非目标
37+ 
38+- 不支持没有原生 Tracing 能力的旧版 vLLM。
39+- 不在 msServiceProfiler 内创建私有 `TracerProvider`
40+- 不将 vLLM 原生 Span 或旧 C++ OTLP 数据转换为 Perfetto。
41+- 不修改、扩展或要求上游适配 vLLM 的进程间通信协议;该约束适用于当前及后续所有版本。
42+- 不保证 YAML 中已删除或改名的业务符号仍可 Hook;缺失符号按现有机制跳过。
43+- 不改变 profiling 的 enable.json、采集、CSV/DB/timeline 和 parse 流程。
44+ 
45+## 4. 总体架构
46+ 
47+```mermaid
48+flowchart LR
49+ Y["内置或用户 YAML"] --> H["现有 Hook Framework"]
50+ H --> A["Tracing around Adapter"]
51+ A --> R["HookTraceRuntime"]
52+ R --> B["OpenTelemetryHookBackend"]
53+ B --> G["vLLM 全局 TracerProvider"]
54+ G --> O["vLLM 原生 OTLP Processor"]
55+ O --> J["Jaeger/OTLP Collector"]
56+ G --> P["可选 PerfettoSpanProcessor"]
57+ P --> S["MSP_PERFETTO_SOCKET"]
58+ S --> F["PerfettoForwarderService"]
59+ F --> C["Chrome Trace JSON"]
60+ 
61+ M["原 MindIE C++ Tracer"] --> L["OTLP_SOCKET"]
62+ L --> X["原 OTLPForwarderService"]
63+```
64+ 
65+两条 Socket 严格隔离:
66+ 
67+- `OTLP_SOCKET`:原 MindIE C++ Trace 的二进制 OTLP 通道,代码和行为不变。
68+- `MSP_PERFETTO_SOCKET`:仅接收 msServiceProfiler Hook Span 的规范化事件包。
69+ 
70+Perfetto 事件包不会进入 OTLP exporter;旧二进制 OTLP 数据也不会进入 Perfetto Forwarder。
71+ 
72+## 5. 关键模块
73+ 
74+| 模块 | 职责 | 边界 |
75+|---|---|---|
76+| `patcher/core/module_hook.py` | 提供一个可选 `around_hook_factory` 扩展点 | 原 context hook 和 profiling handler 逻辑不变 |
77+| `patcher/core/config_loader.py` | 从同一份 YAML 解析 profiling 与 `trace` | 同一解析后符号只允许一个有效 trace |
78+| `patcher/core/trace_hook.py` | 参数/返回值语义提取和 Span 包围逻辑 | 不进入原 profiling context hook 生命周期 |
79+| `tracer/hook_runtime.py` | request context 表、Links 和 Span 公共属性 | 不创建重复 request 根 Span |
80+| `tracer/otel_hook.py` | 通过 OTel 公共 API复用全局 Provider | 不创建、不替换、不关闭 Provider |
81+| `tracer/perfetto_socket.py` | 过滤并异步发送 Hook Span | 不发送 vLLM 原生 Span,不编码 OTLP protobuf |
82+| `tracer/perfetto_forward_service.py` | 独立接收 Hook 包并写 JSON | 不复用或修改原 OTLP Scheduler |
83+| `tracer/perfetto_exporter.py` | Hook 事件包转 Chrome Trace Event | 拒绝旧二进制 OTLP 数据 |
84+| `tracer/otlp_forward_service.py` | 原 MindIE Trace 转发 | 保持原实现,不参与 vLLM Hook Tracing |
85+ 
86+## 6. Hook 组合方式
87+ 
88+### 6.1 最小通用扩展
89+ 
90+现有流程先生成 profiling callable,再由可选的 tracing `around_hook_factory` 包在最外层:
91+ 
92+```text
93+调用方 -> tracing around -> profiling handler(可选)-> 原业务函数
94+```
95+ 
96+不再给通用 `FunctionContext` 增加 `args/kwargs/return_value`,也不重写多个 context manager 的进入退出流程。这样可以把共享 Hook 框架改动限制为一个可选扩展点。
97+ 
98+### 6.2 异常和重复调用
99+ 
100+`TrackableOriginalFunc` 仍负责记录业务函数是否已经执行。若业务函数抛出异常,外层 Hook 必须原样抛出,不作为埋点失败重试;若 tracing 后处理失败,返回已缓存的业务结果,避免第二次执行原函数。
101+ 
102+### 6.3 配置约束
103+ 
104+profiling 和 tracing 是同一符号上的两个独立能力,不是多个 profiling handler:
105+ 
106+```yaml
107+- symbol: vllm.v1.core.sched.scheduler:Scheduler.schedule
108+ handler: ms_service_profiler.patcher.vllm.handlers.v1.scheduler_handlers:schedule
109+ trace:
110+ name: vllm.scheduler.schedule
111+ domain: Schedule
112+ adapter: schedule
113+```
114+ 
115+同一解析后符号只能有一个有效 `trace` 配置。重复项记录告警并忽略后项,避免一层业务调用生成重复 Hook Span。
116+ 
117+## 7. Provider 与请求关联
118+ 
119+### 7.1 Provider 规则
120+ 
121+运行时只接受 `trace.get_tracer_provider()` 返回且支持 `add_span_processor` 的 SDK Provider:
122+ 
123+- Provider 存在:创建 Hook Span。
124+- Provider 不存在:记录一次提示,Hook tracing 变为 no-op,推理继续。
125+- msServiceProfiler 不读取 OTLP endpoint 自建 Provider,不调用 `set_tracer_provider()`,也不关闭 vLLM Provider。
126+ 
127+### 7.2 请求关联
128+ 
129+- 优先从 vLLM 请求参数/对象读取 `request_id``trace_headers`
130+- 使用 W3C `traceparent` 提取真实上下文,并按 `request_id` 暂存。
131+- scheduler/model/output Span 使用 OTel Link 关联已知请求上下文。
132+- 不根据 request ID 伪造 trace ID。
133+- 若进程内拿不到真实上下文,只保留 `request.ids` 属性,不伪造 Link。
134+- 当前及后续版本均不修改 vLLM IPC。跨进程能否获得 `trace_headers` 只取决于当前 vLLM 已公开的数据;获取不到时降级为 `request.ids` 属性,不把修改上游作为补偿方案。
135+- 若未来 vLLM 通过稳定公共接口提供更多 Trace Context,工具可通过版本适配层选择性读取;不能要求 vLLM 为本工具增加字段、消息或反向依赖。
136+ 
137+## 8. 输出模式
138+ 
139+### 8.1 Jaeger
140+ 
141+```bash
142+export MS_TRACE_ENABLE=1
143+export OTEL_SERVICE_NAME=vllm-server
144+export OTEL_EXPORTER_OTLP_TRACES_PROTOCOL=http/protobuf
145+ 
146+vllm serve MODEL \
147+ --otlp-traces-endpoint http://127.0.0.1:4318/v1/traces
148+```
149+ 
150+vLLM 原生 Span 和 Hook Span 共用 Provider,直接发送到 Jaeger/Collector;不需要运行 `python -m ms_service_profiler.trace`
151+ 
152+### 8.2 Jaeger + Perfetto
153+ 
154+```bash
155+python -m ms_service_profiler.trace \
156+ --perfetto-output /tmp/hook_tracing.json &
157+ 
158+export MS_TRACE_ENABLE=1
159+export OTEL_SERVICE_NAME=vllm-server
160+export OTEL_EXPORTER_OTLP_TRACES_PROTOCOL=http/protobuf
161+ 
162+vllm serve MODEL \
163+ --otlp-traces-endpoint http://127.0.0.1:4318/v1/traces
164+```
165+ 
166+`PerfettoSpanProcessor` 只筛选 instrumentation scope 以 `ms_service_profiler.hook` 开头的 Span,因此 JSON 只包含自定义 Hook Span;Jaeger 仍接收 vLLM 原生 Span和 Hook Span。
167+ 
168+### 8.3 不支持的模式
169+ 
170+`python -m ms_service_profiler.trace --perfetto-output ...` 只启动文件接收端,不能自行创建 Span。若 vLLM 未配置 `--otlp-traces-endpoint`,不存在可复用 Provider,Perfetto 文件将保持空数组。这不是 Perfetto-only 模式。
171+ 
172+## 9. 用户接口
173+ 
174+| 配置 | 来源 | 是否新增 | 作用 |
175+|---|---|---:|---|
176+| `MS_TRACE_ENABLE=1` | msServiceProfiler 既有 Trace 开关 | 否 | 开启 Hook tracing |
177+| `PROFILING_SYMBOLS_PATH` | 既有环境变量 | 否 | 可选覆盖同一份 Hook YAML |
178+| `--otlp-traces-endpoint` | vLLM 原生命令参数 | 否 | 必选,初始化 Provider 并配置 OTLP 输出 |
179+| `OTEL_EXPORTER_OTLP_TRACES_PROTOCOL` | OTel 标准环境变量 | 否 | 配置 vLLM exporter 协议 |
180+| `--perfetto-output` | 本特性 CLI 参数 | 是 | 可选增加 Chrome Trace JSON 输出 |
181+ 
182+没有独立 tracing YAML,也没有新增 backend 选择环境变量。默认 YAML 位于 msServiceProfiler;用户 YAML 仍通过 `PROFILING_SYMBOLS_PATH` 覆盖。
183+ 
184+## 10. 兼容性和隔离
185+ 
186+| 场景 | 行为 |
187+|---|---|
188+| vLLM 支持并开启原生 tracing | 正常创建 Hook Span |
189+| vLLM 不支持或未开启原生 tracing | Hook tracing no-op,不属于支持场景 |
190+| vLLM tracing 私有 API 变化 | 不受影响;实现只使用 OTel 公共 API |
191+| YAML 部分符号不存在 | SymbolWatcher 跳过,其他 Hook 继续 |
192+| 新 msServiceProfiler + 旧 vLLM-Ascend | 使用包内 default YAML,存在的符号正常 Hook |
193+| 旧 `libms_service_profiler.so` | 可用;本特性无新增 C/C++ ABI |
194+| profiling 与 tracing 同时开启 | tracing 包裹 profiling;业务函数执行一次 |
195+| Jaeger 不可用 | vLLM exporter 按自身策略处理;Hook fail-open |
196+| Perfetto Forwarder 不可用 | 不注册旁路 Processor;Jaeger不受影响 |
197+ 
198+“兼容旧 vLLM-Ascend”只表示 YAML/符号按现有机制兼容,不表示支持没有 vLLM 原生 Tracing 的 vLLM 核心版本。
199+ 
200+## 11. 性能、安全与可靠性
201+ 
202+- `MS_TRACE_ENABLE` 未开启时不访问 Provider、不创建 Span。
203+- Perfetto 发送使用有界队列和后台线程,推理线程不执行文件 I/O。
204+- request context 表、Link 数量、属性数量和长度均有限制。
205+- 不记录 prompt 正文。
206+- 输出路径检查普通文件、属主、软链接和 `.json` 扩展名。
207+- tracing 与 profiling 同开时,Span 包含少量 profiling handler 处理开销,这是已接受的语义。
208+- 所有 tracing 异常 fail-open;业务异常保持原类型和对象,不触发重复执行。
209+ 
210+## 12. 验证设计
211+ 
212+### 12.1 单元/回归验证
213+ 
214+- 全局 Provider 存在/不存在、Tracing 关闭和 OTel 依赖缺失。
215+- W3C 上下文、request Links、scheduler 显式时间和完成状态。
216+- 同一符号重复 trace 配置只生成一个有效 Hook。
217+- sync/async around Hook 返回值与异常透传。
218+- profiling + tracing 业务函数只执行一次。
219+- Perfetto Processor 只接收 Hook scope。
220+- Perfetto Forwarder 使用独立 Socket;旧 OTLP Forwarder 回归测试不变。
221+- 规范化 Hook 包生成合法 Chrome Complete/Flow Event,旧二进制 OTLP 被拒绝。
222+ 
223+### 12.2 真实 NPU 验证
224+ 
225+必须通过 `vllm bench serve``/v1/completions` 产生真实推理负载;`/health` 只用于就绪检查。验证项包括:
226+ 
227+1. Jaeger 中同时出现 vLLM 原生 Span 和预期 Hook Span。
228+2. 双写时 Perfetto JSON 包含 Hook Span,且 Jaeger 不受影响。
229+3. tracing-only、profiling-only、两者同时开启三种场景推理均成功。
230+4. 停止 Jaeger 或 Perfetto Forwarder 后推理继续。
231+5. 按实际卡号设置 `ASCEND_RT_VISIBLE_DEVICES`,且不得覆盖掉 vLLM-Ascend 平台插件。
232+ 
233+## 13. 评审需确认事项
234+ 
235+1. 接受 vLLM 原生 Tracing 为前置条件,不支持无 Provider 的旧 vLLM 核心。
236+2. 接受 Perfetto 只输出 msServiceProfiler Hook Span,不转换 vLLM 原生 Span或旧 C++ OTLP。
237+3. 接受 tracing + profiling Span 包含少量 profiling handler 开销。
238+4. YAML 符号适配范围和基于现有公开数据能够达到的跨进程 Link 完整性在 NPU 环境验收;不得通过修改 vLLM IPC提高完整性。
239+ 
240+## 14. 回滚
241+ 
242+运行时取消 `MS_TRACE_ENABLE` 即可关闭 Hook tracing。代码回滚可先移除 YAML `trace` 字段,再移除 Python OTel Backend 与独立 Perfetto 旁路;原 C++ Tracer、`OTLP_SOCKET`、OTLP Forwarder 和 profiling 能力不参与回滚。
@@ -149,3 +149,4 @@ nav:
149 - 26.0.0 特性设计: design/MindStudio Service Profiler 26.0.0 特性设计说明书.md149 - 26.0.0 特性设计: design/MindStudio Service Profiler 26.0.0 特性设计说明书.md
150 - ms-service-metric 监控指标设计: design/ms_service_metric_Monitoring_Metrics_Design.md150 - ms-service-metric 监控指标设计: design/ms_service_metric_Monitoring_Metrics_Design.md
151 - ms-service-metric 配置重构设计: design/ms_service_metric_Refactor_Design.md151 - ms-service-metric 配置重构设计: design/ms_service_metric_Refactor_Design.md
152+ - vLLM Hook Tracing 详细设计: design/vLLM_Hook_Tracing_Detailed_Design.md
@@ -0,0 +1,128 @@
1+# -------------------------------------------------------------------------
2+# This file is part of the MindStudio project.
3+# Copyright (c) 2026 Huawei Technologies Co.,Ltd.
4+#
5+# MindStudio is licensed under Mulan PSL v2.
6+# -------------------------------------------------------------------------
7+ 
8+"""Request correlation and lifecycle for OpenTelemetry Hook spans."""
9+ 
10+import os
11+import threading
12+from collections import OrderedDict
13+from dataclasses import dataclass
14+from typing import Iterable, Optional
15+ 
16+from .otel_hook import HookSpanContext, HookTraceSpan, get_hook_tracer_backend, new_noop_hook_span
17+ 
18+ 
19+MAX_INFLIGHT_REQUESTS = 10000
20+MAX_LINKS_PER_SPAN = 128
21+ 
22+ 
23+@dataclass
24+class _RequestTrace:
25+ context: HookSpanContext
26+ 
27+ 
28+class HookTraceRuntime:
29+ def __init__(self, backend=None):
30+ self._backend = backend or get_hook_tracer_backend()
31+ self._requests = OrderedDict()
32+ self._lock = threading.RLock()
33+ 
34+ @property
35+ def enabled(self) -> bool:
36+ return self._backend.enabled
37+ 
38+ def start_span(
39+ self,
40+ name: str,
41+ domain: str,
42+ kind: str = "INTERNAL",
43+ attributes=None,
44+ request_ids: Optional[Iterable[str]] = None,
45+ start_time_ns: Optional[int] = None,
46+ ) -> HookTraceSpan:
47+ if not self.enabled:
48+ return new_noop_hook_span()
49+ normalized_request_ids = [str(item) for item in list(request_ids or [])[:MAX_LINKS_PER_SPAN]]
J
Jjia_ya_nan8月21日

【review】【性能】 ms_service_profiler/tracer/hook_runtime.py 第 49 行

问题:这里先对 request_ids 执行 list(request_ids or []) 再切片,若调用方传入的是生成器或包含大量 request_id 的可迭代对象,会在推理线程中完整消费并分配内存,MAX_LINKS_PER_SPAN 无法限制转换成本,极端情况下会造成明显延迟或内存抖动。

修改建议:使用 itertools.islice 在迭代阶段限制数量,避免先展开完整可迭代对象;request_links 中同类逻辑也建议同步调整。

from itertools import islice

normalized_request_ids = [str(item) for item in islice(request_ids or [], MAX_LINKS_PER_SPAN)]
likedislike
50+ links = self.request_links(normalized_request_ids)
51+ span = self._backend.start_span(name, domain, kind, links=links, start_time_ns=start_time_ns)
52+ if span.is_recording:
53+ span.set_attribute("process.pid", os.getpid())
54+ native_id = getattr(threading, "get_native_id", threading.get_ident)()
55+ span.set_attribute("thread.id", native_id)
56+ if normalized_request_ids:
57+ span.set_attribute("request.ids", normalized_request_ids)
58+ for key, value in (attributes or {}).items():
59+ span.set_attribute(key, value)
60+ return span
61+ 
62+ @staticmethod
63+ def activate(span: HookTraceSpan):
64+ return span.activate() if span is not None else None
65+ 
66+ @staticmethod
67+ def deactivate(token) -> None:
68+ HookTraceSpan.deactivate(token)
69+ 
70+ def _store_request(self, request_id: str, request_trace: _RequestTrace) -> None:
71+ with self._lock:
72+ self._requests.pop(request_id, None)
73+ self._requests[request_id] = request_trace
74+ if len(self._requests) > MAX_INFLIGHT_REQUESTS:
75+ self._requests.popitem(last=False)
76+ 
77+ def register_request_context(self, request_id: str, trace_headers=None) -> bool:
78+ request_id = str(request_id)
79+ context = self._backend.context_from_headers(trace_headers)
80+ if not context.is_valid:
81+ context = self._backend.current_context()
82+ if not context.is_valid:
83+ return False
84+ self._store_request(request_id, _RequestTrace(context=context))
85+ return True
86+ 
87+ def start_request(self, request_id: str, trace_headers=None) -> bool:
88+ """Register vLLM's real request context; never create a duplicate root span."""
89+ return self.register_request_context(request_id, trace_headers)
90+ 
91+ def finish_request(self, request_id: str, success: bool = True, message: str = "") -> bool:
92+ with self._lock:
93+ request_trace = self._requests.pop(str(request_id), None)
94+ return request_trace is not None
95+ 
96+ def request_links(self, request_ids: Iterable[str]):
97+ with self._lock:
98+ request_traces = [
99+ (str(request_id), self._requests.get(str(request_id)))
100+ for request_id in list(request_ids)[:MAX_LINKS_PER_SPAN]
101+ ]
102+ return [
103+ (request_trace.context, request_id)
104+ for request_id, request_trace in request_traces
105+ if request_trace is not None and request_trace.context.is_valid
106+ ]
107+ 
108+ def request_context(self, request_id: str) -> Optional[HookSpanContext]:
109+ with self._lock:
110+ request_trace = self._requests.get(str(request_id))
111+ return request_trace.context if request_trace is not None else None
112+ 
113+ def clear(self) -> None:
114+ with self._lock:
115+ self._requests.clear()
116+ self._backend.shutdown()
117+ 
118+ 
119+_HOOK_TRACE_RUNTIME = HookTraceRuntime()
120+ 
121+ 
122+def get_hook_trace_runtime() -> HookTraceRuntime:
123+ return _HOOK_TRACE_RUNTIME
124+ 
125+ 
126+def _set_hook_trace_runtime_for_test(runtime) -> None:
127+ global _HOOK_TRACE_RUNTIME
128+ _HOOK_TRACE_RUNTIME = runtime
@@ -0,0 +1,280 @@
1+# -------------------------------------------------------------------------
2+# This file is part of the MindStudio project.
3+# Copyright (c) 2026 Huawei Technologies Co.,Ltd.
4+#
5+# MindStudio is licensed under Mulan PSL v2.
6+# -------------------------------------------------------------------------
7+ 
8+"""Hook tracing that reuses only the active vLLM OpenTelemetry provider."""
9+ 
10+import os
11+import threading
12+from dataclasses import dataclass
13+from typing import Any, Iterable, Optional, Tuple
14+ 
15+from ms_service_profiler.tracer.perfetto_socket import HOOK_SCOPE_PREFIX, PerfettoSocketSender, PerfettoSpanProcessor
16+from ms_service_profiler.utils.log import logger
17+ 
18+ 
19+MAX_ATTRIBUTE_COUNT = 32
20+MAX_ATTRIBUTE_KEY_LENGTH = 128
21+MAX_ATTRIBUTE_VALUE_LENGTH = 1024
22+ 
23+ 
24+try:
25+ from opentelemetry import context as otel_context
26+ from opentelemetry import trace
27+ from opentelemetry.trace import Link, NonRecordingSpan, SpanKind, Status, StatusCode
28+ from opentelemetry.trace.propagation.tracecontext import TraceContextTextMapPropagator
29+ 
30+ _OTEL_AVAILABLE = True
31+except ImportError:
32+ otel_context = None
33+ trace = None
34+ Link = None
35+ NonRecordingSpan = None
36+ SpanKind = None
37+ Status = None
38+ StatusCode = None
39+ TraceContextTextMapPropagator = None
40+ _OTEL_AVAILABLE = False
41+ 
42+ 
43+SPAN_KINDS = {
44+ "INTERNAL": "INTERNAL",
45+ "SERVER": "SERVER",
46+ "CLIENT": "CLIENT",
47+ "PRODUCER": "PRODUCER",
48+ "CONSUMER": "CONSUMER",
49+}
50+ 
51+ 
52+@dataclass(frozen=True)
53+class HookSpanContext:
54+ native: Any = None
55+ 
56+ @property
57+ def is_valid(self) -> bool:
58+ return bool(self.native is not None and getattr(self.native, "is_valid", False))
59+ 
60+ @property
61+ def trace_id(self) -> str:
62+ return format(self.native.trace_id, "032x") if self.is_valid else ""
63+ 
64+ @property
65+ def span_id(self) -> str:
66+ return format(self.native.span_id, "016x") if self.is_valid else ""
67+ 
68+ 
69+class HookTraceSpan:
70+ """Fail-open facade over an OTel span owned by vLLM's provider."""
71+ 
72+ def __init__(self, span=None):
73+ self._span = span
74+ native_context = span.get_span_context() if span is not None else None
75+ self.context = HookSpanContext(native_context)
76+ self._ended = False
77+ self._attribute_count = 0
78+ self._lock = threading.Lock()
79+ 
80+ @property
81+ def is_recording(self) -> bool:
82+ return bool(self._span is not None and not self._ended and self._span.is_recording())
83+ 
84+ def set_attribute(self, key, value) -> None:
85+ if not self.is_recording or self._attribute_count >= MAX_ATTRIBUTE_COUNT:
86+ return
87+ if not isinstance(key, str) or not key or len(key) > MAX_ATTRIBUTE_KEY_LENGTH:
88+ return
89+ safe_value = _safe_attribute(value)
90+ if safe_value is None:
91+ return
92+ try:
93+ self._span.set_attribute(key, safe_value)
94+ self._attribute_count += 1
95+ except Exception as exc:
96+ logger.debug("Failed to set Hook tracing attribute: %s", exc)
97+ 
98+ def activate(self):
99+ if not self.is_recording or not _OTEL_AVAILABLE:
100+ return None
101+ try:
102+ return otel_context.attach(trace.set_span_in_context(self._span))
103+ except Exception as exc:
104+ logger.debug("Failed to activate Hook tracing span: %s", exc)
105+ return None
106+ 
107+ @staticmethod
108+ def deactivate(token) -> None:
109+ if token is None or not _OTEL_AVAILABLE:
110+ return
111+ try:
112+ otel_context.detach(token)
113+ except Exception as exc:
114+ logger.debug("Failed to deactivate Hook tracing span: %s", exc)
115+ 
116+ def end(self, success: bool = True, message: str = "", end_time_ns: Optional[int] = None) -> None:
117+ with self._lock:
118+ if self._span is None or self._ended:
119+ return
120+ span = self._span
121+ self._ended = True
122+ try:
123+ status_code = StatusCode.OK if success else StatusCode.ERROR
124+ span.set_status(Status(status_code, str(message)[:MAX_ATTRIBUTE_VALUE_LENGTH] or None))
125+ span.end(end_time=end_time_ns)
126+ except Exception as exc:
127+ logger.debug("Failed to end Hook tracing span: %s", exc)
128+ 
129+ 
130+def new_noop_hook_span() -> HookTraceSpan:
131+ return HookTraceSpan()
132+ 
133+ 
134+def _safe_attribute(value):
135+ if isinstance(value, (str, bool, int, float)):
136+ return value[:MAX_ATTRIBUTE_VALUE_LENGTH] if isinstance(value, str) else value
137+ if isinstance(value, (list, tuple)):
138+ values = []
139+ for item in value[:MAX_ATTRIBUTE_COUNT]:
140+ if not isinstance(item, (str, bool, int, float)):
141+ return None
142+ values.append(item[:MAX_ATTRIBUTE_VALUE_LENGTH] if isinstance(item, str) else item)
143+ return values
144+ return None
145+ 
146+ 
147+class OpenTelemetryHookBackend:
148+ """Use vLLM's provider; never create or shut down a provider."""
149+ 
150+ def __init__(self):
151+ self._lock = threading.RLock()
152+ self._perfetto_providers = set()
153+ self._warned_unavailable = False
154+ self._warned_provider_missing = False
155+ 
156+ @property
157+ def enabled(self) -> bool:
158+ return os.environ.get("MS_TRACE_ENABLE") == "1"
159+ 
160+ @staticmethod
161+ def _active_global_provider():
162+ if not _OTEL_AVAILABLE:
163+ return None
164+ provider = trace.get_tracer_provider()
165+ return provider if callable(getattr(provider, "add_span_processor", None)) else None
166+ 
167+ @property
168+ def reuses_global_provider(self) -> bool:
169+ return self._active_global_provider() is not None
170+ 
171+ def _get_provider(self):
172+ if not self.enabled:
173+ return None
174+ if not _OTEL_AVAILABLE:
175+ if not self._warned_unavailable:
176+ logger.warning("OpenTelemetry is unavailable; vLLM Hook tracing is disabled")
177+ self._warned_unavailable = True
178+ return None
179+ 
180+ provider = self._active_global_provider()
181+ if provider is None:
182+ if not self._warned_provider_missing:
183+ logger.warning("vLLM OpenTelemetry provider is unavailable; start vLLM with --otlp-traces-endpoint")
184+ self._warned_provider_missing = True
185+ return None
186+ 
187+ identity = id(provider)
188+ with self._lock:
189+ perfetto_registered = identity in self._perfetto_providers
190+ try:
191+ perfetto_available = PerfettoSpanProcessor is not None and PerfettoSocketSender.is_available()
192+ except Exception as exc:
193+ logger.debug("Failed to probe Perfetto forwarder: %s", exc)
194+ perfetto_available = False
195+ if not perfetto_registered and perfetto_available:
196+ with self._lock:
197+ if identity not in self._perfetto_providers:
198+ processor = PerfettoSpanProcessor()
199+ try:
200+ provider.add_span_processor(processor)
201+ except Exception as exc:
202+ try:
203+ processor.shutdown()
204+ except Exception as shutdown_exc:
205+ logger.debug("Failed to shut down Perfetto span processor: %s", shutdown_exc)
206+ logger.debug("Failed to register Perfetto span processor: %s", exc)
207+ self._perfetto_providers.add(identity)
J
Jjia_ya_nan8月21日

【review】【错误处理】 ms_service_profiler/tracer/otel_hook.py 第 207 行

问题:即使 provider.add_span_processor(processor) 抛出异常,当前代码仍会把 provider identity 加入 _perfetto_providers,后续同一个 Provider 即使 Perfetto Forwarder 恢复也不会再尝试注册,导致 Perfetto 输出长期缺失且只有 debug 日志可见。

修改建议:仅在 add_span_processor 成功后记录 identity;失败时保持未注册状态,必要时增加退避重试,避免每个 Span 都立即重试。

try:
    provider.add_span_processor(processor)
except Exception as exc:
    processor.shutdown()
    logger.debug("Failed to register Perfetto span processor: %s", exc)
else:
    self._perfetto_providers.add(identity)
likedislike
208+ return provider
209+ 
210+ @staticmethod
211+ def context_from_headers(headers) -> HookSpanContext:
212+ if not _OTEL_AVAILABLE or not headers:
213+ return HookSpanContext()
214+ try:
215+ extracted = TraceContextTextMapPropagator().extract(dict(headers), context=otel_context.Context())
216+ return HookSpanContext(trace.get_current_span(extracted).get_span_context())
217+ except Exception as exc:
218+ logger.debug("Failed to extract Hook tracing context: %s", exc)
219+ return HookSpanContext()
220+ 
221+ @staticmethod
222+ def current_context() -> HookSpanContext:
223+ if not _OTEL_AVAILABLE:
224+ return HookSpanContext()
225+ try:
226+ return HookSpanContext(trace.get_current_span().get_span_context())
227+ except Exception as exc:
228+ logger.debug("Failed to read current Hook tracing context: %s", exc)
229+ return HookSpanContext()
230+ 
231+ def start_span(
232+ self,
233+ name: str,
234+ domain: str,
235+ kind: str = "INTERNAL",
236+ parent: Optional[HookSpanContext] = None,
237+ links: Optional[Iterable[Tuple[HookSpanContext, str]]] = None,
238+ start_time_ns: Optional[int] = None,
239+ ) -> HookTraceSpan:
240+ provider = self._get_provider()
241+ if provider is None:
242+ return new_noop_hook_span()
243+ try:
244+ parent_context = None
245+ if parent is not None and parent.is_valid:
246+ parent_context = trace.set_span_in_context(NonRecordingSpan(parent.native))
247+ native_links = [
248+ Link(link_context.native, {"request.id": str(request_id)})
249+ for link_context, request_id in links or []
250+ if link_context is not None and link_context.is_valid
251+ ]
252+ native_kind = getattr(SpanKind, SPAN_KINDS.get(kind, "INTERNAL"))
253+ scope_name = "{}.{}".format(HOOK_SCOPE_PREFIX, domain or "tracing")
254+ span = provider.get_tracer(scope_name).start_span(
255+ name=name,
256+ context=parent_context,
257+ kind=native_kind,
258+ links=native_links,
259+ start_time=start_time_ns,
260+ )
261+ return HookTraceSpan(span)
262+ except Exception as exc:
263+ logger.debug("Failed to start Hook tracing span %s: %s", name, exc)
264+ return new_noop_hook_span()
265+ 
266+ def shutdown(self) -> None:
267+ # The provider and its processors are owned by vLLM.
268+ return None
269+ 
270+ 
271+_HOOK_TRACER_BACKEND = OpenTelemetryHookBackend()
272+ 
273+ 
274+def get_hook_tracer_backend() -> OpenTelemetryHookBackend:
275+ return _HOOK_TRACER_BACKEND
276+ 
277+ 
278+def _set_hook_tracer_backend_for_test(backend) -> None:
279+ global _HOOK_TRACER_BACKEND
280+ _HOOK_TRACER_BACKEND = backend
@@ -0,0 +1,170 @@
1+# -------------------------------------------------------------------------
2+# This file is part of the MindStudio project.
3+# Copyright (c) 2026 Huawei Technologies Co.,Ltd.
4+#
5+# MindStudio is licensed under Mulan PSL v2.
6+# -------------------------------------------------------------------------
7+ 
8+"""Non-blocking OpenTelemetry Span delivery to the Perfetto forwarder."""
9+ 
10+import json
11+import os
12+import queue
13+import socket
14+import threading
15+from typing import Any, Dict, Optional
16+ 
17+from ms_service_profiler.utils.log import logger
18+ 
19+ 
20+PERFETTO_SOCKET_NAME = "MSP_PERFETTO_SOCKET"
21+PERFETTO_EVENT_MAGIC = b"MSPF_PERFETTO_V1\n"
22+HOOK_SCOPE_PREFIX = "ms_service_profiler.hook"
23+MAX_QUEUE_SIZE = 10000
24+MAX_ATTRIBUTE_COUNT = 64
25+MAX_ATTRIBUTE_LENGTH = 2048
26+ 
27+ 
28+def _safe_value(value: Any):
29+ if isinstance(value, (bool, int, float)) or value is None:
30+ return value
31+ if isinstance(value, str):
32+ return value[:MAX_ATTRIBUTE_LENGTH]
33+ if isinstance(value, (list, tuple)):
34+ return [_safe_value(item) for item in value[:MAX_ATTRIBUTE_COUNT]]
35+ return str(value)[:MAX_ATTRIBUTE_LENGTH]
36+ 
37+ 
38+def _span_context_ids(span_context) -> Dict[str, str]:
39+ if span_context is None or not getattr(span_context, "is_valid", False):
40+ return {"trace_id": "", "span_id": ""}
41+ return {
42+ "trace_id": format(span_context.trace_id, "032x"),
43+ "span_id": format(span_context.span_id, "016x"),
44+ }
45+ 
46+ 
47+def serialize_readable_span(span) -> bytes:
48+ """Serialize only the stable ReadableSpan surface used by Perfetto."""
49+ context_ids = _span_context_ids(span.get_span_context())
50+ parent_ids = _span_context_ids(getattr(span, "parent", None))
51+ attributes = {
52+ str(key)[:MAX_ATTRIBUTE_LENGTH]: _safe_value(value)
53+ for key, value in list((getattr(span, "attributes", None) or {}).items())[:MAX_ATTRIBUTE_COUNT]
54+ }
55+ resource_attributes = {
56+ str(key)[:MAX_ATTRIBUTE_LENGTH]: _safe_value(value)
57+ for key, value in list((getattr(getattr(span, "resource", None), "attributes", None) or {}).items())[
58+ :MAX_ATTRIBUTE_COUNT
59+ ]
60+ }
61+ links = []
62+ for link in list(getattr(span, "links", None) or [])[:MAX_ATTRIBUTE_COUNT]:
63+ link_ids = _span_context_ids(getattr(link, "context", None))
64+ if link_ids["trace_id"] and link_ids["span_id"]:
65+ links.append(link_ids)
66+ instrumentation_scope = getattr(span, "instrumentation_scope", None)
67+ status = getattr(getattr(span, "status", None), "status_code", None)
68+ kind = getattr(span, "kind", None)
69+ payload = {
70+ "name": str(getattr(span, "name", ""))[:MAX_ATTRIBUTE_LENGTH],
71+ "category": str(getattr(instrumentation_scope, "name", "Tracing"))[:MAX_ATTRIBUTE_LENGTH],
72+ "trace_id": context_ids["trace_id"],
73+ "span_id": context_ids["span_id"],
74+ "parent_span_id": parent_ids["span_id"],
75+ "start_time_ns": int(getattr(span, "start_time", 0) or 0),
76+ "end_time_ns": int(getattr(span, "end_time", 0) or 0),
77+ "kind": int(getattr(kind, "value", kind) or 0),
78+ "status": int(getattr(status, "value", status) or 0),
79+ "attributes": attributes,
80+ "resource_attributes": resource_attributes,
81+ "links": links,
82+ }
83+ return PERFETTO_EVENT_MAGIC + json.dumps(payload, ensure_ascii=False, separators=(",", ":")).encode("utf-8")
84+ 
85+ 
86+class PerfettoSocketSender:
87+ """Queue Span packets so tracing failures and socket I/O never block inference."""
88+ 
89+ def __init__(self, socket_name: str = PERFETTO_SOCKET_NAME):
90+ self._address = "\0" + socket_name
91+ self._queue = queue.Queue(MAX_QUEUE_SIZE)
92+ self._stop_event = threading.Event()
93+ self._thread = threading.Thread(target=self._run, name="perfetto-span-sender", daemon=True)
94+ self._thread.start()
95+ 
96+ @classmethod
97+ def is_available(cls, socket_name: str = PERFETTO_SOCKET_NAME) -> bool:
98+ if os.name != "posix" or not hasattr(socket, "AF_UNIX"):
99+ return False
100+ probe = socket.socket(getattr(socket, "AF_UNIX"), socket.SOCK_STREAM)
101+ probe.settimeout(0.05)
102+ try:
103+ probe.connect("\0" + socket_name)
104+ return True
105+ except OSError:
106+ return False
107+ finally:
108+ probe.close()
109+ 
110+ def submit(self, payload: bytes) -> None:
111+ try:
112+ self._queue.put_nowait(payload)
113+ except queue.Full:
114+ logger.warning("Perfetto tracing queue is full; one Span was discarded")
115+ 
116+ def _send(self, payload: bytes) -> None:
117+ client = socket.socket(getattr(socket, "AF_UNIX"), socket.SOCK_STREAM)
118+ client.settimeout(0.2)
119+ try:
120+ client.connect(self._address)
121+ client.sendall(len(payload).to_bytes(4, byteorder="big") + payload)
122+ finally:
123+ client.close()
124+ 
125+ def _run(self) -> None:
126+ while not self._stop_event.is_set() or not self._queue.empty():
127+ try:
128+ payload = self._queue.get(timeout=0.2)
129+ except queue.Empty:
130+ continue
131+ try:
132+ self._send(payload)
133+ except OSError as exc:
134+ logger.debug("Perfetto forwarder is unavailable; discarded one Span: %s", exc)
135+ 
136+ def shutdown(self) -> None:
137+ self._stop_event.set()
138+ if self._thread.is_alive():
139+ self._thread.join(timeout=2)
140+ 
141+ 
142+try:
143+ from opentelemetry.sdk.trace import ReadableSpan, SpanProcessor
144+ 
145+ class PerfettoSpanProcessor(SpanProcessor):
146+ def __init__(self, sender: Optional[PerfettoSocketSender] = None):
147+ self._sender = sender or PerfettoSocketSender()
148+ 
149+ def on_start(self, span, parent_context=None) -> None:
150+ return None
151+ 
152+ def on_end(self, span: ReadableSpan) -> None:
153+ if os.environ.get("MS_TRACE_ENABLE") != "1":
154+ return
155+ scope_name = str(getattr(getattr(span, "instrumentation_scope", None), "name", ""))
156+ if not scope_name.startswith(HOOK_SCOPE_PREFIX):
157+ return
158+ try:
159+ self._sender.submit(serialize_readable_span(span))
160+ except Exception as exc:
161+ logger.debug("Failed to serialize Span for Perfetto: %s", exc)
162+ 
163+ def shutdown(self) -> None:
164+ self._sender.shutdown()
165+ 
166+ def force_flush(self, timeout_millis: int = 30000) -> bool:
167+ return True
168+ 
169+except ImportError:
170+ PerfettoSpanProcessor = None
@@ -0,0 +1,142 @@
1+# -------------------------------------------------------------------------
2+# This file is part of the MindStudio project.
3+# Copyright (c) 2026 Huawei Technologies Co.,Ltd.
4+#
5+# MindStudio is licensed under Mulan PSL v2.
6+# -------------------------------------------------------------------------
7+ 
8+"""Unit tests for the public-API-only OpenTelemetry Hook backend."""
9+ 
10+from unittest.mock import MagicMock, patch
11+ 
12+from opentelemetry.sdk.trace import TracerProvider
13+from opentelemetry.sdk.trace.export import SimpleSpanProcessor
14+from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter
15+from opentelemetry.trace import SpanContext, TraceFlags, TraceState
16+ 
17+from ms_service_profiler.tracer.otel_hook import HookSpanContext, OpenTelemetryHookBackend
18+from ms_service_profiler.tracer.perfetto_socket import (
19+ PERFETTO_EVENT_MAGIC,
20+ PerfettoSpanProcessor,
21+ serialize_readable_span,
22+)
23+ 
24+ 
25+def _remote_context():
26+ return HookSpanContext(
27+ SpanContext(
28+ trace_id=int("1" * 32, 16),
29+ span_id=int("2" * 16, 16),
30+ is_remote=True,
31+ trace_flags=TraceFlags.SAMPLED,
32+ trace_state=TraceState(),
33+ )
34+ )
35+ 
36+ 
37+def test_backend_reuses_active_global_provider_and_exports_custom_span():
38+ provider = TracerProvider()
39+ exporter = InMemorySpanExporter()
40+ provider.add_span_processor(SimpleSpanProcessor(exporter))
41+ backend = OpenTelemetryHookBackend()
42+ 
43+ with (
44+ patch.dict("os.environ", {"MS_TRACE_ENABLE": "1"}, clear=True),
45+ patch.object(backend, "_active_global_provider", return_value=provider),
46+ patch("ms_service_profiler.tracer.otel_hook.PerfettoSocketSender.is_available", return_value=False),
47+ ):
48+ span = backend.start_span(
49+ "vllm.scheduler.schedule",
50+ "Schedule",
51+ links=[(_remote_context(), "request-1")],
52+ )
53+ span.set_attribute("batch.request_count", 1)
54+ span.end(True)
55+ 
56+ finished = exporter.get_finished_spans()
57+ assert len(finished) == 1
58+ assert finished[0].name == "vllm.scheduler.schedule"
59+ assert finished[0].attributes["batch.request_count"] == 1
60+ assert len(finished[0].links) == 1
61+ assert finished[0].links[0].attributes["request.id"] == "request-1"
62+ 
63+ 
64+def test_backend_requires_vllm_global_provider_and_never_creates_private_provider():
65+ backend = OpenTelemetryHookBackend()
66+ 
67+ with (
68+ patch.dict("os.environ", {"MS_TRACE_ENABLE": "1"}, clear=True),
69+ patch.object(backend, "_active_global_provider", return_value=None),
70+ patch("ms_service_profiler.tracer.otel_hook.logger.warning") as warning,
71+ ):
72+ span = backend.start_span("schedule", "Schedule")
73+ 
74+ assert not span.is_recording
75+ warning.assert_called_once()
76+ assert "--otlp-traces-endpoint" in warning.call_args.args[0]
77+ 
78+ 
79+def test_backend_disabled_is_noop_without_touching_provider():
80+ backend = OpenTelemetryHookBackend()
81+ with patch.dict("os.environ", {}, clear=True), patch.object(backend, "_active_global_provider") as active_provider:
82+ span = backend.start_span("schedule", "Schedule")
83+ 
84+ assert not span.is_recording
85+ active_provider.assert_not_called()
86+ 
87+ 
88+def test_perfetto_processor_registration_failure_does_not_disable_jaeger_provider():
89+ backend = OpenTelemetryHookBackend()
90+ provider = MagicMock()
91+ provider.add_span_processor.side_effect = RuntimeError("registration failed")
92+ processor = MagicMock()
93+ 
94+ with (
95+ patch.dict("os.environ", {"MS_TRACE_ENABLE": "1"}, clear=True),
96+ patch.object(backend, "_active_global_provider", return_value=provider),
97+ patch("ms_service_profiler.tracer.otel_hook.PerfettoSocketSender.is_available", return_value=True),
98+ patch("ms_service_profiler.tracer.otel_hook.PerfettoSpanProcessor", return_value=processor),
99+ ):
100+ assert backend._get_provider() is provider
101+ 
102+ processor.shutdown.assert_called_once()
103+ 
104+ 
105+def test_backend_missing_otel_dependency_is_fail_open():
106+ backend = OpenTelemetryHookBackend()
107+ with (
108+ patch.dict("os.environ", {"MS_TRACE_ENABLE": "1"}, clear=True),
109+ patch("ms_service_profiler.tracer.otel_hook._OTEL_AVAILABLE", False),
110+ ):
111+ span = backend.start_span("schedule", "Schedule")
112+ 
113+ assert not span.is_recording
114+ 
115+ 
116+def test_readable_span_serialization_uses_version_neutral_perfetto_packet():
117+ provider = TracerProvider()
118+ exporter = InMemorySpanExporter()
119+ provider.add_span_processor(SimpleSpanProcessor(exporter))
120+ span = provider.get_tracer("Execute").start_span("vllm.model.execute")
121+ span.set_attribute("request.ids", ["request-1"])
122+ span.end()
123+ 
124+ packet = serialize_readable_span(exporter.get_finished_spans()[0])
125+ 
126+ assert packet.startswith(PERFETTO_EVENT_MAGIC)
127+ assert b'"name":"vllm.model.execute"' in packet
128+ 
129+ 
130+def test_perfetto_processor_exports_only_msserviceprofiler_hook_spans():
131+ provider = TracerProvider()
132+ sender = MagicMock()
133+ provider.add_span_processor(PerfettoSpanProcessor(sender))
134+ 
135+ with patch.dict("os.environ", {"MS_TRACE_ENABLE": "1"}, clear=True):
136+ native_span = provider.get_tracer("vllm.engine").start_span("native-vllm")
137+ native_span.end()
138+ hook_span = provider.get_tracer("ms_service_profiler.hook.Schedule").start_span("hook-schedule")
139+ hook_span.end()
140+ 
141+ sender.submit.assert_called_once()
142+ assert b'"name":"hook-schedule"' in sender.submit.call_args.args[0]