Part 4: Ascend Device Plugin 源码超深度分析

310/310P 设备 + 交换机 + DPU + 重置信息管理

分析范围pkg/device/ascend310.gopkg/device/ascend310p.gopkg/device/deviceswitch/ascend_switch.gopkg/device/dpucontrol/dpu_device_find.gopkg/device/dpucontrol/types.gopkg/device/reset_info_mgr.go

源码版本:mind-cluster-v26.0.1

分析者:资深 Go 架构师视角


目录


1. ascend310.go — Ascend 310 设备管理器

1.1 模块定位

业务职责

HwAscend310Manager 是华为 Ascend 310 NPU 芯片的设备管理器实现。它负责:

  1. 设备发现:通过 devmanager 接口枚举节点上的所有 Ascend 310 设备
  2. 设备模式管理:支持普通模式和共享模式两种设备分配方式
  3. 设备状态上报:将健康/不健康设备列表、故障码等信息写入 Kubernetes Node Annotation
  4. Volcano 调度集成:通过 DoWithVolcanoListAndWatch 实现 Volcano 调度器的设备列表与监听
在系统中的位置

外部依赖

Ascend Device Plugin

pkg/device

DevManager 接口

AscendTools 基类

HwAscend310Manager

HwAscend310PManager

HwAscend910Manager

devmanager.DeviceInterface

kubeclient.ClientK8s

pkg/common

Ascend 310 管理器是设备管理器接口的一个具体实现,面向 Ascend 310 芯片场景(边缘推理、轻量训练),在继承 AscendTools 基类能力的基础上,实现了 310 特有的设备发现和状态上报逻辑。

1.2 模块整体结构

类结构

«interface»

DevManager

+GetNPUs()(NpuAllInfo, error)

+DoWithVolcanoListAndWatch(map, int)

+GraceTolerance(ctx, map)

+GetAssociatedLogicIDs(logicID, cardID, deviceID)([]int32, error)

+SetDpu(string, []DpuCMData, map)

AscendTools

-name: string

-unHealthyKey: string

-devCount: int32

-cardInResetMap: map[int32]bool

-resetFailedTimesMap: map[int32]int

-lastUsedChipsContainerMap: map[string]sets.String

+dmgr: DeviceInterface

+client: ClientK8s

HwAscend310Manager

+AscendTools

+GetNPUs()(NpuAllInfo, error)

+DoWithVolcanoListAndWatch(classifyDevs, chipMemory)

-updateDeviceInfo(old, new, devStatusSet) : error

+GraceTolerance(ctx, map)

+GetAssociatedLogicIDs(logicID, cardID, deviceID)([]int32, error)

+SetDpu(string, []DpuCMData, map)

-getNPUsByNormalMode(davinCiDev) : []NpuDevice

核心方法清单
方法名 签名 作用 可见性
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

GetNPUs() 入口

调用 dmgr.GetDeviceList()

err != nil?

返回空 NpuAllInfo + err

devNum > devCount?

返回错误: invalid device num

初始化 allDevices 切片

遍历 devList 中的每个 logicID

调用 getDavinCiDev(logicID) 获取设备详情

err != nil?

返回错误

ShareDev ?

getNPUsByShareMode(davinCiDev)

getNPUsByNormalMode(davinCiDev)

append 到 allDevices

还有更多设备?

返回 NpuAllInfo{AllDevs, AllDevTypes}

逐行解析

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-0
  • assembleNpuDeviceStructAscendTools 基类中的方法,负责将 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 是基类方法,分析所有设备的状态(健康/不健康/网络不健康等),返回 DevStatusSet
  • UpdateNodeDeviceInfo 调用传入的 updateDeviceInfo 回调函数,将状态写入 Node Annotation
  • common.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:
    1. 健康设备:key=HuaweiAscend310,value如"Ascend310-0,Ascend310-1"
    2. 不健康设备:key=Ascend310-Unhealthy
    3. 故障码:key=HuaweiFaultCodeAscend310(注意:源码中此常量值实际为 api.HuaweiAscend310P + "-Fault",这是源码中的一个已知问题)
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 支持更丰富的功能:

  1. 虚拟化设备切分:支持 vNPU 虚拟设备,可将一张物理卡切分为多个虚拟设备
  2. 混合插拔模式Use310PMixedInsert 支持 310P 与其他设备混合使用
  3. 共享模式:支持 ShareDev 模式
  4. AICore 设备:支持 FakeAiCoreDevice 虚拟 AI Core 设备
