已合并
【msserviceprofiler】【需求】【Tracing 3/3】支持vLLM Hook Tracing导出Perfetto及Jaeger双写 #451
【msserviceprofiler】【需求】【Tracing 3/3】支持vLLM Hook Tracing导出Perfetto及Jaeger双写 #451
已合并
ChaseChe77创建于 12 天前
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 
722. 通过配置环境变量支持更灵活的采样控制。722. 通过配置环境变量支持更灵活的采样控制。
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```bash80```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```bash158```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+ 
162options参数说明请参见[参数说明](#section379581401015)。164options参数说明请参见[参数说明](#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>&#8226; debug:调试级别。该级别的日志记录了调试信息,便于开发人员或维护人员定位问题。<br>&#8226; info:正常级别。记录工具正常运行的信息。默认值。<br>&#8226; warning:警告级别。记录工具和预期的状态不一致,但不影响整个进程运行的信息。<br>&#8226; error:一般错误级别。<br>&#8226; critical:严重错误级别。<br>&#8226; fatal:致命错误级别。|否|170|--log-level|设置日志级别,取值为:<br>&#8226; debug:调试级别。该级别的日志记录了调试信息,便于开发人员或维护人员定位问题。<br>&#8226; info:正常级别。记录工具正常运行的信息。默认值。<br>&#8226; warning:警告级别。记录工具和预期的状态不一致,但不影响整个进程运行的信息。<br>&#8226; error:一般错误级别。<br>&#8226; critical:严重错误级别。<br>&#8226; 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)。
175python -m ms_service_profiler.trace178python -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![](figures/可视化结果.png "可视化结果")287![](figures/可视化结果.png "可视化结果")
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.md49 - SGLang 服务化性能采集工具: zh/SGLang_service_oriented_performance_collection_tool.md
50 - 3. Trace 数据链路监测:50 - 3. Trace 数据链路监测:
51 - 服务化Trace数据监测工具: zh/msserviceprofiler_trace_data_monitoring_instruct.md51 - 服务化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.md54 - 服务化性能数据比对工具: zh/ms_service_profiler_compare_tool_instruct.md
54 - 服务化多维度解析工具: zh/msserviceprofiler_multi_analyze_instruct.md55 - 服务化多维度解析工具: 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 
47if __name__ == '__main__':60if __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+ @staticmethod
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+@pytest.mark.parametrize("path", ["trace.txt", "trace.json;"])
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
18from unittest.mock import MagicMock, patch18from 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@pytest.fixture(autouse=True)21@pytest.fixture(autouse=True)
33def mock_span_export_result():22def 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 = 025 mock_result.SUCCESS = 0
37 mock_result.FAILED = 126 mock_result.FAILED = 1
38 mock_result.return_value = mock_result27 mock_result.return_value = mock_result
39- yield mock_result28+ 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+@patch("ms_service_profiler.tracer.perfetto_forward_service.signal.signal")
15+@patch("ms_service_profiler.tracer.perfetto_forward_service.AbstractSocketServer")
16+@patch("ms_service_profiler.tracer.perfetto_forward_service.PerfettoTraceExporter")
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+@patch("ms_service_profiler.tracer.perfetto_forward_service.signal.signal")
38+@patch("ms_service_profiler.tracer.perfetto_forward_service.AbstractSocketServer")
39+@patch("ms_service_profiler.tracer.perfetto_forward_service.PerfettoTraceExporter")
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 @patch('ms_service_profiler.trace.argparse.ArgumentParser')22 @patch('ms_service_profiler.trace.argparse.ArgumentParser')
23 @patch('ms_service_profiler.trace.set_log_level')23 @patch('ms_service_profiler.trace.set_log_level')
24 @patch('ms_service_profiler.trace.OTLPForwarderService')24 @patch('ms_service_profiler.trace.OTLPForwarderService')
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_args29 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_instance31 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+ @patch('ms_service_profiler.trace.argparse.ArgumentParser')
38+ @patch('ms_service_profiler.trace.set_log_level')
39+ @patch('ms_service_profiler.trace.OTLPForwarderService')
40+ @patch('ms_service_profiler.tracer.perfetto_forward_service.PerfettoForwarderService')
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()