已合并
【msserviceprofiler】【需求】【Tracing 3/3】支持vLLM Hook Tracing导出Perfetto及Jaeger双写 #451
ChaseChe77创建于 12 天前
【msserviceprofiler】【需求】【Tracing 3/3】支持vLLM Hook Tracing导出Perfetto及Jaeger双写 #451
已合并
共 11 个文件变更+707-30
| @@ -27,9 +27,9 @@ msServiceProfiler Trace采集MindIE Motor服务中的请求响应时间、响应 | |||
| 27 | |昇腾910系列产品|×| | 27 | |昇腾910系列产品|×| |
| 28 | 28 | ||
| 29 | > [!NOTE] | 29 | > [!NOTE] |
| 30 | -> | 30 | +> |
| 31 | ->针对昇腾A2系列产品,当前仅支持该系列产品中的Atlas 800I A2 推理服务器。 | 31 | +>针对Atlas A2 训练系列产品/Atlas A2 推理系列产品,当前仅支持该系列产品中的Atlas 800I A2 推理服务器。 |
| 32 | ->针对昇腾310P系列产品,当前仅支持该系列产品中的Atlas 300I Duo 推理卡 + A800-3000推理服务器。 | 32 | +>针对Atlas 推理系列产品,当前仅支持该系列产品中的Atlas 300I Duo 推理卡+Atlas 800 推理服务器(型号:3000)。 |
| 33 | 33 | ||
| 34 | ## 使用前准备<a name="ZH-CN_TOPIC_0000002486482024"></a> | 34 | ## 使用前准备<a name="ZH-CN_TOPIC_0000002486482024"></a> |
| 35 | 35 | ||
| @@ -71,11 +71,11 @@ msServiceProfiler Trace转发数据最大支持400并发,超过400并发可能 | |||
| 71 | 71 | ||
| 72 | 2. 通过配置环境变量支持更灵活的采样控制。 | 72 | 2. 通过配置环境变量支持更灵活的采样控制。 |
| 73 | 73 | ||
| 74 | - | 环境变量名 | 说明 | | 74 | + | 环境变量名 | 说明 | |
| 75 | |------------|------| | 75 | |------------|------| |
| 76 | - | `MS_PROFILER_AUTO_TRACE` | 当请求头中没有传递 trace_id 时,是否自动生成 trace_id。设置为 `1` 时开启自动生成;未设置或设置为其他值时不生成。 | | 76 | + | `MS_PROFILER_AUTO_TRACE` | 当请求头中没有传递 trace_id 时,是否自动生成 trace_id。设置为 `1` 时开启自动生成;未设置或设置为其他值时不生成。 | |
| 77 | - | `MS_PROFILER_SAMPLE_RATE` | 设置采样频率,仅对自动生成 `trace_id` 的请求生效。该值为正整数 N,表示每 N 次请求采样 1 次。若未设置或设置为非正整数,则不采样。 | | 77 | + | `MS_PROFILER_SAMPLE_RATE` | 设置采样频率,仅对自动生成 `trace_id` 的请求生效。该值为正整数 N,表示每 N 次请求采样 1 次。若未设置或设置为非正整数,则不采样。 | |
| 78 | - | `MS_PROFILER_SAMPLE_ERROR` | 是否仅上报错误的请求(适用于所有请求)。设置为 `1` 时仅上报错误 Span;未设置或设置为其他值时上报所有请求。 | | 78 | + | `MS_PROFILER_SAMPLE_ERROR` | 是否仅上报错误的请求(适用于所有请求)。设置为 `1` 时仅上报错误 Span;未设置或设置为其他值时上报所有请求。 | |
| 79 | 79 | ||
| 80 | ```bash | 80 | ```bash |
| 81 | # 设置环境变量示例 | 81 | # 设置环境变量示例 |
| @@ -156,9 +156,11 @@ export MS_PROFILER_SAMPLE_ERROR=1 | |||
| 156 | **命令格式<a name="section10872103414491"></a>** | 156 | **命令格式<a name="section10872103414491"></a>** |
| 157 | 157 | ||
| 158 | ```bash | 158 | ```bash |
| 159 | -python -m ms_service_profiler.trace [--log-level] | 159 | +python -m ms_service_profiler.trace [--log-level] [--perfetto-output] |
| 160 | ``` | 160 | ``` |
| 161 | 161 | ||
| 162 | +> 对于 vLLM Hook Tracing,必须由 vLLM 原生 OpenTelemetry Provider 直接导出到 Jaeger。`--perfetto-output` 仅启动 Hook Span 的附加文件出口,不能脱离 vLLM 原生 Tracing 单独使用。MindIE Motor 仍保持本章原有 Forwarder 用法。详见《vLLM Hook Tracing 使用指南》。 | ||
| 163 | + | ||
| 162 | options参数说明请参见[参数说明](#section379581401015)。 | 164 | options参数说明请参见[参数说明](#section379581401015)。 |
| 163 | 165 | ||
| 164 | **参数说明<a name="section379581401015"></a>** | 166 | **参数说明<a name="section379581401015"></a>** |
| @@ -166,6 +168,7 @@ options参数说明请参见[参数说明](#section379581401015)。 | |||
| 166 | |**参数**|说明|**是否必选**| | 168 | |**参数**|说明|**是否必选**| |
| 167 | |--|--|--| | 169 | |--|--|--| |
| 168 | |--log-level|设置日志级别,取值为:<br>• debug:调试级别。该级别的日志记录了调试信息,便于开发人员或维护人员定位问题。<br>• info:正常级别。记录工具正常运行的信息。默认值。<br>• warning:警告级别。记录工具和预期的状态不一致,但不影响整个进程运行的信息。<br>• error:一般错误级别。<br>• critical:严重错误级别。<br>• fatal:致命错误级别。|否| | 170 | |--log-level|设置日志级别,取值为:<br>• debug:调试级别。该级别的日志记录了调试信息,便于开发人员或维护人员定位问题。<br>• info:正常级别。记录工具正常运行的信息。默认值。<br>• warning:警告级别。记录工具和预期的状态不一致,但不影响整个进程运行的信息。<br>• error:一般错误级别。<br>• critical:严重错误级别。<br>• fatal:致命错误级别。|否| |
| 171 | +|--perfetto-output|vLLM Hook Tracing 的可选 Perfetto/Chrome Trace JSON 输出路径。使用时 vLLM 必须同时通过 `--otlp-traces-endpoint` 开启原生 Tracing;该参数不改变 MindIE Motor 原 OTLP Forwarder。|否| | ||
| 169 | 172 | ||
| 170 | **使用示例<a name="section246434914919"></a>** | 173 | **使用示例<a name="section246434914919"></a>** |
| 171 | 174 | ||
| @@ -175,6 +178,12 @@ options参数说明请参见[参数说明](#section379581401015)。 | |||
| 175 | python -m ms_service_profiler.trace | 178 | python -m ms_service_profiler.trace |
| 176 | ``` | 179 | ``` |
| 177 | 180 | ||
| 181 | +vLLM 已通过 `--otlp-traces-endpoint` 开启原生 Tracing 后,如需附加生成 Perfetto 可直接打开的 Chrome Trace JSON,命令如下: | ||
| 182 | + | ||
| 183 | +```bash | ||
| 184 | +python -m ms_service_profiler.trace --perfetto-output /path/to/hook_tracing.json | ||
| 185 | +``` | ||
| 186 | + | ||
| 178 | 启动Trace转发进程使用的用户需要和启动MindIE Motor服务的用户一致,且在同网络命名空间中(同docker或同host)。 | 187 | 启动Trace转发进程使用的用户需要和启动MindIE Motor服务的用户一致,且在同网络命名空间中(同docker或同host)。 |
| 179 | 188 | ||
| 180 | **输出说明<a name="section738017254237"></a>** | 189 | **输出说明<a name="section738017254237"></a>** |
| @@ -274,7 +283,7 @@ curl http://127.0.0.1:1025/v1/chat/completions \ | |||
| 274 | 283 | ||
| 275 | 完成[发送请求](#发送请求)后,可以在支持OTLP协议的开源监测平台(例如Jaeger,须先开启Jaeger平台服务)查看可视化结果,示例如下。 | 284 | 完成[发送请求](#发送请求)后,可以在支持OTLP协议的开源监测平台(例如Jaeger,须先开启Jaeger平台服务)查看可视化结果,示例如下。 |
| 276 | 285 | ||
| 277 | -**图 1** 可视化结果<a name="fig485163113451"></a> | 286 | +**图 1** 可视化结果<a name="fig485163113451"></a> |
| 278 |  | 287 |  |
| 279 | 288 | ||
| 280 | 字段说明如下: | 289 | 字段说明如下: |
| @@ -0,0 +1,220 @@ | |||
| 1 | +# vLLM Hook Tracing 使用指南 | ||
| 2 | + | ||
| 3 | +## 1. 功能边界 | ||
| 4 | + | ||
| 5 | +vLLM Hook Tracing 通过 msServiceProfiler 现有 Hook/YAML 机制,为请求调度、模型执行和输出处理创建 OpenTelemetry Span。 | ||
| 6 | + | ||
| 7 | +- vLLM 原生 Tracing 是前置条件,必须配置 `--otlp-traces-endpoint`。 | ||
| 8 | +- Jaeger:Hook Span 复用 vLLM 的全局 Provider 和 OTLP exporter。 | ||
| 9 | +- Perfetto:可选将 msServiceProfiler Hook Span 额外写成 Chrome Trace JSON。 | ||
| 10 | +- 不支持 Perfetto-only;启动 Perfetto Forwarder 本身不会创建 Span。 | ||
| 11 | +- profiling、metrics、MindIE C++ Trace 和原 OTLP Forwarder 行为不变。 | ||
| 12 | +- 本功能没有新增 C/C++ 接口,不需要替换带新增符号的 `.so`。 | ||
| 13 | +- 当前及后续版本均不修改 vLLM/vLLM-Ascend 源码或 IPC;只读取上游公开的 `request_id`、`trace_headers` 和业务对象字段。 | ||
| 14 | + | ||
| 15 | +Tracing 和 profiling 是两套独立交付流程:Tracing 不需要 `enable.json` 和 `parse`;profiling 仍按原流程生成 CSV、DB 和离线 Chrome Trace 等交付件。 | ||
| 16 | + | ||
| 17 | +## 2. 数据流和地址选择 | ||
| 18 | + | ||
| 19 | +```mermaid | ||
| 20 | +flowchart LR | ||
| 21 | + R[真实推理请求] --> V[vLLM 原生 TracerProvider] | ||
| 22 | + V --> N[vLLM 原生 Span] | ||
| 23 | + V --> H[msServiceProfiler Hook Span] | ||
| 24 | + N --> O[OTLP exporter] | ||
| 25 | + H --> O | ||
| 26 | + O --> J[Jaeger OTLP 4318] | ||
| 27 | + H --> P[Perfetto Span Processor] | ||
| 28 | + P --> F[Perfetto Forwarder] | ||
| 29 | + F --> C[hook_tracing.json] | ||
| 30 | +``` | ||
| 31 | + | ||
| 32 | +Jaeger 和 vLLM 位于不同容器时,地址必须按执行位置选择: | ||
| 33 | + | ||
| 34 | +| 执行位置 | Jaeger查询地址 | OTLP 地址 | | ||
| 35 | +|---|---|---| | ||
| 36 | +| 宿主机 | `http://127.0.0.1:16686` | `http://127.0.0.1:4318/v1/traces` | | ||
| 37 | +| 与 Jaeger 同一 Docker 自定义网络的 vLLM 容器 | `http://ms-trace-jaeger:16686` | `http://ms-trace-jaeger:4318/v1/traces` | | ||
| 38 | +| Jaeger 与 vLLM 同一容器 | `http://127.0.0.1:16686` | `http://127.0.0.1:4318/v1/traces` | | ||
| 39 | + | ||
| 40 | +`127.0.0.1` 始终指向当前进程所在容器或主机,不能跨容器访问另一个服务。 | ||
| 41 | + | ||
| 42 | +## 3. 埋点配置 | ||
| 43 | + | ||
| 44 | +Tracing 复用 profiling 的同一份符号 YAML,不存在第二份 Tracing YAML。未设置 `PROFILING_SYMBOLS_PATH` 时使用 msServiceProfiler 包内 default 配置。 | ||
| 45 | + | ||
| 46 | +```yaml | ||
| 47 | +- symbol: vllm.v1.core.sched.scheduler:Scheduler.schedule | ||
| 48 | + handler: ms_service_profiler.patcher.vllm.handlers.v1.batch_handlers:schedule | ||
| 49 | + trace: | ||
| 50 | + name: vllm.scheduler.schedule | ||
| 51 | + domain: Schedule | ||
| 52 | + adapter: schedule | ||
| 53 | +``` | ||
| 54 | + | ||
| 55 | +同一解析后符号只允许一个有效 `trace` 配置。支持的 Adapter 如下: | ||
| 56 | + | ||
| 57 | +| Adapter | 作用 | | ||
| 58 | +|---|---| | ||
| 59 | +| `call` | 普通函数 Span | | ||
| 60 | +| `request` | 注册请求上下文,不创建重复根 Span | | ||
| 61 | +| `request_context` | 从请求对象读取 `trace_headers` | | ||
| 62 | +| `schedule` | 记录批次请求数、Token 数和请求 Links | | ||
| 63 | +| `model` | 记录模型执行并关联批次 request IDs | | ||
| 64 | +| `output` | 记录输出处理并清理已完成请求上下文 | | ||
| 65 | + | ||
| 66 | +## 4. 启动原则 | ||
| 67 | + | ||
| 68 | +### 4.1 Jaeger-only | ||
| 69 | + | ||
| 70 | +仅验证 Jaeger 时,不需要运行 `python3 -m ms_service_profiler.trace`。在 vLLM 进程环境中开启 Hook tracing,并为 vLLM 配置原生 OTLP endpoint: | ||
| 71 | + | ||
| 72 | +```bash | ||
| 73 | +export MS_TRACE_ENABLE=1 | ||
| 74 | +export OTEL_SERVICE_NAME=vllm-server | ||
| 75 | +export OTEL_EXPORTER_OTLP_TRACES_PROTOCOL=http/protobuf | ||
| 76 | +export OTEL_EXPORTER_OTLP_TRACES_ENDPOINT=http://ms-trace-jaeger:4318/v1/traces | ||
| 77 | + | ||
| 78 | +vllm serve "$MODEL" \ | ||
| 79 | + --host 0.0.0.0 \ | ||
| 80 | + --port "$PORT" \ | ||
| 81 | + --served-model-name "$MODEL_NAME" \ | ||
| 82 | + --otlp-traces-endpoint http://ms-trace-jaeger:4318/v1/traces | ||
| 83 | +``` | ||
| 84 | + | ||
| 85 | +### 4.2 Jaeger + Perfetto 双写 | ||
| 86 | + | ||
| 87 | +Perfetto Forwarder 必须先于 vLLM 启动,确保 Hook backend 初始化时能够注册附加 Processor: | ||
| 88 | + | ||
| 89 | +```bash | ||
| 90 | +export TRACE_VERIFY_DIR=/tmp/ms_trace_verify | ||
| 91 | +mkdir -p "$TRACE_VERIFY_DIR" | ||
| 92 | + | ||
| 93 | +nohup python3 -m ms_service_profiler.trace \ | ||
| 94 | + --log-level debug \ | ||
| 95 | + --perfetto-output "$TRACE_VERIFY_DIR/hook_tracing.json" \ | ||
| 96 | + >"$TRACE_VERIFY_DIR/perfetto_forwarder.log" 2>&1 & | ||
| 97 | + | ||
| 98 | +echo $! >"$TRACE_VERIFY_DIR/perfetto_forwarder.pid" | ||
| 99 | +``` | ||
| 100 | + | ||
| 101 | +然后按 Jaeger-only 场景启动 vLLM。Jaeger 接收 vLLM 原生 Span 和 Hook Span;Perfetto JSON 只接收 instrumentation scope 为 `ms_service_profiler.hook.*` 的 Hook Span。 | ||
| 102 | + | ||
| 103 | +`TRACE_VERIFY_DIR` 和 PID 文件只是验证脚本变量,不是产品配置项。 | ||
| 104 | + | ||
| 105 | +验证时应先确认 Jaeger 与 vLLM 的网络连通性,再按后续章节发送真实推理负载并检查两个数据出口。 | ||
| 106 | + | ||
| 107 | +## 5. 真实推理负载 | ||
| 108 | + | ||
| 109 | +`/health` 只用于等待服务就绪,不能用于验证推理 Trace。必须发送 `/v1/completions` 等真实推理请求,例如: | ||
| 110 | + | ||
| 111 | +```bash | ||
| 112 | +vllm bench serve \ | ||
| 113 | + --backend vllm \ | ||
| 114 | + --host 127.0.0.1 \ | ||
| 115 | + --port "$PORT" \ | ||
| 116 | + --endpoint /v1/completions \ | ||
| 117 | + --model "$MODEL_NAME" \ | ||
| 118 | + --tokenizer "$MODEL" \ | ||
| 119 | + --dataset-name random \ | ||
| 120 | + --random-input-len 128 \ | ||
| 121 | + --random-output-len 16 \ | ||
| 122 | + --num-prompts 8 \ | ||
| 123 | + --max-concurrency 1 \ | ||
| 124 | + --request-rate 1 \ | ||
| 125 | + --ignore-eos | ||
| 126 | +``` | ||
| 127 | + | ||
| 128 | +不要把 `VLLM_PLUGINS` 只设置为 `msserviceprofiler`,否则可能过滤掉 `ascend` 平台插件并导致 vLLM 无法识别 NPU。验证时可 `unset VLLM_PLUGINS`,让 vLLM 自动发现全部插件。 | ||
| 129 | + | ||
| 130 | +## 6. Jaeger 怎么看 | ||
| 131 | + | ||
| 132 | +### 6.1 搜索页 | ||
| 133 | + | ||
| 134 | +1. 打开 `http://<宿主机IP>:16686`。 | ||
| 135 | +2. Service 选择实际服务名,例如 `vllm-server`。 | ||
| 136 | +3. Tags 保持为空,Lookback 选择覆盖压测时间的范围。 | ||
| 137 | +4. 将 Limit Results 调大到 200,避免高频 scheduler/model Trace 挤掉 request/output Trace。 | ||
| 138 | +5. 点击 **Find Traces**。 | ||
| 139 | + | ||
| 140 | +散点图横轴是开始时间,纵轴是 Trace 总时长。明显高于正常分布的点通常是性能异常候选。列表中的 Spans 表示该 Trace 包含的 Span 数量,Duration 表示端到端或当前 Trace 的总时长。 | ||
| 141 | + | ||
| 142 | +### 6.2 Trace 详情页 | ||
| 143 | + | ||
| 144 | +点击一条 Trace 后重点观察: | ||
| 145 | + | ||
| 146 | +| 观察项 | 分析价值 | | ||
| 147 | +|---|---| | ||
| 148 | +| 父子 Span 时间线 | 判断耗时发生在调度、模型执行还是输出处理 | | ||
| 149 | +| Span 自身耗时 | 区分模型计算慢与框架处理慢 | | ||
| 150 | +| 同级 Span 的空白间隔 | 识别排队、同步、IPC 或等待资源的时间 | | ||
| 151 | +| Tags/Attributes | 查看 `request.id(s)`、batch、token、状态和进程/线程信息 | | ||
| 152 | +| References/Links | 查看跨请求批处理或跨进程关联;Link 不等同于父子关系 | | ||
| 153 | +| Error/Status | 定位异常请求和失败阶段 | | ||
| 154 | + | ||
| 155 | +典型埋点含义: | ||
| 156 | + | ||
| 157 | +| Span | 主要定位方向 | | ||
| 158 | +|---|---| | ||
| 159 | +| `vllm.scheduler.schedule` | 调度频率、调度耗时、请求排队、batch 组织 | | ||
| 160 | +| `vllm.model.execute` | Executor 层模型执行耗时 | | ||
| 161 | +| `vllm_ascend.model_runner.execute` | Ascend ModelRunner 实际执行耗时 | | ||
| 162 | +| `vllm.output.process` | 输出处理、请求完成和后处理耗时 | | ||
| 163 | + | ||
| 164 | +### 6.3 常见性能判断 | ||
| 165 | + | ||
| 166 | +- scheduler Span 持续升高、model Span 稳定:重点检查排队、KV Cache 压力、batch 组织和调度策略。 | ||
| 167 | +- model Span 占绝大多数且随 batch 增大明显升高:重点检查模型计算、通信、算子和 NPU 利用率,并与 profiling 结果结合。 | ||
| 168 | +- `vllm.model.execute` 明显大于其内部 `vllm_ascend.model_runner.execute`:Executor 外围可能存在 IPC、同步或数据准备开销。 | ||
| 169 | +- output Span 异常升高:重点检查输出处理、token 后处理、网络回传和请求完成逻辑。 | ||
| 170 | +- 少数 Trace 是长尾:在散点图选中长尾 Trace,与正常 Trace 使用 Compare 对比 Span 时长和属性差异。 | ||
| 171 | +- 大量 Trace 只有 1~2 个 Span:只能证明埋点已导出,不能单独证明完整请求父子链;需要继续检查 Trace 详情中的 Links、`request.ids` 和原生 vLLM 根 Span。 | ||
| 172 | + | ||
| 173 | +Tracing 用于缩小问题所在阶段;若需要分析算子、kernel、HCCL、显存或 NPU 时间线,继续使用原 profiling 采集和解析能力。 | ||
| 174 | + | ||
| 175 | +## 7. 验收层级 | ||
| 176 | + | ||
| 177 | +| 层级 | 判定依据 | | ||
| 178 | +|---|---| | ||
| 179 | +| 数据出口通过 | Jaeger API有 Trace,Perfetto JSON 有合法 Complete Event | | ||
| 180 | +| 埋点展示通过 | Jaeger/Perfetto 能看到预期 Hook Span 名称和耗时 | | ||
| 181 | +| 请求关联通过 | request、schedule、model、output 可通过父子关系、同一 trace ID 或明确 Span Links 关联 | | ||
| 182 | +| 性能分析有效 | 能用 Span 时长和属性区分调度、模型、输出或等待瓶颈,并可下钻 profiling | | ||
| 183 | + | ||
| 184 | +不能用“Jaeger 页面有数据”替代“完整请求链路关联通过”的结论。 | ||
| 185 | + | ||
| 186 | +## 8. Profiling 共存验证 | ||
| 187 | + | ||
| 188 | +Profiling 仍必须按原流程配置 `enable.json`、指定采集目录并执行 `parse`。Tracing 不替代该流程。 | ||
| 189 | + | ||
| 190 | +建议分别验证: | ||
| 191 | + | ||
| 192 | +1. 只开 profiling:交付件和历史版本一致。 | ||
| 193 | +2. 只开 tracing:Jaeger 中出现原生 Span 和 Hook Span;若启动 Perfetto Forwarder,JSON 非空。 | ||
| 194 | +3. profiling + tracing:推理成功,两类交付件均生成;Tracing Span 会包含少量 profiling handler 开销。 | ||
| 195 | + | ||
| 196 | +## 9. 常见问题 | ||
| 197 | + | ||
| 198 | +| 现象 | 原因与处理 | | ||
| 199 | +|---|---| | ||
| 200 | +| `ms-trace-jaeger` 无法解析 | 命令在宿主机执行,或两个容器不在同一 Docker 自定义网络;宿主机改用 `127.0.0.1` | | ||
| 201 | +| `docker exec ... python3 - <<'PY'` 无输出 | heredoc 需要 `docker exec -i` 才会把标准输入传入容器 | | ||
| 202 | +| 路径展开成 `/vllm_tracing.log` | 当前 shell 没有设置 `TRACE_VERIFY_DIR`;先检查变量非空 | | ||
| 203 | +| 容器内看不到另一个容器的 `/tmp` 文件 | 进入了错误容器;始终使用唯一的 `VLLM_CONTAINER` 变量 | | ||
| 204 | +| Python 加载 `/usr/local/Ascend/.../site-packages` | 验证的是已安装包;验证工作区代码时显式设置 `PYTHONPATH=<仓库根目录>` | | ||
| 205 | +| `pkill -9 386` 无法停止 PID 386 | `pkill` 参数是名称模式;按 PID 使用 `kill -TERM 386` | | ||
| 206 | +| Perfetto JSON 是空数组 | Forwarder 启动过晚、未运行,或 vLLM 未开启原生 OTLP tracing | | ||
| 207 | +| Jaeger 可用但 Perfetto 文件不存在 | 两个出口独立;检查 vLLM 容器内 Forwarder 进程和日志 | | ||
| 208 | +| vLLM 无法识别 NPU | 不要将 `VLLM_PLUGINS` 限制为单个 msserviceprofiler 插件;同时检查 `ASCEND_RT_VISIBLE_DEVICES` | | ||
| 209 | +| Jaeger Trace 只有 1~2 个 Span | 检查上下文传播、Span Links 和 `request.ids`;不要直接宣称完整请求链通过 | | ||
| 210 | + | ||
| 211 | +## 10. 异常行为 | ||
| 212 | + | ||
| 213 | +| 场景 | 行为 | | ||
| 214 | +|---|---| | ||
| 215 | +| vLLM 未配置 `--otlp-traces-endpoint` | 无可复用全局 Provider,Hook tracing no-op,并提示开启 vLLM 原生 Tracing | | ||
| 216 | +| OpenTelemetry 依赖缺失 | Hook tracing no-op,推理继续 | | ||
| 217 | +| Jaeger 不可用 | vLLM exporter 按自身策略处理,推理和 Hook 调用继续 | | ||
| 218 | +| Perfetto Forwarder 未启动 | 不注册 Perfetto Processor,Jaeger 不受影响 | | ||
| 219 | +| 部分 YAML 符号不存在 | 跳过该 Hook,其他 Hook 继续 | | ||
| 220 | +| profiling 同时开启 | tracing 包在 profiling 外层,原业务函数只执行一次 | | ||
| @@ -64,6 +64,8 @@ vllm serve Qwen/Qwen2.5-0.5B-Instruct & | |||
| 64 | 64 | ||
| 65 | `service_profiling_symbols.yaml` 为需要导入的埋点配置文件。你也可以选择不设置环境变量 `PROFILING_SYMBOLS_PATH` ,此时将使用默认的配置文件;若你指定的路径下不存在该文件,系统同样会在你指定的路径生成一份配置文件以便后续修改。可参考[点位配置使用指南](#点位配置使用指南)一节进行自定义。 | 65 | `service_profiling_symbols.yaml` 为需要导入的埋点配置文件。你也可以选择不设置环境变量 `PROFILING_SYMBOLS_PATH` ,此时将使用默认的配置文件;若你指定的路径下不存在该文件,系统同样会在你指定的路径生成一份配置文件以便后续修改。可参考[点位配置使用指南](#点位配置使用指南)一节进行自定义。 |
| 66 | 66 | ||
| 67 | +如需通过 Hook 创建自定义链路 Span,并将数据接入 Jaeger 或 Perfetto,请参见 [vLLM Hook Tracing 使用指南](./vLLM_hook_tracing_instruct.md)。 | ||
| 68 | + | ||
| 67 | **2. 开启采集** | 69 | **2. 开启采集** |
| 68 | 70 | ||
| 69 | 将配置文件`ms_service_profiler_config.json`中的 `enable` 字段由 `0` 修改为 `1`,即可开启性能数据采集的开关,可以通过执行下面sed指令完成采集服务的开启: | 71 | 将配置文件`ms_service_profiler_config.json`中的 `enable` 字段由 `0` 修改为 `1`,即可开启性能数据采集的开关,可以通过执行下面sed指令完成采集服务的开启: |
| @@ -49,6 +49,7 @@ nav: | |||
| 49 | - SGLang 服务化性能采集工具: zh/SGLang_service_oriented_performance_collection_tool.md | 49 | - SGLang 服务化性能采集工具: zh/SGLang_service_oriented_performance_collection_tool.md |
| 50 | - 3. Trace 数据链路监测: | 50 | - 3. Trace 数据链路监测: |
| 51 | - 服务化Trace数据监测工具: zh/msserviceprofiler_trace_data_monitoring_instruct.md | 51 | - 服务化Trace数据监测工具: zh/msserviceprofiler_trace_data_monitoring_instruct.md |
| 52 | + - vLLM Hook Tracing 使用指南: zh/vLLM_hook_tracing_instruct.md | ||
| 52 | - 4. 采集数据的比对与多维分析: | 53 | - 4. 采集数据的比对与多维分析: |
| 53 | - 服务化性能数据比对工具: zh/ms_service_profiler_compare_tool_instruct.md | 54 | - 服务化性能数据比对工具: zh/ms_service_profiler_compare_tool_instruct.md |
| 54 | - 服务化多维度解析工具: zh/msserviceprofiler_multi_analyze_instruct.md | 55 | - 服务化多维度解析工具: zh/msserviceprofiler_multi_analyze_instruct.md |
| @@ -27,21 +27,34 @@ def main(): | |||
| 27 | type=str, | 27 | type=str, |
| 28 | default='info', | 28 | default='info', |
| 29 | choices=['debug', 'info', 'warning', 'error', 'fatal', 'critical'], | 29 | choices=['debug', 'info', 'warning', 'error', 'fatal', 'critical'], |
| 30 | - help='Log level to print') | 30 | + help='Log level to print', |
| 31 | + ) | ||
| 32 | + parser.add_argument( | ||
| 33 | + '--perfetto-output', | ||
| 34 | + type=str, | ||
| 35 | + default=None, | ||
| 36 | + help='Optional Chrome Trace JSON file for Hook spans. vLLM OTLP tracing must also be enabled.', | ||
| 37 | + ) | ||
| 31 | args = parser.parse_args() | 38 | args = parser.parse_args() |
| 32 | set_log_level(args.log_level) | 39 | set_log_level(args.log_level) |
| 33 | 40 | ||
| 34 | - if os.name != "nt" and os.getuid() == 0: | 41 | + get_uid = getattr(os, "getuid", lambda: -1) |
| 42 | + if os.name != "nt" and get_uid() == 0: | ||
| 35 | logger.warning( | 43 | logger.warning( |
| 36 | "Security Warning: Running with root privileges may compromise system security. " | 44 | "Security Warning: Running with root privileges may compromise system security. " |
| 37 | "Run the program as the user who runs MindIE." | 45 | "Run the program as the user who runs MindIE." |
| 38 | ) | 46 | ) |
| 39 | 47 | ||
| 40 | try: | 48 | try: |
| 41 | - service = OTLPForwarderService() | 49 | + if args.perfetto_output: |
| 50 | + from ms_service_profiler.tracer.perfetto_forward_service import PerfettoForwarderService | ||
| 51 | + | ||
| 52 | + service = PerfettoForwarderService(args.perfetto_output) | ||
| 53 | + else: | ||
| 54 | + service = OTLPForwarderService() | ||
| 42 | service.start() | 55 | service.start() |
| 43 | except Exception as e: | 56 | except Exception as e: |
| 44 | - logger.error(f"Start OTLPForwarderService failed: {e}") | 57 | + logger.error("Start trace service failed: %s", e) |
| 45 | 58 | ||
| 46 | 59 | ||
| 47 | if __name__ == '__main__': | 60 | if __name__ == '__main__': |
| @@ -0,0 +1,183 @@ | |||
| 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 | +"""Convert Hook span packets into Perfetto-compatible Chrome Trace JSON.""" | ||
| 9 | + | ||
| 10 | +import json | ||
| 11 | +import os | ||
| 12 | +import stat | ||
| 13 | +import threading | ||
| 14 | +from collections import defaultdict | ||
| 15 | +from typing import Dict, Iterable, List | ||
| 16 | + | ||
| 17 | +from ms_service_profiler.utils.file_open_check import is_legal_args_path_string | ||
| 18 | +from ms_service_profiler.utils.log import logger | ||
| 19 | +from ms_service_profiler.tracer.perfetto_socket import PERFETTO_EVENT_MAGIC | ||
| 20 | + | ||
| 21 | + | ||
| 22 | +MAX_TRACKED_SPANS = 100000 | ||
| 23 | + | ||
| 24 | + | ||
| 25 | +def validate_perfetto_output_path(path: str) -> str: | ||
| 26 | + """Return a normalized safe JSON path and create only its parent directory.""" | ||
| 27 | + if not isinstance(path, str) or not path.lower().endswith(".json"): | ||
| 28 | + raise ValueError("--perfetto-output must point to a .json file") | ||
| 29 | + normalized = os.path.abspath(os.path.expanduser(path)) | ||
| 30 | + if not is_legal_args_path_string(normalized): | ||
| 31 | + raise ValueError("--perfetto-output contains unsupported characters") | ||
| 32 | + parent = os.path.dirname(normalized) or os.getcwd() | ||
| 33 | + os.makedirs(parent, mode=0o750, exist_ok=True) | ||
| 34 | + if os.path.islink(normalized): | ||
| 35 | + raise ValueError("--perfetto-output cannot be a symbolic link") | ||
| 36 | + if os.path.exists(normalized): | ||
| 37 | + file_stat = os.stat(normalized) | ||
| 38 | + if not stat.S_ISREG(file_stat.st_mode): | ||
| 39 | + raise ValueError("--perfetto-output must be a regular file") | ||
| 40 | + if hasattr(os, "geteuid") and file_stat.st_uid != os.geteuid(): | ||
| 41 | + raise PermissionError("--perfetto-output must be owned by the current user") | ||
| 42 | + return normalized | ||
| 43 | + | ||
| 44 | + | ||
| 45 | +class PerfettoTraceExporter: | ||
| 46 | + """Append complete Chrome Trace events while keeping the JSON file valid.""" | ||
| 47 | + | ||
| 48 | + def __init__(self, output_path: str): | ||
| 49 | + self.output_path = validate_perfetto_output_path(output_path) | ||
| 50 | + flags = os.O_CREAT | os.O_TRUNC | os.O_RDWR | ||
| 51 | + if hasattr(os, "O_NOFOLLOW"): | ||
| 52 | + flags |= os.O_NOFOLLOW | ||
| 53 | + descriptor = os.open(self.output_path, flags, 0o640) | ||
| 54 | + self._file = os.fdopen(descriptor, "w+", encoding="utf-8") | ||
| 55 | + self._file.write("[]") | ||
| 56 | + self._file.flush() | ||
| 57 | + self._has_events = False | ||
| 58 | + self._closed = False | ||
| 59 | + self._lock = threading.RLock() | ||
| 60 | + self._span_positions = {} | ||
| 61 | + self._pending_links = defaultdict(list) | ||
| 62 | + | ||
| 63 | + def export(self, binary_data: bytes) -> bool: | ||
| 64 | + try: | ||
| 65 | + if not binary_data.startswith(PERFETTO_EVENT_MAGIC): | ||
| 66 | + raise ValueError("unsupported Perfetto packet") | ||
| 67 | + payload = json.loads(binary_data[len(PERFETTO_EVENT_MAGIC) :].decode("utf-8")) | ||
| 68 | + events = self._convert_normalized_span(payload) | ||
| 69 | + self._append_events(events) | ||
| 70 | + return True | ||
| 71 | + except Exception as exc: | ||
| 72 | + logger.warning("Export Hook trace to Perfetto failed: %s", exc) | ||
| 73 | + return False | ||
| 74 | + | ||
| 75 | + def _convert_normalized_span(self, span: Dict) -> List[Dict]: | ||
| 76 | + attributes = dict(span.get("attributes") or {}) | ||
| 77 | + pid = int(attributes.pop("process.pid", os.getpid())) | ||
| 78 | + tid = int(attributes.pop("thread.id", 0)) | ||
| 79 | + start_us = int(span.get("start_time_ns", 0)) / 1000 | ||
| 80 | + duration_us = max(int(span.get("end_time_ns", 0)) - int(span.get("start_time_ns", 0)), 0) / 1000 | ||
| 81 | + trace_id = str(span.get("trace_id", "")) | ||
| 82 | + span_id = str(span.get("span_id", "")) | ||
| 83 | + resource_attributes = dict(span.get("resource_attributes") or {}) | ||
| 84 | + service_name = span.get("service_name") or resource_attributes.get("service.name", "ms_service_profiler") | ||
| 85 | + args = dict(attributes) | ||
| 86 | + args.update( | ||
| 87 | + { | ||
| 88 | + "trace_id": trace_id, | ||
| 89 | + "span_id": span_id, | ||
| 90 | + "parent_span_id": str(span.get("parent_span_id", "")), | ||
| 91 | + "span.kind": int(span.get("kind", 0)), | ||
| 92 | + "span.status": int(span.get("status", 0)), | ||
| 93 | + "service.name": service_name, | ||
| 94 | + } | ||
| 95 | + ) | ||
| 96 | + for key, value in resource_attributes.items(): | ||
| 97 | + args.setdefault("resource.{}".format(key), value) | ||
| 98 | + events = [ | ||
| 99 | + { | ||
| 100 | + "name": span.get("name", ""), | ||
| 101 | + "cat": span.get("category", "Tracing"), | ||
| 102 | + "ph": "X", | ||
| 103 | + "ts": start_us, | ||
| 104 | + "dur": duration_us, | ||
| 105 | + "pid": pid, | ||
| 106 | + "tid": tid, | ||
| 107 | + "args": args, | ||
| 108 | + } | ||
| 109 | + ] | ||
| 110 | + key = (trace_id, span_id) | ||
| 111 | + position = (pid, tid, start_us) | ||
| 112 | + self._remember_position(key, position) | ||
| 113 | + for pending in self._pending_links.pop(key, []): | ||
| 114 | + events.extend(self._make_flow_events(key, position, *pending)) | ||
| 115 | + for link in span.get("links") or []: | ||
| 116 | + source_key = (str(link.get("trace_id", "")), str(link.get("span_id", ""))) | ||
| 117 | + flow_id = "{}:{}:{}".format(source_key[0], source_key[1], span_id) | ||
| 118 | + source_position = self._span_positions.get(source_key) | ||
| 119 | + destination = (pid, tid, start_us, flow_id) | ||
| 120 | + if source_position is None: | ||
| 121 | + if len(self._pending_links) < MAX_TRACKED_SPANS: | ||
| 122 | + self._pending_links[source_key].append(destination) | ||
| 123 | + else: | ||
| 124 | + events.extend(self._make_flow_events(source_key, source_position, *destination)) | ||
| 125 | + return events | ||
| 126 | + | ||
| 127 | + | ||
| 128 | + def _make_flow_events(source_key, source_position, dest_pid, dest_tid, dest_ts, flow_id): | ||
| 129 | + source_pid, source_tid, source_ts = source_position | ||
| 130 | + name = "request.link" | ||
| 131 | + return [ | ||
| 132 | + { | ||
| 133 | + "name": name, | ||
| 134 | + "cat": "Tracing.Flow", | ||
| 135 | + "ph": "s", | ||
| 136 | + "ts": source_ts, | ||
| 137 | + "pid": source_pid, | ||
| 138 | + "tid": source_tid, | ||
| 139 | + "id": flow_id, | ||
| 140 | + "args": {"source.span_id": source_key[1]}, | ||
| 141 | + }, | ||
| 142 | + { | ||
| 143 | + "name": name, | ||
| 144 | + "cat": "Tracing.Flow", | ||
| 145 | + "ph": "f", | ||
| 146 | + "ts": dest_ts, | ||
| 147 | + "pid": dest_pid, | ||
| 148 | + "tid": dest_tid, | ||
| 149 | + "id": flow_id, | ||
| 150 | + "bp": "e", | ||
| 151 | + }, | ||
| 152 | + ] | ||
| 153 | + | ||
| 154 | + def _remember_position(self, key, position) -> None: | ||
| 155 | + if len(self._span_positions) >= MAX_TRACKED_SPANS: | ||
| 156 | + oldest_key = next(iter(self._span_positions)) | ||
| 157 | + self._span_positions.pop(oldest_key, None) | ||
| 158 | + self._span_positions[key] = position | ||
| 159 | + | ||
| 160 | + def _append_events(self, events: Iterable[Dict]) -> None: | ||
| 161 | + encoded = [json.dumps(event, ensure_ascii=False, separators=(",", ":")) for event in events] | ||
| 162 | + if not encoded: | ||
| 163 | + return | ||
| 164 | + with self._lock: | ||
| 165 | + if self._closed: | ||
| 166 | + return | ||
| 167 | + self._file.seek(0, os.SEEK_END) | ||
| 168 | + self._file.seek(self._file.tell() - 1) | ||
| 169 | + if self._has_events: | ||
| 170 | + self._file.write(",") | ||
| 171 | + self._file.write(",".join(encoded)) | ||
| 172 | + self._file.write("]") | ||
| 173 | + self._file.truncate() | ||
| 174 | + self._file.flush() | ||
| 175 | + self._has_events = True | ||
| 176 | + | ||
| 177 | + def close(self) -> None: | ||
| 178 | + with self._lock: | ||
| 179 | + if self._closed: | ||
| 180 | + return | ||
| 181 | + self._closed = True | ||
| 182 | + self._file.flush() | ||
| 183 | + self._file.close() | ||
| @@ -0,0 +1,81 @@ | |||
| 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 | +"""Standalone receiver for Hook span packets; independent of legacy OTLP.""" | ||
| 9 | + | ||
| 10 | +import signal | ||
| 11 | +import threading | ||
| 12 | + | ||
| 13 | +from ms_service_profiler.tracer.perfetto_exporter import PerfettoTraceExporter | ||
| 14 | +from ms_service_profiler.tracer.perfetto_socket import PERFETTO_SOCKET_NAME | ||
| 15 | +from ms_service_profiler.tracer.socket_server import AbstractSocketServer | ||
| 16 | +from ms_service_profiler.utils.log import logger | ||
| 17 | + | ||
| 18 | + | ||
| 19 | +SOCKET_BUFFER_SIZE = 4096 | ||
| 20 | +SOCKET_TIMEOUT = 1 | ||
| 21 | +MAX_LISTEN_NUM = 8 | ||
| 22 | +MAX_QUEUE_SIZE = 100000 | ||
| 23 | +WARNING_QUEUE_SIZE = 10000 | ||
| 24 | +POLL_INTERVAL_SECONDS = 0.05 | ||
| 25 | + | ||
| 26 | + | ||
| 27 | +class PerfettoForwarderService: | ||
| 28 | + """Receive only msserviceprofiler Hook spans and write Chrome Trace JSON.""" | ||
| 29 | + | ||
| 30 | + def __init__(self, output_path: str): | ||
| 31 | + self._stop_event = threading.Event() | ||
| 32 | + self._stopped = False | ||
| 33 | + self._exporter = PerfettoTraceExporter(output_path) | ||
| 34 | + self._socket_server = AbstractSocketServer( | ||
| 35 | + socket_name=PERFETTO_SOCKET_NAME, | ||
| 36 | + buffer_size=SOCKET_BUFFER_SIZE, | ||
| 37 | + max_listen_num=MAX_LISTEN_NUM, | ||
| 38 | + socket_timeout=SOCKET_TIMEOUT, | ||
| 39 | + max_queue_size=MAX_QUEUE_SIZE, | ||
| 40 | + warning_queue_size=WARNING_QUEUE_SIZE, | ||
| 41 | + ) | ||
| 42 | + signal.signal(signal.SIGINT, self._handle_signal) | ||
| 43 | + signal.signal(signal.SIGTERM, self._handle_signal) | ||
| 44 | + | ||
| 45 | + def _handle_signal(self, signum, frame): | ||
| 46 | + logger.info("Receive signal %s, quit...", signum) | ||
| 47 | + self._stop_event.set() | ||
| 48 | + | ||
| 49 | + def _drain(self) -> None: | ||
| 50 | + while True: | ||
| 51 | + data = self._socket_server.get_data() | ||
| 52 | + if not data: | ||
| 53 | + return | ||
| 54 | + self._exporter.export(data) | ||
| 55 | + | ||
| 56 | + def start(self) -> None: | ||
| 57 | + try: | ||
| 58 | + self._socket_server.start() | ||
| 59 | + logger.info("Start PerfettoForwarderService success, running...") | ||
| 60 | + while not self._stop_event.is_set(): | ||
| 61 | + data = self._socket_server.get_data() | ||
| 62 | + if data: | ||
| 63 | + self._exporter.export(data) | ||
| 64 | + else: | ||
| 65 | + self._stop_event.wait(POLL_INTERVAL_SECONDS) | ||
| 66 | + except KeyboardInterrupt: | ||
| 67 | + logger.info("Receive KeyboardInterrupt, quit...") | ||
| 68 | + except Exception as exc: | ||
| 69 | + logger.error("Unexpected Perfetto forwarder error: %s", exc) | ||
| 70 | + finally: | ||
| 71 | + self.stop() | ||
| 72 | + | ||
| 73 | + def stop(self) -> None: | ||
| 74 | + if self._stopped: | ||
| 75 | + return | ||
| 76 | + self._stopped = True | ||
| 77 | + self._stop_event.set() | ||
| 78 | + self._socket_server.stop() | ||
| 79 | + self._drain() | ||
| 80 | + self._exporter.close() | ||
| 81 | + logger.info("Stop PerfettoForwarderService success.") | ||
| @@ -0,0 +1,103 @@ | |||
| 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 | +"""Tests for converting Hook span packets to Perfetto/Chrome Trace.""" | ||
| 9 | + | ||
| 10 | +import json | ||
| 11 | +from unittest.mock import patch | ||
| 12 | + | ||
| 13 | +import pytest | ||
| 14 | + | ||
| 15 | +from ms_service_profiler.tracer.perfetto_exporter import ( | ||
| 16 | + PerfettoTraceExporter, | ||
| 17 | + validate_perfetto_output_path, | ||
| 18 | +) | ||
| 19 | +from ms_service_profiler.tracer.perfetto_socket import PERFETTO_EVENT_MAGIC | ||
| 20 | + | ||
| 21 | + | ||
| 22 | +def _make_packet(name, trace_id, span_id, start_ns, end_ns, link=None): | ||
| 23 | + payload = { | ||
| 24 | + "name": name, | ||
| 25 | + "category": "ms_service_profiler.hook.Tracing", | ||
| 26 | + "trace_id": trace_id, | ||
| 27 | + "span_id": span_id, | ||
| 28 | + "parent_span_id": "", | ||
| 29 | + "start_time_ns": start_ns, | ||
| 30 | + "end_time_ns": end_ns, | ||
| 31 | + "kind": 1, | ||
| 32 | + "status": 1, | ||
| 33 | + "attributes": {"process.pid": 100, "thread.id": 200, "request.id": "request-1"}, | ||
| 34 | + "resource_attributes": {"service.name": "vllm"}, | ||
| 35 | + "links": [{"trace_id": link[0], "span_id": link[1]}] if link else [], | ||
| 36 | + } | ||
| 37 | + return PERFETTO_EVENT_MAGIC + json.dumps(payload).encode("utf-8") | ||
| 38 | + | ||
| 39 | + | ||
| 40 | +def test_perfetto_exporter_writes_complete_events_and_cross_trace_flows(tmp_path): | ||
| 41 | + output = tmp_path / "hook_trace.json" | ||
| 42 | + with patch("ms_service_profiler.tracer.perfetto_exporter.is_legal_args_path_string", return_value=True): | ||
| 43 | + exporter = PerfettoTraceExporter(str(output)) | ||
| 44 | + request_trace_id = "1" * 32 | ||
| 45 | + request_span_id = "2" * 16 | ||
| 46 | + schedule_trace_id = "3" * 32 | ||
| 47 | + schedule_span_id = "4" * 16 | ||
| 48 | + | ||
| 49 | + assert exporter.export( | ||
| 50 | + _make_packet( | ||
| 51 | + "vllm.scheduler.schedule", | ||
| 52 | + schedule_trace_id, | ||
| 53 | + schedule_span_id, | ||
| 54 | + 2_000_000, | ||
| 55 | + 3_000_000, | ||
| 56 | + link=(request_trace_id, request_span_id), | ||
| 57 | + ) | ||
| 58 | + ) | ||
| 59 | + # The linked request may arrive later; pending flow resolution must still | ||
| 60 | + # connect it to the already exported scheduler span. | ||
| 61 | + assert exporter.export(_make_packet("vllm.request", request_trace_id, request_span_id, 1_000_000, 4_000_000)) | ||
| 62 | + exporter.close() | ||
| 63 | + | ||
| 64 | + events = json.loads(output.read_text(encoding="utf-8")) | ||
| 65 | + complete = [event for event in events if event["ph"] == "X"] | ||
| 66 | + flow = [event for event in events if event["ph"] in ("s", "f")] | ||
| 67 | + assert [event["name"] for event in complete] == ["vllm.scheduler.schedule", "vllm.request"] | ||
| 68 | + assert complete[0]["pid"] == 100 | ||
| 69 | + assert complete[0]["tid"] == 200 | ||
| 70 | + assert complete[0]["dur"] == 1000 | ||
| 71 | + assert complete[0]["args"]["trace_id"] == schedule_trace_id | ||
| 72 | + assert len(flow) == 2 | ||
| 73 | + assert flow[0]["id"] == flow[1]["id"] | ||
| 74 | + | ||
| 75 | + | ||
| 76 | +def test_perfetto_exporter_accepts_normalized_hook_packet(tmp_path): | ||
| 77 | + output = tmp_path / "otel_hook_trace.json" | ||
| 78 | + with patch("ms_service_profiler.tracer.perfetto_exporter.is_legal_args_path_string", return_value=True): | ||
| 79 | + exporter = PerfettoTraceExporter(str(output)) | ||
| 80 | + | ||
| 81 | + assert exporter.export(_make_packet("vllm.model.execute", "1" * 32, "2" * 16, 1_000_000, 2_000_000)) | ||
| 82 | + exporter.close() | ||
| 83 | + | ||
| 84 | + events = json.loads(output.read_text(encoding="utf-8")) | ||
| 85 | + assert events[0]["name"] == "vllm.model.execute" | ||
| 86 | + assert events[0]["pid"] == 100 | ||
| 87 | + assert events[0]["dur"] == 1000 | ||
| 88 | + | ||
| 89 | + | ||
| 90 | +def test_perfetto_exporter_rejects_legacy_binary_otlp_payload(tmp_path): | ||
| 91 | + output = tmp_path / "hook_trace.json" | ||
| 92 | + with patch("ms_service_profiler.tracer.perfetto_exporter.is_legal_args_path_string", return_value=True): | ||
| 93 | + exporter = PerfettoTraceExporter(str(output)) | ||
| 94 | + | ||
| 95 | + assert exporter.export(b"binary-otlp") is False | ||
| 96 | + exporter.close() | ||
| 97 | + assert json.loads(output.read_text(encoding="utf-8")) == [] | ||
| 98 | + | ||
| 99 | + | ||
| 100 | + | ||
| 101 | +def test_perfetto_output_rejects_invalid_path(path): | ||
| 102 | + with pytest.raises(ValueError): | ||
| 103 | + validate_perfetto_output_path(path) | ||
| @@ -18,17 +18,6 @@ import pytest | |||
| 18 | from unittest.mock import MagicMock, patch | 18 | from unittest.mock import MagicMock, patch |
| 19 | 19 | ||
| 20 | 20 | ||
| 21 | -mock_modules = { | ||
| 22 | - 'opentelemetry': MagicMock(), | ||
| 23 | - 'opentelemetry.proto.collector.trace.v1.trace_service_pb2': MagicMock(), | ||
| 24 | - 'opentelemetry.exporter.otlp.proto.http.trace_exporter': MagicMock(), | ||
| 25 | - 'opentelemetry.exporter.otlp.proto.grpc.trace_exporter': MagicMock(), | ||
| 26 | - 'opentelemetry.sdk.trace.export': MagicMock() | ||
| 27 | -} | ||
| 28 | -patch_obj = patch.dict('sys.modules', mock_modules) | ||
| 29 | -patch_obj.start() | ||
| 30 | - | ||
| 31 | - | ||
| 32 | 21 | ||
| 33 | def mock_span_export_result(): | 22 | def mock_span_export_result(): |
| 34 | """Fixture: Mock SpanExportResult enum""" | 23 | """Fixture: Mock SpanExportResult enum""" |
| @@ -36,4 +25,4 @@ def mock_span_export_result(): | |||
| 36 | mock_result.SUCCESS = 0 | 25 | mock_result.SUCCESS = 0 |
| 37 | mock_result.FAILED = 1 | 26 | mock_result.FAILED = 1 |
| 38 | mock_result.return_value = mock_result | 27 | mock_result.return_value = mock_result |
| 39 | - yield mock_result | 28 | + yield mock_result |
| @@ -0,0 +1,59 @@ | |||
| 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 | +import signal | ||
| 9 | +from unittest.mock import call, patch | ||
| 10 | + | ||
| 11 | +from ms_service_profiler.tracer.perfetto_forward_service import PerfettoForwarderService | ||
| 12 | + | ||
| 13 | + | ||
| 14 | + | ||
| 15 | + | ||
| 16 | + | ||
| 17 | +def test_perfetto_forwarder_uses_independent_socket(mock_exporter, mock_socket, mock_signal): | ||
| 18 | + service = PerfettoForwarderService("hook_trace.json") | ||
| 19 | + | ||
| 20 | + mock_exporter.assert_called_once_with("hook_trace.json") | ||
| 21 | + mock_socket.assert_called_once_with( | ||
| 22 | + socket_name="MSP_PERFETTO_SOCKET", | ||
| 23 | + buffer_size=4096, | ||
| 24 | + max_listen_num=8, | ||
| 25 | + socket_timeout=1, | ||
| 26 | + max_queue_size=100000, | ||
| 27 | + warning_queue_size=10000, | ||
| 28 | + ) | ||
| 29 | + mock_signal.assert_has_calls( | ||
| 30 | + [ | ||
| 31 | + call(signal.SIGINT, service._handle_signal), | ||
| 32 | + call(signal.SIGTERM, service._handle_signal), | ||
| 33 | + ] | ||
| 34 | + ) | ||
| 35 | + | ||
| 36 | + | ||
| 37 | + | ||
| 38 | + | ||
| 39 | + | ||
| 40 | +def test_perfetto_forwarder_exports_hook_packets_and_closes(mock_exporter, mock_socket, _): | ||
| 41 | + socket_instance = mock_socket.return_value | ||
| 42 | + exporter_instance = mock_exporter.return_value | ||
| 43 | + service = PerfettoForwarderService("hook_trace.json") | ||
| 44 | + packet = b"hook-packet" | ||
| 45 | + reads = iter([packet, None]) | ||
| 46 | + | ||
| 47 | + def get_data(): | ||
| 48 | + value = next(reads) | ||
| 49 | + if value == packet: | ||
| 50 | + service._stop_event.set() | ||
| 51 | + return value | ||
| 52 | + | ||
| 53 | + socket_instance.get_data.side_effect = get_data | ||
| 54 | + service.start() | ||
| 55 | + | ||
| 56 | + socket_instance.start.assert_called_once() | ||
| 57 | + exporter_instance.export.assert_called_once_with(packet) | ||
| 58 | + socket_instance.stop.assert_called_once() | ||
| 59 | + exporter_instance.close.assert_called_once() | ||
| @@ -22,13 +22,30 @@ class TestMain: | |||
| 22 | 22 | ||
| 23 | 23 | ||
| 24 | 24 | ||
| 25 | - def test_main_success(self, mock_otlp_service, mock_set_log_level, mock_arg_parser): | 25 | + def test_main_without_perfetto_keeps_legacy_otlp_service( |
| 26 | - """Test the behavior of the main function in a successful scenario""" | 26 | + self, mock_otlp_service, mock_set_log_level, mock_arg_parser |
| 27 | - mock_args = MagicMock() | 27 | + ): |
| 28 | - mock_args.log_level = 'info' | 28 | + mock_args = MagicMock(log_level='info', perfetto_output=None) |
| 29 | mock_arg_parser.return_value.parse_args.return_value = mock_args | 29 | mock_arg_parser.return_value.parse_args.return_value = mock_args |
| 30 | mock_service_instance = MagicMock() | 30 | mock_service_instance = MagicMock() |
| 31 | mock_otlp_service.return_value = mock_service_instance | 31 | mock_otlp_service.return_value = mock_service_instance |
| 32 | main() | 32 | main() |
| 33 | mock_set_log_level.assert_called_once_with('info') | 33 | mock_set_log_level.assert_called_once_with('info') |
| 34 | - mock_service_instance.start.assert_called_once() | 34 | + mock_otlp_service.assert_called_once_with() |
| 35 | + mock_service_instance.start.assert_called_once() | ||
| 36 | + | ||
| 37 | + | ||
| 38 | + | ||
| 39 | + | ||
| 40 | + | ||
| 41 | + def test_main_with_perfetto_uses_independent_hook_service( | ||
| 42 | + self, mock_perfetto_service, mock_otlp_service, mock_set_log_level, mock_arg_parser | ||
| 43 | + ): | ||
| 44 | + mock_args = MagicMock(log_level='debug', perfetto_output='/tmp/hook_trace.json') | ||
| 45 | + mock_arg_parser.return_value.parse_args.return_value = mock_args | ||
| 46 | + main() | ||
| 47 | + | ||
| 48 | + mock_set_log_level.assert_called_once_with('debug') | ||
| 49 | + mock_perfetto_service.assert_called_once_with('/tmp/hook_trace.json') | ||
| 50 | + mock_perfetto_service.return_value.start.assert_called_once() | ||
| 51 | + mock_otlp_service.assert_not_called() | ||