与 310 的关键差异
特性 Ascend 310 Ascend 310P
虚拟化切分
混合插拔
AICore 设备
最大设备数 256 (64×4) 100 (MaxDevicesNum)
设备名 Ascend310 / davinci-mini Ascend310P
FD 模式前缀 支持 不支持

2.2 模块整体结构

«interface»

DevManager

+GetNPUs()

+DoWithVolcanoListAndWatch()

+GraceTolerance()

+GetAssociatedLogicIDs()

+SetDpu()

AscendTools

-name: string

-devCount: int32

-cardInResetMap: map

-resetFailedTimesMap: map

-lastUsedChipsContainerMap: map

HwAscend310PManager

+AscendTools

+GetNPUs()(NpuAllInfo, error)

+DoWithVolcanoListAndWatch(classifyDevs, chipMemory)

-updateDeviceInfo(old, new, devStatusSet) : error

+GraceTolerance(ctx, map)

+GetAssociatedLogicIDs(logicID, cardID, deviceID)

+SetDpu(string, []DpuCMData, map)

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 — 核心复杂逻辑

GetNPUs() 入口

dmgr.GetDeviceList()

err != nil?

返回错误

devNum > devCount?

返回错误

初始化 allDevices, aiCoreDevices, allDeviceTypes

遍历 devList (i=0..devNum-1)

getDavinCiDev(devList[i])

err?

返回错误

Use310PMixedInsert?

assemble310PMixedPhyDevices

continue 下一个设备

getVirtualDevice(devList[i])

vDevNum > MaxVirtualDeviceNum?

返回错误: invalid virtual device count

vDevNum > 0 && ShareDev?

返回错误: virtual device + shareDev 冲突

!PresetVDevice?

FakeAiCoreDevice()

vDevNum > 0?

assembleVirtualDevices()

ShareDev ?

assembleShareModeDevices()

assemblePhyDevices()

continue

还有设备?

removeDuplicate(allDeviceTypes)

返回 NpuAllInfo{AllDevs, AICoreDevs, AllDevTypes}

逐行解析

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 相同,GraceToleranceGetAssociatedLogicIDsSetDpu 均为空实现,310P 同样不支持这些功能。


3. ascend_switch.go — 交换机故障管理

3.1 模块定位

业务职责

SwitchDevManager 负责管理 Ascend 910A3 交换机芯片的故障事件。这是整个 device-plugin 中唯一涉及 CGO 交互的模块,功能包括:

  1. 驱动库加载:通过 dlopen 加载 liblingqu-dcmi.so 动态库
  2. 故障订阅:通过 lq_dcmi_subscribe_fault_event 订阅交换机故障事件
  3. 故障轮询:通过 lq_dcmi_get_fault_info 定时查询故障信息
  4. 故障码组装:将 C 结构体转换为 Go 的 SwitchFaultEvent 并组装标准故障码
  5. 故障等级映射:将故障码映射到不同处理等级(不处理、亚健康、重启请求、预隔离、隔离)
在系统中的位置

C 动态库

Device Plugin 进程

pkg/common

pkg/device/deviceswitch

CGO dlopen

SwitchDevManager

UpdateSwitchFaultLevel()

fault_code.go

device.go

SwitchFaultCode (全局)

SwitchFaultLevelMap (全局)

liblingqu-dcmi.so

lq_dcmi_init

lq_dcmi_subscribe_fault_event

lq_dcmi_get_fault_info

3.2 模块整体结构

CGO 架构

动态库

C 层 (CGO)

Go 层

dlopen

dlsym

dlsym

dlsym

//export

SwitchDevManager

goFaultEventHandler

convertFaultEvent

