已合并
【msserviceprofiler】【需求】【Tracing 1/3】支持通过Hook机制实现自定义tracing埋点 #446
ChaseChe77创建于 8月21日
【msserviceprofiler】【需求】【Tracing 1/3】支持通过Hook机制实现自定义tracing埋点 #446
已合并
共 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 特性设计说明书.md | 149 | - 26.0.0 特性设计: design/MindStudio Service Profiler 26.0.0 特性设计说明书.md |
| 150 | - ms-service-metric 监控指标设计: design/ms_service_metric_Monitoring_Metrics_Design.md | 150 | - ms-service-metric 监控指标设计: design/ms_service_metric_Monitoring_Metrics_Design.md |
| 151 | - ms-service-metric 配置重构设计: design/ms_service_metric_Refactor_Design.md | 151 | - 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 | + | ||
| 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 | + | ||
| 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 | |||
| 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 | + | ||
| 63 | + def activate(span: HookTraceSpan): | ||
| 64 | + return span.activate() if span is not None else None | ||
| 65 | + | ||
| 66 | + | ||
| 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 | + | ||
| 53 | +class HookSpanContext: | ||
| 54 | + native: Any = None | ||
| 55 | + | ||
| 56 | + | ||
| 57 | + def is_valid(self) -> bool: | ||
| 58 | + return bool(self.native is not None and getattr(self.native, "is_valid", False)) | ||
| 59 | + | ||
| 60 | + | ||
| 61 | + def trace_id(self) -> str: | ||
| 62 | + return format(self.native.trace_id, "032x") if self.is_valid else "" | ||
| 63 | + | ||
| 64 | + | ||
| 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 | + | ||
| 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 | + | ||
| 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 | + | ||
| 157 | + def enabled(self) -> bool: | ||
| 158 | + return os.environ.get("MS_TRACE_ENABLE") == "1" | ||
| 159 | + | ||
| 160 | + | ||
| 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 | + | ||
| 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 【review】【错误处理】 问题:即使 provider.add_span_processor(processor) 抛出异常,当前代码仍会把 provider identity 加入 _perfetto_providers,后续同一个 Provider 即使 Perfetto Forwarder 恢复也不会再尝试注册,导致 Perfetto 输出长期缺失且只有 debug 日志可见。 修改建议:仅在 add_span_processor 成功后记录 identity;失败时保持未注册状态,必要时增加退避重试,避免每个 Span 都立即重试。
![]() ![]() | |||
| 208 | + return provider | ||
| 209 | + | ||
| 210 | + | ||
| 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 | + | ||
| 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 | + | ||
| 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] | ||


【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)]