【NPU】Ascend Device Plugin 源码超深度分析 — Part 4 310/310P 设备 + 交换机 + DPU + 重置信息管理之一
Part 4: Ascend Device Plugin 源码超深度分析
310/310P 设备 + 交换机 + DPU + 重置信息管理
分析范围:
pkg/device/ascend310.go、pkg/device/ascend310p.go、pkg/device/deviceswitch/ascend_switch.go、pkg/device/dpucontrol/dpu_device_find.go、pkg/device/dpucontrol/types.go、pkg/device/reset_info_mgr.go源码版本:mind-cluster-v26.0.1
分析者:资深 Go 架构师视角
目录
- 1. ascend310.go — Ascend 310 设备管理器
- 2. ascend310p.go — Ascend 310P 设备管理器
- 3. ascend_switch.go — 交换机故障管理
- 4. dpu_device_find.go — DPU 设备发现与过滤
- 5. types.go — DPU 控制类型定义
- 6. reset_info_mgr.go — NPU 重置信息管理器
- 7. 模块间协作关系总览
1. ascend310.go — Ascend 310 设备管理器
1.1 模块定位
业务职责
HwAscend310Manager 是华为 Ascend 310 NPU 芯片的设备管理器实现。它负责:
- 设备发现:通过 devmanager 接口枚举节点上的所有 Ascend 310 设备
- 设备模式管理:支持普通模式和共享模式两种设备分配方式
- 设备状态上报:将健康/不健康设备列表、故障码等信息写入 Kubernetes Node Annotation
- Volcano 调度集成:通过
DoWithVolcanoListAndWatch实现 Volcano 调度器的设备列表与监听
在系统中的位置
Ascend 310 管理器是设备管理器接口的一个具体实现,面向 Ascend 310 芯片场景(边缘推理、轻量训练),在继承 AscendTools 基类能力的基础上,实现了 310 特有的设备发现和状态上报逻辑。
1.2 模块整体结构
类结构
核心方法清单
| 方法名 | 签名 | 作用 | 可见性 |
|---|---|---|---|
NewHwAscend310Manager |
() *HwAscend310Manager |
工厂函数,创建 310 管理器实例 | 公开 |
GetNPUs |
() (common.NpuAllInfo, error) |
发现所有 310 设备 | 公开 |
getNPUsByNormalMode |
(davinCiDev) []NpuDevice |
普通模式下组装设备信息 | 私有 |
DoWithVolcanoListAndWatch |
(classifyDevs, chipMemory) |
Volcano 调度设备监听 | 公开 |
updateDeviceInfo |
(old, new, devStatusSet) error |
更新 Node Annotation | 私有 |
GraceTolerance |
(ctx, map) |
优雅容错(310 不支持,空实现) | 公开 |
GetAssociatedLogicIDs |
(logicID, cardID, deviceID) ([]int32, error) |
获取关联逻辑ID(310 不支持,空实现) | 公开 |
SetDpu |
(string, []DpuCMData, map) |
DPU 写入(310 不支持,空实现) | 公开 |
1.3 核心业务逻辑深度解析
1.3.1 构造函数 NewHwAscend310Manager
func NewHwAscend310Manager() *HwAscend310Manager {
name := api.Ascend310 // 默认设备名 "Ascend310"
if common.ParamOption.GetFdFlag {
name = common.AscendfdPrefix // 若启用 FD 模式,使用 "davinci-mini"
}
return &HwAscend310Manager{
AscendTools: AscendTools{
name: name,
unHealthyKey: common.HuaweiUnHealthAscend310, // "Ascend310-Unhealthy"
devCount: common.MaxCardNum * common.MaxDevNumInCard, // 64 * 4 = 256
cardInResetMap: make(map[int32]bool, common.GeneralMapSize), // 初始容量 8
resetFailedTimesMap: make(map[int32]int, common.GeneralMapSize),
lastUsedChipsContainerMap: make(map[string]sets.String),
},
}
}
逐行解析:
name:设备类型名称。api.Ascend310是标准名称(如"Ascend310"),若GetFdFlag为 true 则改为davinci-mini,这是为了兼容 FD(Fd Device Plugin)模式下的设备命名。unHealthyKey:不健康设备在 Node Annotation 中的 key,值为"Ascend310-Unhealthy"。devCount:最大设备数 =MaxCardNum(64)×MaxDevNumInCard(4)= 256。这是 310 芯片在单节点上的理论上限。cardInResetMap:记录哪些逻辑ID的卡正在重置中,map[int32]bool,预分配容量 8(GeneralMapSize)。resetFailedTimesMap:记录每张卡重置失败的次数。lastUsedChipsContainerMap:记录容器最后使用的芯片集合,map[string]sets.String,用于追踪设备使用状态。
1.3.2 设备发现 GetNPUs
逐行解析:
func (hnm *HwAscend310Manager) GetNPUs() (common.NpuAllInfo, error) {
// 1. 调用 devmanager 接口获取设备列表
devNum, devList, err := hnm.dmgr.GetDeviceList()
if err != nil {
return common.NpuAllInfo{}, err
}
// 2. 设备数量合法性校验,防止异常值导致后续处理溢出
if devNum > hnm.devCount {
return common.NpuAllInfo{}, fmt.Errorf("invalid device num: %d", devNum)
}
// 3. 初始化设备结果切片
var allDevices = make([]common.NpuDevice, 0)
// 4. 遍历每个逻辑设备
for logicIDIdx := 0; logicIDIdx < len(devList); logicIDIdx++ {
// 4.1 获取 DavinCi 设备详细信息(PhyID, CardID, DeviceID 等)
davinCiDev, err := hnm.getDavinCiDev(devList[logicIDIdx])
if err != nil {
return common.NpuAllInfo{}, err
}
// 4.2 根据是否启用共享模式选择不同的设备组装方式
normalDevices := hnm.getNPUsByNormalMode(davinCiDev)
if common.ShareDev() {
normalDevices = hnm.getNPUsByShareMode(davinCiDev)
}
// 4.3 将组装好的设备追加到结果集
allDevices = append(allDevices, normalDevices...)
}
// 5. 返回所有设备和设备类型(310 只有一种类型)
return common.NpuAllInfo{AllDevs: allDevices, AllDevTypes: []string{hnm.name}}, err
}
设计意图:
- 310 是较简单的芯片,不支持虚拟化切分,因此
GetNPUs逻辑相对简单——每个物理设备映射为一个NpuDevice AllDevTypes只包含一种类型(hnm.name),因为 310 不涉及 AICore 等子设备- 共享模式(ShareDev)允许一张卡被多个 Pod 使用,通过
getNPUsByShareMode生成不同的设备名
1.3.3 普通模式设备组装 getNPUsByNormalMode
func (hnm *HwAscend310Manager) getNPUsByNormalMode(davinCiDev common.DavinCiDev) []common.NpuDevice {
// 构造设备名:如 "Ascend310-0"、"Ascend310-1"
deviceName := fmt.Sprintf("%s-%d", hnm.name, davinCiDev.PhyID)
// 组装并返回单个 NpuDevice 结构
return []common.NpuDevice{hnm.assembleNpuDeviceStruct(hnm.name, deviceName, davinCiDev)}
}
解析:
deviceName格式为<设备类型>-<物理ID>,如Ascend310-0assembleNpuDeviceStruct是AscendTools基类中的方法,负责将 DavinCiDev 信息填充到 NpuDevice 结构中- 返回切片(而非单个值)是为了与共享模式接口统一(共享模式可能返回多个虚拟设备)
1.3.4 Volcano 调度监听 DoWithVolcanoListAndWatch
func (hnm *HwAscend310Manager) DoWithVolcanoListAndWatch(classifyDevs map[string][]*common.NpuDevice, chipMemory int) {
// 1. 获取设备状态集合(健康/不健康/故障等)
devStatusSet := hnm.getDevStatesDevSet(classifyDevs, chipMemory)
// 2. 更新 Node 上的设备信息 Annotation
if err := hnm.UpdateNodeDeviceInfo(devStatusSet, common.DpuInfo{}, hnm.updateDeviceInfo); err != nil {
hwlog.RunLog.Errorf("update device info failed, err: %v", err)
}
}
解析:
classifyDevs是按设备类型分类的设备列表chipMemory是芯片内存大小,用于计算设备资源getDevStatesDevSet是基类方法,分析所有设备的状态(健康/不健康/网络不健康等),返回DevStatusSetUpdateNodeDeviceInfo调用传入的updateDeviceInfo回调函数,将状态写入 Node Annotationcommon.DpuInfo{}传入空 DPU 信息,因为 310 不涉及 DPU
1.3.5 设备信息更新回调 updateDeviceInfo
func (hnm *HwAscend310Manager) updateDeviceInfo(_, newDeviceInfo map[string]string,
devStatusSet common.DevStatusSet) error {
// 1. 参数校验
if newDeviceInfo == nil {
return fmt.Errorf("invalid new device info")
}
// 2. 写入健康设备列表
// key="Ascend310", value=逗号分隔的健康设备名列表
newDeviceInfo[api.HuaweiAscend310] = common.ToString(
devStatusSet.FreeHealthyDevice[hnm.name], common.CommaSepDev)
// 3. 写入不健康设备列表
// key="Ascend310-Unhealthy", value=逗号分隔的不健康设备名列表
newDeviceInfo[hnm.unHealthyKey] = common.ToString(
devStatusSet.UnHealthyDevice, common.CommaSepDev)
// 4. 序列化故障码并写入
var data []byte
if data = common.MarshalData(devStatusSet.DeviceFault); len(data) == 0 {
return fmt.Errorf("device fault code marshal failed")
}
// key="Ascend310P-Fault" (注意:这里用的是 HuaweiFaultCodeAscend310 常量)
newDeviceInfo[common.HuaweiFaultCodeAscend310] = string(data)
return nil
}
解析:
- 第一个参数
_(oldDeviceInfo)被忽略,因为 310 采用全量覆盖策略 - 三类信息写入 Annotation:
- 健康设备:key=
HuaweiAscend310,value如"Ascend310-0,Ascend310-1" - 不健康设备:key=
Ascend310-Unhealthy - 故障码:key=
HuaweiFaultCodeAscend310(注意:源码中此常量值实际为api.HuaweiAscend310P + "-Fault",这是源码中的一个已知问题)
- 健康设备:key=
1.3.6 空实现方法
// GraceTolerance — 优雅容错,310 不支持
func (hnm *HwAscend310Manager) GraceTolerance(context.Context, map[string][]*common.NpuDevice) {
return // 直接返回,无操作
}
// GetAssociatedLogicIDs — 获取关联逻辑ID,310 不支持
func (hnm *HwAscend310Manager) GetAssociatedLogicIDs(logicID, cardID, deviceID int32) ([]int32, error) {
return nil, nil // 返回空
}
// SetDpu — DPU 写入,310 不支持
func (hnm *HwAscend310Manager) SetDpu(string, []common.DpuCMData, map[string][]string) {
return // 直接返回,无操作
}
设计意图:这三个方法是 DevManager 接口的一部分,310 芯片不需要这些功能(没有交换机互联、没有 DPU),但仍需实现接口,因此使用空实现(null object 模式)。
2. ascend310p.go — Ascend 310P 设备管理器
2.1 模块定位
业务职责
HwAscend310PManager 管理华为 Ascend 310P NPU 芯片。与 310 相比,310P 支持更丰富的功能:
- 虚拟化设备切分:支持 vNPU 虚拟设备,可将一张物理卡切分为多个虚拟设备
- 混合插拔模式:
Use310PMixedInsert支持 310P 与其他设备混合使用 - 共享模式:支持 ShareDev 模式
- AICore 设备:支持 FakeAiCoreDevice 虚拟 AI Core 设备
与 310 的关键差异
| 特性 | Ascend 310 | Ascend 310P |
|---|---|---|
| 虚拟化切分 | ❌ | ✅ |
| 混合插拔 | ❌ | ✅ |
| AICore 设备 | ❌ | ✅ |
| 最大设备数 | 256 (64×4) | 100 (MaxDevicesNum) |
| 设备名 | Ascend310 / davinci-mini | Ascend310P |
| FD 模式前缀 | 支持 | 不支持 |
2.2 模块整体结构
2.3 核心业务逻辑深度解析
2.3.1 构造函数 NewHwAscend310PManager
func NewHwAscend310PManager() *HwAscend310PManager {
return &HwAscend310PManager{
AscendTools: AscendTools{
name: api.Ascend310P, // "Ascend310P"
unHealthyKey: common.HuaweiUnHealthAscend310P, // "Ascend310P-Unhealthy"
devCount: common.MaxDevicesNum, // 100
cardInResetMap: make(map[int32]bool, common.GeneralMapSize),
resetFailedTimesMap: make(map[int32]int, common.GeneralMapSize),
lastUsedChipsContainerMap: make(map[string]sets.String),
},
}
}
与 310 构造函数的差异:
- 无 FD 模式判断:310P 不需要支持
davinci-mini前缀 devCount:使用MaxDevicesNum(100)而非MaxCardNum * MaxDevNumInCard(256),因为 310P 的设备拓扑结构不同- 其余字段结构与 310 一致
2.3.2 设备发现 GetNPUs — 核心复杂逻辑
逐行解析:
func (hnm *HwAscend310PManager) GetNPUs() (common.NpuAllInfo, error) {
// 1. 获取设备列表
devNum, devList, err := hnm.dmgr.GetDeviceList()
if err != nil {
return common.NpuAllInfo{}, err
}
// 2. 设备数校验
if devNum > hnm.devCount {
return common.NpuAllInfo{}, fmt.Errorf("invalid device num: %d", devNum)
}
// 3. 初始化结果变量
var allDevices []common.NpuDevice // 所有设备(物理+虚拟+共享)
var aiCoreDevices []*common.NpuDevice // AI Core 虚拟设备
var allDeviceTypes = make([]string, 0) // 所有设备类型名(可能多种)
// 4. 遍历每个设备
for i := int32(0); i < devNum; i++ {
// 4.1 获取 DavinCi 设备详情
davinCiDev, err := hnm.getDavinCiDev(devList[i])
if err != nil {
return common.NpuAllInfo{}, err
}
// 4.2 混合插拔模式分支
if common.ParamOption.Use310PMixedInsert {
if err = hnm.assemble310PMixedPhyDevices(davinCiDev, &allDevices, &allDeviceTypes); err != nil {
hwlog.RunLog.Errorf("assemble mixed phy devices failed: %v", err)
}
continue // 混合模式处理完直接跳到下一个设备
}
// 4.3 获取虚拟设备信息
vDevInfos, err := hnm.getVirtualDevice(devList[i])
if err != nil {
// 虚拟设备查询失败不致命,记录日志继续
hwlog.RunLog.Errorf("The virtual device is considered not exist, please check the error: %v", err)
}
// 4.4 虚拟设备数量校验
if vDevInfos.TotalResource.VDevNum > common.MaxVirtualDeviceNum {
return common.NpuAllInfo{}, fmt.Errorf("invalid virtual device count")
}
// 4.5 虚拟设备与共享模式互斥校验
if vDevInfos.TotalResource.VDevNum > 0 && common.ShareDev() {
return common.NpuAllInfo{}, fmt.Errorf("virtual device is exist, shareDevCount should be 1")
}
// 4.6 若未预设 VDevice,则创建 Fake AI Core 设备
if !common.ParamOption.PresetVDevice {
common.FakeAiCoreDevice(davinCiDev, &aiCoreDevices)
}
// 4.7 虚拟设备分支
if vDevInfos.TotalResource.VDevNum > 0 {
hnm.assembleVirtualDevices(davinCiDev, vDevInfos, &allDevices, &allDeviceTypes)
continue
}
// 4.8 共享模式分支
if common.ShareDev() {
hnm.assembleShareModeDevices(davinCiDev, &allDevices, &allDeviceTypes)
} else {
// 4.9 普通物理设备分支
hnm.assemblePhyDevices(davinCiDev, &allDevices, &allDeviceTypes)
}
}
// 5. 去重设备类型
allDeviceTypes = hnm.removeDuplicate(&allDeviceTypes)
// 6. 返回完整信息(包含 AICoreDevs,310 没有此字段)
return common.NpuAllInfo{AllDevs: allDevices, AICoreDevs: aiCoreDevices, AllDevTypes: allDeviceTypes}, nil
}
关键分支分析:
| 分支条件 | 调用方法 | 说明 |
|---|---|---|
Use310PMixedInsert=true |
assemble310PMixedPhyDevices |
310P 混合插拔场景 |
vDevNum > 0 |
assembleVirtualDevices |
有虚拟设备,组装虚拟设备 |
ShareDev()=true |
assembleShareModeDevices |
共享模式,一张卡多 Pod |
| 默认 | assemblePhyDevices |
普通物理设备模式 |
互斥校验:虚拟设备(vDevNum > 0)与共享模式(ShareDev)不能同时使用,否则报错。
2.3.3 Volcano 监听与信息更新
func (hnm *HwAscend310PManager) DoWithVolcanoListAndWatch(classifyDevs map[string][]*common.NpuDevice, chipMemory int) {
devStatusSet := hnm.getDevStatesDevSet(classifyDevs, chipMemory)
if err := hnm.UpdateNodeDeviceInfo(devStatusSet, common.DpuInfo{}, hnm.updateDeviceInfo); err != nil {
hwlog.RunLog.Errorf("update device info failed, err: %v", err)
}
}
与 310 的实现结构完全一致,区别在于 updateDeviceInfo 回调中使用的 key 不同:
func (hnm *HwAscend310PManager) updateDeviceInfo(_, newDevInfo map[string]string,
devStatusSet common.DevStatusSet) error {
if newDevInfo == nil {
return fmt.Errorf("invalid new device info")
}
// 健康设备 key="Ascend310P"
newDevInfo[api.HuaweiAscend310P] = common.ToString(
devStatusSet.FreeHealthyDevice[hnm.name], common.CommaSepDev)
// 不健康设备 key="Ascend310P-Unhealthy"
newDevInfo[hnm.unHealthyKey] = common.ToString(
devStatusSet.UnHealthyDevice, common.CommaSepDev)
// 故障码 key="Ascend310P-Fault"
var data []byte
if data = common.MarshalData(devStatusSet.DeviceFault); len(data) == 0 {
return fmt.Errorf("device fault code marshal failed")
}
newDevInfo[common.HuaweiFaultCodeAscend310P] = string(data)
return nil
}
2.3.4 空实现方法
与 310 相同,GraceTolerance、GetAssociatedLogicIDs、SetDpu 均为空实现,310P 同样不支持这些功能。
3. ascend_switch.go — 交换机故障管理
3.1 模块定位
业务职责
SwitchDevManager 负责管理 Ascend 910A3 交换机芯片的故障事件。这是整个 device-plugin 中唯一涉及 CGO 交互的模块,功能包括:
- 驱动库加载:通过
dlopen加载liblingqu-dcmi.so动态库 - 故障订阅:通过
lq_dcmi_subscribe_fault_event订阅交换机故障事件 - 故障轮询:通过
lq_dcmi_get_fault_info定时查询故障信息 - 故障码组装:将 C 结构体转换为 Go 的
SwitchFaultEvent并组装标准故障码 - 故障等级映射:将故障码映射到不同处理等级(不处理、亚健康、重启请求、预隔离、隔离)
在系统中的位置
3.2 模块整体结构
CGO 架构
核心方法清单
| 方法名 | 签名 | 作用 |
|---|---|---|
UpdateSwitchFaultLevel |
() |
更新故障码到等级的全局映射表 |
NewSwitchDevManager |
() *SwitchDevManager |
创建交换机管理器 |
InitSwitchDev |
() error |
初始化驱动库 |
ShutDownSwitch |
() |
关闭驱动库 |
GetSwitchFaultCodeByInterval |
(ctx, interval) |
定时轮询故障码 |
SubscribeSwitchFaults |
() error |
订阅故障事件 |
GetSwitchFaults |
() ([]SwitchFaultEvent, error) |
查询当前所有故障 |
goFaultEventHandler |
(event *C.struct_LqDcmiEvent) |
CGO 回调入口 |
convertFaultEvent |
(event) SwitchFaultEvent |
C→Go 结构转换 |
setExtraFaultInfo |
(event) |
组装标准故障码 |
isPortLevelFault |
(switchPortId) bool |
判断是否端口级故障 |
isFaultRecoveredEvent |
(fault, recover) bool |
判断故障恢复事件 |
3.3 核心业务逻辑深度解析
3.3.1 CGO C 代码解析
文件顶部通过 import "C" 嵌入了 C 代码,这是 Go 与驱动库交互的桥梁:
// 定义函数指针,用于动态加载 .so 中的符号
static int (*lq_dcmi_init_func)();
static int (*lq_dcmi_get_fault_info_func)(unsigned int listLen, unsigned int *eventListLen,
struct LqDcmiEvent *eventList);
static int (*lq_dcmi_subscribe_fault_event_func)(struct lq_dcmi_event_filter filter,
LqDcmiFaultEventCallback handler);
// 事件回调中间层:C 函数 → Go 函数
static void event_handler(struct LqDcmiEvent *fault_event) {
goFaultEventHandler(fault_event); // 调用 Go 侧的 export 函数
}
// dlopen 加载动态库并绑定函数符号
static int dcmiInit_lq(const char* dcmiLibPath) {
dcmiHandle = dlopen(dcmiLibPath, RTLD_LAZY | RTLD_GLOBAL);
lq_dcmi_init_func = dlsym(dcmiHandle, "lq_dcmi_init");
lq_dcmi_subscribe_fault_event_func = dlsym(dcmiHandle, "lq_dcmi_subscribe_fault_event");
lq_dcmi_get_fault_info_func = dlsym(dcmiHandle, "lq_dcmi_get_fault_info");
return SUCCESS;
}
// dlclose 关闭动态库
static int lqDcmiShutDown(void) {
if (dcmiHandle == NULL) return SUCCESS;
return (dlclose(dcmiHandle) ? ERROR_UNKNOWN : SUCCESS);
}
设计意图:
- 使用
dlopen/dlsym而非直接链接,是为了在没有驱动库的环境下也能编译通过 RTLD_LAZY延迟解析符号,RTLD_GLOBAL使符号可被后续加载的库使用event_handler是 C 侧的回调函数,它转发到 Go 侧的goFaultEventHandler(通过//export导出)
3.3.2 故障等级映射 UpdateSwitchFaultLevel
func UpdateSwitchFaultLevel() {
// 定义故障码分组:每组对应一个处理等级
faultCodeGroups := []struct {
codes []string
level int
}{
{common.NotHandleFaultCodes, common.NotHandleFaultLevel}, // 0 - 不处理
{common.SubHealthFaultCodes, common.SubHealthFaultLevel}, // 1 - 亚健康
{common.RestartRequestFaultCodes, common.RestartRequestFaultLevel}, // 2 - 重启请求
{common.PreSeparateFaultCodes, common.PreSeparateFaultLevel}, // 3 - 预隔离
{common.SeparateFaultCodes, common.SeparateFaultLevel}, // 4 - 隔离
}
// 加锁保护全局映射表
common.SwitchFaultLevelMapLock.Lock()
defer common.SwitchFaultLevelMapLock.Unlock()
// 重建映射表
common.SwitchFaultLevelMap = make(map[string]int, common.GeneralMapSize)
for _, group := range faultCodeGroups {
for _, code := range group.codes {
common.SwitchFaultLevelMap[code] = group.level
}
}
}
故障等级体系:
解析:
- 每次调用时完全重建
SwitchFaultLevelMap,而非增量更新,确保一致性 - 使用全局锁
SwitchFaultLevelMapLock保护并发访问 - 故障码来源是
common包中的全局变量(从配置文件加载)
3.3.3 驱动库初始化 InitSwitchDev
func (sdm *SwitchDevManager) InitSwitchDev() error {
dcmiLibName := "liblingqu-dcmi.so"
// 1. 查找驱动库路径(可能在多个位置)
dcmiLibPath, err := utils.GetDriverLibPath(dcmiLibName)
if err != nil {
return fmt.Errorf("failed to find switch library so, err: %s", err.Error())
}
// 2. Go string → C string(需手动释放)
cDcmiTemplateName := C.CString(dcmiLibPath)
defer C.free(unsafe.Pointer(cDcmiTemplateName))
// 3. 加载动态库并绑定函数符号
if retCode := C.dcmiInit_lq(cDcmiTemplateName); retCode != C.SUCCESS {
return fmt.Errorf("dcmi lib load failed, error code: %d", int32(retCode))
}
// 4. 调用驱动初始化函数
if retCode := C.dcmi_init_lq(); retCode != C.SUCCESS {
return fmt.Errorf("dcmi init call failed, error code: %d", int32(retCode))
}
hwlog.RunLog.Info("init switch library succeeded")
return nil
}
3.3.4 故障事件回调 goFaultEventHandler
这是 CGO 回调的核心入口,当驱动上报故障事件时被调用:
//export goFaultEventHandler
func goFaultEventHandler(event *C.struct_LqDcmiEvent) {
// 1. 确保触发更新通知(defer 保证即使 panic 也会执行)
defer func() {
common.TriggerUpdate("switch fault occur")
}()
// 2. C 结构体 → Go 结构体转换
faultEvent := convertFaultEvent(event)
hwlog.RunLog.Warnf("switch subscribe got fault:%s, faultCode:%v", ...)
// 3. 分两种情况处理
if int8(faultEvent.Assertion) == devmanagercommon.FaultRecover {
// 3.1 故障恢复事件:从当前故障列表中移除已恢复的故障
newFaultCodes := make([]common.SwitchFaultEvent, 0)
for _, errInfo := range common.GetSwitchFaultCode() {
if !isFaultRecoveredEvent(errInfo, faultEvent) {
newFaultCodes = append(newFaultCodes, errInfo) // 未恢复的保留
}
}
common.SetSwitchFaultCode(newFaultCodes)
return
}
// 3.2 新故障事件:追加到当前故障列表
currentFault := common.GetSwitchFaultCode()
common.SetSwitchFaultCode(append(currentFault, faultEvent))
}
关键设计:
//export goFaultEventHandler是 Go 的 CGO 导出指令,使 C 代码可以调用此 Go 函数defer TriggerUpdate确保任何故障事件都会触发设备状态更新- 故障恢复逻辑通过
isFaultRecoveredEvent比较关键字段来判断是否是同一个故障的恢复事件
3.3.5 故障码组装 setExtraFaultInfo
func setExtraFaultInfo(event *common.SwitchFaultEvent) {
// 1. 判断对端设备类型
PeerDeviceType, PeerDeviceName := int(event.PeerPortDevice), ""
if isPortLevelFault(int(event.SwitchPortId)) {
PeerDeviceName = getPeerDeviceName(PeerDeviceType) // cpu/npu/L2
} else {
PeerDeviceName = common.PeerDeviceNAPortName // "na"
}
// 2. 组装故障码
alarmID, faultID := event.EventType, event.SubType
if faultID == invalidNum { // 0xFFFFFFFF
// linkdown 特殊处理:faultID 为 "na"
event.FaultID = common.PeerDeviceNAPortName
event.AssembledFaultCode = fmt.Sprintf("[0x%08x,na,%s,na]", alarmID, PeerDeviceName)
} else {
event.FaultID = strconv.Itoa(int(faultID))
event.AssembledFaultCode = fmt.Sprintf("[0x%08x,%d,%s,na]", alarmID, faultID, PeerDeviceName)
}
// 3. 补充告警时间
if event.AlarmRaisedTime == int64(0) {
event.AlarmRaisedTime = time.Now().UnixMilli()
}
}
故障码格式:[AlarmID, FaultID, PeerDeviceName, na]
| 字段 | 说明 | 示例 |
|---|---|---|
| AlarmID | 告警ID,十六进制 | 0x00f1ff09 |
| FaultID | 故障子ID,数字或"na" | 155912 或 na |
| PeerDeviceName | 对端设备类型 | cpu/npu/L2/na |
| na | 保留位 | na |
对端设备类型映射:
3.3.6 定时故障查询 GetSwitchFaultCodeByInterval
func (sdm *SwitchDevManager) GetSwitchFaultCodeByInterval(ctx context.Context, interval time.Duration) {
// 1. 初始查询
hwlog.RunLog.Info("performing initial query of switch fault codes")
runtime.LockOSThread() // 锁定 OS 线程(CGO 需要在同一线程)
updateSwitchFaultCode(true)
// 2. 定时轮询
ticker := time.NewTicker(interval)
defer func() {
ticker.Stop()
hwlog.RunLog.Info("query switch fault by interval stopped")
}()
for {
select {
case _, ok := <-ctx.Done():
if !ok {
hwlog.RunLog.Info("catch stop signal channel closed")
}
hwlog.RunLog.Infof("received stop signal: %v", ctx.Err())
return
case <-ticker.C:\n updateSwitchFaultCode(false)\n }\n }\n}\n```\n\n**`updateSwitchFaultCode` 内部逻辑**:
```go
func updateSwitchFaultCode(isInit bool) {
// 非初始化时,若无故障码则跳过
if !isInit {
switchFaultCodes := common.GetSwitchFaultCode()
if len(switchFaultCodes) == 0 {
hwlog.RunLog.Info("no switch fault codes to query, skip this cycle")
return
}
}
// 调用 C 接口查询故障
errCodes, err := GetSwitchFaults()
if err != nil {
hwlog.RunLog.Errorf("failed to query switch fault codes: %v", err)
return
}
common.SetSwitchFaultCode(errCodes)
}
设计意图:
runtime.LockOSThread()确保整个轮询在同一个 OS 线程上运行,因为 CGO 调用依赖线程局部状态- 初始查询(
isInit=true)无论是否有故障码都会执行,确保启动时获取完整状态 - 后续轮询(
isInit=false)如果当前没有故障码则跳过,减少不必要的 CGO 调用
3.3.7 批量故障查询 GetSwitchFaults
func GetSwitchFaults() ([]common.SwitchFaultEvent, error) {
var errCount C.uint
var errInfoArray [maxFaultNum]C.struct_LqDcmiEvent // 128 个事件槽位
// 1. 调用 C 接口批量获取故障
if retCode := C.lq_dcmi_get_fault_info(C.uint(maxFaultNum), &errCount, &errInfoArray[0]);
int32(retCode) != devmanagercommon.Success {
return []common.SwitchFaultEvent{}, fmt.Errorf("failed to get switch device errorcodes, errCode:%v", retCode)
}
// 2. 校验故障数量
if int32(errCount) < 0 || int32(errCount) > maxFaultNum {
return []common.SwitchFaultEvent{}, fmt.Errorf("failed to get switch device errcodes, cause errcodes nums %v is illegal", errCount)
}
// 3. 遍历故障事件,过滤恢复事件
errorCodes := make([]string, 0)
retErrorInfo := make([]common.SwitchFaultEvent, 0)
for i := 0; i < int(errCount); i++ {
faultEvent := convertFaultEvent(&errInfoArray[i])
if int8(faultEvent.Assertion) == devmanagercommon.FaultRecover {
continue // 恢复事件不加入结果
}
errorCodes = append(errorCodes, faultEvent.AssembledFaultCode)
retErrorInfo = append(retErrorInfo, faultEvent)
}
if len(errorCodes) > 0 {
hwlog.RunLog.Warnf("switch of 910A3 get fault codes: %#v", errorCodes)
}
return retErrorInfo, nil
}
3.3.8 故障恢复判断 isFaultRecoveredEvent
func isFaultRecoveredEvent(faultEvent, recoverEvent common.SwitchFaultEvent) bool {
// 1. recoverEvent 必须是恢复类型,且与 faultEvent 的 Assertion 不同
if int8(recoverEvent.Assertion) != devmanagercommon.FaultRecover ||
recoverEvent.Assertion == faultEvent.Assertion {
return false
}
// 2. 比较所有关键字段是否一致
faultEventInfo := fmt.Sprintf("EventType:%v,FaultID:%v,AssembledFaultCode:%v,PeerPortDevice:%v,PeerPortId:%v,SwitchChipId:%v,SwitchPortId:%v",
faultEvent.EventType, faultEvent.SubType, ...)
recoveredEventInfo := fmt.Sprintf("EventType:%v,FaultID:%v,AssembledFaultCode:%v,PeerPortDevice:%v,PeerPortId:%v,SwitchChipId:%v,SwitchPortId:%v",
recoverEvent.EventType, recoverEvent.SubType, ...)
// 3. 字符串比较:完全一致则认为是同一故障的恢复
return faultEventInfo == recoveredEventInfo
}
设计意图:通过将所有关键标识字段拼接为字符串后比较,来判断故障事件和恢复事件是否匹配。虽然字符串拼接比较不是最高效的方式,但逻辑清晰且不易遗漏字段。
3.3.9 C 结构体转换 convertFaultEvent
func convertFaultEvent(event *C.struct_LqDcmiEvent) common.SwitchFaultEvent {
fault := common.SwitchFaultEvent{
EventType: uint(event.eventType), // 告警类型
SubType: uint(event.subType), // 告警子类型
PeerPortDevice: uint(event.peerportDevice), // 对端端口设备类型
PeerPortId: uint(event.peerportId), // 对端端口ID
SwitchChipId: uint(event.switchChipid), // 交换芯片ID
SwitchPortId: uint(event.switchPortid), // 交换端口ID
Severity: uint(event.severity), // 严重程度
Assertion: uint(event.assertion), // 断言类型(发生/恢复)
EventSerialNum: int(event.eventSerialNum), // 事件序列号
NotifySerialNum: int(event.notifySerialNum), // 通知序列号
AlarmRaisedTime: int64(event.alarmRaisedTime), // 告警时间
}
setExtraFaultInfo(&fault) // 组装额外字段(FaultID, AssembledFaultCode 等)
return fault
}
4. dpu_device_find.go — DPU 设备发现与过滤
4.1 模块定位
业务职责
DpuFilter 负责在节点上发现和识别 DPU(Data Processing Unit)设备,并建立 NPU 与 DPU 的映射关系。功能包括:
- 配置加载:从
/user/mindx-dl/dpu/dpu-config.json加载用户 DPU 配置 - 两种总线模式:
- UB 模式(
busTypeUb):通过/sys/class/net目录扫描 DPU 网卡 - PCIe 模式(
busTypePcie):通过 NPU 的 PCIe 总线信息反向查找连接的 DPU
- UB 模式(
- 设备过滤:按 Vendor、DeviceId、DeviceName 三维度过滤
- NPU-DPU 映射:建立 NPU 与 DPU 的对应关系
在系统中的位置
4.2 模块整体结构
4.3 核心业务逻辑深度解析
4.3.1 主入口 SaveDpuConfToNode
func (df *DpuFilter) SaveDpuConfToNode(dmgr devmanager.DeviceInterface) error {
// 1. 检查配置文件是否存在,不存在则跳过(DPU 功能未启用)
if !utils.IsExist(DpuConfigPath) {
hwlog.RunLog.Infof("%s file %s not found", api.DpuLogPrefix, DpuConfigPath)
return nil
}
// 2. 加载配置文件
err := df.loadDpuConfigFromFile()
if err != nil {
return fmt.Errorf("dpu devices find is not enable,err:%v", err)
}
// 3. 根据总线类型选择发现方式
switch df.UserConfig.BusType {
case busTypeUb:
// 3.1 UB 模式:直接扫描 /sys/class/net
entries, err := os.ReadDir(netPath)
if err != nil {
return fmt.Errorf("read dpu file dir err:%v", err)
}
df.entries = entries
dpuInfos, err := df.filterDpu()
if err != nil {
return fmt.Errorf("filter dpu err:%v", err)
}
df.dpuInfos = dpuInfos
err = df.getNpuCorrespDpuInfo()
if err != nil {
return fmt.Errorf("build npu correspond dpu infos err:%v", err)
}
case busTypePcie:
// 3.2 PCIe 模式:通过 NPU 反向查找 DPU
err = df.getDpuWithNpuPcieSwitch(dmgr)
if err != nil {
return fmt.Errorf("get dpu by npu error: %v", err)
}
default:
return fmt.Errorf("unsupported busType: %s", df.UserConfig.BusType)
}
// 4. 校验结果非空
if len(df.NpuWithDpuInfos) == 0 {
return errors.New("filter dpu infos result is nil")
}
hwlog.RunLog.Infof("%s successfully get DPU infos: %v", api.DpuLogPrefix, df.NpuWithDpuInfos)
return nil
}
4.3.2 配置加载 loadDpuConfigFromFile
func (df *DpuFilter) loadDpuConfigFromFile() error {
// 1. 读取 JSON 配置文件
jsonContent, err := utils.LoadFile(DpuConfigPath)
if err != nil {
return fmt.Errorf("load config from file error:%v", err)
}
// 2. JSON 反序列化
var configList ConfigList
if err = json.Unmarshal(jsonContent, &configList); err != nil {
return fmt.Errorf("parse config from file error:%v", err)
}
// 3. 参数校验
userConfigList := configList.UserDpuConfigList
if len(userConfigList) == 0 || userConfigList[0].Selectors == nil || (userConfigList[0].BusType == "") {
return errors.New("config missing parameter, dpu devices find is not enable")
}
// 4. 只取第一个配置(当前版本只支持单个 DPU 配置)
userConfig := userConfigList[0]
busType := userConfig.BusType
if busType != busTypeUb && busType != busTypePcie {
return fmt.Errorf("invalid busType: %s", busType)
}
// 5. 至少需要 Vendor 或 DeviceIds 中的一个
selectors := userConfig.Selectors
if len(selectors.Vendor) == 0 && len(selectors.DeviceIds) == 0 {
return errors.New("no vendor and deviceIds found, dpu devices find is not enable")
}
df.UserConfig = userConfig
return nil
}
配置文件格式:
{
"configList": [
{
"busType": "ub",
"selectors": {
"vendor": ["0x15b3"],
"deviceIds": ["0x1021"],
"devices": ["eth0"]
}
}
]
}
4.3.3 PCIe 模式 DPU 发现 getDpuWithNpuPcieSwitch
func (df *DpuFilter) getDpuWithNpuPcieSwitch(dcMgr devmanager.DeviceInterface) error {
// 1. 获取所有 NPU 逻辑 ID 列表
devNum, logicIDList, err := dcMgr.GetDeviceList()
if err != nil || devNum == 0 {
return fmt.Errorf("get device list error: %v", err)
}
// 2. PCIe Switch ID 去重 map
pcieSwIds := make(map[string]struct{})
for _, logicID := range logicIDList {
// 3. 通过逻辑 ID 获取 PCIe 总线信息
pcieBusInfo, err := dcMgr.GetPCIeBusInfo(logicID)
if err != nil {
return fmt.Errorf("get pcie bus info of logicID %v failed: %v", logicID, err)
}
// 4. 通过 PCIe 总线信息查找 DPU
pcieSwId, dpuInfo, err := df.getDpuByPcieBusInfo(pcieBusInfo)
if err != nil {
hwlog.RunLog.Errorf(...)
continue
}
// 5. 去重:同一个 PCIe Switch 下的第一个 NPU 对应第一张 DPU,
// 第二个 NPU 对应第二张 DPU
if _, ok := pcieSwIds[pcieSwId]; !ok {
pcieSwIds[pcieSwId] = struct{}{}
df.addDpuByNpuId(logicID, dpuIndexFir, dpuInfo) // dpuIndexFir=0
continue
}
df.addDpuByNpuId(logicID, dpuIndexSec, dpuInfo) // dpuIndexSec=1
}
return nil
}
设计意图:
- 每个 PCIe Switch 下通常连接两台 DPU,通过去重判断第一个 NPU 对应 DPU0,第二个 NPU 对应 DPU1
addDpuByNpuId中有特殊处理:如果只有一台 DPU(onlyOneDpu=1),则不分索引直接添加
4.3.4 PCIe 总线信息查找 getDpuByPcieBusInfo
func (df *DpuFilter) getDpuByPcieBusInfo(pcieBusInfo string) (string, []BaseDpuInfo, error) {
// 1. 通过 PCIe Bus ID 找到 PCIe Switch 路径
pcieSw, err := df.getPcieswByBusId(pcieBusInfo)
if err != nil {
return "", []BaseDpuInfo{}, err
}
// 2. 找到挂在同一 PCIe Switch 下的所有网卡
nics, err := df.getNicsByPcieSw(pcieSw)
if err != nil {
return pcieSw, []BaseDpuInfo{}, err
}
// 3. 过滤出 DPU 设备
df.entries = nics
dpuInfos, err := df.filterDpu()
if err != nil {
return pcieSw, []BaseDpuInfo{}, err
}
if len(dpuInfos) == 0 {
return pcieSw, []BaseDpuInfo{}, fmt.Errorf("filter dpu infos is nil")
}
return pcieSw, dpuInfos, nil
}
4.3.5 PCIe Switch 路径解析 getPcieswByBusId
func (df *DpuFilter) getPcieswByBusId(busId string) (string, error) {
// 1. 读取 /sys/bus/pci/devices/{busId} 的符号链接
targetPath, err := os.Readlink(filepath.Join(pcieSwitchDir, busId))
if err != nil {
return "", err
}
// 2. 转为绝对路径
abs, err := filepath.Abs(targetPath)
if err != nil {
return "", err
}
// 3. 标准化路径,确保以 /sys 开头
absPath := filepath.Join("/sys", strings.TrimPrefix(abs, "/"))
// 4. 截取前 4 级目录作为 PCIe Switch 根路径
// 例如 /sys/devices/pci0000:00/0000:00:1c.0 → 取前4部分
parts := strings.Split(absPath, "/")
if len(parts) < pcieDirLen { // pcieDirLen = 4
return "", fmt.Errorf(...)
}
return strings.Join(parts[:pcieDirLen], "/"), nil
}
路径解析示例:
输入 busId: "0000:00:1c.0"
Readlink 结果: "../../../devices/pci0000:00/0000:00:1c.0"
绝对路径: /sys/devices/pci0000:00/0000:00:1c.0
截取前4级: /sys/devices/pci0000:00
→ 这就是 PCIe Switch 的根路径
4.3.6 DPU 过滤 filterDpu
func (df *DpuFilter) filterDpu() ([]BaseDpuInfo, error) {
configVendors := df.UserConfig.Selectors.Vendor
configDeviceIDs := df.UserConfig.Selectors.DeviceIds
configDeviceNames := df.UserConfig.Selectors.DeviceNames
// 1. 校验 entries 长度
if len(df.entries) == 0 || len(df.entries) > math.MaxInt32 {
return []BaseDpuInfo{}, fmt.Errorf("the length of df.entries is invalid: %v", len(df.entries))
}
var dpuInfos []BaseDpuInfo
for _, entry := range df.entries {
ifaceName := entry.Name()
dpuPath := filepath.Join(netPath, ifaceName)
// 2. 读取网卡的符号链接,获取 PCI 设备路径
dpuDir, err := os.Readlink(dpuPath)
if err != nil {
hwlog.RunLog.Errorf(...)
continue
}
ifacePath := filepath.Join(filepath.Dir(dpuPath), dpuDir)
dpuDeviceDirPath := filepath.Join(ifacePath, deviceDir) // /sys/class/net/{iface}/device
// 3. 三维度过滤:Vendor、DeviceID、DeviceName
isVendorFiltered, vendorValue := df.shouldFilterByVendor(dpuDeviceDirPath, configVendors)
isDeviceIDFiltered, deviceIDValue := df.shouldFilterByDeviceID(dpuDeviceDirPath, configDeviceIDs)
isDeviceNameFiltered := df.shouldFilterByDeviceName(ifaceName, configDeviceNames)
if isVendorFiltered || isDeviceIDFiltered || isDeviceNameFiltered {
continue // 任一条件命中则跳过
}
// 4. 获取网卡 IP 地址
ips := getInterfaceIPs(ifaceName)
// 5. 构造 DPU 信息
dpuInfos = append(dpuInfos, BaseDpuInfo{
DeviceName: ifaceName,
DpuIP: ips,
Vendor: vendorValue,
DeviceId: deviceIDValue,
Operstate: api.DpuStatusDown, // 初始状态为 down
})
}
return dpuInfos, nil
}
4.3.7 通用字段过滤 shouldFilterByField
func (df *DpuFilter) shouldFilterByField(basePath, fileName string, allowed []string) (bool, string) {
// 1. 读取 sysfs 中的文件内容(如 /sys/class/net/eth0/device/vendor)
value, err := readFileContent(filepath.Join(basePath, fileName))
if err != nil {
hwlog.RunLog.Errorf(...)
return true, "" // 读取失败也过滤掉
}
// 2. 如果 allowed 列表为空,不过滤(接受所有)
if len(allowed) > 0 && !slices.Contains(allowed, value) {
return true, "" // 不在允许列表中,过滤掉
}
return false, value // 通过过滤
}
设计意图:
shouldFilterByVendor:读取vendor文件,与配置的 Vendor 列表比对shouldFilterByDeviceID:读取device文件,与配置的 DeviceIds 列表比对shouldFilterByDeviceName:直接比对网卡名,无需读取文件- 当
allowed为空时不过滤,表示该维度不启用过滤
4.3.8 UB 模式 NPU-DPU 映射 getNpuCorrespDpuInfo
func (df *DpuFilter) getNpuCorrespDpuInfo() error {
// NPU 0-3 对应 slot 1 和 slot 9 的 DPU
// NPU 4-7 对应 slot 2 和 slot 10 的 DPU
for npuId := 0; npuId < api.NpuCountPerNode; npuId++ {
if npuId < npuIdxCorrespDpuRangeMiddle { // 4
dpuInfos := df.getDpuPair(dpuSlotIdx1, dpuSlotIdx9) // slot "1" 和 "9"
if len(dpuInfos) != dpuIpAddrsLen { // 必须正好 2 台 DPU
return fmt.Errorf("get npu %d correspond dpuinfos error", npuId)
}
df.NpuWithDpuInfos = append(df.NpuWithDpuInfos, NpuWithDpuInfo{
NpuId: int32(npuId),
DpuInfo: dpuInfos,
})
}
if npuId >= npuIdxCorrespDpuRangeMiddle {
dpuInfos := df.getDpuPair(dpuSlotIdx2, dpuSlotIdx10) // slot "2" 和 "10"
if len(dpuInfos) != dpuIpAddrsLen {
return fmt.Errorf("get npu %d correspond dpuinfos error", npuId)
}
df.NpuWithDpuInfos = append(df.NpuWithDpuInfos, NpuWithDpuInfo{
NpuId: int32(npuId),
DpuInfo: dpuInfos,
})
}
}
return nil
}
NPU-DPU 映射关系(UB 模式):
4.3.9 Slot ID 获取 getSlotId
func (df *DpuFilter) getSlotId(ifaceName string) (string, error) {
// 1. 读取网卡符号链接
dpuPath := filepath.Join(netPath, ifaceName)
dpuDir, err := os.Readlink(dpuPath)
if err != nil {
return "", fmt.Errorf("readlink %s error:%v", dpuPath, err)
}
// 2. 构造 device 目录路径
ifacePath := filepath.Join(filepath.Dir(dpuPath), dpuDir)
dpuDeviceDirPath := filepath.Join(ifacePath, deviceDir)
// 3. 只有 UB 模式才能读取 slot_id
if df.UserConfig.BusType == busTypeUb {
slotID, err := readFileContent(filepath.Join(dpuDeviceDirPath, slotIdFile))
if err != nil {
return "", fmt.Errorf("dpu %s read slot_id error:%v", ifaceName, err)
}
return slotID, nil
}
return "", fmt.Errorf("busType is %s not ub", df.UserConfig.BusType)
}
4.3.10 IP 地址获取 getInterfaceIPs
func getInterfaceIPs(ifaceName string) string {
// 1. 通过 net.InterfaceByName 获取网卡接口
iface, err := net.InterfaceByName(ifaceName)
if err != nil {
hwlog.RunLog.Errorf(...)
return ""
}
// 2. 获取所有地址
addrs, err := iface.Addrs()
if err != nil || len(addrs) == 0 {
hwlog.RunLog.Errorf(...)
return ""
}
// 3. 取第一个地址的 IP 部分
ipNet, ok := addrs[0].(*net.IPNet)
if !ok {
hwlog.RunLog.Errorf(...)
return ""
}
return ipNet.IP.String()
}
设计意图:只取第一个 IP 地址,DPU 通常只有一个管理 IP。如果获取失败返回空字符串(不阻断流程)。
鲲鹏昇腾开发者社区是面向全社会开放的“联接全球计算开发者,聚合华为+生态”的社区,内容涵盖鲲鹏、昇腾资源,帮助开发者快速获取所需的知识、经验、软件、工具、算力,支撑开发者易学、好用、成功,成为核心开发者。
更多推荐


所有评论(0)