setExtraFaultInfo

dcmiInit_lq

dcmi_init_lq

lq_dcmi_subscribe_fault_event

lq_dcmi_get_fault_info

event_handler

lqDcmiShutDown

liblingqu-dcmi.so

核心方法清单
方法名 签名 作用
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
        }
    }
}

故障等级体系

Level 0: NotHandle
不处理

Level 1: SubHealth
亚健康

Level 2: RestartRequest
重启请求

Level 3: PreSeparate
预隔离

Level 4: Separate
隔离

解析

  • 每次调用时完全重建 SwitchFaultLevelMap,而非增量更新,确保一致性
  • 使用全局锁 SwitchFaultLevelMapLock 保护并发访问
  • 故障码来源是 common 包中的全局变量(从配置文件加载)
3.3.3 驱动库初始化 InitSwitchDev

InitSwitchDev()

查找 liblingqu-dcmi.so 路径

找到?

返回错误: failed to find switch library

C.CString 转换路径

defer C.free 释放

C.dcmiInit_lq 加载库

retCode == SUCCESS?

返回错误: dcmi lib load failed

C.dcmi_init_lq 初始化驱动

retCode == SUCCESS?

返回错误: dcmi init call failed

日志: init switch library succeeded

返回 nil

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 回调的核心入口,当驱动上报故障事件时被调用:

是: 恢复事件

是: 匹配恢复

否: 不匹配

否: 新故障

C event_handler 回调

goFaultEventHandler(event)

defer TriggerUpdate('switch fault occur')

convertFaultEvent(event)

日志记录故障详情

Assertion == FaultRecover?

遍历当前故障列表

isFaultRecoveredEvent?

从列表中移除

保留在列表中

SetSwitchFaultCode(过滤后列表)

获取当前故障列表

append 新故障

SetSwitchFaultCode(追加后列表)

//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" 155912na
PeerDeviceName 对端设备类型 cpu/npu/L2/na
na 保留位 na

对端设备类型映射

ChipOrCpu

NpuPort

L2Port

Default

PeerPortDevice=0

cpu

PeerPortDevice=1

npu

PeerPortDevice=2

L2

其他

na

3.3.6 定时故障查询 GetSwitchFaultCodeByInterval

ctx.Done

ticker.C

GetSwitchFaultCodeByInterval(ctx, interval)

runtime.LockOSThread()

updateSwitchFaultCode(true) 初始查询

创建 ticker(interval)

defer ticker.Stop()

select
循环

日志: received stop signal

return 退出

updateSwitchFaultCode(false)

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 的映射关系。功能包括:

  1. 配置加载:从 /user/mindx-dl/dpu/dpu-config.json 加载用户 DPU 配置
  2. 两种总线模式
    • UB 模式busTypeUb):通过 /sys/class/net 目录扫描 DPU 网卡
    • PCIe 模式busTypePcie):通过 NPU 的 PCIe 总线信息反向查找连接的 DPU
  3. 设备过滤:按 Vendor、DeviceId、DeviceName 三维度过滤
  4. NPU-DPU 映射:建立 NPU 与 DPU 的对应关系
在系统中的位置

devmanager

调用方

DPU 发现流程

dpu-config.json

DpuFilter

UB 模式
/sys/class/net

PCIe 模式
/sys/bus/pci/devices

NpuWithDpuInfos

pkg/server/dpu.go

DeviceInterface

4.2 模块整体结构

DpuFilter

+NpuWithDpuInfos: []NpuWithDpuInfo

+UserConfig: UserDpuConfig

-entries: []os.DirEntry

-dpuInfos: []BaseDpuInfo

+SaveDpuConfToNode(dmgr) : error

-getDpuWithNpuPcieSwitch(dmgr) : error

-addDpuByNpuId(logicID, dpuIndex, dpuInfo)

-getDpuByPcieBusInfo(busInfo)(string, []BaseDpuInfo, error)

-getPcieswByBusId(busId)(string, error)

-getNicsByPcieSw(busId)([]DirEntry, error)

-filterDpu()([]BaseDpuInfo, error)

-shouldFilterByField(basePath, fileName, allowed)(bool, string)

-shouldFilterByVendor(path, vendors)(bool, string)

-shouldFilterByDeviceID(path, deviceIDs)(bool, string)

-shouldFilterByDeviceName(ifaceName, deviceNames) : bool

-getNpuCorrespDpuInfo() : error

-getDpuPair(slot1, slot2) : []BaseDpuInfo

-loadDpuConfigFromFile() : error

-getSlotId(ifaceName)(string, error)

BaseDpuInfo

+Operstate: string

+DeviceName: string

+DpuIP: string

+Vendor: string

+DeviceId: string

NpuWithDpuInfo

+NpuId: int32

+DpuInfo: []BaseDpuInfo

UserDpuConfig

+BusType: string

+Selectors: *DeviceSelectors

DeviceSelectors

+Vendor: []string

+DeviceIds: []string

+DeviceNames: []string

4.3 核心业务逻辑深度解析

4.3.1 主入口 SaveDpuConfToNode

ub

pcie

其他

SaveDpuConfToNode(dmgr)

DpuConfigPath 存在?

返回 nil (DPU 未配置)

loadDpuConfigFromFile()

err?

返回错误

BusType?

读取 /sys/class/net 目录

filterDpu() 过滤 DPU

getNpuCorrespDpuInfo() 建立映射

getDpuWithNpuPcieSwitch(dmgr)

返回错误: unsupported busType

NpuWithDpuInfos 为空?

返回错误: filter result is nil

日志: successfully get DPU infos

返回 nil

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

否: 第一张 DPU

是: 第二张 DPU

getDpuWithNpuPcieSwitch(dcMgr)

dcMgr.GetDeviceList() 获取所有 NPU

err 或 devNum==0?

返回错误

初始化 pcieSwIds 去重 map

遍历 logicIDList

dcMgr.GetPCIeBusInfo(logicID)

err?

返回错误

getDpuByPcieBusInfo(busInfo)

err?

日志错误, continue

pcieSwId 已存在?

addDpuByNpuId(logicID, 0, dpuInfo)

addDpuByNpuId(logicID, 1, dpuInfo)

还有更多?

返回 nil

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

getDpuByPcieBusInfo(pcieBusInfo)

getPcieswByBusId(busId)

os.Readlink(/sys/bus/pci/devices/{busId})

解析绝对路径

截取前4级目录作为 PCIe Switch 路径

err?

返回错误

getNicsByPcieSw(pcieSw)

遍历 /sys/class/net

os.Readlink 检查每个网卡

路径包含 pcieSw?

加入 nics 列表

跳过

还有更多?

filterDpu() 过滤

返回 (pcieSw, dpuInfos, nil)

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

filterDpu()

entries 长度合法?

返回错误

遍历 entries

读取网卡符号链接

构造 device 目录路径

shouldFilterByVendor?

shouldFilterByDeviceID?

shouldFilterByDeviceName?

任一过滤命中?

跳过此设备

getInterfaceIPs 获取 IP

构造 BaseDpuInfo

还有更多?

返回 dpuInfos

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

getNpuCorrespDpuInfo()

遍历 npuId 0..7

npuId < 4?

getDpuPair(slot1='1', slot9='9')

getDpuPair(slot2='2', slot10='10')

len == 2?

返回错误

append NpuWithDpuInfo

还有更多?

返回 nil

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 模式):

DPU

NPU 0-3

NPU 0

NPU 1

NPU 2

NPU 3

DPU slot=1

DPU slot=9

DPU

NPU 4-7

NPU 4

NPU 5

NPU 6

NPU 7

DPU slot=2

DPU slot=10

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。如果获取失败返回空字符串(不阻断流程)。


Logo

鲲鹏昇腾开发者社区是面向全社会开放的“联接全球计算开发者,聚合华为+生态”的社区,内容涵盖鲲鹏、昇腾资源,帮助开发者快速获取所需的知识、经验、软件、工具、算力,支撑开发者易学、好用、成功,成为核心开发者。

更多推荐